package gateway import ( "context" "errors" "sync" ) var ErrQueueClosed = errors.New("gateway queue closed") // BoundedQueue is deliberately fixed-size. Media uses PushLatest so a slow // client drops old frames instead of allowing provider output to accumulate. type BoundedQueue[T any] struct { mu sync.Mutex items []T limit int dropped uint64 closed bool wake chan struct{} } func NewBoundedQueue[T any](limit int) *BoundedQueue[T] { if limit < 1 { limit = 1 } return &BoundedQueue[T]{limit: limit, wake: make(chan struct{}, 1)} } func (q *BoundedQueue[T]) PushLatest(item T) error { q.mu.Lock() defer q.mu.Unlock() if q.closed { return ErrQueueClosed } if len(q.items) == q.limit { var zero T q.items[0] = zero q.items = q.items[1:] q.dropped++ } q.items = append(q.items, item) select { case q.wake <- struct{}{}: default: } return nil } func (q *BoundedQueue[T]) Pop(ctx context.Context) (T, error) { for { q.mu.Lock() if len(q.items) > 0 { item := q.items[0] q.items[0] = *new(T) q.items = q.items[1:] q.mu.Unlock() return item, nil } if q.closed { q.mu.Unlock() var zero T return zero, ErrQueueClosed } q.mu.Unlock() select { case <-ctx.Done(): var zero T return zero, ctx.Err() case <-q.wake: } } } func (q *BoundedQueue[T]) Close() { q.mu.Lock() if q.closed { q.mu.Unlock() return } q.closed = true select { case q.wake <- struct{}{}: default: } q.mu.Unlock() } func (q *BoundedQueue[T]) Dropped() uint64 { q.mu.Lock() defer q.mu.Unlock() return q.dropped } func (q *BoundedQueue[T]) Len() int { q.mu.Lock() defer q.mu.Unlock() return len(q.items) }