Files
VerseVDI-Data-Plane/gateway/queue.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)
}