From 7617aeadda0fefd777c66ca22526310295e4fcc6 Mon Sep 17 00:00:00 2001 From: 9seconds Date: Tue, 16 Mar 2021 14:54:17 +0300 Subject: [PATCH] Test event stream --- events/event_stream.go | 8 +-- events/event_stream_test.go | 122 ++++++++++++++++++++++++++++++++++++ events/init.go | 2 +- events/init_test.go | 2 +- mtglib/events.go | 4 +- 5 files changed, 130 insertions(+), 8 deletions(-) create mode 100644 events/event_stream_test.go diff --git a/events/event_stream.go b/events/event_stream.go index 0a053b4..be860f6 100644 --- a/events/event_stream.go +++ b/events/event_stream.go @@ -18,12 +18,10 @@ type eventStream struct { func (e eventStream) Send(ctx context.Context, evt mtglib.Event) { var chanNo uint32 - streamID := evt.StreamID() - - if streamID == "" { - chanNo = rand.Uint32() - } else { + if streamID := evt.StreamID(); streamID != "" { chanNo = xxhash.ChecksumString32(streamID) + } else { + chanNo = rand.Uint32() } select { diff --git a/events/event_stream_test.go b/events/event_stream_test.go new file mode 100644 index 0000000..b86b68b --- /dev/null +++ b/events/event_stream_test.go @@ -0,0 +1,122 @@ +package events_test + +import ( + "context" + "net" + "testing" + "time" + + "github.com/9seconds/mtg/v2/events" + "github.com/9seconds/mtg/v2/mtglib" + "github.com/stretchr/testify/mock" + "github.com/stretchr/testify/suite" +) + +type EventStreamTestSuite struct { + suite.Suite + + ctx context.Context + ctxCancel context.CancelFunc + observerMock1 *ObserverMock + observerMock2 *ObserverMock + stream mtglib.EventStream +} + +func (suite *EventStreamTestSuite) SetupTest() { + suite.ctx, suite.ctxCancel = context.WithCancel(context.Background()) + + suite.observerMock1 = &ObserverMock{} + suite.observerMock2 = &ObserverMock{} + + suite.observerMock1.On("Shutdown") + suite.observerMock2.On("Shutdown") + + factories := make([]events.ObserverFactory, 2) + factories[0] = func() events.Observer { return suite.observerMock1 } + factories[1] = func() events.Observer { return suite.observerMock2 } + + suite.stream = events.NewEventStream(factories) +} + +func (suite *EventStreamTestSuite) TestEventStartOk() { + evt := mtglib.EventStart{ + CreatedAt: time.Now(), + ConnID: "connID", + RemoteIP: net.ParseIP("10.0.0.1"), + } + + for _, v := range []*ObserverMock{suite.observerMock1, suite.observerMock2} { + v. + On("EventStart", mock.Anything). + Once(). + Run(func(args mock.Arguments) { + caught := args.Get(0).(mtglib.EventStart) + + suite.Equal(evt.CreatedAt, caught.CreatedAt) + suite.Equal(evt.ConnID, caught.ConnID) + suite.Equal(evt.RemoteIP.String(), caught.RemoteIP.String()) + suite.Equal(evt.StreamID(), caught.StreamID()) + }) + } + + suite.stream.Send(suite.ctx, evt) + time.Sleep(100 * time.Millisecond) +} + +func (suite *EventStreamTestSuite) TestEventFinishOk() { + evt := mtglib.EventFinish{ + CreatedAt: time.Now(), + ConnID: "connID", + } + + for _, v := range []*ObserverMock{suite.observerMock1, suite.observerMock2} { + v. + On("EventFinish", mock.Anything). + Once(). + Run(func(args mock.Arguments) { + caught := args.Get(0).(mtglib.EventFinish) + + suite.Equal(evt.CreatedAt, caught.CreatedAt) + suite.Equal(evt.ConnID, caught.ConnID) + suite.Equal(evt.StreamID(), caught.StreamID()) + }) + } + + suite.stream.Send(suite.ctx, evt) + time.Sleep(100 * time.Millisecond) +} + +func (suite *EventStreamTestSuite) TestEventConcurrencyLimitedOk() { + evt := mtglib.EventConcurrencyLimited{ + CreatedAt: time.Now(), + } + + for _, v := range []*ObserverMock{suite.observerMock1, suite.observerMock2} { + v. + On("EventConcurrencyLimited", mock.Anything). + Once(). + Run(func(args mock.Arguments) { + caught := args.Get(0).(mtglib.EventConcurrencyLimited) + + suite.Equal(evt.CreatedAt, caught.CreatedAt) + }) + } + + suite.stream.Send(suite.ctx, evt) + time.Sleep(100 * time.Millisecond) +} + +func (suite *EventStreamTestSuite) TearDownTest() { + suite.stream.Shutdown() + suite.ctxCancel() + + time.Sleep(100 * time.Millisecond) + + suite.observerMock1.AssertExpectations(suite.T()) + suite.observerMock2.AssertExpectations(suite.T()) +} + +func TestEventStream(t *testing.T) { + t.Parallel() + suite.Run(t, &EventStreamTestSuite{}) +} diff --git a/events/init.go b/events/init.go index 9a632f4..0e348f5 100644 --- a/events/init.go +++ b/events/init.go @@ -5,7 +5,7 @@ import "github.com/9seconds/mtg/v2/mtglib" type Observer interface { EventStart(mtglib.EventStart) EventFinish(mtglib.EventFinish) - EventConcurrencyLimited(mtglib.EventConcurrencyLimited) + EventConcurrencyLimited(mtglib.EventConcurrencyLimited) Shutdown() } diff --git a/events/init_test.go b/events/init_test.go index c67efa9..e75e5e9 100644 --- a/events/init_test.go +++ b/events/init_test.go @@ -13,7 +13,7 @@ func (o *ObserverMock) EventStart(evt mtglib.EventStart) { o.Called(evt) } -func (o *ObserverMock) EventFinish(evt mtglib.EventStart) { +func (o *ObserverMock) EventFinish(evt mtglib.EventFinish) { o.Called(evt) } diff --git a/mtglib/events.go b/mtglib/events.go index 61f282a..e7ea557 100644 --- a/mtglib/events.go +++ b/mtglib/events.go @@ -24,7 +24,9 @@ func (e EventFinish) StreamID() string { return e.ConnID } -type EventConcurrencyLimited struct{} +type EventConcurrencyLimited struct { + CreatedAt time.Time +} func (e EventConcurrencyLimited) StreamID() string { return ""