diff --git a/listener.go b/listener.go new file mode 100644 index 0000000..fb93759 --- /dev/null +++ b/listener.go @@ -0,0 +1,52 @@ +package events + +import ( + "sync" +) + +// ListenFunc is a type for listener functions must implement +type ListenFunc func(event Event) + +// ListenerSink passes events through ListenFunc before draining through the +// next sink below +type ListenerSink struct { + dst Sink + mu sync.Mutex + closed bool + ListenFunc +} + +// NewListenerSink returns a sink that will pass events to a listener +func NewListenerSink(dst Sink, fn ListenFunc) *ListenerSink { + ls := &ListenerSink{ + dst: dst, + ListenFunc: fn, + closed: false, + } + return ls +} + +// Write passes the event to a listener and then forwards to the next sink +func (ls *ListenerSink) Write(event Event) error { + ls.mu.Lock() + defer ls.mu.Unlock() + + if ls.closed { + return ErrSinkClosed + } + + ls.ListenFunc(event) + + return ls.dst.Write(event) +} + +// Close calls Close() on the sink below +func (ls *ListenerSink) Close() error { + if ls.closed { + return nil + } + + // set closed flag + ls.closed = true + return ls.dst.Close() +} diff --git a/listener_test.go b/listener_test.go new file mode 100644 index 0000000..eada0ee --- /dev/null +++ b/listener_test.go @@ -0,0 +1,67 @@ +package events + +import ( + "fmt" + "sync" + "testing" +) + +func TestListener(t *testing.T) { + const nevents = 100 + tm := &testMetrics{0, 0, 0} + + ts := newTestSink(t, nevents) + es := NewListenerSink(ts, tm.egress) + eq := NewQueue(es) + is := NewListenerSink(eq, tm.ingress) + + var wg sync.WaitGroup + for i := 1; i <= nevents; i++ { + wg.Add(1) + go func(event Event) { + if err := is.Write(event); err != nil { + t.Fatalf("error writing event: %v", err) + } + wg.Done() + }("event-" + fmt.Sprint(i)) + } + wg.Wait() + checkClose(t, is) + + ts.mu.Lock() + defer ts.mu.Unlock() + + t.Logf("%#v", tm) + + if len(ts.events) != nevents { + t.Fatalf("events did not make it to the sink: %d != %d", len(ts.events), 1000) + } + + if tm.events != nevents && tm.incoming != nevents && tm.outgoing != nevents { + t.Fatalf("events, incoming, outgoing should all == %d, %#v", nevents, tm) + } + + if !ts.closed { + t.Fatalf("sink should have been closed") + } + // t.Fatalf("%#v", tm) +} + +type endSink struct { + Sink +} + +type testMetrics struct { + events int + incoming int + outgoing int +} + +func (m *testMetrics) ingress(event Event) { + m.events++ + m.incoming++ +} + +func (m *testMetrics) egress(event Event) { + m.outgoing++ +}