package msg_test import ( "context" "testing" "time" "gitea.dwysokinski.me/twhelp/core/internal/msg/internal/model" "github.com/stretchr/testify/assert" "github.com/ThreeDotsLabs/watermill/message/subscriber" "github.com/stretchr/testify/require" "gitea.dwysokinski.me/twhelp/core/internal/domain" "gitea.dwysokinski.me/twhelp/core/internal/msg" ) func TestServerPublisher_CmdRefresh(t *testing.T) { t.Parallel() marshaler := msg.GobMarshaler{} pubSub := newPubSub(t) payloads := []domain.RefreshServersCmdPayload{ { URL: "https://host1.com", VersionCode: "vc1", }, { URL: "https://host2.com", VersionCode: "vc2", }, } msgs, err := pubSub.Subscribe(context.Background(), "servers.cmd.refresh") require.NoError(t, err) require.NoError(t, msg.NewServerPublisher(pubSub, marshaler).CmdRefresh(context.Background(), payloads...)) receivedMsgs, _ := subscriber.BulkRead(msgs, len(payloads), time.Second) require.Len(t, receivedMsgs, len(payloads)) for _, m := range receivedMsgs { var received model.RefreshServersCmdPayload assert.NoError(t, marshaler.Unmarshal(m, &received)) found := false for _, payload := range payloads { if payload.VersionCode == received.VersionCode && payload.URL == received.URL { found = true break } } assert.True(t, found) } }