From e009d05a906f34c7d50cf9de7fd907b8afada33a Mon Sep 17 00:00:00 2001 From: 9seconds Date: Tue, 16 Mar 2021 11:28:34 +0300 Subject: [PATCH] Add tests for noop event stream --- events/event_stream.go | 25 ++++++++++++-- events/init.go | 1 + events/init_test.go | 4 +++ events/multi_observer.go | 15 ++++++++ events/noop.go | 27 +++++++++++++++ events/noop_test.go | 75 ++++++++++++++++++++++++++++++++++++++++ mtglib/events.go | 24 ++++++++----- mtglib/init.go | 2 +- 8 files changed, 160 insertions(+), 13 deletions(-) create mode 100644 events/noop.go create mode 100644 events/noop_test.go diff --git a/events/event_stream.go b/events/event_stream.go index 577e713..0a053b4 100644 --- a/events/event_stream.go +++ b/events/event_stream.go @@ -2,6 +2,7 @@ package events import ( "context" + "math/rand" "runtime" "github.com/9seconds/mtg/v2/mtglib" @@ -15,12 +16,20 @@ type eventStream struct { } func (e eventStream) Send(ctx context.Context, evt mtglib.Event) { - chanNo := int(xxhash.ChecksumString32(evt.ConnectionID())) % len(e.chans) + var chanNo uint32 + + streamID := evt.StreamID() + + if streamID == "" { + chanNo = rand.Uint32() + } else { + chanNo = xxhash.ChecksumString32(streamID) + } select { case <-ctx.Done(): case <-e.ctx.Done(): - case e.chans[chanNo] <- evt: + case e.chans[int(chanNo)%len(e.chans)] <- evt: } } @@ -29,6 +38,10 @@ func (e eventStream) Shutdown() { } func NewEventStream(observerFactories []ObserverFactory) mtglib.EventStream { + if len(observerFactories) == 0 { + observerFactories = append(observerFactories, NewNoopObserver) + } + ctx, cancel := context.WithCancel(context.Background()) rv := eventStream{ ctx: ctx, @@ -39,7 +52,11 @@ func NewEventStream(observerFactories []ObserverFactory) mtglib.EventStream { for i := 0; i < runtime.NumCPU(); i++ { rv.chans[i] = make(chan mtglib.Event, 1) - go eventStreamProcessor(ctx, rv.chans[i], newMultiObserver(observerFactories)) + if len(observerFactories) == 1 { + go eventStreamProcessor(ctx, rv.chans[i], observerFactories[0]()) + } else { + go eventStreamProcessor(ctx, rv.chans[i], newMultiObserver(observerFactories)) + } } return rv @@ -58,6 +75,8 @@ func eventStreamProcessor(ctx context.Context, eventChan <-chan mtglib.Event, ob observer.EventStart(typedEvt) case mtglib.EventFinish: observer.EventFinish(typedEvt) + case mtglib.EventConcurrencyLimited: + observer.EventConcurrencyLimited(typedEvt) } } } diff --git a/events/init.go b/events/init.go index f84b03b..9a632f4 100644 --- a/events/init.go +++ b/events/init.go @@ -5,6 +5,7 @@ import "github.com/9seconds/mtg/v2/mtglib" type Observer interface { EventStart(mtglib.EventStart) EventFinish(mtglib.EventFinish) + EventConcurrencyLimited(mtglib.EventConcurrencyLimited) Shutdown() } diff --git a/events/init_test.go b/events/init_test.go index 829ca5c..c67efa9 100644 --- a/events/init_test.go +++ b/events/init_test.go @@ -17,6 +17,10 @@ func (o *ObserverMock) EventFinish(evt mtglib.EventStart) { o.Called(evt) } +func (o *ObserverMock) EventConcurrencyLimited(evt mtglib.EventConcurrencyLimited) { + o.Called(evt) +} + func (o *ObserverMock) Shutdown() { o.Called() } diff --git a/events/multi_observer.go b/events/multi_observer.go index 9370fb0..698d87c 100644 --- a/events/multi_observer.go +++ b/events/multi_observer.go @@ -40,6 +40,21 @@ func (m multiObserver) EventFinish(evt mtglib.EventFinish) { wg.Wait() } +func (m multiObserver) EventConcurrencyLimited(evt mtglib.EventConcurrencyLimited) { + wg := &sync.WaitGroup{} + wg.Add(len(m.observers)) + + for _, v := range m.observers { + go func(obs Observer) { + defer wg.Done() + + obs.EventConcurrencyLimited(evt) + }(v) + } + + wg.Wait() +} + func (m multiObserver) Shutdown() { for _, v := range m.observers { v.Shutdown() diff --git a/events/noop.go b/events/noop.go new file mode 100644 index 0000000..7eb4041 --- /dev/null +++ b/events/noop.go @@ -0,0 +1,27 @@ +package events + +import ( + "context" + + "github.com/9seconds/mtg/v2/mtglib" +) + +type noop struct{} + +func (n noop) Send(ctx context.Context, evt mtglib.Event) {} +func (n noop) Shutdown() {} + +func NewNoopStream() mtglib.EventStream { + return noop{} +} + +type noopObserver struct{} + +func (n noopObserver) EventStart(_ mtglib.EventStart) {} +func (n noopObserver) EventFinish(_ mtglib.EventFinish) {} +func (n noopObserver) EventConcurrencyLimited(_ mtglib.EventConcurrencyLimited) {} +func (n noopObserver) Shutdown() {} + +func NewNoopObserver() Observer { + return noopObserver{} +} diff --git a/events/noop_test.go b/events/noop_test.go new file mode 100644 index 0000000..b5120f3 --- /dev/null +++ b/events/noop_test.go @@ -0,0 +1,75 @@ +package events_test + +import ( + "context" + "net" + "testing" + "time" + + "github.com/9seconds/mtg/v2/events" + "github.com/9seconds/mtg/v2/mtglib" + "github.com/stretchr/testify/suite" +) + +type NoopTestSuite struct { + suite.Suite + + testData map[string]mtglib.Event + ctx context.Context +} + +func (suite *NoopTestSuite) SetupSuite() { + suite.testData = map[string]mtglib.Event{ + "start": mtglib.EventStart{ + CreatedAt: time.Now(), + ConnID: "connID", + RemoteIP: net.ParseIP("127.0.0.1"), + }, + "finish": mtglib.EventFinish{ + CreatedAt: time.Now(), + ConnID: "connID", + }, + "concurrency-limited": mtglib.EventConcurrencyLimited{}, + } + suite.ctx = context.Background() +} + +func (suite *NoopTestSuite) TestStream() { + stream := events.NewNoopStream() + + for name, v := range suite.testData { + value := v + + suite.T().Run(name, func(t *testing.T) { + stream.Send(suite.ctx, value) + }) + } + + stream.Shutdown() +} + +func (suite *NoopTestSuite) TestObserver() { + observer := events.NewNoopObserver() + + for name, v := range suite.testData { + value := v + + suite.T().Run(name, func(t *testing.T) { + switch typedEvt := value.(type) { + case mtglib.EventStart: + observer.EventStart(typedEvt) + case mtglib.EventFinish: + observer.EventFinish(typedEvt) + case mtglib.EventConcurrencyLimited: + observer.EventConcurrencyLimited(typedEvt) + } + }) + } + + observer.Shutdown() +} + +func TestNoop(t *testing.T) { + t.Parallel() + suite.Run(t, &NoopTestSuite{}) +} diff --git a/mtglib/events.go b/mtglib/events.go index ef361a2..61f282a 100644 --- a/mtglib/events.go +++ b/mtglib/events.go @@ -5,21 +5,27 @@ import ( "time" ) -type eventBase struct { +type EventStart struct { + CreatedAt time.Time + ConnID string + RemoteIP net.IP +} + +func (e EventStart) StreamID() string { + return e.ConnID +} + +type EventFinish struct { CreatedAt time.Time ConnID string } -func (e eventBase) ConnectionID() string { +func (e EventFinish) StreamID() string { return e.ConnID } -type EventStart struct { - eventBase +type EventConcurrencyLimited struct{} - RemoteIP net.IP -} - -type EventFinish struct { - eventBase +func (e EventConcurrencyLimited) StreamID() string { + return "" } diff --git a/mtglib/init.go b/mtglib/init.go index 8a3b764..22c5472 100644 --- a/mtglib/init.go +++ b/mtglib/init.go @@ -28,7 +28,7 @@ type IPBlocklist interface { } type Event interface { - ConnectionID() string + StreamID() string } type EventStream interface {