streampump

package
v0.1.0 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Oct 5, 2026 License: Apache-2.0 Imports: 5 Imported by: 0

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

View Source
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.

func (EventPump[T]) Recv

func (p EventPump[T]) Recv(ctx context.Context) (lipapi.Event, error)

Recv returns the next canonical event, draining the pending queue before reading wire input.

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.

func (*PendingEventQueue) PopFront

func (q *PendingEventQueue) PopFront() (lipapi.Event, bool)

PopFront removes and returns the oldest event. The second result is false when empty.

func (*PendingEventQueue) Push

func (q *PendingEventQueue) Push(ev lipapi.Event) error

Push appends an event to the tail.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL