99 lines
1.6 KiB
Go
99 lines
1.6 KiB
Go
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)
|
|
}
|