package ws import ( "context" "encoding/json" "net/http" "net/http/httptest" "strings" "testing" "time" "github.com/alicebob/miniredis/v2" "github.com/gorilla/websocket" "github.com/redis/go-redis/v9" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" "github.com/ThreeDotsLabs/watermill/message" wspkg "github.com/gochat/gochat/internal/ws" ) // === NewHub tests === func TestNewHub_WithSubsystems(t *testing.T) { mr := miniredis.RunT(t) defer mr.Close() rdb := redis.NewClient(&redis.Options{Addr: mr.Addr()}) relay := wspkg.NewBroadcastRelay(rdb, nil) typing := wspkg.NewTypingTracker(rdb, relay) presence := wspkg.NewPresenceTracker(rdb, relay) pm := wspkg.NewPresenceManager(presence, wspkg.DefaultHeartbeatConfig()) hub := NewHub(relay, typing, presence, pm) require.NotNil(t, hub) assert.NotNil(t, hub.clients) assert.NotNil(t, hub.rooms) assert.NotNil(t, hub.accounts) assert.NotNil(t, hub.commandChan) assert.Same(t, relay, hub.relay) assert.Same(t, typing, hub.typing) assert.Same(t, presence, hub.presence) assert.Same(t, pm, hub.presenceMgr) } // === Hub.Register with presenceMgr (contact path) === func TestHub_Register_ContactWithPresence(t *testing.T) { mr := miniredis.RunT(t) defer mr.Close() rdb := redis.NewClient(&redis.Options{Addr: mr.Addr()}) relay := wspkg.NewBroadcastRelay(rdb, nil) presence := wspkg.NewPresenceTracker(rdb, relay) pm := wspkg.NewPresenceManager(presence, wspkg.DefaultHeartbeatConfig()) hub := NewHub(relay, nil, presence, pm) client := NewClient(0, 10, nil, hub) client.IsContact = true client.ContactID = 5 client.PubsubToken = "token123" hub.Register(client) assert.Contains(t, hub.clients, client.ID) // Contacts only receive events addressed to their own pubsub token. assert.False(t, client.SubscribedRooms[accountRoomName(10)]) assert.True(t, client.SubscribedRooms[pubsubTokenRoomName("token123")]) hub.Unregister(client) } func TestHub_Register_AgentWithPresence(t *testing.T) { mr := miniredis.RunT(t) defer mr.Close() rdb := redis.NewClient(&redis.Options{Addr: mr.Addr()}) relay := wspkg.NewBroadcastRelay(rdb, nil) presence := wspkg.NewPresenceTracker(rdb, relay) pm := wspkg.NewPresenceManager(presence, wspkg.DefaultHeartbeatConfig()) hub := NewHub(relay, nil, presence, pm) client := NewClient(1, 10, nil, hub) hub.Register(client) assert.Contains(t, hub.clients, client.ID) assert.NotNil(t, client.CancelPresence) hub.Unregister(client) assert.NotContains(t, hub.clients, client.ID) } // === Hub.processCommand tests === func TestProcessCommand_UnknownCommand(t *testing.T) { hub := NewHubSimple() client := NewClient(1, 10, nil, hub) hub.Register(client) defer hub.Unregister(client) hub.processCommand(&ClientCommand{ Client: client, Cmd: wspkg.WSCommand{Command: "unknown_command"}, }) // Should not panic, no message sent select { case <-client.Send: t.Fatal("unexpected message for unknown command") default: } } // === Hub.handleSubscribe tests === func TestHubHandleSubscribe_AccountChannel(t *testing.T) { hub := NewHubSimple() client := NewClient(1, 10, nil, hub) hub.Register(client) defer hub.Unregister(client) data, _ := json.Marshal(wspkg.SubscribeData{ Channel: wspkg.ChannelAccount, AccountID: 10, }) hub.processCommand(&ClientCommand{ Client: client, Cmd: wspkg.WSCommand{Command: "subscribe", Data: string(data)}, }) // Should receive subscribe confirmation select { case msg := <-client.Send: var resp wspkg.WSMessage err := json.Unmarshal(msg, &resp) require.NoError(t, err) assert.Equal(t, wspkg.EventSubscribeConfirm, resp.Event) default: t.Fatal("expected subscribe confirmation message") } // Client should be subscribed to account room roomName := accountRoomName(10) assert.True(t, client.SubscribedRooms[roomName]) } func TestHubHandleSubscribe_ConversationChannel(t *testing.T) { hub := NewHubSimple() client := NewClient(1, 10, nil, hub) hub.Register(client) defer hub.Unregister(client) data, _ := json.Marshal(wspkg.SubscribeData{ Channel: wspkg.ChannelConversation, AccountID: 10, ConversationID: 42, }) hub.processCommand(&ClientCommand{ Client: client, Cmd: wspkg.WSCommand{Command: "subscribe", Data: string(data)}, }) select { case msg := <-client.Send: var resp wspkg.WSMessage err := json.Unmarshal(msg, &resp) require.NoError(t, err) assert.Equal(t, wspkg.EventSubscribeConfirm, resp.Event) default: t.Fatal("expected subscribe confirmation message") } roomName := conversationRoomName(10, 42) assert.True(t, client.SubscribedRooms[roomName]) } func TestHubHandleSubscribe_AccountMismatch(t *testing.T) { hub := NewHubSimple() client := NewClient(1, 10, nil, hub) hub.Register(client) defer hub.Unregister(client) data, _ := json.Marshal(wspkg.SubscribeData{ Channel: wspkg.ChannelAccount, AccountID: 999, // mismatch }) hub.processCommand(&ClientCommand{ Client: client, Cmd: wspkg.WSCommand{Command: "subscribe", Data: string(data)}, }) select { case msg := <-client.Send: var resp wspkg.WSMessage err := json.Unmarshal(msg, &resp) require.NoError(t, err) assert.Equal(t, wspkg.EventSubscribeReject, resp.Event) default: t.Fatal("expected subscribe reject message") } } func TestHubHandleSubscribe_ConversationNoID(t *testing.T) { hub := NewHubSimple() client := NewClient(1, 10, nil, hub) hub.Register(client) defer hub.Unregister(client) data, _ := json.Marshal(wspkg.SubscribeData{ Channel: wspkg.ChannelConversation, AccountID: 10, // ConversationID omitted }) hub.processCommand(&ClientCommand{ Client: client, Cmd: wspkg.WSCommand{Command: "subscribe", Data: string(data)}, }) select { case msg := <-client.Send: var resp wspkg.WSMessage err := json.Unmarshal(msg, &resp) require.NoError(t, err) assert.Equal(t, wspkg.EventSubscribeReject, resp.Event) default: t.Fatal("expected subscribe reject message") } } func TestHubHandleSubscribe_UnknownChannel(t *testing.T) { hub := NewHubSimple() client := NewClient(1, 10, nil, hub) hub.Register(client) defer hub.Unregister(client) data, _ := json.Marshal(wspkg.SubscribeData{ Channel: "UnknownChannel", AccountID: 10, }) hub.processCommand(&ClientCommand{ Client: client, Cmd: wspkg.WSCommand{Command: "subscribe", Data: string(data)}, }) select { case msg := <-client.Send: var resp wspkg.WSMessage err := json.Unmarshal(msg, &resp) require.NoError(t, err) assert.Equal(t, wspkg.EventSubscribeReject, resp.Event) default: t.Fatal("expected subscribe reject message") } } func TestHubHandleSubscribe_InvalidData(t *testing.T) { hub := NewHubSimple() client := NewClient(1, 10, nil, hub) hub.Register(client) defer hub.Unregister(client) hub.processCommand(&ClientCommand{ Client: client, Cmd: wspkg.WSCommand{Command: "subscribe", Data: "invalid json"}, }) // Should not send any message, just return select { case <-client.Send: t.Fatal("unexpected message for invalid subscribe data") default: } } // === Hub.handleUnsubscribe tests === func TestHubHandleUnsubscribe_AccountChannel(t *testing.T) { hub := NewHubSimple() client := NewClient(1, 10, nil, hub) hub.Register(client) defer hub.Unregister(client) // First subscribe data, _ := json.Marshal(wspkg.SubscribeData{ Channel: wspkg.ChannelAccount, AccountID: 10, }) hub.processCommand(&ClientCommand{ Client: client, Cmd: wspkg.WSCommand{Command: "subscribe", Data: string(data)}, }) <-client.Send // drain confirm // Now unsubscribe hub.processCommand(&ClientCommand{ Client: client, Cmd: wspkg.WSCommand{Command: "unsubscribe", Data: string(data)}, }) select { case msg := <-client.Send: var resp wspkg.WSMessage err := json.Unmarshal(msg, &resp) require.NoError(t, err) assert.Equal(t, wspkg.EventUnsubscribeConfirm, resp.Event) default: t.Fatal("expected unsubscribe confirmation message") } } func TestHubHandleUnsubscribe_ConversationChannel(t *testing.T) { hub := NewHubSimple() client := NewClient(1, 10, nil, hub) hub.Register(client) defer hub.Unregister(client) data, _ := json.Marshal(wspkg.SubscribeData{ Channel: wspkg.ChannelConversation, AccountID: 10, ConversationID: 42, }) hub.processCommand(&ClientCommand{ Client: client, Cmd: wspkg.WSCommand{Command: "subscribe", Data: string(data)}, }) <-client.Send // drain confirm hub.processCommand(&ClientCommand{ Client: client, Cmd: wspkg.WSCommand{Command: "unsubscribe", Data: string(data)}, }) select { case msg := <-client.Send: var resp wspkg.WSMessage err := json.Unmarshal(msg, &resp) require.NoError(t, err) assert.Equal(t, wspkg.EventUnsubscribeConfirm, resp.Event) default: t.Fatal("expected unsubscribe confirmation message") } } func TestHubHandleUnsubscribe_UnknownChannel(t *testing.T) { hub := NewHubSimple() client := NewClient(1, 10, nil, hub) hub.Register(client) defer hub.Unregister(client) data, _ := json.Marshal(wspkg.SubscribeData{ Channel: "UnknownChannel", AccountID: 10, }) hub.processCommand(&ClientCommand{ Client: client, Cmd: wspkg.WSCommand{Command: "unsubscribe", Data: string(data)}, }) // Unknown channel — should return without sending select { case <-client.Send: t.Fatal("unexpected message for unknown channel unsubscribe") default: } } func TestHubHandleUnsubscribe_InvalidData(t *testing.T) { hub := NewHubSimple() client := NewClient(1, 10, nil, hub) hub.Register(client) defer hub.Unregister(client) hub.processCommand(&ClientCommand{ Client: client, Cmd: wspkg.WSCommand{Command: "unsubscribe", Data: "invalid json"}, }) select { case <-client.Send: t.Fatal("unexpected message for invalid unsubscribe data") default: } } // === Hub.handlePing tests === func TestHubHandlePing(t *testing.T) { hub := NewHubSimple() client := NewClient(1, 10, nil, hub) hub.Register(client) defer hub.Unregister(client) hub.processCommand(&ClientCommand{ Client: client, Cmd: wspkg.WSCommand{Command: "ping", Data: ""}, }) select { case msg := <-client.Send: var resp wspkg.WSMessage err := json.Unmarshal(msg, &resp) require.NoError(t, err) assert.Equal(t, wspkg.EventPingResponse, resp.Event) default: t.Fatal("expected ping response message") } } // === Hub.handleTypingOn tests === func TestHubHandleTypingOn_NilTyping(t *testing.T) { hub := NewHubSimple() // no typing tracker client := NewClient(1, 10, nil, hub) hub.Register(client) defer hub.Unregister(client) data, _ := json.Marshal(wspkg.TypingData{ ConversationID: 42, AccountID: 10, }) hub.processCommand(&ClientCommand{ Client: client, Cmd: wspkg.WSCommand{Command: "typing_on", Data: string(data)}, }) // Should return early without panic } func TestHubHandleTypingOn_WithTypingTracker(t *testing.T) { mr := miniredis.RunT(t) defer mr.Close() rdb := redis.NewClient(&redis.Options{Addr: mr.Addr()}) relay := wspkg.NewBroadcastRelay(rdb, nil) typing := wspkg.NewTypingTracker(rdb, relay) hub := NewHub(nil, typing, nil, nil) client := NewClient(1, 10, nil, hub) hub.Register(client) defer hub.Unregister(client) data, _ := json.Marshal(wspkg.TypingData{ ConversationID: 42, AccountID: 10, }) hub.processCommand(&ClientCommand{ Client: client, Cmd: wspkg.WSCommand{Command: "typing_on", Data: string(data)}, }) // Should not panic } func TestHubHandleTypingOn_InvalidData(t *testing.T) { mr := miniredis.RunT(t) defer mr.Close() rdb := redis.NewClient(&redis.Options{Addr: mr.Addr()}) relay := wspkg.NewBroadcastRelay(rdb, nil) typing := wspkg.NewTypingTracker(rdb, relay) hub := NewHub(nil, typing, nil, nil) client := NewClient(1, 10, nil, hub) hub.Register(client) defer hub.Unregister(client) hub.processCommand(&ClientCommand{ Client: client, Cmd: wspkg.WSCommand{Command: "typing_on", Data: "invalid json"}, }) // Should not panic } func TestHubHandleTypingOn_ContactClient(t *testing.T) { mr := miniredis.RunT(t) defer mr.Close() rdb := redis.NewClient(&redis.Options{Addr: mr.Addr()}) relay := wspkg.NewBroadcastRelay(rdb, nil) typing := wspkg.NewTypingTracker(rdb, relay) hub := NewHub(nil, typing, nil, nil) client := NewClient(0, 10, nil, hub) client.IsContact = true client.ContactID = 5 hub.Register(client) defer hub.Unregister(client) data, _ := json.Marshal(wspkg.TypingData{ ConversationID: 42, AccountID: 10, }) hub.processCommand(&ClientCommand{ Client: client, Cmd: wspkg.WSCommand{Command: "typing_on", Data: string(data)}, }) } // === Hub.handleTypingOff tests === func TestHubHandleTypingOff_NilTyping(t *testing.T) { hub := NewHubSimple() client := NewClient(1, 10, nil, hub) hub.Register(client) defer hub.Unregister(client) data, _ := json.Marshal(wspkg.TypingData{ ConversationID: 42, AccountID: 10, }) hub.processCommand(&ClientCommand{ Client: client, Cmd: wspkg.WSCommand{Command: "typing_off", Data: string(data)}, }) } func TestHubHandleTypingOff_WithTypingTracker(t *testing.T) { mr := miniredis.RunT(t) defer mr.Close() rdb := redis.NewClient(&redis.Options{Addr: mr.Addr()}) relay := wspkg.NewBroadcastRelay(rdb, nil) typing := wspkg.NewTypingTracker(rdb, relay) hub := NewHub(nil, typing, nil, nil) client := NewClient(1, 10, nil, hub) hub.Register(client) defer hub.Unregister(client) data, _ := json.Marshal(wspkg.TypingData{ ConversationID: 42, AccountID: 10, }) hub.processCommand(&ClientCommand{ Client: client, Cmd: wspkg.WSCommand{Command: "typing_off", Data: string(data)}, }) } func TestHubHandleTypingOff_InvalidData(t *testing.T) { mr := miniredis.RunT(t) defer mr.Close() rdb := redis.NewClient(&redis.Options{Addr: mr.Addr()}) relay := wspkg.NewBroadcastRelay(rdb, nil) typing := wspkg.NewTypingTracker(rdb, relay) hub := NewHub(nil, typing, nil, nil) client := NewClient(1, 10, nil, hub) hub.Register(client) defer hub.Unregister(client) hub.processCommand(&ClientCommand{ Client: client, Cmd: wspkg.WSCommand{Command: "typing_off", Data: "invalid json"}, }) } func TestHubHandleTypingOff_ContactClient(t *testing.T) { mr := miniredis.RunT(t) defer mr.Close() rdb := redis.NewClient(&redis.Options{Addr: mr.Addr()}) relay := wspkg.NewBroadcastRelay(rdb, nil) typing := wspkg.NewTypingTracker(rdb, relay) hub := NewHub(nil, typing, nil, nil) client := NewClient(0, 10, nil, hub) client.IsContact = true client.ContactID = 5 hub.Register(client) defer hub.Unregister(client) data, _ := json.Marshal(wspkg.TypingData{ ConversationID: 42, AccountID: 10, }) hub.processCommand(&ClientCommand{ Client: client, Cmd: wspkg.WSCommand{Command: "typing_off", Data: string(data)}, }) } // === Hub.handleUpdatePresence tests === func TestHubHandleUpdatePresence_NilPresence(t *testing.T) { hub := NewHubSimple() client := NewClient(1, 10, nil, hub) hub.Register(client) defer hub.Unregister(client) data, _ := json.Marshal(wspkg.PresenceData{Status: "online"}) hub.processCommand(&ClientCommand{ Client: client, Cmd: wspkg.WSCommand{Command: "update_presence", Data: string(data)}, }) } func TestHubHandleUpdatePresence_Online(t *testing.T) { mr := miniredis.RunT(t) defer mr.Close() rdb := redis.NewClient(&redis.Options{Addr: mr.Addr()}) relay := wspkg.NewBroadcastRelay(rdb, nil) presence := wspkg.NewPresenceTracker(rdb, relay) hub := NewHub(nil, nil, presence, nil) client := NewClient(1, 10, nil, hub) hub.Register(client) defer hub.Unregister(client) data, _ := json.Marshal(wspkg.PresenceData{Status: "online"}) hub.processCommand(&ClientCommand{ Client: client, Cmd: wspkg.WSCommand{Command: "update_presence", Data: string(data)}, }) } func TestHubHandleUpdatePresence_Busy(t *testing.T) { mr := miniredis.RunT(t) defer mr.Close() rdb := redis.NewClient(&redis.Options{Addr: mr.Addr()}) relay := wspkg.NewBroadcastRelay(rdb, nil) presence := wspkg.NewPresenceTracker(rdb, relay) hub := NewHub(nil, nil, presence, nil) client := NewClient(1, 10, nil, hub) hub.Register(client) defer hub.Unregister(client) data, _ := json.Marshal(wspkg.PresenceData{Status: "busy"}) hub.processCommand(&ClientCommand{ Client: client, Cmd: wspkg.WSCommand{Command: "update_presence", Data: string(data)}, }) } func TestHubHandleUpdatePresence_Offline(t *testing.T) { mr := miniredis.RunT(t) defer mr.Close() rdb := redis.NewClient(&redis.Options{Addr: mr.Addr()}) relay := wspkg.NewBroadcastRelay(rdb, nil) presence := wspkg.NewPresenceTracker(rdb, relay) hub := NewHub(nil, nil, presence, nil) client := NewClient(1, 10, nil, hub) hub.Register(client) defer hub.Unregister(client) data, _ := json.Marshal(wspkg.PresenceData{Status: "offline"}) hub.processCommand(&ClientCommand{ Client: client, Cmd: wspkg.WSCommand{Command: "update_presence", Data: string(data)}, }) } func TestHubHandleUpdatePresence_UnknownStatus(t *testing.T) { mr := miniredis.RunT(t) defer mr.Close() rdb := redis.NewClient(&redis.Options{Addr: mr.Addr()}) relay := wspkg.NewBroadcastRelay(rdb, nil) presence := wspkg.NewPresenceTracker(rdb, relay) hub := NewHub(nil, nil, presence, nil) client := NewClient(1, 10, nil, hub) hub.Register(client) defer hub.Unregister(client) data, _ := json.Marshal(wspkg.PresenceData{Status: "unknown"}) hub.processCommand(&ClientCommand{ Client: client, Cmd: wspkg.WSCommand{Command: "update_presence", Data: string(data)}, }) } func TestHubHandleUpdatePresence_InvalidData(t *testing.T) { mr := miniredis.RunT(t) defer mr.Close() rdb := redis.NewClient(&redis.Options{Addr: mr.Addr()}) relay := wspkg.NewBroadcastRelay(rdb, nil) presence := wspkg.NewPresenceTracker(rdb, relay) hub := NewHub(nil, nil, presence, nil) client := NewClient(1, 10, nil, hub) hub.Register(client) defer hub.Unregister(client) hub.processCommand(&ClientCommand{ Client: client, Cmd: wspkg.WSCommand{Command: "update_presence", Data: "invalid json"}, }) } // === Hub.Run with command processing === func TestHub_Run_ProcessesCommands(t *testing.T) { hub := NewHubSimple() ctx, cancel := context.WithCancel(context.Background()) go hub.Run(ctx) client := NewClient(1, 10, nil, hub) hub.Register(client) // Submit a ping command via SubmitCommand (goes through the event loop) hub.SubmitCommand(client, wspkg.WSCommand{Command: "ping", Data: ""}) // Wait for the ping response select { case msg := <-client.Send: var resp wspkg.WSMessage err := json.Unmarshal(msg, &resp) require.NoError(t, err) assert.Equal(t, wspkg.EventPingResponse, resp.Event) case <-time.After(2 * time.Second): t.Fatal("did not receive ping response from event loop") } hub.Unregister(client) cancel() } // === Hub.shutdown with clients (with CancelPresence) === func TestHub_Shutdown_WithCancelPresence(t *testing.T) { mr := miniredis.RunT(t) defer mr.Close() rdb := redis.NewClient(&redis.Options{Addr: mr.Addr()}) relay := wspkg.NewBroadcastRelay(rdb, nil) presence := wspkg.NewPresenceTracker(rdb, relay) pm := wspkg.NewPresenceManager(presence, wspkg.DefaultHeartbeatConfig()) hub := NewHub(relay, nil, presence, pm) // Register an agent client — this will start the presence refresh loop // Use a dummy websocket connection to avoid nil panic in shutdown() server, clientConn := newTestWSConn(t) defer server.Close() client := NewClient(1, 10, clientConn, hub) hub.Register(client) // CancelPresence is set by Register when presenceMgr != nil require.NotNil(t, client.CancelPresence) // Shutdown should call CancelPresence and close Send hub.Shutdown(context.Background()) assert.Empty(t, hub.clients) } // === SendToAccount with slow client (drop path) === func TestHub_SendToAccount_SlowClient(t *testing.T) { hub := NewHubSimple() client := NewClient(1, 10, nil, hub) client.Identifier = `{"channel":"AccountChannel","account_id":10}` hub.Register(client) // Fill the Send channel buffer for i := 0; i < SendChannelSize; i++ { client.Send <- []byte("fill") } // Now send — should be dropped, not block hub.SendToAccount(10, []byte(`{"event":"test"}`)) // If we reach here, it didn't block } func TestHub_SendToAccountConversation_SlowClient(t *testing.T) { hub := NewHubSimple() client := NewClient(1, 10, nil, hub) client.Identifier = `{"channel":"ConversationChannel","account_id":10,"conversation_id":5}` hub.Register(client) room := conversationRoomName(10, 5) hub.subscribeClient(client.ID, room) client.SubscribedRooms[room] = true // Fill the Send channel for i := 0; i < SendChannelSize; i++ { client.Send <- []byte("fill") } hub.SendToAccountConversation(10, 5, []byte(`{"event":"test"}`)) } func TestHub_SendToRoom_SlowClient(t *testing.T) { hub := NewHubSimple() client := NewClient(1, 10, nil, hub) client.Identifier = `{"channel":"RoomChannel"}` hub.Register(client) hub.subscribeClient(client.ID, "custom_room") client.SubscribedRooms["custom_room"] = true for i := 0; i < SendChannelSize; i++ { client.Send <- []byte("fill") } hub.SendToRoom("custom_room", []byte(`{"event":"test"}`)) } func TestHub_SendToClient_SlowClient(t *testing.T) { hub := NewHubSimple() client := NewClient(1, 10, nil, hub) client.Identifier = `{"channel":"AccountChannel"}` hub.Register(client) for i := 0; i < SendChannelSize; i++ { client.Send <- []byte("fill") } hub.SendToClient(client.ID, []byte(`{"event":"test"}`)) } // === SendToClient with no identifier === func TestHub_SendToClient_NoIdentifier(t *testing.T) { hub := NewHubSimple() client := NewClient(1, 10, nil, hub) // No Identifier set hub.Register(client) hub.SendToClient(client.ID, []byte(`{"event":"test"}`)) // Should be skipped (wrapActionCableMessage returns nil) select { case <-client.Send: // Could have the welcome or other messages, but the SendToClient should not add default: } } // === Subscriber tests === func TestNewSubscriber(t *testing.T) { mr := miniredis.RunT(t) defer mr.Close() rdb := redis.NewClient(&redis.Options{Addr: mr.Addr()}) hub := NewHubSimple() sub, err := NewSubscriber(hub, rdb) require.NoError(t, err) require.NotNil(t, sub) assert.NotNil(t, sub.hub) assert.NotNil(t, sub.subscriber) assert.NotNil(t, sub.router) assert.NotNil(t, sub.redisClient) } func TestNewSubscriber_InvalidRedis(t *testing.T) { // Use a non-existent redis address — the subscriber creation may still succeed // because watermill is lazy, but let's test the path rdb := redis.NewClient(&redis.Options{Addr: "127.0.0.1:1"}) // port 1 won't connect hub := NewHubSimple() // NewSubscriber should still succeed — it doesn't connect until Run sub, err := NewSubscriber(hub, rdb) if err != nil { // If it fails, that's also acceptable — we're testing the error path assert.Nil(t, sub) return } require.NotNil(t, sub) } func TestSubscriber_Running(t *testing.T) { mr := miniredis.RunT(t) defer mr.Close() rdb := redis.NewClient(&redis.Options{Addr: mr.Addr()}) hub := NewHubSimple() sub, err := NewSubscriber(hub, rdb) require.NoError(t, err) // Not running yet assert.False(t, sub.Running()) } func TestSubscriber_Close(t *testing.T) { mr := miniredis.RunT(t) defer mr.Close() rdb := redis.NewClient(&redis.Options{Addr: mr.Addr()}) hub := NewHubSimple() sub, err := NewSubscriber(hub, rdb) require.NoError(t, err) // Run the router in a goroutine so Close can work properly ctx, cancel := context.WithCancel(context.Background()) defer cancel() go func() { _ = sub.Run(ctx) }() // Give it a moment to start time.Sleep(100 * time.Millisecond) // Close should work done := make(chan error, 1) go func() { done <- sub.Close() }() select { case err := <-done: _ = err case <-time.After(5 * time.Second): t.Fatal("subscriber.Close() timed out") } } func TestSubscriber_Close_Twice(t *testing.T) { mr := miniredis.RunT(t) defer mr.Close() rdb := redis.NewClient(&redis.Options{Addr: mr.Addr()}) hub := NewHubSimple() sub, err := NewSubscriber(hub, rdb) require.NoError(t, err) // Run the router so Close works ctx, cancel := context.WithCancel(context.Background()) defer cancel() go func() { _ = sub.Run(ctx) }() time.Sleep(100 * time.Millisecond) // First close done := make(chan error, 1) go func() { done <- sub.Close() }() select { case <-done: case <-time.After(5 * time.Second): t.Fatal("first Close timed out") } // Second close should handle "already closed" gracefully done2 := make(chan error, 1) go func() { done2 <- sub.Close() }() select { case err := <-done2: assert.NoError(t, err) case <-time.After(5 * time.Second): t.Fatal("second Close timed out") } } // === extractPayload tests === func TestExtractPayload_Valid(t *testing.T) { payloadJSON := `{"account_id":10,"conversation_id":42,"data":{"key":"value"}}` msg := message.NewMessage("test-id", []byte(payloadJSON)) payload, err := extractPayload(msg) require.NoError(t, err) require.NotNil(t, payload) assert.Equal(t, uint(10), payload.AccountID) assert.Equal(t, uint(42), payload.ConversationID) } func TestExtractPayload_InvalidJSON(t *testing.T) { msg := message.NewMessage("test-id", []byte("invalid json")) payload, err := extractPayload(msg) require.Error(t, err) assert.Nil(t, payload) } func TestExtractPayload_Empty(t *testing.T) { msg := message.NewMessage("test-id", []byte("{}")) payload, err := extractPayload(msg) require.NoError(t, err) assert.Equal(t, uint(0), payload.AccountID) } // === forwardToAccountAndConversation tests === func TestForwardToAccountAndConversation_Valid(t *testing.T) { hub := NewHubSimple() client := NewClient(1, 10, nil, hub) client.Identifier = `{"channel":"AccountChannel","account_id":10}` hub.Register(client) // Also subscribe to conversation room convRoom := conversationRoomName(10, 42) hub.subscribeClient(client.ID, convRoom) client.SubscribedRooms[convRoom] = true sub := &Subscriber{hub: hub} payloadJSON := `{"account_id":10,"conversation_id":42,"data":{"message":"hello"}}` msg := message.NewMessage("test-id", []byte(payloadJSON)) handler := sub.forwardToAccountAndConversation(EventMessageCreated) err := handler(msg) require.NoError(t, err) // Should receive message on account room <-client.Send // may have welcome or prior messages, drain // Try to read the forwarded message // Since the client.Send channel is buffered, we should get messages } func TestForwardToAccountAndConversation_MissingAccountID(t *testing.T) { hub := NewHubSimple() sub := &Subscriber{hub: hub} payloadJSON := `{"account_id":0,"conversation_id":42,"data":{}}` msg := message.NewMessage("test-id", []byte(payloadJSON)) handler := sub.forwardToAccountAndConversation(EventMessageCreated) err := handler(msg) require.NoError(t, err) // returns nil even on missing account_id } func TestForwardToAccountAndConversation_InvalidPayload(t *testing.T) { hub := NewHubSimple() sub := &Subscriber{hub: hub} msg := message.NewMessage("test-id", []byte("invalid json")) handler := sub.forwardToAccountAndConversation(EventMessageCreated) err := handler(msg) require.NoError(t, err) // returns nil on bad payload } func TestForwardToAccountAndConversation_NoConversationID(t *testing.T) { hub := NewHubSimple() client := NewClient(1, 10, nil, hub) client.Identifier = `{"channel":"AccountChannel","account_id":10}` hub.Register(client) sub := &Subscriber{hub: hub} payloadJSON := `{"account_id":10,"data":{"key":"value"}}` msg := message.NewMessage("test-id", []byte(payloadJSON)) handler := sub.forwardToAccountAndConversation(EventMessageCreated) err := handler(msg) require.NoError(t, err) } // === forwardToAccount tests === func TestForwardToAccount_Valid(t *testing.T) { hub := NewHubSimple() client := NewClient(1, 10, nil, hub) client.Identifier = `{"channel":"AccountChannel","account_id":10}` hub.Register(client) sub := &Subscriber{hub: hub} payloadJSON := `{"account_id":10,"data":{"key":"value"}}` msg := message.NewMessage("test-id", []byte(payloadJSON)) handler := sub.forwardToAccount(EventConversationCreated) err := handler(msg) require.NoError(t, err) } func TestForwardToAccount_MissingAccountID(t *testing.T) { hub := NewHubSimple() sub := &Subscriber{hub: hub} payloadJSON := `{"account_id":0,"data":{}}` msg := message.NewMessage("test-id", []byte(payloadJSON)) handler := sub.forwardToAccount(EventConversationCreated) err := handler(msg) require.NoError(t, err) } func TestForwardToAccount_InvalidPayload(t *testing.T) { hub := NewHubSimple() sub := &Subscriber{hub: hub} msg := message.NewMessage("test-id", []byte("invalid json")) handler := sub.forwardToAccount(EventConversationCreated) err := handler(msg) require.NoError(t, err) } // === forwardToConversation tests === func TestForwardToConversation_Valid(t *testing.T) { hub := NewHubSimple() client := NewClient(1, 10, nil, hub) client.Identifier = `{"channel":"ConversationChannel","account_id":10,"conversation_id":42}` hub.Register(client) convRoom := conversationRoomName(10, 42) hub.subscribeClient(client.ID, convRoom) client.SubscribedRooms[convRoom] = true sub := &Subscriber{hub: hub} payloadJSON := `{"account_id":10,"conversation_id":42,"data":{"key":"value"}}` msg := message.NewMessage("test-id", []byte(payloadJSON)) handler := sub.forwardToConversation(EventAgentTypingOn) err := handler(msg) require.NoError(t, err) } func TestForwardToConversation_MissingAccountID(t *testing.T) { hub := NewHubSimple() sub := &Subscriber{hub: hub} payloadJSON := `{"account_id":0,"conversation_id":42,"data":{}}` msg := message.NewMessage("test-id", []byte(payloadJSON)) handler := sub.forwardToConversation(EventAgentTypingOn) err := handler(msg) require.NoError(t, err) } func TestForwardToConversation_MissingConversationID(t *testing.T) { hub := NewHubSimple() sub := &Subscriber{hub: hub} payloadJSON := `{"account_id":10,"conversation_id":0,"data":{}}` msg := message.NewMessage("test-id", []byte(payloadJSON)) handler := sub.forwardToConversation(EventAgentTypingOn) err := handler(msg) require.NoError(t, err) } func TestForwardToConversation_InvalidPayload(t *testing.T) { hub := NewHubSimple() sub := &Subscriber{hub: hub} msg := message.NewMessage("test-id", []byte("invalid json")) handler := sub.forwardToConversation(EventAgentTypingOn) err := handler(msg) require.NoError(t, err) } // === watermillZapAdapter tests === func TestWatermillZapAdapter_Error(t *testing.T) { adapter := &watermillZapAdapter{} adapter.Error("test error", assert.AnError, nil) adapter.Error("test no error", nil, nil) } func TestWatermillZapAdapter_Info(t *testing.T) { adapter := &watermillZapAdapter{} adapter.Info("test info", nil) } func TestWatermillZapAdapter_Debug(t *testing.T) { adapter := &watermillZapAdapter{} adapter.Debug("test debug", nil) } func TestWatermillZapAdapter_Trace(t *testing.T) { adapter := &watermillZapAdapter{} adapter.Trace("test trace", nil) } func TestWatermillZapAdapter_With(t *testing.T) { adapter := &watermillZapAdapter{} newAdapter := adapter.With(nil) assert.NotNil(t, newAdapter) } // === Handler.handleUnsubscribe with invalid identifier === // newTestWSConn creates a pair of connected websocket connections for testing. // Returns the server side and the client side. func newTestWSConn(t *testing.T) (*websocket.Conn, *websocket.Conn) { t.Helper() upgrader := websocket.Upgrader{ CheckOrigin: func(r *http.Request) bool { return true }, } serverConn := make(chan *websocket.Conn, 1) srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { conn, err := upgrader.Upgrade(w, r, nil) if err != nil { return } serverConn <- conn })) t.Cleanup(srv.Close) wsURL := "ws" + strings.TrimPrefix(srv.URL, "http") dialer := websocket.Dialer{HandshakeTimeout: 2 * time.Second} clientConn, _, err := dialer.Dial(wsURL, nil) require.NoError(t, err) select { case conn := <-serverConn: return conn, clientConn case <-time.After(2 * time.Second): clientConn.Close() t.Fatal("server connection was not established") return nil, nil } } func TestHandlerHandleUnsubscribe_InvalidIdentifier(t *testing.T) { hub := NewHubSimple() h := &Handler{hub: hub} client := NewClient(1, 10, nil, hub) hub.Register(client) defer hub.Unregister(client) // Call handleUnsubscribe with invalid identifier JSON h.handleUnsubscribe(client, CommandFrame{ Command: CommandUnsubscribe, Identifier: "invalid json", }) // Should return without sending confirmation }