Skip to content

Commit d9e5899

Browse files
committed
fix(web): scope inbound thread updates
Resolve inbound events to their authenticated message before clearing unread state so activity from another conversation cannot mark the open thread as read. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: c4b0e509-8bcb-4ace-82f4-3e64f6fcc0a6
1 parent 33fe141 commit d9e5899

6 files changed

Lines changed: 137 additions & 15 deletions

File tree

api/pkg/handlers/message_thread_handler_test.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -93,7 +93,7 @@ func TestMessageThreadHandlerUpdate_ReturnsNotFoundWhenThreadIsMissing(t *testin
9393
func TestMessageThreadHandlerUpdate_RejectsLegacyIsReadPayload(t *testing.T) {
9494
logger := &messageThreadHandlerNoopLogger{}
9595
tracer := telemetry.NewOtelLogger("test", logger)
96-
service := services.NewMessageThreadService(logger, tracer, &messageThreadHandlerRepositoryStub{}, nil, nil)
96+
service := services.NewMessageThreadService(logger, tracer, &messageThreadHandlerRepositoryStub{}, nil, nil, nil)
9797
handler := NewMessageThreadHandler(logger, tracer, validators.NewMessageThreadHandlerValidator(logger, tracer), service)
9898

9999
app := fiber.New()

api/pkg/listeners/websocket_listener.go

Lines changed: 10 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,10 @@ type WebsocketListener struct {
1919
client *pusher.Client
2020
}
2121

22+
type websocketMessagePayload struct {
23+
MessageID string `json:"message_id"`
24+
}
25+
2226
// NewWebsocketListener creates a new instance of WebsocketListener
2327
func NewWebsocketListener(
2428
logger telemetry.Logger,
@@ -49,7 +53,9 @@ func (listener *WebsocketListener) onMessageCallMissed(ctx context.Context, even
4953
return listener.tracer.WrapErrorSpan(span, stacktrace.Propagatef(err, "cannot decode [%s] into [%T]", event.Data(), payload))
5054
}
5155

52-
if err := listener.client.Trigger(payload.UserID.String(), event.Type(), event.ID()); err != nil {
56+
if err := listener.client.Trigger(payload.UserID.String(), event.Type(), websocketMessagePayload{
57+
MessageID: payload.MessageID.String(),
58+
}); err != nil {
5359
return listener.tracer.WrapErrorSpan(span, stacktrace.Propagatef(err, "cannot trigger websocket [%s] event with ID [%s] for user with ID [%s]", event.Type(), event.ID(), payload.UserID))
5460
}
5561
return nil
@@ -82,7 +88,9 @@ func (listener *WebsocketListener) onMessagePhoneReceived(ctx context.Context, e
8288
return listener.tracer.WrapErrorSpan(span, stacktrace.Propagatef(err, "cannot decode [%s] into [%T]", event.Data(), payload))
8389
}
8490

85-
if err := listener.client.Trigger(payload.UserID.String(), event.Type(), event.ID()); err != nil {
91+
if err := listener.client.Trigger(payload.UserID.String(), event.Type(), websocketMessagePayload{
92+
MessageID: payload.MessageID.String(),
93+
}); err != nil {
8694
return listener.tracer.WrapErrorSpan(span, stacktrace.Propagatef(err, "cannot trigger websocket [%s] event with ID [%s] for user with ID [%s]", event.Type(), event.ID(), payload.UserID))
8795
}
8896

api/pkg/listeners/websocket_listener_test.go

Lines changed: 81 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,12 +1,22 @@
11
package listeners
22

33
import (
4+
"context"
5+
"encoding/json"
6+
"io"
7+
"net/http"
8+
"net/http/httptest"
9+
"strings"
410
"testing"
511

12+
"github.com/NdoleStudio/httpsms/pkg/entities"
613
"github.com/NdoleStudio/httpsms/pkg/events"
714
"github.com/NdoleStudio/httpsms/pkg/telemetry"
15+
cloudevents "github.com/cloudevents/sdk-go/v2"
16+
"github.com/google/uuid"
817
"github.com/pusher/pusher-http-go/v5"
918
"github.com/stretchr/testify/assert"
19+
"github.com/stretchr/testify/require"
1020
)
1121

1222
func TestWebsocketListenerRegistersMissedCalls(t *testing.T) {
@@ -16,3 +26,74 @@ func TestWebsocketListenerRegistersMissedCalls(t *testing.T) {
1626

1727
assert.Contains(t, routes, events.MessageCallMissed)
1828
}
29+
30+
func TestWebsocketListenerPublishesReceivedMessageID(t *testing.T) {
31+
userID := "user-id"
32+
messageID := uuid.New()
33+
event := cloudevents.NewEvent()
34+
event.SetID(uuid.NewString())
35+
event.SetType(events.EventTypeMessagePhoneReceived)
36+
require.NoError(t, event.SetData(cloudevents.ApplicationJSON, events.MessagePhoneReceivedPayload{
37+
MessageID: messageID,
38+
UserID: entities.UserID(userID),
39+
}))
40+
41+
payload := captureWebsocketPayload(t, event, userID)
42+
43+
assert.Equal(t, messageID.String(), payload.MessageID)
44+
}
45+
46+
func TestWebsocketListenerPublishesMissedCallMessageID(t *testing.T) {
47+
userID := "user-id"
48+
messageID := uuid.New()
49+
event := cloudevents.NewEvent()
50+
event.SetID(uuid.NewString())
51+
event.SetType(events.MessageCallMissed)
52+
require.NoError(t, event.SetData(cloudevents.ApplicationJSON, events.MessageCallMissedPayload{
53+
MessageID: messageID,
54+
UserID: entities.UserID(userID),
55+
}))
56+
57+
payload := captureWebsocketPayload(t, event, userID)
58+
59+
assert.Equal(t, messageID.String(), payload.MessageID)
60+
}
61+
62+
func captureWebsocketPayload(t *testing.T, event cloudevents.Event, userID string) websocketMessagePayload {
63+
t.Helper()
64+
65+
requestBody := make(chan []byte, 1)
66+
server := httptest.NewServer(http.HandlerFunc(func(writer http.ResponseWriter, request *http.Request) {
67+
body, err := io.ReadAll(request.Body)
68+
require.NoError(t, err)
69+
requestBody <- body
70+
writer.Header().Set("Content-Type", "application/json")
71+
_, err = writer.Write([]byte(`{}`))
72+
require.NoError(t, err)
73+
}))
74+
t.Cleanup(server.Close)
75+
76+
logger := &noopListenerLogger{}
77+
tracer := telemetry.NewOtelLogger("test", logger)
78+
client := &pusher.Client{
79+
AppID: "app-id",
80+
Key: "key",
81+
Secret: "secret",
82+
Host: strings.TrimPrefix(server.URL, "http://"),
83+
HTTPClient: server.Client(),
84+
}
85+
_, routes := NewWebsocketListener(logger, tracer, client)
86+
87+
require.NoError(t, routes[event.Type()](context.Background(), event))
88+
89+
var trigger struct {
90+
Channels []string `json:"channels"`
91+
Data string `json:"data"`
92+
}
93+
require.NoError(t, json.Unmarshal(<-requestBody, &trigger))
94+
assert.Equal(t, []string{userID}, trigger.Channels)
95+
96+
var payload websocketMessagePayload
97+
require.NoError(t, json.Unmarshal([]byte(trigger.Data), &payload))
98+
return payload
99+
}

api/pkg/services/message_thread_service_test.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -125,7 +125,7 @@ func newMessageThreadServiceForTest(repository repositories.MessageThreadReposit
125125
func newMessageThreadServiceWithPhoneForTest(repository repositories.MessageThreadRepository, phoneRepository repositories.PhoneRepository) *MessageThreadService {
126126
logger := &noopLogger{}
127127
tracer := telemetry.NewOtelLogger("test", logger)
128-
return NewMessageThreadService(logger, tracer, repository, phoneRepository, nil)
128+
return NewMessageThreadService(logger, tracer, repository, phoneRepository, nil, nil)
129129
}
130130

131131
func TestUpdateThreadPassesUnreadWatermarkForInboundActivity(t *testing.T) {

web/app/pages/threads/[id]/index.vue

Lines changed: 36 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -62,6 +62,10 @@ const formMessageRules = [
6262
6363
let webhookChannel: Channel | null = null
6464
65+
interface WebsocketMessageEvent {
66+
message_id: string
67+
}
68+
6569
const contactIsPhoneNumber = computed(() => {
6670
const thread = currentThread.value
6771
if (!thread) return false
@@ -137,6 +141,30 @@ async function resetCurrentThreadUnreadCount(force = false) {
137141
}
138142
}
139143
144+
async function handleInboundMessage(event: WebsocketMessageEvent) {
145+
if (loadingMessages.value) return
146+
147+
try {
148+
const message = await messagesStore.getMessage(event.message_id)
149+
await threadsStore.loadThreads()
150+
151+
const thread = currentThread.value
152+
if (
153+
!thread ||
154+
message.owner !== thread.owner ||
155+
message.contact !== thread.contact ||
156+
loadingMessages.value
157+
) {
158+
return
159+
}
160+
161+
await resetCurrentThreadUnreadCount(true)
162+
loadMessages(false, false)
163+
} catch (error) {
164+
console.error(error)
165+
}
166+
}
167+
140168
function loadMessages(hide = true, resetUnreadCount = true) {
141169
loadingMessages.value = true
142170
const threadId = route.params.id as string
@@ -257,17 +285,14 @@ onMounted(async () => {
257285
webhookChannel.bind('message.send.failed', () => {
258286
if (!loadingMessages.value) loadMessages(false)
259287
})
260-
webhookChannel.bind('message.phone.received', () => {
261-
if (!loadingMessages.value) {
262-
void resetCurrentThreadUnreadCount(true)
263-
loadMessages(false, false)
264-
}
265-
})
266-
webhookChannel.bind('message.call.missed', () => {
267-
if (!loadingMessages.value) {
268-
void resetCurrentThreadUnreadCount(true)
269-
loadMessages(false, false)
270-
}
288+
webhookChannel.bind(
289+
'message.phone.received',
290+
(event: WebsocketMessageEvent) => {
291+
void handleInboundMessage(event)
292+
},
293+
)
294+
webhookChannel.bind('message.call.missed', (event: WebsocketMessageEvent) => {
295+
void handleInboundMessage(event)
271296
})
272297
})
273298

web/app/stores/messages.ts

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -48,6 +48,13 @@ export const useMessagesStore = defineStore('messages', () => {
4848
})
4949
}
5050

51+
async function getMessage(messageId: string): Promise<EntitiesMessage> {
52+
const response = await apiFetch<{ data: EntitiesMessage }>(
53+
`/v1/messages/${messageId}`,
54+
)
55+
return response.data
56+
}
57+
5158
async function searchMessages(
5259
payload: SearchMessagesRequest,
5360
): Promise<EntitiesMessage[]> {
@@ -88,6 +95,7 @@ export const useMessagesStore = defineStore('messages', () => {
8895
return {
8996
sendMessage,
9097
deleteMessage,
98+
getMessage,
9199
searchMessages,
92100
sendBulkMessages,
93101
fetchBulkMessageOrders,

0 commit comments

Comments
 (0)