Skip to content
Closed
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
40 changes: 35 additions & 5 deletions queue.go
Original file line number Diff line number Diff line change
Expand Up @@ -11,11 +11,18 @@ import (
// by a sink. It is unbounded and thread safe but the sink must be reliable or
// events will be dropped.
type Queue struct {
dst Sink
events *list.List
cond *sync.Cond
mu sync.Mutex
closed bool
dst Sink
events *list.List
listeners []QueueListener
cond *sync.Cond
mu sync.Mutex
closed bool
}

// QueueListener is called when various events happen on the queue.
type QueueListener interface {
Ingress(event Event)
Egress(event Event)
}

// NewQueue returns a queue to the provided Sink dst.
Expand All @@ -30,6 +37,21 @@ func NewQueue(dst Sink) *Queue {
return &eq
}

// NewListenerQueue returns a queue to the provided Sink dst. If the updater
// is non-nil, it will be called to update pending metrics on ingress and
// egress.
func NewListenerQueue(dst Sink, listeners ...QueueListener) *Queue {
eq := Queue{
dst: dst,
events: list.New(),
listeners: listeners,
}

eq.cond = sync.NewCond(&eq.mu)
go eq.run()
return &eq
}

// Write accepts the events into the queue, only failing if the queue has
// been closed.
func (eq *Queue) Write(event Event) error {
Expand All @@ -40,6 +62,10 @@ func (eq *Queue) Write(event Event) error {
return ErrSinkClosed
}

for _, listener := range eq.listeners {
listener.Ingress(event)
}

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

Expand Down Expand Up @@ -84,6 +110,10 @@ func (eq *Queue) run() {
"sink": eq.dst,
}).WithError(err).Debug("eventqueue: dropped event")
}

for _, listener := range eq.listeners {
listener.Egress(event)
}
}
}

Expand Down
99 changes: 98 additions & 1 deletion queue_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,7 @@ func TestQueue(t *testing.T) {
Sink: ts,
delay: time.Millisecond * 1,
})
time.Sleep(10 * time.Millisecond) // let's queue settle to wait conidition.
time.Sleep(10 * time.Millisecond) // let's queue settle to wait condition.

var wg sync.WaitGroup
for i := 1; i <= nevents; i++ {
Expand All @@ -44,3 +44,100 @@ func TestQueue(t *testing.T) {
t.Fatalf("sink should have been closed")
}
}

func TestListenerQueue(t *testing.T) {
const nevents = 1000
ts := newTestSink(t, nevents)
metrics := newSafeMetrics()
eq := NewListenerQueue(
// delayed sync simulates destination slower than channel comms
&delayedSink{
Sink: ts,
delay: time.Millisecond * 1,
}, metrics.eventQueueListener())
time.Sleep(10 * time.Millisecond) // let's queue settle to wait condition.

var wg sync.WaitGroup
for i := 1; i <= nevents; i++ {
wg.Add(1)
go func(event Event) {
if err := eq.Write(event); err != nil {
t.Fatalf("error writing event: %v", err)
}
wg.Done()
}("event-" + fmt.Sprint(i))
}

wg.Wait()
checkClose(t, eq)

ts.mu.Lock()
defer ts.mu.Unlock()
metrics.Lock()
defer metrics.Unlock()

if len(ts.events) != nevents {
t.Fatalf("events did not make it to the sink: %d != %d", len(ts.events), nevents)
}

if !ts.closed {
t.Fatalf("sink should have been closed")
}

if metrics.Events != nevents {
t.Fatalf("unexpected ingress count: %d != %d", metrics.Events, nevents)
}

if metrics.Pending != 0 {
t.Fatalf("unexpected egress count: %d != %d", metrics.Pending, 0)
}
}

type testMetrics struct {
Pending int // events pending in queue
Events int // total events incoming
Successes int // total events written successfully
Failures int // total events failed
Errors int // total events errored
Statuses map[string]int // status code histogram, per call event
}

// safeMetrics guards the metrics implementation with a lock and provides a
// safe update function.
type safeMetrics struct {
testMetrics
sync.Mutex // protects statuses map
}

// newSafeMetrics returns safeMetrics with map allocated.
func newSafeMetrics() *safeMetrics {
var sm safeMetrics
sm.Statuses = make(map[string]int)
return &sm
}

// eventQueueListener returns a listener that maintains queue related counters.
func (sm *safeMetrics) eventQueueListener() QueueListener {
return &testMetricsEventQueueListener{
safeMetrics: sm,
}
}

// testMetricsEventQueueListener maintains the incoming events counter and
// the queues pending count.
type testMetricsEventQueueListener struct {
*safeMetrics
}

func (eqc *testMetricsEventQueueListener) Ingress(event Event) {
eqc.Lock()
defer eqc.Unlock()
eqc.Events++
eqc.Pending++
}

func (eqc *testMetricsEventQueueListener) Egress(event Event) {
eqc.Lock()
defer eqc.Unlock()
eqc.Pending--
}