diff --git a/queue.go b/queue.go index 4bb770a..2e52c72 100644 --- a/queue.go +++ b/queue.go @@ -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. @@ -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 { @@ -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 @@ -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) + } } } diff --git a/queue_test.go b/queue_test.go index 7bfe0e3..029df49 100644 --- a/queue_test.go +++ b/queue_test.go @@ -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++ { @@ -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-- +}