Documentation
¶
Overview ¶
Package streampump provides a bounded pending-event queue and generic Recv loop shared by core stream adapters and connector backends.
Index ¶
Constants ¶
This section is empty.
Variables ¶
var ErrPendingQueueFull = errors.New("streampump: pending event queue capacity exceeded")
ErrPendingQueueFull is returned when PendingEventQueue.Push would exceed a configured max length.
Functions ¶
func DrainPending ¶
func DrainPending(q *PendingEventQueue) []lipapi.Event
DrainPending pops every queued event in order and returns them. The queue is empty afterward.
Types ¶
type EventPump ¶
type EventPump[T any] struct { Lock *sync.Mutex Pending *PendingEventQueue IsClosed func() bool Read func() (T, bool, error) Handle func(T) error OnEOF func() (bool, error) }
EventPump owns the common "pending queue before wire read" loop used by stream adapters. Callers keep provider-specific reading and mapping in Read/Handle.
type PendingEventQueue ¶
type PendingEventQueue struct {
// contains filtered or unexported fields
}
PendingEventQueue buffers canonical events for adapters that translate one wire chunk into zero or more lipapi.Event values. It avoids slice-prefix dequeue (pending = pending[1:]) which retains a large backing array over long streams.
When constructed with a positive max length (see NewPendingEventQueue), Push returns ErrPendingQueueFull once the queue would exceed that cap. When max length is zero (default), the queue is unbounded until other request limits apply.
func NewPendingEventQueue ¶
func NewPendingEventQueue(maxLen int) PendingEventQueue
NewPendingEventQueue returns a queue with the given max pending events (0 = unlimited).
func (*PendingEventQueue) Len ¶
func (q *PendingEventQueue) Len() int
Len returns the number of queued events.