Dawid Wysokiński
0f0e0cd3ee
All checks were successful
continuous-integration/drone/push Build is passing
Reviewed-on: twhelp/core#40
85 lines
2.5 KiB
Go
85 lines
2.5 KiB
Go
package msg_test
|
|
|
|
import (
|
|
"context"
|
|
"testing"
|
|
"time"
|
|
|
|
"github.com/ThreeDotsLabs/watermill"
|
|
|
|
"gitea.dwysokinski.me/twhelp/core/internal/msg/internal/model"
|
|
|
|
"github.com/ThreeDotsLabs/watermill/message/subscriber"
|
|
|
|
"gitea.dwysokinski.me/twhelp/core/internal/msg"
|
|
"gitea.dwysokinski.me/twhelp/core/internal/msg/internal/mock"
|
|
"github.com/stretchr/testify/assert"
|
|
"github.com/stretchr/testify/require"
|
|
)
|
|
|
|
func TestPlayerConsumer_refresh(t *testing.T) {
|
|
t.Parallel()
|
|
|
|
t.Run("OK", func(t *testing.T) {
|
|
t.Parallel()
|
|
|
|
marshaler := msg.GobMarshaler{}
|
|
pubSub := newPubSub(t)
|
|
|
|
var numPlayers int64 = 12344
|
|
playerSvc := &mock.FakePlayerService{}
|
|
playerSvc.RefreshReturns(numPlayers, nil)
|
|
|
|
runRouter(t, msg.NewPlayerConsumer(marshaler, pubSub, pubSub, watermill.NopLogger{}, playerSvc))
|
|
|
|
msgs, err := pubSub.Subscribe(context.Background(), "players.event.refreshed")
|
|
require.NoError(t, err)
|
|
|
|
serverRefreshedPayload := model.ServerRefreshedEvPayload{
|
|
Key: "pl151",
|
|
URL: "https://pl151.plemiona.pl",
|
|
Open: true,
|
|
VersionCode: "pl",
|
|
}
|
|
ev, err := marshaler.Marshal(serverRefreshedPayload)
|
|
require.NoError(t, err)
|
|
require.NoError(t, pubSub.Publish("servers.event.refreshed", ev))
|
|
|
|
receivedMsgs, _ := subscriber.BulkRead(msgs, 1, 10*time.Second)
|
|
require.Len(t, receivedMsgs, 1)
|
|
var playersRefreshedPayload model.PlayersRefreshedEvPayload
|
|
assert.NoError(t, marshaler.Unmarshal(receivedMsgs[0], &playersRefreshedPayload))
|
|
assert.Equal(t, serverRefreshedPayload.Key, playersRefreshedPayload.Key)
|
|
assert.Equal(t, serverRefreshedPayload.URL, playersRefreshedPayload.URL)
|
|
assert.Equal(t, serverRefreshedPayload.VersionCode, playersRefreshedPayload.VersionCode)
|
|
assert.Equal(t, numPlayers, playersRefreshedPayload.NumPlayers)
|
|
})
|
|
|
|
t.Run("OK: only open servers should be updated", func(t *testing.T) {
|
|
t.Parallel()
|
|
|
|
marshaler := msg.GobMarshaler{}
|
|
pubSub := newPubSub(t)
|
|
playerSvc := &mock.FakePlayerService{}
|
|
|
|
waitMiddleware, wait := newWaitForHandlerMiddleware("PlayerConsumer.refresh")
|
|
runRouter(
|
|
t,
|
|
waitMiddleware,
|
|
msg.NewPlayerConsumer(marshaler, pubSub, pubSub, watermill.NopLogger{}, playerSvc),
|
|
)
|
|
|
|
ev, err := marshaler.Marshal(model.ServerRefreshedEvPayload{Open: false})
|
|
require.NoError(t, err)
|
|
require.NoError(t, pubSub.Publish("servers.event.refreshed", ev))
|
|
|
|
select {
|
|
case <-wait:
|
|
case <-time.After(time.Second):
|
|
t.Fatal("timeout")
|
|
}
|
|
|
|
assert.Equal(t, 0, playerSvc.RefreshCallCount())
|
|
})
|
|
}
|