Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions errors.go
Original file line number Diff line number Diff line change
Expand Up @@ -7,4 +7,9 @@ var (
// closed. If encountered, the error should be considered terminal and
// retries will not be successful.
ErrSinkClosed = fmt.Errorf("events: sink closed")

// ErrQueueFull is returned if a write is issued to a queue that does not
// have enough space to store an additional event. If encountered, further
// replies may be successful if any of the queue elements was consumed.
ErrQueueFull = fmt.Errorf("events: queue full")
)
13 changes: 12 additions & 1 deletion queue.go
Original file line number Diff line number Diff line change
Expand Up @@ -13,16 +13,23 @@ import (
type Queue struct {
dst Sink
events *list.List
limit int
cond *sync.Cond
mu sync.Mutex
closed bool
}

// NewQueue returns a queue to the provided Sink dst.
// NewQueue returns an infinite queue to the provided Sink dst.
func NewQueue(dst Sink) *Queue {
return NewSizedQueue(dst, 0)
}

// NewSizedQueue returns a sized queue to the provided Sink dst.
func NewSizedQueue(dst Sink, limit int) *Queue {
eq := Queue{
dst: dst,
events: list.New(),
limit: limit,
}

eq.cond = sync.NewCond(&eq.mu)
Expand All @@ -40,6 +47,10 @@ func (eq *Queue) Write(event Event) error {
return ErrSinkClosed
}

if eq.limit > 0 && eq.events.Len() >= eq.limit {
return ErrQueueFull
}

eq.events.PushBack(event)
eq.cond.Signal() // signal waiters

Expand Down