diff --git a/api/docs/docs.go b/api/docs/docs.go index c7fcd1145..8cd408cb1 100644 --- a/api/docs/docs.go +++ b/api/docs/docs.go @@ -871,7 +871,7 @@ const docTemplate = `{ "200": { "description": "OK", "schema": { - "$ref": "#/definitions/responses.PhoneResponse" + "$ref": "#/definitions/responses.MessageThreadResponse" } }, "400": { @@ -886,6 +886,12 @@ const docTemplate = `{ "$ref": "#/definitions/responses.Unauthorized" } }, + "404": { + "description": "Not Found", + "schema": { + "$ref": "#/definitions/responses.NotFound" + } + }, "422": { "description": "Unprocessable Entity", "schema": { @@ -3691,6 +3697,7 @@ const docTemplate = `{ "created_at", "id", "is_archived", + "is_read", "last_message_content", "last_message_id", "order_timestamp", @@ -3720,6 +3727,10 @@ const docTemplate = `{ "type": "boolean", "example": false }, + "is_read": { + "type": "boolean", + "example": true + }, "last_message_content": { "type": "string", "example": "This is a sample message content" @@ -4388,13 +4399,14 @@ const docTemplate = `{ }, "requests.MessageThreadUpdate": { "type": "object", - "required": [ - "is_archived" - ], "properties": { "is_archived": { "type": "boolean", "example": true + }, + "is_read": { + "type": "boolean", + "example": true } } }, @@ -4904,6 +4916,27 @@ const docTemplate = `{ } } }, + "responses.MessageThreadResponse": { + "type": "object", + "required": [ + "data", + "message", + "status" + ], + "properties": { + "data": { + "$ref": "#/definitions/entities.MessageThread" + }, + "message": { + "type": "string", + "example": "Request handled successfully" + }, + "status": { + "type": "string", + "example": "success" + } + } + }, "responses.MessageThreadsResponse": { "type": "object", "required": [ diff --git a/api/docs/swagger.json b/api/docs/swagger.json index 25db6fc31..ac9c12c98 100644 --- a/api/docs/swagger.json +++ b/api/docs/swagger.json @@ -868,7 +868,7 @@ "200": { "description": "OK", "schema": { - "$ref": "#/definitions/responses.PhoneResponse" + "$ref": "#/definitions/responses.MessageThreadResponse" } }, "400": { @@ -883,6 +883,12 @@ "$ref": "#/definitions/responses.Unauthorized" } }, + "404": { + "description": "Not Found", + "schema": { + "$ref": "#/definitions/responses.NotFound" + } + }, "422": { "description": "Unprocessable Entity", "schema": { @@ -3688,6 +3694,7 @@ "created_at", "id", "is_archived", + "is_read", "last_message_content", "last_message_id", "order_timestamp", @@ -3717,6 +3724,10 @@ "type": "boolean", "example": false }, + "is_read": { + "type": "boolean", + "example": true + }, "last_message_content": { "type": "string", "example": "This is a sample message content" @@ -4385,13 +4396,14 @@ }, "requests.MessageThreadUpdate": { "type": "object", - "required": [ - "is_archived" - ], "properties": { "is_archived": { "type": "boolean", "example": true + }, + "is_read": { + "type": "boolean", + "example": true } } }, @@ -4901,6 +4913,27 @@ } } }, + "responses.MessageThreadResponse": { + "type": "object", + "required": [ + "data", + "message", + "status" + ], + "properties": { + "data": { + "$ref": "#/definitions/entities.MessageThread" + }, + "message": { + "type": "string", + "example": "Request handled successfully" + }, + "status": { + "type": "string", + "example": "success" + } + } + }, "responses.MessageThreadsResponse": { "type": "object", "required": [ diff --git a/api/docs/swagger.yaml b/api/docs/swagger.yaml index 892c9ae45..ddc6ae700 100644 --- a/api/docs/swagger.yaml +++ b/api/docs/swagger.yaml @@ -319,6 +319,9 @@ definitions: is_archived: example: false type: boolean + is_read: + example: true + type: boolean last_message_content: example: This is a sample message content type: string @@ -346,6 +349,7 @@ definitions: - created_at - id - is_archived + - is_read - last_message_content - last_message_id - order_timestamp @@ -853,8 +857,9 @@ definitions: is_archived: example: true type: boolean - required: - - is_archived + is_read: + example: true + type: boolean type: object requests.PhoneAPIKeyStoreRequest: properties: @@ -1227,6 +1232,21 @@ definitions: - message - status type: object + responses.MessageThreadResponse: + properties: + data: + $ref: '#/definitions/entities.MessageThread' + message: + example: Request handled successfully + type: string + status: + example: success + type: string + required: + - data + - message + - status + type: object responses.MessageThreadsResponse: properties: data: @@ -2182,7 +2202,7 @@ paths: "200": description: OK schema: - $ref: '#/definitions/responses.PhoneResponse' + $ref: '#/definitions/responses.MessageThreadResponse' "400": description: Bad Request schema: @@ -2191,6 +2211,10 @@ paths: description: Unauthorized schema: $ref: '#/definitions/responses.Unauthorized' + "404": + description: Not Found + schema: + $ref: '#/definitions/responses.NotFound' "422": description: Unprocessable Entity schema: diff --git a/api/pkg/entities/message_thread.go b/api/pkg/entities/message_thread.go index 2eb1e2d33..3766acc77 100644 --- a/api/pkg/entities/message_thread.go +++ b/api/pkg/entities/message_thread.go @@ -12,6 +12,8 @@ type MessageThread struct { Owner string `json:"owner" example:"+18005550199"` Contact string `json:"contact" example:"+18005550100"` IsArchived bool `json:"is_archived" example:"false"` + IsRead bool `json:"is_read" gorm:"not null;default:true" example:"true"` + LastReadAt time.Time `json:"-" gorm:"not null;default:CURRENT_TIMESTAMP"` UserID UserID `json:"user_id" example:"WB7DRDWrJZRGbYrv2CKGkqbzvqdC"` Color string `json:"color" example:"indigo"` Status MessageStatus `json:"status" example:"PENDING"` diff --git a/api/pkg/entities/message_thread_test.go b/api/pkg/entities/message_thread_test.go new file mode 100644 index 000000000..1409dff29 --- /dev/null +++ b/api/pkg/entities/message_thread_test.go @@ -0,0 +1,25 @@ +package entities + +import ( + "reflect" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func TestMessageThreadReadFieldsHaveBackwardCompatibleDefaults(t *testing.T) { + threadType := reflect.TypeOf(MessageThread{}) + + isRead, ok := threadType.FieldByName("IsRead") + require.True(t, ok) + assert.Contains(t, isRead.Tag.Get("gorm"), "not null") + assert.Contains(t, isRead.Tag.Get("gorm"), "default:true") + assert.Equal(t, "is_read", isRead.Tag.Get("json")) + + lastReadAt, ok := threadType.FieldByName("LastReadAt") + require.True(t, ok) + assert.Contains(t, lastReadAt.Tag.Get("gorm"), "not null") + assert.Contains(t, lastReadAt.Tag.Get("gorm"), "default:CURRENT_TIMESTAMP") + assert.Equal(t, "-", lastReadAt.Tag.Get("json")) +} diff --git a/api/pkg/handlers/message_thread_handler.go b/api/pkg/handlers/message_thread_handler.go index 1a00f7737..9520b6634 100644 --- a/api/pkg/handlers/message_thread_handler.go +++ b/api/pkg/handlers/message_thread_handler.go @@ -100,9 +100,10 @@ func (h *MessageThreadHandler) Index(c fiber.Ctx) error { // @Produce json // @Param messageThreadID path string true "ID of the message thread" default(32343a19-da5e-4b1b-a767-3298a73703ca) // @Param payload body requests.MessageThreadUpdate true "Payload of message thread details to update" -// @Success 200 {object} responses.PhoneResponse +// @Success 200 {object} responses.MessageThreadResponse // @Failure 400 {object} responses.BadRequest // @Failure 401 {object} responses.Unauthorized +// @Failure 404 {object} responses.NotFound // @Failure 422 {object} responses.UnprocessableEntity // @Failure 500 {object} responses.InternalServerError // @Router /message-threads/{messageThreadID} [put] @@ -123,6 +124,9 @@ func (h *MessageThreadHandler) Update(c fiber.Ctx) error { } thread, err := h.service.UpdateStatus(ctx, request.ToUpdateParams(h.userIDFomContext(c))) + if stacktrace.GetCode(err) == repositories.ErrCodeNotFound { + return h.responseNotFound(c, fmt.Sprintf("cannot find message thread with ID [%s]", request.MessageThreadID)) + } if err != nil { ctxLogger.Error(stacktrace.Propagate(err, "cannot update message thread with params [%+#v]", request)) return h.responseInternalServerError(c) diff --git a/api/pkg/handlers/message_thread_handler_test.go b/api/pkg/handlers/message_thread_handler_test.go new file mode 100644 index 000000000..ab2d75962 --- /dev/null +++ b/api/pkg/handlers/message_thread_handler_test.go @@ -0,0 +1,110 @@ +package handlers + +import ( + "bytes" + "context" + "encoding/json" + "net/http" + "net/http/httptest" + "testing" + "time" + + "github.com/NdoleStudio/httpsms/pkg/entities" + "github.com/NdoleStudio/httpsms/pkg/middlewares" + "github.com/NdoleStudio/httpsms/pkg/repositories" + "github.com/NdoleStudio/httpsms/pkg/services" + "github.com/NdoleStudio/httpsms/pkg/telemetry" + "github.com/NdoleStudio/httpsms/pkg/validators" + "github.com/gofiber/fiber/v3" + "github.com/google/uuid" + "github.com/palantir/stacktrace" + "github.com/stretchr/testify/require" + "go.opentelemetry.io/otel/trace" + "gorm.io/gorm" +) + +type messageThreadHandlerRepositoryStub struct{} + +func (stub *messageThreadHandlerRepositoryStub) Store(context.Context, *entities.MessageThread) error { + return nil +} + +func (stub *messageThreadHandlerRepositoryStub) UpdateActivity(context.Context, repositories.MessageThreadActivityUpdate) error { + return nil +} + +func (stub *messageThreadHandlerRepositoryStub) UpdateStatus(context.Context, entities.UserID, uuid.UUID, repositories.MessageThreadStatusUpdate) (*entities.MessageThread, error) { + return nil, stacktrace.PropagateWithCode(gorm.ErrRecordNotFound, repositories.ErrCodeNotFound, "not found") +} + +func (stub *messageThreadHandlerRepositoryStub) LoadByOwnerContact(context.Context, entities.UserID, string, string) (*entities.MessageThread, error) { + return nil, nil +} + +func (stub *messageThreadHandlerRepositoryStub) Load(context.Context, entities.UserID, uuid.UUID) (*entities.MessageThread, error) { + return nil, nil +} + +func (stub *messageThreadHandlerRepositoryStub) Index(context.Context, entities.UserID, string, bool, repositories.IndexParams) (*[]entities.MessageThread, error) { + return nil, nil +} + +func (stub *messageThreadHandlerRepositoryStub) UpdateAfterDeletedMessage(context.Context, repositories.MessageThreadDeletedUpdate) error { + return nil +} + +func (stub *messageThreadHandlerRepositoryStub) Delete(context.Context, entities.UserID, uuid.UUID) error { + return nil +} + +func (stub *messageThreadHandlerRepositoryStub) DeleteAllForUser(context.Context, entities.UserID) error { + return nil +} + +func TestMessageThreadHandlerUpdate_ReturnsNotFoundWhenThreadIsMissing(t *testing.T) { + logger := &messageThreadHandlerNoopLogger{} + tracer := telemetry.NewOtelLogger("test", logger) + service := services.NewMessageThreadService(logger, tracer, &messageThreadHandlerRepositoryStub{}, nil, nil) + handler := NewMessageThreadHandler(logger, tracer, validators.NewMessageThreadHandlerValidator(logger, tracer), service) + + app := fiber.New() + app.Use(func(c fiber.Ctx) error { + c.Locals(middlewares.ContextKeyAuthUserID, entities.AuthContext{ID: entities.UserID("user-id"), Email: "user@example.com"}) + return c.Next() + }) + handler.RegisterRoutes(app) + + messageThreadID := uuid.New() + req := httptest.NewRequest(http.MethodPut, "/v1/message-threads/"+messageThreadID.String(), bytes.NewBufferString(`{"is_read":true}`)) + req.Header.Set("Content-Type", "application/json") + + resp, err := app.Test(req, fiber.TestConfig{Timeout: time.Second}) + + require.NoError(t, err) + require.Equal(t, http.StatusNotFound, resp.StatusCode) + + var payload struct { + Message string `json:"message"` + } + require.NoError(t, json.NewDecoder(resp.Body).Decode(&payload)) + require.Equal(t, "cannot find message thread with ID ["+messageThreadID.String()+"]", payload.Message) +} + +type messageThreadHandlerNoopLogger struct{} + +var _ telemetry.Logger = (*messageThreadHandlerNoopLogger)(nil) + +func (logger *messageThreadHandlerNoopLogger) Error(_ error) {} +func (logger *messageThreadHandlerNoopLogger) WithService(_ string) telemetry.Logger { return logger } + +func (logger *messageThreadHandlerNoopLogger) WithString(_, _ string) telemetry.Logger { return logger } + +func (logger *messageThreadHandlerNoopLogger) WithSpan(_ trace.SpanContext) telemetry.Logger { + return logger +} +func (logger *messageThreadHandlerNoopLogger) Trace(_ string) {} +func (logger *messageThreadHandlerNoopLogger) Info(_ string) {} +func (logger *messageThreadHandlerNoopLogger) Warn(_ error) {} +func (logger *messageThreadHandlerNoopLogger) Debug(_ string) {} +func (logger *messageThreadHandlerNoopLogger) Fatal(_ error) {} +func (logger *messageThreadHandlerNoopLogger) Printf(_ string, _ ...interface{}) {} diff --git a/api/pkg/listeners/message_thread_listener.go b/api/pkg/listeners/message_thread_listener.go index e5f90ca39..47621296e 100644 --- a/api/pkg/listeners/message_thread_listener.go +++ b/api/pkg/listeners/message_thread_listener.go @@ -40,6 +40,7 @@ func NewMessageThreadListener( events.EventTypeMessagePhoneDelivered: l.OnMessagePhoneDelivered, events.EventTypeMessageSendFailed: l.OnMessagePhoneFailed, events.EventTypeMessagePhoneReceived: l.OnMessagePhoneReceived, + events.MessageCallMissed: l.OnMessageCallMissed, events.EventTypeMessageNotificationScheduled: l.onMessageNotificationScheduled, events.EventTypeMessageSendExpired: l.onMessageExpired, events.UserAccountDeleted: l.onUserAccountDeleted, @@ -209,13 +210,15 @@ func (listener *MessageThreadListener) OnMessagePhoneReceived(ctx context.Contex } updateParams := services.MessageThreadUpdateParams{ - Owner: payload.Owner, - Contact: payload.Contact, - Timestamp: payload.Timestamp, - UserID: payload.UserID, - Status: entities.MessageStatusReceived, - Content: payload.Content, - MessageID: payload.MessageID, + Owner: payload.Owner, + Contact: payload.Contact, + Timestamp: payload.Timestamp, + UserID: payload.UserID, + Status: entities.MessageStatusReceived, + Content: payload.Content, + MessageID: payload.MessageID, + MarkAsUnread: true, + EventTimestamp: event.Time(), } if err := listener.service.UpdateThread(ctx, updateParams); err != nil { @@ -225,6 +228,33 @@ func (listener *MessageThreadListener) OnMessagePhoneReceived(ctx context.Contex return nil } +// OnMessageCallMissed handles the events.MessageCallMissed event +func (listener *MessageThreadListener) OnMessageCallMissed(ctx context.Context, event cloudevents.Event) error { + ctx, span := listener.tracer.Start(ctx) + defer span.End() + + var payload events.MessageCallMissedPayload + if err := event.DataAs(&payload); err != nil { + return listener.tracer.WrapErrorSpan(span, stacktrace.Propagate(err, "cannot decode [%s] into [%T]", event.Data(), payload)) + } + + params := services.MessageThreadUpdateParams{ + Owner: payload.Owner, + Contact: payload.Contact, + UserID: payload.UserID, + Status: entities.MessageStatusReceived, + Timestamp: payload.Timestamp, + Content: "Missed phone call", + MessageID: payload.MessageID, + MarkAsUnread: true, + EventTimestamp: event.Time(), + } + if err := listener.service.UpdateThread(ctx, params); err != nil { + return listener.tracer.WrapErrorSpan(span, stacktrace.Propagate(err, "cannot update thread for missed call [%s] on event [%s]", payload.MessageID, event.ID())) + } + return nil +} + // onMessageNotificationScheduled handles the events.EventTypeMessageNotificationScheduled event func (listener *MessageThreadListener) onMessageNotificationScheduled(ctx context.Context, event cloudevents.Event) error { ctx, span := listener.tracer.Start(ctx) diff --git a/api/pkg/listeners/message_thread_listener_test.go b/api/pkg/listeners/message_thread_listener_test.go new file mode 100644 index 000000000..39a444cce --- /dev/null +++ b/api/pkg/listeners/message_thread_listener_test.go @@ -0,0 +1,72 @@ +package listeners + +import ( + "context" + "testing" + "time" + + "github.com/NdoleStudio/httpsms/pkg/entities" + "github.com/NdoleStudio/httpsms/pkg/events" + "github.com/NdoleStudio/httpsms/pkg/services" + "github.com/NdoleStudio/httpsms/pkg/telemetry" + cloudevents "github.com/cloudevents/sdk-go/v2" + "github.com/google/uuid" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func TestMessageThreadListenerMarksInboundMessageUnread(t *testing.T) { + repository, routes := newMessageThreadListenerForTest() + event := cloudevents.NewEvent() + event.SetID(uuid.NewString()) + event.SetSource("/v1/messages/phone-received") + event.SetType(events.EventTypeMessagePhoneReceived) + event.SetTime(time.Date(2026, 7, 18, 7, 0, 0, 0, time.UTC)) + require.NoError(t, event.SetData(cloudevents.ApplicationJSON, events.MessagePhoneReceivedPayload{ + MessageID: uuid.New(), + UserID: entities.UserID("user-id"), + Owner: "+18005550199", + Contact: "+18005550100", + Content: "hello", + Timestamp: time.Date(2026, 7, 18, 6, 59, 0, 0, time.UTC), + })) + + err := routes[events.EventTypeMessagePhoneReceived](context.Background(), event) + + require.NoError(t, err) + assert.True(t, repository.activity.MarkAsUnread) + assert.Equal(t, event.Time(), repository.activity.EventTimestamp) +} + +func TestMessageThreadListenerMarksMissedCallUnread(t *testing.T) { + repository, routes := newMessageThreadListenerForTest() + event := cloudevents.NewEvent() + event.SetID(uuid.NewString()) + event.SetSource("/v1/messages/call-missed") + event.SetType(events.MessageCallMissed) + event.SetTime(time.Date(2026, 7, 18, 7, 0, 0, 0, time.UTC)) + require.NoError(t, event.SetData(cloudevents.ApplicationJSON, events.MessageCallMissedPayload{ + MessageID: uuid.New(), + UserID: entities.UserID("user-id"), + Owner: "+18005550199", + Contact: "+18005550100", + Timestamp: time.Date(2026, 7, 18, 6, 59, 0, 0, time.UTC), + })) + require.Contains(t, routes, events.MessageCallMissed) + + err := routes[events.MessageCallMissed](context.Background(), event) + + require.NoError(t, err) + assert.True(t, repository.activity.MarkAsUnread) + assert.Equal(t, "Missed phone call", repository.activity.Content) + assert.Equal(t, event.Time(), repository.activity.EventTimestamp) +} + +func newMessageThreadListenerForTest() (*listenerMessageThreadRepository, map[string]events.EventListener) { + repository := &listenerMessageThreadRepository{} + logger := &noopListenerLogger{} + tracer := telemetry.NewOtelLogger("test", logger) + service := services.NewMessageThreadService(logger, tracer, repository, nil, nil) + _, routes := NewMessageThreadListener(logger, tracer, service) + return repository, routes +} diff --git a/api/pkg/listeners/read_receipts_test_helpers_test.go b/api/pkg/listeners/read_receipts_test_helpers_test.go new file mode 100644 index 000000000..60cd27e65 --- /dev/null +++ b/api/pkg/listeners/read_receipts_test_helpers_test.go @@ -0,0 +1,66 @@ +package listeners + +import ( + "context" + + "github.com/NdoleStudio/httpsms/pkg/entities" + "github.com/NdoleStudio/httpsms/pkg/repositories" + "github.com/NdoleStudio/httpsms/pkg/telemetry" + "github.com/google/uuid" + "go.opentelemetry.io/otel/trace" +) + +type noopListenerLogger struct{} + +func (logger *noopListenerLogger) Error(error) {} +func (logger *noopListenerLogger) WithService(string) telemetry.Logger { return logger } +func (logger *noopListenerLogger) WithString(string, string) telemetry.Logger { return logger } +func (logger *noopListenerLogger) WithSpan(trace.SpanContext) telemetry.Logger { return logger } +func (logger *noopListenerLogger) Trace(string) {} +func (logger *noopListenerLogger) Info(string) {} +func (logger *noopListenerLogger) Warn(error) {} +func (logger *noopListenerLogger) Debug(string) {} +func (logger *noopListenerLogger) Fatal(error) {} +func (logger *noopListenerLogger) Printf(string, ...interface{}) {} + +type listenerMessageThreadRepository struct { + activity repositories.MessageThreadActivityUpdate +} + +func (repository *listenerMessageThreadRepository) Store(context.Context, *entities.MessageThread) error { + return nil +} + +func (repository *listenerMessageThreadRepository) UpdateActivity(_ context.Context, params repositories.MessageThreadActivityUpdate) error { + repository.activity = params + return nil +} + +func (repository *listenerMessageThreadRepository) UpdateStatus(_ context.Context, _ entities.UserID, threadID uuid.UUID, _ repositories.MessageThreadStatusUpdate) (*entities.MessageThread, error) { + return &entities.MessageThread{ID: threadID}, nil +} + +func (repository *listenerMessageThreadRepository) UpdateAfterDeletedMessage(context.Context, repositories.MessageThreadDeletedUpdate) error { + return nil +} + +func (repository *listenerMessageThreadRepository) LoadByOwnerContact(context.Context, entities.UserID, string, string) (*entities.MessageThread, error) { + return &entities.MessageThread{ID: uuid.New()}, nil +} + +func (repository *listenerMessageThreadRepository) Load(context.Context, entities.UserID, uuid.UUID) (*entities.MessageThread, error) { + return &entities.MessageThread{}, nil +} + +func (repository *listenerMessageThreadRepository) Index(context.Context, entities.UserID, string, bool, repositories.IndexParams) (*[]entities.MessageThread, error) { + threads := []entities.MessageThread{} + return &threads, nil +} + +func (repository *listenerMessageThreadRepository) Delete(context.Context, entities.UserID, uuid.UUID) error { + return nil +} + +func (repository *listenerMessageThreadRepository) DeleteAllForUser(context.Context, entities.UserID) error { + return nil +} diff --git a/api/pkg/listeners/websocket_listener.go b/api/pkg/listeners/websocket_listener.go index 3108ed671..8aa9054f0 100644 --- a/api/pkg/listeners/websocket_listener.go +++ b/api/pkg/listeners/websocket_listener.go @@ -36,9 +36,25 @@ func NewWebsocketListener( events.EventTypeMessagePhoneSent: l.onMessagePhoneSent, events.EventTypeMessageSendFailed: l.onMessagePhoneFailed, events.EventTypeMessagePhoneReceived: l.onMessagePhoneReceived, + events.MessageCallMissed: l.onMessageCallMissed, } } +func (listener *WebsocketListener) onMessageCallMissed(ctx context.Context, event cloudevents.Event) error { + ctx, span, _ := listener.tracer.StartWithLogger(ctx, listener.logger) + defer span.End() + + var payload events.MessageCallMissedPayload + if err := event.DataAs(&payload); err != nil { + return listener.tracer.WrapErrorSpan(span, stacktrace.Propagate(err, "cannot decode [%s] into [%T]", event.Data(), payload)) + } + + if err := listener.client.Trigger(payload.UserID.String(), event.Type(), event.ID()); err != nil { + return listener.tracer.WrapErrorSpan(span, stacktrace.Propagate(err, "cannot trigger websocket [%s] event with ID [%s] for user with ID [%s]", event.Type(), event.ID(), payload.UserID)) + } + return nil +} + // onMessagePhoneSent handles the events.EventTypeMessagePhoneSent event func (listener *WebsocketListener) onMessagePhoneSent(ctx context.Context, event cloudevents.Event) error { ctx, span, _ := listener.tracer.StartWithLogger(ctx, listener.logger) diff --git a/api/pkg/listeners/websocket_listener_test.go b/api/pkg/listeners/websocket_listener_test.go new file mode 100644 index 000000000..fdaf1c6d5 --- /dev/null +++ b/api/pkg/listeners/websocket_listener_test.go @@ -0,0 +1,18 @@ +package listeners + +import ( + "testing" + + "github.com/NdoleStudio/httpsms/pkg/events" + "github.com/NdoleStudio/httpsms/pkg/telemetry" + "github.com/pusher/pusher-http-go/v5" + "github.com/stretchr/testify/assert" +) + +func TestWebsocketListenerRegistersMissedCalls(t *testing.T) { + logger := &noopListenerLogger{} + tracer := telemetry.NewOtelLogger("test", logger) + _, routes := NewWebsocketListener(logger, tracer, &pusher.Client{}) + + assert.Contains(t, routes, events.MessageCallMissed) +} diff --git a/api/pkg/repositories/gorm_message_thread_repository.go b/api/pkg/repositories/gorm_message_thread_repository.go index 7bc07f7b3..071f5cc79 100644 --- a/api/pkg/repositories/gorm_message_thread_repository.go +++ b/api/pkg/repositories/gorm_message_thread_repository.go @@ -35,6 +35,48 @@ func NewGormMessageThreadRepository( } } +func messageThreadActivityUpdates(params MessageThreadActivityUpdate) map[string]any { + updates := map[string]any{ + "order_timestamp": params.Timestamp, + "last_message_id": params.MessageID, + "last_message_content": params.Content, + "status": params.Status, + } + if params.Unarchive { + updates["is_archived"] = false + } + if params.MarkAsUnread { + updates["is_read"] = gorm.Expr( + "CASE WHEN last_read_at < ? THEN ? ELSE is_read END", + params.EventTimestamp, + false, + ) + } + return updates +} + +func messageThreadDeletedUpdates(params MessageThreadDeletedUpdate) map[string]any { + return map[string]any{ + "last_message_id": params.LastMessageID, + "last_message_content": params.LastMessageContent, + "status": params.LastMessageStatus, + } +} + +func messageThreadStatusUpdates(params MessageThreadStatusUpdate) map[string]any { + updates := make(map[string]any) + if params.IsArchived != nil { + updates["is_archived"] = *params.IsArchived + } + if params.IsRead != nil { + updates["is_read"] = *params.IsRead + if *params.IsRead { + updates["last_read_at"] = params.ReadAt + } + } + return updates +} + func (repository *gormMessageThreadRepository) DeleteAllForUser(ctx context.Context, userID entities.UserID) error { ctx, span := repository.tracer.Start(ctx) defer span.End() @@ -60,20 +102,17 @@ func (repository *gormMessageThreadRepository) Delete(ctx context.Context, userI } // UpdateAfterDeletedMessage updates a thread after the original message has been deleted -func (repository *gormMessageThreadRepository) UpdateAfterDeletedMessage(ctx context.Context, userID entities.UserID, messageID uuid.UUID) error { +func (repository *gormMessageThreadRepository) UpdateAfterDeletedMessage(ctx context.Context, params MessageThreadDeletedUpdate) error { ctx, span := repository.tracer.Start(ctx) defer span.End() - err := repository.db.WithContext(ctx).Model(&entities.MessageThread{}). - Where("user_id = ?", userID). - Where("last_message_id = ?", messageID). - Updates(map[string]any{ - "last_message_id": nil, - "last_message_content": nil, - "status": entities.MessageStatusDeleted, - }).Error - if err != nil { - return repository.tracer.WrapErrorSpan(span, stacktrace.Propagate(err, "cannot update thread after message is deleted with userID [%s] and messageID [%s]", userID, messageID)) + result := repository.db.WithContext(ctx). + Model(&entities.MessageThread{}). + Where("user_id = ?", params.UserID). + Where("id = ?", params.MessageThreadID). + Updates(messageThreadDeletedUpdates(params)) + if result.Error != nil { + return repository.tracer.WrapErrorSpan(span, stacktrace.Propagate(result.Error, "cannot update deleted-message metadata for thread [%s]", params.MessageThreadID)) } return nil @@ -84,25 +123,77 @@ func (repository *gormMessageThreadRepository) Store(ctx context.Context, thread ctx, span := repository.tracer.Start(ctx) defer span.End() - if err := repository.db.WithContext(ctx).Clauses(clause.OnConflict{DoNothing: true}).Create(thread).Error; err != nil { + isRead := thread.IsRead + err := repository.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { + result := tx.Clauses(clause.OnConflict{DoNothing: true}).Create(thread) + thread.IsRead = isRead + if result.Error != nil { + return result.Error + } + if result.RowsAffected == 0 || isRead { + return nil + } + + return tx.Model(&entities.MessageThread{}). + Where("user_id = ?", thread.UserID). + Where("id = ?", thread.ID). + UpdateColumn("is_read", false). + Error + }) + if err != nil { return repository.tracer.WrapErrorSpan(span, stacktrace.Propagate(err, "cannot save message thread with ID [%s]", thread.ID)) } return nil } -// Update a new entities.MessageThread -func (repository *gormMessageThreadRepository) Update(ctx context.Context, thread *entities.MessageThread) error { +// UpdateActivity persists the last-message activity fields for a thread +func (repository *gormMessageThreadRepository) UpdateActivity(ctx context.Context, params MessageThreadActivityUpdate) error { ctx, span := repository.tracer.Start(ctx) defer span.End() - if err := repository.db.WithContext(ctx).Save(thread).Error; err != nil { - return repository.tracer.WrapErrorSpan(span, stacktrace.Propagate(err, "cannot update message thread thread with ID [%s]", thread.ID)) + result := repository.db.WithContext(ctx). + Model(&entities.MessageThread{}). + Where("user_id = ?", params.UserID). + Where("id = ?", params.MessageThreadID). + Updates(messageThreadActivityUpdates(params)) + if result.Error != nil { + return repository.tracer.WrapErrorSpan(span, stacktrace.Propagate(result.Error, "cannot update message activity for thread [%s]", params.MessageThreadID)) + } + if result.RowsAffected == 0 { + return repository.tracer.WrapErrorSpan(span, stacktrace.PropagateWithCode(gorm.ErrRecordNotFound, ErrCodeNotFound, "thread with id [%s] not found", params.MessageThreadID)) } return nil } +// UpdateStatus persists archive/read status fields for a thread +func (repository *gormMessageThreadRepository) UpdateStatus( + ctx context.Context, + userID entities.UserID, + messageThreadID uuid.UUID, + params MessageThreadStatusUpdate, +) (*entities.MessageThread, error) { + ctx, span := repository.tracer.Start(ctx) + defer span.End() + + thread := new(entities.MessageThread) + result := repository.db.WithContext(ctx). + Model(thread). + Clauses(clause.Returning{}). + Where("user_id = ?", userID). + Where("id = ?", messageThreadID). + Updates(messageThreadStatusUpdates(params)) + if result.Error != nil { + return nil, repository.tracer.WrapErrorSpan(span, stacktrace.Propagate(result.Error, "cannot update status for thread [%s] and user [%s]", messageThreadID, userID)) + } + if result.RowsAffected == 0 { + return nil, repository.tracer.WrapErrorSpan(span, stacktrace.PropagateWithCode(gorm.ErrRecordNotFound, ErrCodeNotFound, "thread with id [%s] not found for user with ID [%s]", messageThreadID, userID)) + } + + return thread, nil +} + // LoadByOwnerContact a thread between 2 users func (repository *gormMessageThreadRepository) LoadByOwnerContact(ctx context.Context, userID entities.UserID, owner string, contact string) (*entities.MessageThread, error) { ctx, span := repository.tracer.Start(ctx) diff --git a/api/pkg/repositories/gorm_message_thread_repository_test.go b/api/pkg/repositories/gorm_message_thread_repository_test.go new file mode 100644 index 000000000..31dc43ed5 --- /dev/null +++ b/api/pkg/repositories/gorm_message_thread_repository_test.go @@ -0,0 +1,204 @@ +package repositories + +import ( + "context" + "database/sql" + "database/sql/driver" + "errors" + "strings" + "testing" + "time" + + "github.com/NdoleStudio/httpsms/pkg/entities" + "github.com/NdoleStudio/httpsms/pkg/telemetry" + "github.com/google/uuid" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "go.opentelemetry.io/otel/trace" + "gorm.io/driver/postgres" + "gorm.io/gorm" +) + +type messageThreadTestStatement struct { + query string + args []any +} + +type messageThreadTestConnPool struct { + statements []messageThreadTestStatement +} + +func (messageThreadTestConnPool) PrepareContext(context.Context, string) (*sql.Stmt, error) { + return nil, errors.New("unexpected PrepareContext") +} + +func (pool *messageThreadTestConnPool) ExecContext(_ context.Context, query string, args ...any) (sql.Result, error) { + pool.statements = append(pool.statements, messageThreadTestStatement{ + query: query, + args: append([]any(nil), args...), + }) + return driver.RowsAffected(1), nil +} + +func (messageThreadTestConnPool) QueryContext(context.Context, string, ...any) (*sql.Rows, error) { + return nil, errors.New("unexpected QueryContext") +} + +func (messageThreadTestConnPool) QueryRowContext(context.Context, string, ...any) *sql.Row { + return &sql.Row{} +} + +func (pool *messageThreadTestConnPool) BeginTx(context.Context, *sql.TxOptions) (gorm.ConnPool, error) { + return pool, nil +} + +func (*messageThreadTestConnPool) Commit() error { + return nil +} + +func (*messageThreadTestConnPool) Rollback() error { + return nil +} + +type messageThreadTestLogger struct{} + +func (logger *messageThreadTestLogger) Error(error) {} +func (logger *messageThreadTestLogger) WithService(string) telemetry.Logger { return logger } + +func (logger *messageThreadTestLogger) WithString(string, string) telemetry.Logger { return logger } + +func (logger *messageThreadTestLogger) WithSpan(trace.SpanContext) telemetry.Logger { return logger } +func (logger *messageThreadTestLogger) Trace(string) {} +func (logger *messageThreadTestLogger) Info(string) {} +func (logger *messageThreadTestLogger) Warn(error) {} +func (logger *messageThreadTestLogger) Debug(string) {} +func (logger *messageThreadTestLogger) Fatal(error) {} +func (logger *messageThreadTestLogger) Printf(string, ...interface{}) {} + +func TestMessageThreadStorePreservesExplicitUnreadState(t *testing.T) { + pool := &messageThreadTestConnPool{} + db, err := gorm.Open( + postgres.New(postgres.Config{ + Conn: pool, + WithoutReturning: true, + }), + &gorm.Config{DisableAutomaticPing: true}, + ) + require.NoError(t, err) + + logger := &messageThreadTestLogger{} + repository := NewGormMessageThreadRepository(logger, telemetry.NewOtelLogger("test", logger), db) + thread := &entities.MessageThread{ + ID: uuid.New(), + IsRead: false, + } + + require.NoError(t, repository.Store(context.Background(), thread)) + assert.False(t, thread.IsRead) + + require.NotEmpty(t, pool.statements) + update := pool.statements[len(pool.statements)-1] + assert.True(t, strings.HasPrefix(update.query, `UPDATE "message_threads"`)) + assert.Contains(t, update.query, `"is_read"=$1`) + assert.Contains(t, update.args, false) +} + +func TestMessageThreadActivityUpdatesOwnOnlyMessageColumns(t *testing.T) { + messageID := uuid.New() + updates := messageThreadActivityUpdates(MessageThreadActivityUpdate{ + Timestamp: time.Date(2026, 7, 18, 7, 0, 0, 0, time.UTC), + MessageID: messageID, + Content: "hello", + Status: entities.MessageStatusReceived, + }) + + assert.Equal(t, map[string]any{ + "order_timestamp": time.Date(2026, 7, 18, 7, 0, 0, 0, time.UTC), + "last_message_id": messageID, + "last_message_content": "hello", + "status": entities.MessageStatus(entities.MessageStatusReceived), + }, updates) + assert.NotContains(t, updates, "is_read") + assert.NotContains(t, updates, "is_archived") + assert.NotContains(t, updates, "last_read_at") +} + +func TestUpdateActivityMarksUnreadWithOneQuery(t *testing.T) { + pool := &messageThreadTestConnPool{} + db, err := gorm.Open( + postgres.New(postgres.Config{ + Conn: pool, + WithoutReturning: true, + }), + &gorm.Config{DisableAutomaticPing: true}, + ) + require.NoError(t, err) + + logger := &messageThreadTestLogger{} + repository := NewGormMessageThreadRepository(logger, telemetry.NewOtelLogger("test", logger), db) + + err = repository.UpdateActivity(context.Background(), MessageThreadActivityUpdate{ + MessageThreadID: uuid.New(), + UserID: entities.UserID("user-id"), + Timestamp: time.Date(2026, 7, 19, 10, 0, 0, 0, time.UTC), + MessageID: uuid.New(), + Content: "hello", + Status: entities.MessageStatusReceived, + MarkAsUnread: true, + EventTimestamp: time.Date(2026, 7, 19, 10, 0, 1, 0, time.UTC), + }) + + require.NoError(t, err) + var updates []messageThreadTestStatement + for _, statement := range pool.statements { + if strings.HasPrefix(statement.query, `UPDATE "message_threads"`) { + updates = append(updates, statement) + } + } + require.Len(t, updates, 1) + assert.Contains(t, updates[0].query, `"is_read"=CASE WHEN last_read_at <`) +} + +func TestMessageThreadDeletedUpdatesPreserveStatusType(t *testing.T) { + messageID := uuid.New() + content := "previous message" + updates := messageThreadDeletedUpdates(MessageThreadDeletedUpdate{ + LastMessageID: &messageID, + LastMessageContent: &content, + LastMessageStatus: entities.MessageStatusDelivered, + }) + + assert.Equal(t, map[string]any{ + "last_message_id": &messageID, + "last_message_content": &content, + "status": entities.MessageStatus(entities.MessageStatusDelivered), + }, updates) +} + +func TestMessageThreadStatusUpdatesReadOnly(t *testing.T) { + isRead := true + readAt := time.Date(2026, 7, 18, 7, 1, 0, 0, time.UTC) + + updates := messageThreadStatusUpdates(MessageThreadStatusUpdate{ + IsRead: &isRead, + ReadAt: readAt, + }) + + assert.Equal(t, map[string]any{ + "is_read": true, + "last_read_at": readAt, + }, updates) + assert.NotContains(t, updates, "is_archived") +} + +func TestMessageThreadStatusUpdatesArchiveOnly(t *testing.T) { + isArchived := true + + updates := messageThreadStatusUpdates(MessageThreadStatusUpdate{ + IsArchived: &isArchived, + }) + + assert.Equal(t, map[string]any{"is_archived": true}, updates) + assert.NotContains(t, updates, "is_read") + assert.NotContains(t, updates, "last_read_at") +} diff --git a/api/pkg/repositories/message_thread_repository.go b/api/pkg/repositories/message_thread_repository.go index 542fd519b..e60931419 100644 --- a/api/pkg/repositories/message_thread_repository.go +++ b/api/pkg/repositories/message_thread_repository.go @@ -2,19 +2,50 @@ package repositories import ( "context" + "time" "github.com/google/uuid" "github.com/NdoleStudio/httpsms/pkg/entities" ) +type MessageThreadActivityUpdate struct { + MessageThreadID uuid.UUID + UserID entities.UserID + // Timestamp controls thread activity ordering; EventTimestamp is the server-side unread watermark. + Timestamp time.Time + MessageID uuid.UUID + Content string + Status entities.MessageStatus + MarkAsUnread bool + EventTimestamp time.Time + Unarchive bool +} + +type MessageThreadStatusUpdate struct { + IsArchived *bool + IsRead *bool + ReadAt time.Time +} + +type MessageThreadDeletedUpdate struct { + MessageThreadID uuid.UUID + UserID entities.UserID + LastMessageID *uuid.UUID + LastMessageContent *string + LastMessageStatus entities.MessageStatus +} + // MessageThreadRepository loads and persists an entities.MessageThread type MessageThreadRepository interface { // Store a new entities.MessageThread Store(ctx context.Context, thread *entities.MessageThread) error - // Update a new entities.MessageThread - Update(ctx context.Context, thread *entities.MessageThread) error + // UpdateActivity persists the last-message activity fields for a thread + UpdateActivity(ctx context.Context, params MessageThreadActivityUpdate) error + + // UpdateStatus persists archive/read status fields for a thread + UpdateStatus(ctx context.Context, userID entities.UserID, messageThreadID uuid.UUID, params MessageThreadStatusUpdate) (*entities.MessageThread, error) // LoadByOwnerContact fetches a thread between owner and contact LoadByOwnerContact(ctx context.Context, userID entities.UserID, owner string, contact string) (*entities.MessageThread, error) @@ -26,7 +57,7 @@ type MessageThreadRepository interface { Index(ctx context.Context, userID entities.UserID, owner string, archived bool, params IndexParams) (*[]entities.MessageThread, error) // UpdateAfterDeletedMessage updates a thread after the original message has been deleted - UpdateAfterDeletedMessage(ctx context.Context, userID entities.UserID, messageID uuid.UUID) error + UpdateAfterDeletedMessage(ctx context.Context, params MessageThreadDeletedUpdate) error // Delete an entities.MessageThread by ID Delete(ctx context.Context, userID entities.UserID, messageThreadID uuid.UUID) error diff --git a/api/pkg/requests/message_thread_update_request.go b/api/pkg/requests/message_thread_update_request.go index d64d75c59..3309fc91d 100644 --- a/api/pkg/requests/message_thread_update_request.go +++ b/api/pkg/requests/message_thread_update_request.go @@ -10,7 +10,8 @@ import ( // MessageThreadUpdate is the payload for updating a message thread type MessageThreadUpdate struct { request - IsArchived bool `json:"is_archived" example:"true"` + IsArchived *bool `json:"is_archived,omitempty" example:"true"` + IsRead *bool `json:"is_read,omitempty" example:"true"` MessageThreadID string `json:"messageThreadID" swaggerignore:"true"` // used internally for validation } @@ -21,5 +22,6 @@ func (input *MessageThreadUpdate) ToUpdateParams(userID entities.UserID) service UserID: userID, MessageThreadID: uuid.MustParse(input.MessageThreadID), IsArchived: input.IsArchived, + IsRead: input.IsRead, } } diff --git a/api/pkg/requests/message_thread_update_request_test.go b/api/pkg/requests/message_thread_update_request_test.go new file mode 100644 index 000000000..9f9579fda --- /dev/null +++ b/api/pkg/requests/message_thread_update_request_test.go @@ -0,0 +1,25 @@ +package requests + +import ( + "testing" + + "github.com/NdoleStudio/httpsms/pkg/entities" + "github.com/google/uuid" + "github.com/stretchr/testify/assert" +) + +func TestMessageThreadUpdateToUpdateParamsPreservesOptionalFields(t *testing.T) { + threadID := uuid.New() + isRead := true + input := MessageThreadUpdate{ + MessageThreadID: threadID.String(), + IsRead: &isRead, + } + + params := input.ToUpdateParams(entities.UserID("user-id")) + + assert.Equal(t, threadID, params.MessageThreadID) + assert.Equal(t, entities.UserID("user-id"), params.UserID) + assert.Nil(t, params.IsArchived) + assert.Same(t, &isRead, params.IsRead) +} diff --git a/api/pkg/responses/message_thead_responses.go b/api/pkg/responses/message_thead_responses.go index 917706a47..e47fec287 100644 --- a/api/pkg/responses/message_thead_responses.go +++ b/api/pkg/responses/message_thead_responses.go @@ -7,3 +7,9 @@ type MessageThreadsResponse struct { response Data []entities.MessageThread `json:"data"` } + +// MessageThreadResponse is the payload containing entities.MessageThread +type MessageThreadResponse struct { + response + Data entities.MessageThread `json:"data"` +} diff --git a/api/pkg/services/message_thread_service.go b/api/pkg/services/message_thread_service.go index c77458cbf..e3dc882c7 100644 --- a/api/pkg/services/message_thread_service.go +++ b/api/pkg/services/message_thread_service.go @@ -50,7 +50,10 @@ type MessageThreadUpdateParams struct { Content string UserID entities.UserID MessageID uuid.UUID - Timestamp time.Time + // Timestamp controls thread activity ordering; EventTimestamp is the server-side unread watermark. + Timestamp time.Time + MarkAsUnread bool + EventTimestamp time.Time } // shouldCheckUnarchive reports whether a thread update is a new inbound message @@ -101,27 +104,39 @@ func (service *MessageThreadService) UpdateThread(ctx context.Context, params Me return nil } + activity := repositories.MessageThreadActivityUpdate{ + MessageThreadID: thread.ID, + UserID: params.UserID, + Timestamp: params.Timestamp, + MessageID: params.MessageID, + Content: params.Content, + Status: params.Status, + MarkAsUnread: params.MarkAsUnread, + EventTimestamp: params.EventTimestamp, + } + if service.shouldCheckUnarchive(thread, params) { phone, phoneErr := service.phoneRepository.Load(ctx, params.UserID, params.Owner) if phoneErr != nil { ctxLogger.Warn(stacktrace.Propagate(phoneErr, "cannot load phone [%s] for user [%s] to resolve UnarchiveThread; leaving thread [%s] archived", params.Owner, params.UserID, thread.ID)) } else if phone.UnarchiveThread { - thread.UpdateArchive(false) + activity.Unarchive = true ctxLogger.Info(fmt.Sprintf("unarchiving thread [%s] after inbound message [%s]", thread.ID, params.MessageID)) } } - if err = service.repository.Update(ctx, thread.Update(params.Timestamp, params.MessageID, params.Content, params.Status)); err != nil { - return service.tracer.WrapErrorSpan(span, stacktrace.Propagate(err, "cannot update message thread with id [%s] after adding message [%s]", thread.ID, params.MessageID)) + if err = service.repository.UpdateActivity(ctx, activity); err != nil { + return service.tracer.WrapErrorSpan(span, stacktrace.PropagateWithCode(err, stacktrace.GetCode(err), "cannot update message thread with id [%s] after adding message [%s]", thread.ID, params.MessageID)) } - ctxLogger.Info(fmt.Sprintf("thread with id [%s] updated with last message [%s] and status [%s]", thread.ID, thread.LastMessageID, thread.Status)) + ctxLogger.Info(fmt.Sprintf("thread with id [%s] updated with last message [%s] and status [%s]", thread.ID, params.MessageID, params.Status)) return nil } // MessageThreadStatusParams are parameters for updating a thread status type MessageThreadStatusParams struct { - IsArchived bool + IsArchived *bool + IsRead *bool UserID entities.UserID MessageThreadID uuid.UUID } @@ -131,18 +146,16 @@ func (service *MessageThreadService) UpdateStatus(ctx context.Context, params Me ctx, span := service.tracer.Start(ctx) defer span.End() - ctxLogger := service.tracer.CtxLogger(service.logger, span) - - thread, err := service.repository.Load(ctx, params.UserID, params.MessageThreadID) - if err != nil { - return nil, service.tracer.WrapErrorSpan(span, stacktrace.Propagate(err, "cannot find thread with id [%s]", params.MessageThreadID)) + update := repositories.MessageThreadStatusUpdate{ + IsArchived: params.IsArchived, + IsRead: params.IsRead, + ReadAt: time.Now().UTC(), } - - if err = service.repository.Update(ctx, thread.UpdateArchive(params.IsArchived)); err != nil { - return nil, service.tracer.WrapErrorSpan(span, stacktrace.Propagate(err, "cannot update message thread with id [%s] with archive status [%t]", thread.ID, params.IsArchived)) + thread, err := service.repository.UpdateStatus(ctx, params.UserID, params.MessageThreadID, update) + if err != nil { + return nil, service.tracer.WrapErrorSpan(span, stacktrace.PropagateWithCode(err, stacktrace.GetCode(err), "cannot update message thread with ID [%s] for user [%s]", params.MessageThreadID, params.UserID)) } - ctxLogger.Info(fmt.Sprintf("thread with id [%s] updated with archive status [%t]", thread.ID, thread.IsArchived)) return thread, nil } @@ -172,12 +185,13 @@ func (service *MessageThreadService) UpdateAfterDeletedMessage(ctx context.Conte return nil } - thread.LastMessageContent = payload.PreviousMessageContent - thread.LastMessageID = payload.PreviousMessageID - thread.Status = *payload.PreviousMessageStatus - thread.UpdatedAt = time.Now().UTC() - - if err = service.repository.Update(ctx, thread); err != nil { + if err = service.repository.UpdateAfterDeletedMessage(ctx, repositories.MessageThreadDeletedUpdate{ + MessageThreadID: thread.ID, + UserID: thread.UserID, + LastMessageID: payload.PreviousMessageID, + LastMessageContent: payload.PreviousMessageContent, + LastMessageStatus: *payload.PreviousMessageStatus, + }); err != nil { return service.tracer.WrapErrorSpan(span, stacktrace.Propagate(err, "cannot update thread with ID [%s] for user with ID [%s]", thread.ID, thread.UserID)) } @@ -191,18 +205,21 @@ func (service *MessageThreadService) createThread(ctx context.Context, params Me ctxLogger := service.tracer.CtxLogger(service.logger, span) + now := time.Now().UTC() thread := &entities.MessageThread{ ID: uuid.New(), Owner: params.Owner, Contact: params.Contact, UserID: params.UserID, IsArchived: false, + IsRead: !params.MarkAsUnread, + LastReadAt: now, Color: service.getColor(), LastMessageContent: ¶ms.Content, Status: params.Status, LastMessageID: ¶ms.MessageID, - CreatedAt: time.Now().UTC(), - UpdatedAt: time.Now().UTC(), + CreatedAt: now, + UpdatedAt: now, OrderTimestamp: params.Timestamp, } diff --git a/api/pkg/services/message_thread_service_test.go b/api/pkg/services/message_thread_service_test.go index 7e48ab433..fc1e84313 100644 --- a/api/pkg/services/message_thread_service_test.go +++ b/api/pkg/services/message_thread_service_test.go @@ -1,12 +1,226 @@ package services import ( + "context" "testing" + "time" "github.com/NdoleStudio/httpsms/pkg/entities" + "github.com/NdoleStudio/httpsms/pkg/repositories" + "github.com/NdoleStudio/httpsms/pkg/telemetry" + "github.com/google/uuid" + "github.com/palantir/stacktrace" "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "gorm.io/gorm" ) +type messageThreadRepositoryStub struct { + loadByOwnerContact func(context.Context, entities.UserID, string, string) (*entities.MessageThread, error) + load func(context.Context, entities.UserID, uuid.UUID) (*entities.MessageThread, error) + store func(context.Context, *entities.MessageThread) error + updateActivity func(context.Context, repositories.MessageThreadActivityUpdate) error + updateStatus func(context.Context, entities.UserID, uuid.UUID, repositories.MessageThreadStatusUpdate) (*entities.MessageThread, error) +} + +func (stub *messageThreadRepositoryStub) Store(ctx context.Context, thread *entities.MessageThread) error { + if stub.store != nil { + return stub.store(ctx, thread) + } + return nil +} + +func (stub *messageThreadRepositoryStub) UpdateActivity(ctx context.Context, params repositories.MessageThreadActivityUpdate) error { + if stub.updateActivity != nil { + return stub.updateActivity(ctx, params) + } + return nil +} + +func (stub *messageThreadRepositoryStub) UpdateStatus(ctx context.Context, userID entities.UserID, threadID uuid.UUID, params repositories.MessageThreadStatusUpdate) (*entities.MessageThread, error) { + if stub.updateStatus != nil { + return stub.updateStatus(ctx, userID, threadID, params) + } + return &entities.MessageThread{ID: threadID}, nil +} + +func (stub *messageThreadRepositoryStub) UpdateAfterDeletedMessage(context.Context, repositories.MessageThreadDeletedUpdate) error { + return nil +} + +func (stub *messageThreadRepositoryStub) LoadByOwnerContact(ctx context.Context, userID entities.UserID, owner string, contact string) (*entities.MessageThread, error) { + return stub.loadByOwnerContact(ctx, userID, owner, contact) +} + +func (stub *messageThreadRepositoryStub) Load(ctx context.Context, userID entities.UserID, id uuid.UUID) (*entities.MessageThread, error) { + return stub.load(ctx, userID, id) +} + +func (stub *messageThreadRepositoryStub) Index(context.Context, entities.UserID, string, bool, repositories.IndexParams) (*[]entities.MessageThread, error) { + threads := []entities.MessageThread{} + return &threads, nil +} + +func (stub *messageThreadRepositoryStub) Delete(context.Context, entities.UserID, uuid.UUID) error { + return nil +} + +func (stub *messageThreadRepositoryStub) DeleteAllForUser(context.Context, entities.UserID) error { + return nil +} + +func newMessageThreadServiceForTest(repository repositories.MessageThreadRepository) *MessageThreadService { + logger := &noopLogger{} + tracer := telemetry.NewOtelLogger("test", logger) + return NewMessageThreadService(logger, tracer, repository, nil, nil) +} + +func TestUpdateThreadPassesUnreadWatermarkForInboundActivity(t *testing.T) { + threadID := uuid.New() + eventTimestamp := time.Date(2026, 7, 18, 7, 0, 0, 0, time.UTC) + var captured repositories.MessageThreadActivityUpdate + repository := &messageThreadRepositoryStub{ + loadByOwnerContact: func(context.Context, entities.UserID, string, string) (*entities.MessageThread, error) { + return &entities.MessageThread{ID: threadID}, nil + }, + updateActivity: func(_ context.Context, params repositories.MessageThreadActivityUpdate) error { + captured = params + return nil + }, + } + + service := newMessageThreadServiceForTest(repository) + err := service.UpdateThread(context.Background(), MessageThreadUpdateParams{ + UserID: entities.UserID("user-id"), + Owner: "+18005550199", + Contact: "+18005550100", + MessageID: uuid.New(), + Content: "hello", + Status: entities.MessageStatusReceived, + Timestamp: eventTimestamp, + MarkAsUnread: true, + EventTimestamp: eventTimestamp, + }) + + require.NoError(t, err) + assert.True(t, captured.MarkAsUnread) + assert.Equal(t, eventTimestamp, captured.EventTimestamp) +} + +func TestUpdateThreadPreservesReadStateForOutboundActivity(t *testing.T) { + var captured repositories.MessageThreadActivityUpdate + repository := &messageThreadRepositoryStub{ + loadByOwnerContact: func(context.Context, entities.UserID, string, string) (*entities.MessageThread, error) { + return &entities.MessageThread{ID: uuid.New(), IsRead: false}, nil + }, + updateActivity: func(_ context.Context, params repositories.MessageThreadActivityUpdate) error { + captured = params + return nil + }, + } + + service := newMessageThreadServiceForTest(repository) + err := service.UpdateThread(context.Background(), MessageThreadUpdateParams{ + UserID: entities.UserID("user-id"), + Owner: "+18005550199", + Contact: "+18005550100", + MessageID: uuid.New(), + Content: "outbound", + Status: entities.MessageStatusSent, + Timestamp: time.Now().UTC(), + }) + + require.NoError(t, err) + assert.False(t, captured.MarkAsUnread) +} + +func TestCreateThreadSetsReadStateFromActivityDirection(t *testing.T) { + tests := []struct { + name string + marksUnread bool + wantRead bool + }{ + {name: "inbound", marksUnread: true, wantRead: false}, + {name: "outbound", marksUnread: false, wantRead: true}, + } + + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + var stored *entities.MessageThread + repository := &messageThreadRepositoryStub{ + loadByOwnerContact: func(context.Context, entities.UserID, string, string) (*entities.MessageThread, error) { + return nil, stacktrace.PropagateWithCode(gorm.ErrRecordNotFound, repositories.ErrCodeNotFound, "not found") + }, + store: func(_ context.Context, thread *entities.MessageThread) error { + stored = thread + return nil + }, + } + + service := newMessageThreadServiceForTest(repository) + err := service.UpdateThread(context.Background(), MessageThreadUpdateParams{ + UserID: entities.UserID("user-id"), + Owner: "+18005550199", + Contact: "+18005550100", + MessageID: uuid.New(), + Content: "hello", + Status: entities.MessageStatusReceived, + Timestamp: time.Now().UTC(), + MarkAsUnread: test.marksUnread, + }) + + require.NoError(t, err) + require.NotNil(t, stored) + assert.Equal(t, test.wantRead, stored.IsRead) + assert.False(t, stored.LastReadAt.IsZero()) + }) + } +} + +func TestUpdateStatusChangesOnlyRequestedState(t *testing.T) { + threadID := uuid.New() + isRead := false + var captured repositories.MessageThreadStatusUpdate + repository := &messageThreadRepositoryStub{ + updateStatus: func(_ context.Context, _ entities.UserID, _ uuid.UUID, params repositories.MessageThreadStatusUpdate) (*entities.MessageThread, error) { + captured = params + return &entities.MessageThread{ID: threadID, IsArchived: true, IsRead: false}, nil + }, + } + + service := newMessageThreadServiceForTest(repository) + thread, err := service.UpdateStatus(context.Background(), MessageThreadStatusParams{ + UserID: entities.UserID("user-id"), + MessageThreadID: threadID, + IsRead: &isRead, + }) + + require.NoError(t, err) + assert.Nil(t, captured.IsArchived) + assert.Same(t, &isRead, captured.IsRead) + assert.False(t, captured.ReadAt.IsZero()) + assert.True(t, thread.IsArchived) + assert.False(t, thread.IsRead) +} + +func TestUpdateStatusPreservesNotFoundCode(t *testing.T) { + repository := &messageThreadRepositoryStub{ + updateStatus: func(context.Context, entities.UserID, uuid.UUID, repositories.MessageThreadStatusUpdate) (*entities.MessageThread, error) { + return nil, stacktrace.PropagateWithCode(gorm.ErrRecordNotFound, repositories.ErrCodeNotFound, "not found") + }, + } + + service := newMessageThreadServiceForTest(repository) + isRead := true + _, err := service.UpdateStatus(context.Background(), MessageThreadStatusParams{ + UserID: entities.UserID("user-id"), + MessageThreadID: uuid.New(), + IsRead: &isRead, + }) + + assert.Equal(t, repositories.ErrCodeNotFound, stacktrace.GetCode(err)) +} + func TestShouldCheckUnarchive(t *testing.T) { service := &MessageThreadService{} diff --git a/api/pkg/validators/message_thread_handler_validator.go b/api/pkg/validators/message_thread_handler_validator.go index a4e5cd647..72a64194f 100644 --- a/api/pkg/validators/message_thread_handler_validator.go +++ b/api/pkg/validators/message_thread_handler_validator.go @@ -72,5 +72,13 @@ func (validator *MessageThreadHandlerValidator) ValidateUpdate(_ context.Context }, }) - return v.ValidateStruct() + errors := v.ValidateStruct() + if request.IsArchived == nil && request.IsRead == nil { + if errors == nil { + errors = url.Values{} + } + errors.Add("payload", "at least one of is_archived or is_read is required") + } + + return errors } diff --git a/api/pkg/validators/message_thread_handler_validator_test.go b/api/pkg/validators/message_thread_handler_validator_test.go new file mode 100644 index 000000000..e5ec7c1b7 --- /dev/null +++ b/api/pkg/validators/message_thread_handler_validator_test.go @@ -0,0 +1,34 @@ +package validators + +import ( + "context" + "testing" + + "github.com/NdoleStudio/httpsms/pkg/requests" + "github.com/google/uuid" + "github.com/stretchr/testify/assert" +) + +func TestValidateUpdateRequiresAtLeastOneStatusField(t *testing.T) { + validator := &MessageThreadHandlerValidator{} + request := requests.MessageThreadUpdate{ + MessageThreadID: uuid.NewString(), + } + + errors := validator.ValidateUpdate(context.Background(), request) + + assert.NotEmpty(t, errors.Get("payload")) +} + +func TestValidateUpdateAcceptsReadOnlyUpdate(t *testing.T) { + validator := &MessageThreadHandlerValidator{} + isRead := true + request := requests.MessageThreadUpdate{ + MessageThreadID: uuid.NewString(), + IsRead: &isRead, + } + + errors := validator.ValidateUpdate(context.Background(), request) + + assert.Empty(t, errors) +} diff --git a/docs/superpowers/plans/2026-07-18-read-receipts.md b/docs/superpowers/plans/2026-07-18-read-receipts.md new file mode 100644 index 000000000..1887a1511 --- /dev/null +++ b/docs/superpowers/plans/2026-07-18-read-receipts.md @@ -0,0 +1,1760 @@ +# Message Thread Read Receipts Implementation Plan + +> **For agentic workers:** REQUIRED SUB-SKILL: Use superpowers:subagent-driven-development (recommended) or superpowers:executing-plans to implement this plan task-by-task. Steps use checkbox (`- [ ]`) syntax for tracking. + +**Goal:** Persist message-thread read state, mark inbound SMS and missed calls unread, automatically mark opened threads read, and highlight unread threads in the web UI. + +**Architecture:** Store `is_read` plus an internal `last_read_at` watermark on each message thread. Replace full-row thread saves with field-scoped GORM updates so concurrent event listeners cannot overwrite a newer read action; inbound updates compare CloudEvent creation time against the stored watermark. Extend the existing thread update endpoint with optional archive/read fields, then have the Nuxt store call it automatically when the thread page opens or refreshes after inbound realtime events. + +**Tech Stack:** Go 1.25, Fiber v3, GORM/PostgreSQL, CloudEvents, Pusher, Nuxt 4 SPA, Vue 3, Pinia, Vuetify 4, TypeScript. + +## Global Constraints + +- Work only in `C:\Users\achoa\Work\NdoleStudio\httpsms-read-receipts` on branch `feat/read-receipts`, based on `origin/main`. +- Existing message threads must migrate as read. +- Only incoming SMS and missed-call activity marks a thread unread. +- Outbound messages and delivery/status events preserve the existing read state. +- Opening a thread marks it read automatically; inbound activity while it is open must leave it read. +- Do not add a new endpoint or a manual read/unread control. +- Use the existing `PUT /v1/message-threads/{messageThreadID}` endpoint with optional `is_archived` and `is_read`. +- Use GORM query builders with `WithContext(ctx)`; do not use raw SQL. +- Wrap API errors with `stacktrace.Propagate`/`PropagateWithCode`. +- Web code uses single quotes, no semicolons, and 2-space indentation. +- Every commit must end with: + +```text +Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> +Copilot-Session: bf00a0ac-e11f-4015-b295-3cdd9b491229 +``` + +- Add end-to-end coverage in the existing `tests/` integration project. + +## File Structure + +- `api/pkg/entities/message_thread.go`: persisted read fields. +- `api/pkg/entities/message_thread_test.go`: schema-default regression tests. +- `api/pkg/repositories/message_thread_repository.go`: field-scoped update parameter types and interface methods. +- `api/pkg/repositories/gorm_message_thread_repository.go`: transactional activity updates, conditional unread transition, partial user-status update, and partial deleted-message update. +- `api/pkg/repositories/gorm_message_thread_repository_test.go`: update-map ownership tests. +- `api/pkg/services/message_thread_service.go`: read-state business rules and repository orchestration. +- `api/pkg/services/message_thread_service_test.go`: service behavior with a repository stub. +- `api/pkg/requests/message_thread_update_request.go`: optional archive/read request fields. +- `api/pkg/requests/message_thread_update_request_test.go`: request-to-service conversion tests. +- `api/pkg/validators/message_thread_handler_validator.go`: require at least one update field. +- `api/pkg/validators/message_thread_handler_validator_test.go`: empty-payload and valid-payload tests. +- `api/pkg/responses/message_thead_responses.go`: single-thread Swagger response. +- `api/pkg/handlers/message_thread_handler.go`: not-found mapping and corrected Swagger response. +- `api/pkg/listeners/message_thread_listener.go`: incoming SMS and missed-call unread updates. +- `api/pkg/listeners/websocket_listener.go`: missed-call Pusher publication. +- `api/pkg/listeners/read_receipts_test_helpers_test.go`: listener logger and repository test doubles. +- `api/pkg/listeners/message_thread_listener_test.go`: listener route and payload tests. +- `api/pkg/listeners/websocket_listener_test.go`: missed-call route registration test. +- `api/docs/docs.go`, `api/docs/swagger.json`, `api/docs/swagger.yaml`: regenerated API documentation. +- `web/shared/types/api.ts`: regenerated thread/request types. +- `web/app/stores/threads.ts`: automatic read update and local state replacement. +- `web/app/pages/threads/[id]/index.vue`: non-blocking read calls and missed-call realtime refresh. +- `web/app/components/MessageThread.vue`: unread visual treatment. +- `tests/read_receipts_test.go`: Docker-stack read/unread lifecycle coverage. +- `tests/README.md`: integration coverage documentation. + +--- + +### Task 1: Add the read-state API contract + +**Files:** + +- Modify: `api/pkg/entities/message_thread.go` +- Create: `api/pkg/entities/message_thread_test.go` +- Modify: `api/pkg/requests/message_thread_update_request.go` +- Create: `api/pkg/requests/message_thread_update_request_test.go` +- Modify: `api/pkg/validators/message_thread_handler_validator.go` +- Create: `api/pkg/validators/message_thread_handler_validator_test.go` +- Modify: `api/pkg/responses/message_thead_responses.go` +- Modify: `api/pkg/services/message_thread_service.go` + +**Interfaces:** + +- Produces: `MessageThread.IsRead bool`, `MessageThread.LastReadAt time.Time` +- Produces: `MessageThreadUpdate.IsArchived *bool`, `MessageThreadUpdate.IsRead *bool` +- Produces: `MessageThreadStatusParams.IsArchived *bool`, `MessageThreadStatusParams.IsRead *bool` +- Produces: `responses.MessageThreadResponse` + +- [ ] **Step 1: Write failing entity schema tests** + +Create `api/pkg/entities/message_thread_test.go`: + +```go +package entities + +import ( + "reflect" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func TestMessageThreadReadFieldsHaveBackwardCompatibleDefaults(t *testing.T) { + threadType := reflect.TypeOf(MessageThread{}) + + isRead, ok := threadType.FieldByName("IsRead") + require.True(t, ok) + assert.Contains(t, isRead.Tag.Get("gorm"), "not null") + assert.Contains(t, isRead.Tag.Get("gorm"), "default:true") + assert.Equal(t, "is_read", isRead.Tag.Get("json")) + + lastReadAt, ok := threadType.FieldByName("LastReadAt") + require.True(t, ok) + assert.Contains(t, lastReadAt.Tag.Get("gorm"), "not null") + assert.Contains(t, lastReadAt.Tag.Get("gorm"), "default:CURRENT_TIMESTAMP") + assert.Equal(t, "-", lastReadAt.Tag.Get("json")) +} +``` + +- [ ] **Step 2: Write failing request and validator tests** + +Create `api/pkg/requests/message_thread_update_request_test.go`: + +```go +package requests + +import ( + "testing" + + "github.com/NdoleStudio/httpsms/pkg/entities" + "github.com/google/uuid" + "github.com/stretchr/testify/assert" +) + +func TestMessageThreadUpdateToUpdateParamsPreservesOptionalFields(t *testing.T) { + threadID := uuid.New() + isRead := true + input := MessageThreadUpdate{ + MessageThreadID: threadID.String(), + IsRead: &isRead, + } + + params := input.ToUpdateParams(entities.UserID("user-id")) + + assert.Equal(t, threadID, params.MessageThreadID) + assert.Equal(t, entities.UserID("user-id"), params.UserID) + assert.Nil(t, params.IsArchived) + assert.Same(t, &isRead, params.IsRead) +} +``` + +Create `api/pkg/validators/message_thread_handler_validator_test.go`: + +```go +package validators + +import ( + "context" + "testing" + + "github.com/NdoleStudio/httpsms/pkg/requests" + "github.com/google/uuid" + "github.com/stretchr/testify/assert" +) + +func TestValidateUpdateRequiresAtLeastOneStatusField(t *testing.T) { + validator := &MessageThreadHandlerValidator{} + request := requests.MessageThreadUpdate{ + MessageThreadID: uuid.NewString(), + } + + errors := validator.ValidateUpdate(context.Background(), request) + + assert.NotEmpty(t, errors.Get("payload")) +} + +func TestValidateUpdateAcceptsReadOnlyUpdate(t *testing.T) { + validator := &MessageThreadHandlerValidator{} + isRead := true + request := requests.MessageThreadUpdate{ + MessageThreadID: uuid.NewString(), + IsRead: &isRead, + } + + errors := validator.ValidateUpdate(context.Background(), request) + + assert.Empty(t, errors) +} +``` + +- [ ] **Step 3: Run the focused tests and confirm they fail** + +Run: + +```powershell +Set-Location 'C:\Users\achoa\Work\NdoleStudio\httpsms-read-receipts\api' +go test ./pkg/entities ./pkg/requests ./pkg/validators +``` + +Expected: compile failures because the read fields and optional request fields do not exist. + +- [ ] **Step 4: Add the persisted fields and optional request contract** + +Add to `entities.MessageThread` after `IsArchived`: + +```go +IsRead bool `json:"is_read" gorm:"not null;default:true" example:"true"` +LastReadAt time.Time `json:"-" gorm:"not null;default:CURRENT_TIMESTAMP"` +``` + +Change `requests.MessageThreadUpdate` and its conversion: + +```go +type MessageThreadUpdate struct { + request + IsArchived *bool `json:"is_archived,omitempty" example:"true"` + IsRead *bool `json:"is_read,omitempty" example:"true"` + + MessageThreadID string `json:"messageThreadID" swaggerignore:"true"` +} + +func (input *MessageThreadUpdate) ToUpdateParams(userID entities.UserID) services.MessageThreadStatusParams { + return services.MessageThreadStatusParams{ + UserID: userID, + MessageThreadID: uuid.MustParse(input.MessageThreadID), + IsArchived: input.IsArchived, + IsRead: input.IsRead, + } +} +``` + +Change the service parameter type: + +```go +type MessageThreadStatusParams struct { + IsArchived *bool + IsRead *bool + UserID entities.UserID + MessageThreadID uuid.UUID +} +``` + +After `v.ValidateStruct()` in `ValidateUpdate`, add an explicit payload check: + +```go +errors := v.ValidateStruct() +if request.IsArchived == nil && request.IsRead == nil { + errors.Add("payload", "at least one of is_archived or is_read is required") +} +return errors +``` + +Add a single-thread response beside `MessageThreadsResponse`: + +```go +// MessageThreadResponse is the payload containing entities.MessageThread +type MessageThreadResponse struct { + response + Data entities.MessageThread `json:"data"` +} +``` + +- [ ] **Step 5: Run focused tests** + +Run: + +```powershell +go test ./pkg/entities ./pkg/requests ./pkg/validators +``` + +Expected: PASS. + +- [ ] **Step 6: Commit the contract** + +```powershell +git add api/pkg/entities/message_thread.go api/pkg/entities/message_thread_test.go api/pkg/requests/message_thread_update_request.go api/pkg/requests/message_thread_update_request_test.go api/pkg/validators/message_thread_handler_validator.go api/pkg/validators/message_thread_handler_validator_test.go api/pkg/responses/message_thead_responses.go api/pkg/services/message_thread_service.go +@' +feat(api): add thread read state contract + +Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> +Copilot-Session: bf00a0ac-e11f-4015-b295-3cdd9b491229 +'@ | git commit -F - +``` + +### Task 2: Add atomic thread persistence operations + +**Files:** + +- Modify: `api/pkg/repositories/message_thread_repository.go` +- Modify: `api/pkg/repositories/gorm_message_thread_repository.go` +- Create: `api/pkg/repositories/gorm_message_thread_repository_test.go` + +**Interfaces:** + +- Produces: `repositories.MessageThreadActivityUpdate` +- Produces: `repositories.MessageThreadStatusUpdate` +- Produces: `repositories.MessageThreadDeletedUpdate` +- Produces: `MessageThreadRepository.UpdateActivity`, `UpdateStatus`, `UpdateAfterDeletedMessage` +- Removes: `MessageThreadRepository.Update` + +- [ ] **Step 1: Write failing update-map ownership tests** + +Create `api/pkg/repositories/gorm_message_thread_repository_test.go`: + +```go +package repositories + +import ( + "testing" + "time" + + "github.com/NdoleStudio/httpsms/pkg/entities" + "github.com/google/uuid" + "github.com/stretchr/testify/assert" +) + +func TestMessageThreadActivityUpdatesOwnOnlyMessageColumns(t *testing.T) { + messageID := uuid.New() + updates := messageThreadActivityUpdates(MessageThreadActivityUpdate{ + Timestamp: time.Date(2026, 7, 18, 7, 0, 0, 0, time.UTC), + MessageID: messageID, + Content: "hello", + Status: entities.MessageStatusReceived, + }) + + assert.Equal(t, map[string]any{ + "order_timestamp": time.Date(2026, 7, 18, 7, 0, 0, 0, time.UTC), + "last_message_id": messageID, + "last_message_content": "hello", + "status": entities.MessageStatusReceived, + }, updates) + assert.NotContains(t, updates, "is_read") + assert.NotContains(t, updates, "is_archived") + assert.NotContains(t, updates, "last_read_at") +} + +func TestMessageThreadStatusUpdatesReadOnly(t *testing.T) { + isRead := true + readAt := time.Date(2026, 7, 18, 7, 1, 0, 0, time.UTC) + + updates := messageThreadStatusUpdates(MessageThreadStatusUpdate{ + IsRead: &isRead, + ReadAt: readAt, + }) + + assert.Equal(t, map[string]any{ + "is_read": true, + "last_read_at": readAt, + }, updates) + assert.NotContains(t, updates, "is_archived") +} + +func TestMessageThreadStatusUpdatesArchiveOnly(t *testing.T) { + isArchived := true + + updates := messageThreadStatusUpdates(MessageThreadStatusUpdate{ + IsArchived: &isArchived, + }) + + assert.Equal(t, map[string]any{"is_archived": true}, updates) + assert.NotContains(t, updates, "is_read") + assert.NotContains(t, updates, "last_read_at") +} +``` + +- [ ] **Step 2: Run the repository test and confirm it fails** + +Run: + +```powershell +go test ./pkg/repositories -run MessageThread -v +``` + +Expected: compile failures because the parameter types and helper functions do not exist. + +- [ ] **Step 3: Define repository update parameters and methods** + +In `message_thread_repository.go`, add: + +```go +type MessageThreadActivityUpdate struct { + MessageThreadID uuid.UUID + UserID entities.UserID + Timestamp time.Time + MessageID uuid.UUID + Content string + Status entities.MessageStatus + MarksUnread bool + EventTimestamp time.Time +} + +type MessageThreadStatusUpdate struct { + IsArchived *bool + IsRead *bool + ReadAt time.Time +} + +type MessageThreadDeletedUpdate struct { + MessageThreadID uuid.UUID + UserID entities.UserID + LastMessageID *uuid.UUID + LastMessageContent *string + LastMessageStatus entities.MessageStatus +} +``` + +Add the `time` import. Replace `Update` and the old deleted-message method in the interface with: + +```go +UpdateActivity(ctx context.Context, params MessageThreadActivityUpdate) error +UpdateStatus(ctx context.Context, userID entities.UserID, messageThreadID uuid.UUID, params MessageThreadStatusUpdate) error +UpdateAfterDeletedMessage(ctx context.Context, params MessageThreadDeletedUpdate) error +``` + +- [ ] **Step 4: Implement field-scoped GORM updates** + +Add these helpers to `gorm_message_thread_repository.go`: + +```go +func messageThreadActivityUpdates(params MessageThreadActivityUpdate) map[string]any { + return map[string]any{ + "order_timestamp": params.Timestamp, + "last_message_id": params.MessageID, + "last_message_content": params.Content, + "status": params.Status, + } +} + +func messageThreadStatusUpdates(params MessageThreadStatusUpdate) map[string]any { + updates := make(map[string]any) + if params.IsArchived != nil { + updates["is_archived"] = *params.IsArchived + } + if params.IsRead != nil { + updates["is_read"] = *params.IsRead + if *params.IsRead { + updates["last_read_at"] = params.ReadAt + } + } + return updates +} +``` + +Replace the full-row `Update` implementation with: + +```go +func (repository *gormMessageThreadRepository) UpdateActivity(ctx context.Context, params MessageThreadActivityUpdate) error { + ctx, span := repository.tracer.Start(ctx) + defer span.End() + + err := repository.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { + query := tx.Model(&entities.MessageThread{}). + Where("user_id = ?", params.UserID). + Where("id = ?", params.MessageThreadID) + + result := query.Updates(messageThreadActivityUpdates(params)) + if result.Error != nil { + return result.Error + } + if result.RowsAffected == 0 { + return stacktrace.PropagateWithCode( + gorm.ErrRecordNotFound, + ErrCodeNotFound, + fmt.Sprintf("thread with id [%s] not found", params.MessageThreadID), + ) + } + + if !params.MarksUnread { + return nil + } + + return tx.Model(&entities.MessageThread{}). + Where("user_id = ?", params.UserID). + Where("id = ?", params.MessageThreadID). + Where("last_read_at < ?", params.EventTimestamp). + Update("is_read", false). + Error + }) + if err != nil { + msg := fmt.Sprintf("cannot update message activity for thread [%s]", params.MessageThreadID) + return repository.tracer.WrapErrorSpan(span, stacktrace.PropagateWithCode(err, stacktrace.GetCode(err), msg)) + } + return nil +} + +func (repository *gormMessageThreadRepository) UpdateStatus( + ctx context.Context, + userID entities.UserID, + messageThreadID uuid.UUID, + params MessageThreadStatusUpdate, +) error { + ctx, span := repository.tracer.Start(ctx) + defer span.End() + + result := repository.db.WithContext(ctx). + Model(&entities.MessageThread{}). + Where("user_id = ?", userID). + Where("id = ?", messageThreadID). + Updates(messageThreadStatusUpdates(params)) + if result.Error != nil { + msg := fmt.Sprintf("cannot update status for thread [%s]", messageThreadID) + return repository.tracer.WrapErrorSpan(span, stacktrace.Propagate(result.Error, msg)) + } + if result.RowsAffected == 0 { + msg := fmt.Sprintf("thread with id [%s] not found", messageThreadID) + return repository.tracer.WrapErrorSpan(span, stacktrace.PropagateWithCode(gorm.ErrRecordNotFound, ErrCodeNotFound, msg)) + } + return nil +} + +func (repository *gormMessageThreadRepository) UpdateAfterDeletedMessage(ctx context.Context, params MessageThreadDeletedUpdate) error { + ctx, span := repository.tracer.Start(ctx) + defer span.End() + + result := repository.db.WithContext(ctx). + Model(&entities.MessageThread{}). + Where("user_id = ?", params.UserID). + Where("id = ?", params.MessageThreadID). + Updates(map[string]any{ + "last_message_id": params.LastMessageID, + "last_message_content": params.LastMessageContent, + "status": params.LastMessageStatus, + }) + if result.Error != nil { + msg := fmt.Sprintf("cannot update deleted-message metadata for thread [%s]", params.MessageThreadID) + return repository.tracer.WrapErrorSpan(span, stacktrace.Propagate(result.Error, msg)) + } + return nil +} +``` + +Delete the old `Save(thread)` method and the old unused `UpdateAfterDeletedMessage(userID, messageID)` implementation. + +- [ ] **Step 5: Run repository tests** + +Run: + +```powershell +go test ./pkg/repositories -run MessageThread -v +``` + +Expected: PASS. + +- [ ] **Step 6: Commit atomic persistence** + +```powershell +git add api/pkg/repositories/message_thread_repository.go api/pkg/repositories/gorm_message_thread_repository.go api/pkg/repositories/gorm_message_thread_repository_test.go +@' +refactor(api): make thread updates atomic + +Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> +Copilot-Session: bf00a0ac-e11f-4015-b295-3cdd9b491229 +'@ | git commit -F - +``` + +### Task 3: Apply read rules in services and listeners + +**Files:** + +- Modify: `api/pkg/services/message_thread_service.go` +- Create: `api/pkg/services/message_thread_service_test.go` +- Modify: `api/pkg/listeners/message_thread_listener.go` +- Create: `api/pkg/listeners/read_receipts_test_helpers_test.go` +- Create: `api/pkg/listeners/message_thread_listener_test.go` +- Modify: `api/pkg/listeners/websocket_listener.go` +- Create: `api/pkg/listeners/websocket_listener_test.go` + +**Interfaces:** + +- Consumes: repository update types and methods from Task 2. +- Produces: `MessageThreadUpdateParams.MarksUnread bool` +- Produces: `MessageThreadUpdateParams.EventTimestamp time.Time` +- Produces: missed-call thread and websocket listeners. + +- [ ] **Step 1: Write failing service behavior tests** + +Create `api/pkg/services/message_thread_service_test.go`. Implement a repository stub with all interface methods returning zero values by default and function hooks for the methods under test: + +```go +package services + +import ( + "context" + "testing" + "time" + + "github.com/NdoleStudio/httpsms/pkg/entities" + "github.com/NdoleStudio/httpsms/pkg/repositories" + "github.com/NdoleStudio/httpsms/pkg/telemetry" + "github.com/google/uuid" + "github.com/palantir/stacktrace" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "gorm.io/gorm" +) + +type messageThreadRepositoryStub struct { + loadByOwnerContact func(context.Context, entities.UserID, string, string) (*entities.MessageThread, error) + load func(context.Context, entities.UserID, uuid.UUID) (*entities.MessageThread, error) + store func(context.Context, *entities.MessageThread) error + updateActivity func(context.Context, repositories.MessageThreadActivityUpdate) error + updateStatus func(context.Context, entities.UserID, uuid.UUID, repositories.MessageThreadStatusUpdate) error +} + +func (stub *messageThreadRepositoryStub) Store(ctx context.Context, thread *entities.MessageThread) error { + if stub.store != nil { + return stub.store(ctx, thread) + } + return nil +} + +func (stub *messageThreadRepositoryStub) UpdateActivity(ctx context.Context, params repositories.MessageThreadActivityUpdate) error { + if stub.updateActivity != nil { + return stub.updateActivity(ctx, params) + } + return nil +} + +func (stub *messageThreadRepositoryStub) UpdateStatus(ctx context.Context, userID entities.UserID, threadID uuid.UUID, params repositories.MessageThreadStatusUpdate) error { + if stub.updateStatus != nil { + return stub.updateStatus(ctx, userID, threadID, params) + } + return nil +} + +func (stub *messageThreadRepositoryStub) UpdateAfterDeletedMessage(context.Context, repositories.MessageThreadDeletedUpdate) error { + return nil +} + +func (stub *messageThreadRepositoryStub) LoadByOwnerContact(ctx context.Context, userID entities.UserID, owner string, contact string) (*entities.MessageThread, error) { + return stub.loadByOwnerContact(ctx, userID, owner, contact) +} + +func (stub *messageThreadRepositoryStub) Load(ctx context.Context, userID entities.UserID, id uuid.UUID) (*entities.MessageThread, error) { + return stub.load(ctx, userID, id) +} + +func (stub *messageThreadRepositoryStub) Index(context.Context, entities.UserID, string, bool, repositories.IndexParams) (*[]entities.MessageThread, error) { + threads := []entities.MessageThread{} + return &threads, nil +} + +func (stub *messageThreadRepositoryStub) Delete(context.Context, entities.UserID, uuid.UUID) error { + return nil +} + +func (stub *messageThreadRepositoryStub) DeleteAllForUser(context.Context, entities.UserID) error { + return nil +} + +func newMessageThreadServiceForTest(repository repositories.MessageThreadRepository) *MessageThreadService { + logger := &noopLogger{} + tracer := telemetry.NewOtelLogger("test", logger) + return NewMessageThreadService(logger, tracer, repository, nil) +} +``` + +Add these tests below the stub: + +```go +func TestUpdateThreadPassesUnreadWatermarkForInboundActivity(t *testing.T) { + threadID := uuid.New() + eventTimestamp := time.Date(2026, 7, 18, 7, 0, 0, 0, time.UTC) + var captured repositories.MessageThreadActivityUpdate + repository := &messageThreadRepositoryStub{ + loadByOwnerContact: func(context.Context, entities.UserID, string, string) (*entities.MessageThread, error) { + return &entities.MessageThread{ID: threadID}, nil + }, + updateActivity: func(_ context.Context, params repositories.MessageThreadActivityUpdate) error { + captured = params + return nil + }, + } + + service := newMessageThreadServiceForTest(repository) + err := service.UpdateThread(context.Background(), MessageThreadUpdateParams{ + UserID: entities.UserID("user-id"), + Owner: "+18005550199", + Contact: "+18005550100", + MessageID: uuid.New(), + Content: "hello", + Status: entities.MessageStatusReceived, + Timestamp: eventTimestamp, + MarksUnread: true, + EventTimestamp: eventTimestamp, + }) + + require.NoError(t, err) + assert.True(t, captured.MarksUnread) + assert.Equal(t, eventTimestamp, captured.EventTimestamp) +} + +func TestUpdateThreadPreservesReadStateForOutboundActivity(t *testing.T) { + var captured repositories.MessageThreadActivityUpdate + repository := &messageThreadRepositoryStub{ + loadByOwnerContact: func(context.Context, entities.UserID, string, string) (*entities.MessageThread, error) { + return &entities.MessageThread{ID: uuid.New(), IsRead: false}, nil + }, + updateActivity: func(_ context.Context, params repositories.MessageThreadActivityUpdate) error { + captured = params + return nil + }, + } + + service := newMessageThreadServiceForTest(repository) + err := service.UpdateThread(context.Background(), MessageThreadUpdateParams{ + UserID: entities.UserID("user-id"), + Owner: "+18005550199", + Contact: "+18005550100", + MessageID: uuid.New(), + Content: "outbound", + Status: entities.MessageStatusSent, + Timestamp: time.Now().UTC(), + }) + + require.NoError(t, err) + assert.False(t, captured.MarksUnread) +} + +func TestCreateThreadSetsReadStateFromActivityDirection(t *testing.T) { + tests := []struct { + name string + marksUnread bool + wantRead bool + }{ + {name: "inbound", marksUnread: true, wantRead: false}, + {name: "outbound", marksUnread: false, wantRead: true}, + } + + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + var stored *entities.MessageThread + repository := &messageThreadRepositoryStub{ + loadByOwnerContact: func(context.Context, entities.UserID, string, string) (*entities.MessageThread, error) { + return nil, stacktrace.PropagateWithCode(gorm.ErrRecordNotFound, repositories.ErrCodeNotFound, "not found") + }, + store: func(_ context.Context, thread *entities.MessageThread) error { + stored = thread + return nil + }, + } + + service := newMessageThreadServiceForTest(repository) + err := service.UpdateThread(context.Background(), MessageThreadUpdateParams{ + UserID: entities.UserID("user-id"), + Owner: "+18005550199", + Contact: "+18005550100", + MessageID: uuid.New(), + Content: "hello", + Status: entities.MessageStatusReceived, + Timestamp: time.Now().UTC(), + MarksUnread: test.marksUnread, + }) + + require.NoError(t, err) + require.NotNil(t, stored) + assert.Equal(t, test.wantRead, stored.IsRead) + assert.False(t, stored.LastReadAt.IsZero()) + }) + } +} + +func TestUpdateStatusChangesOnlyRequestedState(t *testing.T) { + threadID := uuid.New() + isRead := true + var captured repositories.MessageThreadStatusUpdate + repository := &messageThreadRepositoryStub{ + updateStatus: func(_ context.Context, _ entities.UserID, _ uuid.UUID, params repositories.MessageThreadStatusUpdate) error { + captured = params + return nil + }, + load: func(context.Context, entities.UserID, uuid.UUID) (*entities.MessageThread, error) { + return &entities.MessageThread{ID: threadID, IsArchived: true, IsRead: true}, nil + }, + } + + service := newMessageThreadServiceForTest(repository) + thread, err := service.UpdateStatus(context.Background(), MessageThreadStatusParams{ + UserID: entities.UserID("user-id"), + MessageThreadID: threadID, + IsRead: &isRead, + }) + + require.NoError(t, err) + assert.Nil(t, captured.IsArchived) + assert.Same(t, &isRead, captured.IsRead) + assert.False(t, captured.ReadAt.IsZero()) + assert.True(t, thread.IsArchived) +} + +func TestUpdateStatusPreservesNotFoundCode(t *testing.T) { + repository := &messageThreadRepositoryStub{ + updateStatus: func(context.Context, entities.UserID, uuid.UUID, repositories.MessageThreadStatusUpdate) error { + return stacktrace.PropagateWithCode(gorm.ErrRecordNotFound, repositories.ErrCodeNotFound, "not found") + }, + } + + service := newMessageThreadServiceForTest(repository) + isRead := true + _, err := service.UpdateStatus(context.Background(), MessageThreadStatusParams{ + UserID: entities.UserID("user-id"), + MessageThreadID: uuid.New(), + IsRead: &isRead, + }) + + assert.Equal(t, repositories.ErrCodeNotFound, stacktrace.GetCode(err)) +} +``` + +- [ ] **Step 2: Run service tests and confirm they fail** + +Run: + +```powershell +go test ./pkg/services -run 'Test(UpdateThread|CreateThread|UpdateStatus)' -v +``` + +Expected: compile failures because the new update parameters and repository calls are not wired. + +- [ ] **Step 3: Implement service read rules** + +Extend `MessageThreadUpdateParams`: + +```go +MarksUnread bool +EventTimestamp time.Time +``` + +Replace the existing full-row update in `UpdateThread` with: + +```go +if err = service.repository.UpdateActivity(ctx, repositories.MessageThreadActivityUpdate{ + MessageThreadID: thread.ID, + UserID: params.UserID, + Timestamp: params.Timestamp, + MessageID: params.MessageID, + Content: params.Content, + Status: params.Status, + MarksUnread: params.MarksUnread, + EventTimestamp: params.EventTimestamp, +}); err != nil { + msg := fmt.Sprintf("cannot update message thread with id [%s] after adding message [%s]", thread.ID, params.MessageID) + return service.tracer.WrapErrorSpan(span, stacktrace.PropagateWithCode(err, stacktrace.GetCode(err), msg)) +} +``` + +In `createThread`, create one `now := time.Now().UTC()` and initialize: + +```go +IsRead: !params.MarksUnread, +LastReadAt: now, +CreatedAt: now, +UpdatedAt: now, +``` + +Replace `UpdateStatus` with a partial repository update followed by a reload: + +```go +func (service *MessageThreadService) UpdateStatus(ctx context.Context, params MessageThreadStatusParams) (*entities.MessageThread, error) { + ctx, span := service.tracer.Start(ctx) + defer span.End() + + update := repositories.MessageThreadStatusUpdate{ + IsArchived: params.IsArchived, + IsRead: params.IsRead, + ReadAt: time.Now().UTC(), + } + if err := service.repository.UpdateStatus(ctx, params.UserID, params.MessageThreadID, update); err != nil { + msg := fmt.Sprintf("cannot update message thread with id [%s]", params.MessageThreadID) + return nil, service.tracer.WrapErrorSpan(span, stacktrace.PropagateWithCode(err, stacktrace.GetCode(err), msg)) + } + + thread, err := service.repository.Load(ctx, params.UserID, params.MessageThreadID) + if err != nil { + msg := fmt.Sprintf("cannot reload message thread with id [%s]", params.MessageThreadID) + return nil, service.tracer.WrapErrorSpan(span, stacktrace.PropagateWithCode(err, stacktrace.GetCode(err), msg)) + } + return thread, nil +} +``` + +In `UpdateAfterDeletedMessage`, replace the full-row update with: + +```go +if err = service.repository.UpdateAfterDeletedMessage(ctx, repositories.MessageThreadDeletedUpdate{ + MessageThreadID: thread.ID, + UserID: thread.UserID, + LastMessageID: payload.PreviousMessageID, + LastMessageContent: payload.PreviousMessageContent, + LastMessageStatus: *payload.PreviousMessageStatus, +}); err != nil { + msg := fmt.Sprintf("cannot update thread with ID [%s] for user with ID [%s]", thread.ID, thread.UserID) + return service.tracer.WrapErrorSpan(span, stacktrace.Propagate(err, msg)) +} +``` + +- [ ] **Step 4: Add inbound and missed-call listener behavior** + +In `NewMessageThreadListener`, register: + +```go +events.MessageCallMissed: l.OnMessageCallMissed, +``` + +In `OnMessagePhoneReceived`, add: + +```go +MarksUnread: true, +EventTimestamp: event.Time(), +``` + +Add: + +```go +func (listener *MessageThreadListener) OnMessageCallMissed(ctx context.Context, event cloudevents.Event) error { + ctx, span := listener.tracer.Start(ctx) + defer span.End() + + var payload events.MessageCallMissedPayload + if err := event.DataAs(&payload); err != nil { + msg := fmt.Sprintf("cannot decode [%s] into [%T]", event.Data(), payload) + return listener.tracer.WrapErrorSpan(span, stacktrace.Propagate(err, msg)) + } + + params := services.MessageThreadUpdateParams{ + Owner: payload.Owner, + Contact: payload.Contact, + UserID: payload.UserID, + Status: entities.MessageStatusReceived, + Timestamp: payload.Timestamp, + Content: "Missed phone call", + MessageID: payload.MessageID, + MarksUnread: true, + EventTimestamp: event.Time(), + } + if err := listener.service.UpdateThread(ctx, params); err != nil { + msg := fmt.Sprintf("cannot update thread for missed call [%s] on event [%s]", payload.MessageID, event.ID()) + return listener.tracer.WrapErrorSpan(span, stacktrace.Propagate(err, msg)) + } + return nil +} +``` + +In `NewWebsocketListener`, register: + +```go +events.MessageCallMissed: l.onMessageCallMissed, +``` + +Add: + +```go +func (listener *WebsocketListener) onMessageCallMissed(ctx context.Context, event cloudevents.Event) error { + ctx, span, _ := listener.tracer.StartWithLogger(ctx, listener.logger) + defer span.End() + + var payload events.MessageCallMissedPayload + if err := event.DataAs(&payload); err != nil { + msg := fmt.Sprintf("cannot decode [%s] into [%T]", event.Data(), payload) + return listener.tracer.WrapErrorSpan(span, stacktrace.Propagate(err, msg)) + } + + if err := listener.client.Trigger(payload.UserID.String(), event.Type(), event.ID()); err != nil { + msg := fmt.Sprintf("cannot trigger websocket [%s] event with ID [%s] for user with ID [%s]", event.Type(), event.ID(), payload.UserID) + return listener.tracer.WrapErrorSpan(span, stacktrace.Propagate(err, msg)) + } + return nil +} +``` + +- [ ] **Step 5: Add listener tests** + +Create `api/pkg/listeners/read_receipts_test_helpers_test.go`: + +```go +package listeners + +import ( + "context" + + "github.com/NdoleStudio/httpsms/pkg/entities" + "github.com/NdoleStudio/httpsms/pkg/repositories" + "github.com/NdoleStudio/httpsms/pkg/telemetry" + "github.com/google/uuid" + "go.opentelemetry.io/otel/trace" +) + +type noopListenerLogger struct{} + +func (logger *noopListenerLogger) Error(error) {} +func (logger *noopListenerLogger) WithService(string) telemetry.Logger { return logger } +func (logger *noopListenerLogger) WithString(string, string) telemetry.Logger { return logger } +func (logger *noopListenerLogger) WithSpan(trace.SpanContext) telemetry.Logger { return logger } +func (logger *noopListenerLogger) Trace(string) {} +func (logger *noopListenerLogger) Info(string) {} +func (logger *noopListenerLogger) Warn(error) {} +func (logger *noopListenerLogger) Debug(string) {} +func (logger *noopListenerLogger) Fatal(error) {} +func (logger *noopListenerLogger) Printf(string, ...interface{}) {} + +type listenerMessageThreadRepository struct { + activity repositories.MessageThreadActivityUpdate +} + +func (repository *listenerMessageThreadRepository) Store(context.Context, *entities.MessageThread) error { + return nil +} + +func (repository *listenerMessageThreadRepository) UpdateActivity(_ context.Context, params repositories.MessageThreadActivityUpdate) error { + repository.activity = params + return nil +} + +func (repository *listenerMessageThreadRepository) UpdateStatus(context.Context, entities.UserID, uuid.UUID, repositories.MessageThreadStatusUpdate) error { + return nil +} + +func (repository *listenerMessageThreadRepository) UpdateAfterDeletedMessage(context.Context, repositories.MessageThreadDeletedUpdate) error { + return nil +} + +func (repository *listenerMessageThreadRepository) LoadByOwnerContact(context.Context, entities.UserID, string, string) (*entities.MessageThread, error) { + return &entities.MessageThread{ID: uuid.New()}, nil +} + +func (repository *listenerMessageThreadRepository) Load(context.Context, entities.UserID, uuid.UUID) (*entities.MessageThread, error) { + return &entities.MessageThread{}, nil +} + +func (repository *listenerMessageThreadRepository) Index(context.Context, entities.UserID, string, bool, repositories.IndexParams) (*[]entities.MessageThread, error) { + threads := []entities.MessageThread{} + return &threads, nil +} + +func (repository *listenerMessageThreadRepository) Delete(context.Context, entities.UserID, uuid.UUID) error { + return nil +} + +func (repository *listenerMessageThreadRepository) DeleteAllForUser(context.Context, entities.UserID) error { + return nil +} +``` + +Create `api/pkg/listeners/message_thread_listener_test.go`: + +```go +package listeners + +import ( + "context" + "testing" + "time" + + "github.com/NdoleStudio/httpsms/pkg/entities" + "github.com/NdoleStudio/httpsms/pkg/events" + "github.com/NdoleStudio/httpsms/pkg/services" + "github.com/NdoleStudio/httpsms/pkg/telemetry" + cloudevents "github.com/cloudevents/sdk-go/v2" + "github.com/google/uuid" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func TestMessageThreadListenerMarksMissedCallUnread(t *testing.T) { + repository := &listenerMessageThreadRepository{} + logger := &noopListenerLogger{} + tracer := telemetry.NewOtelLogger("test", logger) + service := services.NewMessageThreadService(logger, tracer, repository, nil) + _, routes := NewMessageThreadListener(logger, tracer, service) + + event := cloudevents.NewEvent() + event.SetID(uuid.NewString()) + event.SetSource("/v1/messages/call-missed") + event.SetType(events.MessageCallMissed) + event.SetTime(time.Date(2026, 7, 18, 7, 0, 0, 0, time.UTC)) + require.NoError(t, event.SetData(cloudevents.ApplicationJSON, events.MessageCallMissedPayload{ + MessageID: uuid.New(), + UserID: entities.UserID("user-id"), + Owner: "+18005550199", + Contact: "+18005550100", + Timestamp: time.Date(2026, 7, 18, 6, 59, 0, 0, time.UTC), + })) + + err := routes[events.MessageCallMissed](context.Background(), event) + + require.NoError(t, err) + assert.True(t, repository.activity.MarksUnread) + assert.Equal(t, "Missed phone call", repository.activity.Content) + assert.Equal(t, event.Time(), repository.activity.EventTimestamp) +} +``` + +Create `api/pkg/listeners/websocket_listener_test.go`: + +```go +package listeners + +import ( + "testing" + + "github.com/NdoleStudio/httpsms/pkg/events" + "github.com/NdoleStudio/httpsms/pkg/telemetry" + "github.com/pusher/pusher-http-go/v5" + "github.com/stretchr/testify/assert" +) + +func TestWebsocketListenerRegistersMissedCalls(t *testing.T) { + logger := &noopListenerLogger{} + tracer := telemetry.NewOtelLogger("test", logger) + _, routes := NewWebsocketListener(logger, tracer, &pusher.Client{}) + + assert.Contains(t, routes, events.MessageCallMissed) +} +``` + +- [ ] **Step 6: Run service and listener tests** + +Run: + +```powershell +go test ./pkg/services ./pkg/listeners -run 'MessageThread|MissedCall|UpdateStatus|CreateThread' -v +``` + +Expected: PASS. + +- [ ] **Step 7: Commit service and realtime behavior** + +```powershell +git add api/pkg/services/message_thread_service.go api/pkg/services/message_thread_service_test.go api/pkg/listeners/message_thread_listener.go api/pkg/listeners/read_receipts_test_helpers_test.go api/pkg/listeners/message_thread_listener_test.go api/pkg/listeners/websocket_listener.go api/pkg/listeners/websocket_listener_test.go +@' +feat(api): track unread inbound threads + +Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> +Copilot-Session: bf00a0ac-e11f-4015-b295-3cdd9b491229 +'@ | git commit -F - +``` + +### Task 4: Finish the HTTP behavior and regenerate Swagger + +**Files:** + +- Modify: `api/pkg/handlers/message_thread_handler.go` +- Modify: `api/docs/docs.go` +- Modify: `api/docs/swagger.json` +- Modify: `api/docs/swagger.yaml` + +**Interfaces:** + +- Consumes: `responses.MessageThreadResponse` from Task 1. +- Produces: 404 behavior for missing thread updates. +- Produces: Swagger schema fields `is_read`, optional `is_archived`, optional `is_read`. + +- [ ] **Step 1: Update the handler error mapping and annotations** + +Change the update annotation: + +```go +// @Success 200 {object} responses.MessageThreadResponse +// @Failure 404 {object} responses.NotFound +``` + +Change the error block: + +```go +thread, err := h.service.UpdateStatus(ctx, request.ToUpdateParams(h.userIDFomContext(c))) +if stacktrace.GetCode(err) == repositories.ErrCodeNotFound { + return h.responseNotFound(c, fmt.Sprintf("cannot find message thread with ID [%s]", request.MessageThreadID)) +} +if err != nil { + msg := fmt.Sprintf("cannot update message thread with params [%+#v]", request) + ctxLogger.Error(stacktrace.Propagate(err, msg)) + return h.responseInternalServerError(c) +} +``` + +- [ ] **Step 2: Regenerate Swagger** + +Run: + +```powershell +Set-Location 'C:\Users\achoa\Work\NdoleStudio\httpsms-read-receipts\api' +swag init --requiredByDefault --parseDependency --parseInternal +``` + +Expected: `api/docs/docs.go`, `api/docs/swagger.json`, and `api/docs/swagger.yaml` update successfully. + +- [ ] **Step 3: Run the full API suite** + +Run: + +```powershell +go test -vet=off ./pkg/entities ./pkg/repositories ./pkg/requests ./pkg/services ./pkg/validators ./pkg/listeners +go test -vet=off ./pkg/handlers -run '^$' +``` + +Expected: PASS. Vet is disabled because clean `origin/main` has unrelated Go +1.25 non-constant-format vet failures. Handler tests are compile-checked only +because the pre-existing suite expects an API server on `localhost:8000`; Task +7 provides the HTTP integration coverage against the Docker stack. + +- [ ] **Step 4: Commit HTTP and documentation changes** + +```powershell +git add api/pkg/handlers/message_thread_handler.go api/docs/docs.go api/docs/swagger.json api/docs/swagger.yaml +@' +docs(api): publish thread read state + +Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> +Copilot-Session: bf00a0ac-e11f-4015-b295-3cdd9b491229 +'@ | git commit -F - +``` + +### Task 5: Add automatic read updates to the web store + +**Files:** + +- Modify: `web/shared/types/api.ts` +- Modify: `web/app/stores/threads.ts` + +**Interfaces:** + +- Consumes: API `PUT /v1/message-threads/{id}` with `{ is_read: true }`. +- Produces: `markThreadRead(threadId: string, force?: boolean): Promise`. + +- [ ] **Step 1: Regenerate TypeScript API models** + +Run: + +```powershell +Set-Location 'C:\Users\achoa\Work\NdoleStudio\httpsms-read-receipts\web' +pnpm api:models +``` + +Expected: `EntitiesMessageThread.is_read: boolean` and optional fields on `RequestsMessageThreadUpdate`. + +- [ ] **Step 2: Add a typed local replacement helper** + +In `web/app/stores/threads.ts`, add: + +```ts +function replaceThread(updatedThread: EntitiesMessageThread) { + const index = threads.value.findIndex( + (thread) => thread.id === updatedThread.id, + ); + if (index !== -1) threads.value[index] = updatedThread; +} +``` + +- [ ] **Step 3: Implement non-silent automatic read persistence** + +Add: + +```ts +async function markThreadRead(threadId: string, force = false) { + const thread = threads.value.find((item) => item.id === threadId); + if (!thread) throw new Error(`Cannot find thread with id ${threadId}`); + if (!force && thread.is_read) return; + + try { + const response = await apiFetch<{ data: EntitiesMessageThread }>( + `/v1/message-threads/${threadId}`, + { + method: "PUT", + body: { is_read: true }, + }, + ); + replaceThread(response.data); + } catch (error) { + notificationsStore.addNotification({ + message: "The message thread could not be marked as read", + type: "error", + }); + await loadThreads(); + throw error; + } +} +``` + +Add `markThreadRead` to the returned store API. + +Keep archive updates independent: + +```ts +body: { is_archived: payload.isArchived }, +``` + +Do not send `is_read` from `updateThread`. + +- [ ] **Step 4: Run web static checks** + +Run: + +```powershell +pnpm lint:js +pnpm lint:prettier +``` + +Expected: PASS. + +- [ ] **Step 5: Commit store and generated types** + +```powershell +git add web/shared/types/api.ts web/app/stores/threads.ts +@' +feat(web): persist opened threads as read + +Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> +Copilot-Session: bf00a0ac-e11f-4015-b295-3cdd9b491229 +'@ | git commit -F - +``` + +### Task 6: Mark opened threads read and highlight unread threads + +**Files:** + +- Modify: `web/app/pages/threads/[id]/index.vue` +- Modify: `web/app/components/MessageThread.vue` + +**Interfaces:** + +- Consumes: `threadsStore.markThreadRead(threadId, force)`. +- Produces: automatic initial and realtime read updates. +- Produces: shared unread list styling on mobile and desktop. + +- [ ] **Step 1: Make read updates independent from message loading** + +In `web/app/pages/threads/[id]/index.vue`, add: + +```ts +async function markCurrentThreadRead(force = false) { + const threadId = route.params.id as string; + try { + await threadsStore.markThreadRead(threadId, force); + } catch (error) { + console.error(error); + } +} +``` + +At the start of `loadMessages`, after computing `threadId`, call without awaiting: + +```ts +void markCurrentThreadRead(); +``` + +This lets message loading continue even if the read update fails. + +- [ ] **Step 2: Force a newer read watermark after inbound realtime events** + +Change the incoming SMS binding: + +```ts +webhookChannel.bind("message.phone.received", () => { + if (!loadingMessages.value) { + void markCurrentThreadRead(true); + loadMessages(false); + } +}); +``` + +Add the missed-call binding with the same behavior: + +```ts +webhookChannel.bind("message.call.missed", () => { + if (!loadingMessages.value) { + void markCurrentThreadRead(true); + loadMessages(false); + } +}); +``` + +The forced call is required because the local thread may still say `is_read: +true` when the concurrent backend listener has not committed its unread update. + +- [ ] **Step 3: Add unread presentation** + +In `MessageThread.vue`, add the class binding to `v-list-item`: + +```vue +:class="{ 'message-thread--unread': !thread.is_read }" +``` + +Change title and subtitle class bindings: + +```vue + + {{ formatPhoneNumber(thread.contact) }} + + +``` + +Add: + +```vue + +``` + +- [ ] **Step 4: Run web validation** + +Run: + +```powershell +Set-Location 'C:\Users\achoa\Work\NdoleStudio\httpsms-read-receipts\web' +pnpm lint +pnpm run generate +``` + +Expected: both commands PASS. + +- [ ] **Step 5: Commit UI behavior** + +```powershell +git add web/app/pages/threads/[id]/index.vue web/app/components/MessageThread.vue +@' +feat(web): highlight unread message threads + +Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> +Copilot-Session: bf00a0ac-e11f-4015-b295-3cdd9b491229 +'@ | git commit -F - +``` + +### Task 7: Add integration coverage + +**Files:** + +- Create: `tests/read_receipts_test.go` +- Modify: `tests/README.md` + +**Interfaces:** + +- Consumes: thread index and update HTTP endpoints from Tasks 1-4. +- Consumes: incoming SMS, missed-call, and outbound message flows from Task 3. +- Produces: `TestMessageThreadReadReceipts`. + +- [ ] **Step 1: Write the end-to-end read-receipts test** + +Create `tests/read_receipts_test.go`: + +```go +package tests + +import ( + "bytes" + "context" + "encoding/json" + "fmt" + "io" + "net/http" + "net/url" + "testing" + "time" + + httpsms "github.com/NdoleStudio/httpsms-go" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +type integrationMessageThread struct { + ID string `json:"id"` + Contact string `json:"contact"` + IsRead bool `json:"is_read"` + LastMessageContent *string `json:"last_message_content"` +} + +func requestJSON( + ctx context.Context, + t *testing.T, + method string, + path string, + apiKey string, + payload any, + expectedStatus int, + output any, +) { + t.Helper() + + var body io.Reader + if payload != nil { + encoded, err := json.Marshal(payload) + require.NoError(t, err) + body = bytes.NewReader(encoded) + } + + request, err := http.NewRequestWithContext(ctx, method, apiBaseURL+path, body) + require.NoError(t, err) + request.Header.Set("x-api-key", apiKey) + request.Header.Set("Content-Type", "application/json") + + response, err := http.DefaultClient.Do(request) + require.NoError(t, err) + defer response.Body.Close() + + responseBody, err := io.ReadAll(response.Body) + require.NoError(t, err) + require.Equal(t, expectedStatus, response.StatusCode, "response: %s", string(responseBody)) + + if output != nil { + require.NoError(t, json.Unmarshal(responseBody, output)) + } +} + +func fetchMessageThreads(ctx context.Context, t *testing.T, owner string) []integrationMessageThread { + t.Helper() + + var response struct { + Data []integrationMessageThread `json:"data"` + } + path := fmt.Sprintf( + "/v1/message-threads?owner=%s&skip=0&limit=20&is_archived=false", + url.QueryEscape(owner), + ) + requestJSON(ctx, t, http.MethodGet, path, userAPIKey, nil, http.StatusOK, &response) + return response.Data +} + +func waitForMessageThread( + ctx context.Context, + t *testing.T, + owner string, + contact string, + timeout time.Duration, + matches func(integrationMessageThread) bool, +) integrationMessageThread { + t.Helper() + + deadline := time.Now().Add(timeout) + for time.Now().Before(deadline) { + for _, thread := range fetchMessageThreads(ctx, t, owner) { + if thread.Contact == contact && matches(thread) { + return thread + } + } + time.Sleep(500 * time.Millisecond) + } + + t.Fatalf("thread %s -> %s did not reach the expected state within %v", owner, contact, timeout) + return integrationMessageThread{} +} + +func markMessageThreadRead(ctx context.Context, t *testing.T, threadID string) integrationMessageThread { + t.Helper() + + var response struct { + Data integrationMessageThread `json:"data"` + } + requestJSON( + ctx, + t, + http.MethodPut, + "/v1/message-threads/"+threadID, + userAPIKey, + map[string]any{"is_read": true}, + http.StatusOK, + &response, + ) + return response.Data +} + +func TestMessageThreadReadReceipts(t *testing.T) { + ctx := context.Background() + phone := setupPhone(ctx, t, 60) + contact := randomPhoneNumber() + + requestJSON( + ctx, + t, + http.MethodPost, + "/v1/messages/receive", + phone.PhoneAPIKey, + map[string]any{ + "from": contact, + "to": phone.PhoneNumber, + "content": "Unread inbound message", + "encrypted": false, + "sim": "SIM1", + "timestamp": time.Now().UTC().Format(time.RFC3339), + }, + http.StatusOK, + nil, + ) + + thread := waitForMessageThread(ctx, t, phone.PhoneNumber, contact, 20*time.Second, func(thread integrationMessageThread) bool { + return !thread.IsRead + }) + assert.False(t, thread.IsRead) + + updated := markMessageThreadRead(ctx, t, thread.ID) + assert.True(t, updated.IsRead) + waitForMessageThread(ctx, t, phone.PhoneNumber, contact, 10*time.Second, func(thread integrationMessageThread) bool { + return thread.IsRead + }) + + requestJSON( + ctx, + t, + http.MethodPost, + "/v1/messages/calls/missed", + phone.PhoneAPIKey, + map[string]any{ + "from": contact, + "to": phone.PhoneNumber, + "sim": "SIM1", + "timestamp": time.Now().UTC().Format(time.RFC3339), + }, + http.StatusOK, + nil, + ) + + thread = waitForMessageThread(ctx, t, phone.PhoneNumber, contact, 20*time.Second, func(thread integrationMessageThread) bool { + return !thread.IsRead && + thread.LastMessageContent != nil && + *thread.LastMessageContent == "Missed phone call" + }) + assert.False(t, thread.IsRead) + + outboundContent := "Outbound activity preserves unread" + client := newAPIClient() + _, response, err := client.Messages.Send(ctx, &httpsms.MessageSendParams{ + From: phone.PhoneNumber, + To: contact, + Content: outboundContent, + }) + require.NoError(t, err) + require.Equal(t, http.StatusOK, response.HTTPResponse.StatusCode) + + thread = waitForMessageThread(ctx, t, phone.PhoneNumber, contact, 20*time.Second, func(thread integrationMessageThread) bool { + return thread.LastMessageContent != nil && + *thread.LastMessageContent == outboundContent + }) + assert.False(t, thread.IsRead, "outbound activity must not clear unread state") +} +``` + +- [ ] **Step 2: Run the test before implementation is complete** + +With the Docker integration stack running, execute: + +```powershell +Set-Location 'C:\Users\achoa\Work\NdoleStudio\httpsms-read-receipts\tests' +go test -v -timeout 120s -run TestMessageThreadReadReceipts ./... +``` + +Expected before the feature implementation: FAIL because thread responses do not expose/persist `is_read`. + +- [ ] **Step 3: Update integration-test documentation** + +Add to the Test Coverage checklist in `tests/README.md`: + +```markdown +- [x] **Message thread read receipts E2E** — Incoming SMS and missed calls mark a thread unread, the existing thread update endpoint marks it read, and outbound activity preserves unread state +``` + +- [ ] **Step 4: Run the integration test against the completed stack** + +Generate credentials, start the Docker stack, run the focused test, and tear down: + +```powershell +Set-Location 'C:\Users\achoa\Work\NdoleStudio\httpsms-read-receipts\tests' +& 'C:\Program Files\Git\bin\bash.exe' generate-firebase-credentials.sh +$env:FIREBASE_CREDENTIALS = (Get-Content firebase-credentials.json -Raw | ConvertFrom-Json | ConvertTo-Json -Compress) +docker compose up -d --build --wait +docker compose wait seed +Start-Sleep -Seconds 2 +try { + go test -v -timeout 120s -run TestMessageThreadReadReceipts ./... +} finally { + docker compose down -v +} +``` + +Expected: PASS. + +- [ ] **Step 5: Commit integration coverage** + +```powershell +git add tests/read_receipts_test.go tests/README.md +@' +test: cover message thread read receipts + +Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> +Copilot-Session: bf00a0ac-e11f-4015-b295-3cdd9b491229 +'@ | git commit -F - +``` + +### Task 8: Verify the complete feature + +**Files:** + +- Review: all files changed since `origin/main` + +**Interfaces:** + +- Consumes: all previous tasks. +- Produces: a clean, verified feature branch. + +- [ ] **Step 1: Run API formatting and tests** + +Run: + +```powershell +Set-Location 'C:\Users\achoa\Work\NdoleStudio\httpsms-read-receipts\api' +go-fumpt -w pkg/entities/message_thread.go pkg/entities/message_thread_test.go pkg/repositories/message_thread_repository.go pkg/repositories/gorm_message_thread_repository.go pkg/repositories/gorm_message_thread_repository_test.go pkg/services/message_thread_service.go pkg/services/message_thread_service_test.go pkg/requests/message_thread_update_request.go pkg/requests/message_thread_update_request_test.go pkg/validators/message_thread_handler_validator.go pkg/validators/message_thread_handler_validator_test.go pkg/responses/message_thead_responses.go pkg/handlers/message_thread_handler.go pkg/listeners/message_thread_listener.go pkg/listeners/read_receipts_test_helpers_test.go pkg/listeners/message_thread_listener_test.go pkg/listeners/websocket_listener.go pkg/listeners/websocket_listener_test.go +go test -vet=off ./pkg/entities ./pkg/repositories ./pkg/requests ./pkg/services ./pkg/validators ./pkg/listeners +go test -vet=off ./pkg/handlers -run '^$' +``` + +Expected: formatting completes and all API tests PASS. + +- [ ] **Step 2: Run complete web validation** + +Run: + +```powershell +Set-Location 'C:\Users\achoa\Work\NdoleStudio\httpsms-read-receipts\web' +pnpm lint +pnpm run generate +``` + +Expected: PASS. + +- [ ] **Step 3: Inspect the final diff** + +Run: + +```powershell +Set-Location 'C:\Users\achoa\Work\NdoleStudio\httpsms-read-receipts' +git --no-pager diff --check origin/main...HEAD +git --no-pager diff --stat origin/main...HEAD +git --no-pager status --short --branch +``` + +Expected: no whitespace errors and no uncommitted implementation changes. + +- [ ] **Step 4: Perform manual acceptance checks against a migrated local database** + +Verify: + +1. Existing rows return `"is_read": true`. +2. A closed thread becomes highlighted after an incoming SMS. +3. A closed thread becomes highlighted after a missed call. +4. Opening the thread removes the highlight and refresh keeps it read. +5. Incoming activity while the thread is open does not leave it unread. +6. Sending and delivery-status changes do not clear an unread thread. +7. Archiving and unarchiving preserve read state. + +- [ ] **Step 5: Commit any formatting-only corrections** + +Skip this step when Step 1 produced no changes. Otherwise: + +```powershell +git add api web +@' +style: format read receipts changes + +Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> +Copilot-Session: bf00a0ac-e11f-4015-b295-3cdd9b491229 +'@ | git commit -F - +``` diff --git a/docs/superpowers/specs/2026-07-18-read-receipts-design.md b/docs/superpowers/specs/2026-07-18-read-receipts-design.md new file mode 100644 index 000000000..6462cbaa5 --- /dev/null +++ b/docs/superpowers/specs/2026-07-18-read-receipts-design.md @@ -0,0 +1,286 @@ +# Message Thread Read Receipts + +- Date: 2026-07-18 +- Status: Approved (design) +- Scope: `api/` Go backend and `web/` Nuxt frontend. Android is unchanged. +- Branch: `feat/read-receipts`, based on `origin/main` + +## Problem + +The message-thread list does not distinguish conversations with unseen inbound +activity from conversations the user has already opened. Users need unread +threads to stand out and become read automatically when opened. + +The change must be backward compatible: every thread that exists when the new +schema is deployed starts as read. + +## Decisions + +- Unread state is stored per message thread. +- Existing threads are read after migration. +- A new incoming SMS or missed call marks its thread unread. +- Outbound messages and later delivery/status updates do not change read state. +- Opening a thread in the web UI marks it read automatically. +- A thread that receives inbound activity while already open remains read. +- There is no manual "Mark as read/unread" UI. +- Unread threads use bold contact/preview text and a subtle primary-tinted + background. +- The existing `PUT /v1/message-threads/{messageThreadID}` endpoint handles the + read update; no new endpoint is added. + +## Design + +### 1. Thread persistence + +Add these fields to `api/pkg/entities/message_thread.go`: + +```go +IsRead bool `json:"is_read" gorm:"not null;default:true" example:"true"` +LastReadAt time.Time `json:"-" gorm:"not null;default:CURRENT_TIMESTAMP"` +``` + +`IsRead` is the public binary state used by API clients and the web UI. +`LastReadAt` is internal ordering metadata used to make concurrent event +processing deterministic. + +GORM `AutoMigrate` adds both columns. The database defaults make all existing +rows read and give them a migration-time read watermark. New-thread creation +sets `IsRead` explicitly according to whether the triggering activity is +inbound. + +Existing-row updates must not use the repository's current full-row `Save` +method. A stale entity loaded by one event handler could otherwise overwrite a +newer read action from the web request. + +### 2. Atomic repository updates + +Add field-scoped repository operations for the two update paths: + +- **Message activity update:** update only `order_timestamp`, + `last_message_id`, `last_message_content`, and `status`. For inbound activity, + conditionally set `is_read = false` only where the stored `last_read_at` is + older than the CloudEvent creation time. +- **User status update:** update only the supplied `is_archived` and `is_read` + fields. Marking read writes `is_read = true` and `last_read_at = now` in the + same database update. + +The message activity update runs its metadata and conditional unread statements +inside one GORM transaction. Both operations remain scoped by `user_id` and +thread ID and use GORM's query builder with context propagation. + +After a user status update, reload the thread for the API response. Do not use +`Save` in either path, because it writes unrelated columns and reintroduces lost +updates. + +The stored timestamp condition prevents this race: + +1. An inbound event is created. +2. The websocket listener notifies the open web page. +3. The web page marks the thread read. +4. The message-thread listener finishes later. + +Because CloudEvent creation time precedes the read request, the late listener +checks the newer database watermark and does not overwrite the read action. If +the message listener commits first, the later read update wins normally. + +### 3. Thread update parameters + +Extend `services.MessageThreadUpdateParams` with: + +```go +MarksUnread bool +EventTimestamp time.Time +``` + +Inbound SMS and missed-call listeners set `MarksUnread = true` and use +`event.Time()` as `EventTimestamp`. Other message events leave `MarksUnread` +false, so sending, scheduling, delivery, failure, and expiry updates preserve +the current read state. + +For an existing thread, the service delegates the field-scoped message update +and optional conditional unread transition to the repository. + +For a newly created thread: + +- inbound SMS or missed call: `IsRead = false`; +- outbound activity: `IsRead = true`. + +### 4. Missed-call thread updates + +`MessageThreadListener` currently handles received SMS events but not +`message.call.missed`. Register the missed-call event and update/create the +corresponding thread with: + +- owner and contact from `MessageCallMissedPayload`; +- message ID and event timestamp; +- received status; +- a concise preview such as `Missed phone call`; +- `MarksUnread = true`. + +This ensures a missed call is visible as thread activity even when no auto-reply +is configured. + +### 5. Existing update endpoint + +Change `requests.MessageThreadUpdate` to use optional fields: + +```go +IsArchived *bool `json:"is_archived,omitempty" example:"true"` +IsRead *bool `json:"is_read,omitempty" example:"true"` +``` + +The validator requires at least one supported field. Pointer fields distinguish +"not supplied" from `false`, preventing a read-only update from unarchiving a +thread and an archive-only update from changing read state. + +`MessageThreadStatusParams` receives the optional fields. `UpdateStatus`: + +- loads the authenticated user's thread; +- builds an update map containing only supplied fields; +- writes `last_read_at = time.Now().UTC()` alongside `is_read = true`; +- supports `IsRead = false` at the API level without adding a web control; +- persists one combined partial update and reloads the updated thread. + +The handler preserves repository error codes so a missing thread returns 404. +Invalid IDs or empty/unsupported update bodies return validation errors. + +### 6. Web store behavior + +Regenerate the API model so `EntitiesMessageThread` includes `is_read` and +`RequestsMessageThreadUpdate` has optional `is_archived` and `is_read`. + +In `web/app/stores/threads.ts`: + +- keep archive updates on the existing endpoint with + `{ is_archived: value }`; +- add `markThreadRead(threadId)` using `{ is_read: true }`; +- replace the matching thread in local state with the successful API response; +- make the operation idempotent and skip the request when the local thread is + already read unless it is called after an inbound realtime event. + +Opening a thread invokes `markThreadRead` automatically as part of the thread +loading flow. Message fetching remains independent: if the read update fails, +messages still display, the thread remains/reloads as unread, and the existing +notification system reports the failure. + +### 7. Realtime behavior + +The message detail page already reloads messages for incoming SMS websocket +events. That reload also marks the currently selected thread read. + +Add `message.call.missed` to: + +- `WebsocketListener` registrations in the API; +- the Pusher bindings on the thread detail page. + +The websocket payload can remain the event ID. The selected thread is marked +read idempotently after any inbound realtime refresh; a different thread's +backend state remains unread. + +### 8. Thread-list presentation + +In `web/app/components/MessageThread.vue`, unread list items receive: + +- bold contact title; +- bold message preview; +- a subtle background using the Vuetify primary theme color at low opacity. + +Read threads retain the current appearance. Active-route styling remains +visible and takes precedence where Vuetify applies it. + +The highlight is applied in the shared thread-list component, so it works in +both the mobile thread page and desktop navigation drawer, including archived +threads. + +### 9. Swagger and generated types + +After changing the entity and request annotations: + +1. Run `swag init --requiredByDefault --parseDependency --parseInternal` in + `api/`. +2. Run `pnpm api:models` in `web/`. + +Generated Swagger files and `web/shared/types/api.ts` are committed with the +implementation. + +## Error Handling + +- Repository and service errors continue to use `stacktrace.Propagate`. +- The update handler returns 404 for a thread that does not belong to the + authenticated user. +- A failed automatic read update is not swallowed and does not block message + display. +- Local UI state changes only from a successful response or a subsequent thread + reload; the UI does not claim a persisted read state after an API failure. +- Realtime callbacks report failures through the existing notification/error + path rather than using empty catches. + +## Testing + +### API + +Add targeted tests for: + +- entity schema defaults and new-thread read-state initialization; +- an inbound SMS making an existing thread unread; +- a missed call creating or updating an unread thread; +- outbound and delivery/status events preserving read state; +- a read action winning when its timestamp is newer than an inbound event; +- concurrent message/read updates changing only their owned columns; +- archive-only updates preserving read state; +- read-only updates preserving archive state; +- missing threads preserving the not-found error code; +- existing-row migration defaults (`is_read = true`). + +Run: + +```bash +cd api +go test ./... +``` + +### Web + +No frontend unit-test runner is configured. Validate the generated types, +Pinia/Vue changes, styles, and production output with: + +```bash +cd web +pnpm lint +pnpm run generate +``` + +### Integration + +Add `tests/read_receipts_test.go` to exercise the feature against the Docker +integration stack: + +- an inbound SMS creates or updates an unread thread; +- marking the thread read through the existing update endpoint persists; +- a missed call makes the thread unread again; +- subsequent outbound activity does not clear that unread state. + +Update `tests/README.md` coverage and run: + +```bash +cd tests +go test -v -timeout 120s -run TestMessageThreadReadReceipts ./... +``` + +Manual acceptance checks: + +1. Existing threads appear read immediately after migration. +2. A new incoming SMS highlights its closed thread. +3. A missed call highlights its closed thread. +4. Opening an unread thread removes the highlight and persists across refresh. +5. Incoming activity in the currently open thread does not leave it unread. +6. Sending or receiving delivery updates does not clear another unread state. +7. Archiving/unarchiving does not change read state. + +## Out of Scope + +- Per-message read receipts. +- Android read/unread UI. +- Unread counts or badges outside the thread list. +- Manual read/unread controls. +- Cross-user/group-conversation receipt tracking. diff --git a/tests/README.md b/tests/README.md index 345d84d1e..45bbb278f 100644 --- a/tests/README.md +++ b/tests/README.md @@ -53,6 +53,7 @@ The API's Firebase SDK is configured (via `FCM_ENDPOINT` env var) to redirect al - [x] **Send SMS E2E** — Full send lifecycle: API → FCM push → emulator responds with SENT/DELIVERED events → message reaches `delivered` status - [x] **Receive SMS E2E** — Phone submits received message to API → message is stored and retrievable via GET endpoint +- [x] **Message thread read receipts E2E** — Incoming SMS and missed calls mark a thread unread, the existing thread update endpoint marks it read, and outbound activity preserves unread state - [x] **Unarchive Thread on Receive E2E** — Archived thread returns to the inbox on inbound message when the phone's `unarchive_thread` setting is enabled, and stays archived when disabled ## Prerequisites diff --git a/tests/read_receipts_test.go b/tests/read_receipts_test.go new file mode 100644 index 000000000..a65e7eb73 --- /dev/null +++ b/tests/read_receipts_test.go @@ -0,0 +1,206 @@ +package tests + +import ( + "bytes" + "context" + "encoding/json" + "fmt" + "io" + "net/http" + "net/url" + "testing" + "time" + + httpsms "github.com/NdoleStudio/httpsms-go" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +type integrationMessageThread struct { + ID string `json:"id"` + Contact string `json:"contact"` + IsRead bool `json:"is_read"` + LastMessageContent *string `json:"last_message_content"` +} + +func requestJSON( + ctx context.Context, + t *testing.T, + method string, + path string, + apiKey string, + payload any, + expectedStatus int, + output any, +) { + t.Helper() + + var encoded []byte + if payload != nil { + var err error + encoded, err = json.Marshal(payload) + require.NoError(t, err) + } + + deadline := time.Now().Add(20 * time.Second) + for { + request, err := http.NewRequestWithContext(ctx, method, apiBaseURL+path, bytes.NewReader(encoded)) + require.NoError(t, err) + request.Header.Set("x-api-key", apiKey) + request.Header.Set("Content-Type", "application/json") + + response, err := http.DefaultClient.Do(request) + require.NoError(t, err) + + responseBody, err := io.ReadAll(response.Body) + response.Body.Close() + require.NoError(t, err) + + if response.StatusCode == http.StatusUnauthorized && + apiKey != userAPIKey && + time.Now().Before(deadline) { + time.Sleep(500 * time.Millisecond) + continue + } + + require.Equal(t, expectedStatus, response.StatusCode, "response: %s", string(responseBody)) + if output != nil { + require.NoError(t, json.Unmarshal(responseBody, output)) + } + return + } +} + +func fetchMessageThreads(ctx context.Context, t *testing.T, owner string) []integrationMessageThread { + t.Helper() + + var response struct { + Data []integrationMessageThread `json:"data"` + } + path := fmt.Sprintf( + "/v1/message-threads?owner=%s&skip=0&limit=20&is_archived=false", + url.QueryEscape(owner), + ) + requestJSON(ctx, t, http.MethodGet, path, userAPIKey, nil, http.StatusOK, &response) + return response.Data +} + +func waitForMessageThread( + ctx context.Context, + t *testing.T, + owner string, + contact string, + timeout time.Duration, + matches func(integrationMessageThread) bool, +) integrationMessageThread { + t.Helper() + + deadline := time.Now().Add(timeout) + for time.Now().Before(deadline) { + for _, thread := range fetchMessageThreads(ctx, t, owner) { + if thread.Contact == contact && matches(thread) { + return thread + } + } + time.Sleep(500 * time.Millisecond) + } + + t.Fatalf("thread %s -> %s did not reach the expected state within %v", owner, contact, timeout) + return integrationMessageThread{} +} + +func markMessageThreadRead(ctx context.Context, t *testing.T, threadID string) integrationMessageThread { + t.Helper() + + var response struct { + Data integrationMessageThread `json:"data"` + } + requestJSON( + ctx, + t, + http.MethodPut, + "/v1/message-threads/"+threadID, + userAPIKey, + map[string]any{"is_read": true}, + http.StatusOK, + &response, + ) + return response.Data +} + +func TestMessageThreadReadReceipts(t *testing.T) { + ctx := context.Background() + phone := setupPhone(ctx, t, 60) + contact := randomPhoneNumber() + + requestJSON( + ctx, + t, + http.MethodPost, + "/v1/messages/receive", + phone.PhoneAPIKey, + map[string]any{ + "from": contact, + "to": phone.PhoneNumber, + "content": "Unread inbound message", + "encrypted": false, + "sim": "SIM1", + "timestamp": time.Now().UTC().Format(time.RFC3339Nano), + }, + http.StatusOK, + nil, + ) + + thread := waitForMessageThread(ctx, t, phone.PhoneNumber, contact, 20*time.Second, func(thread integrationMessageThread) bool { + return !thread.IsRead + }) + assert.False(t, thread.IsRead) + + updated := markMessageThreadRead(ctx, t, thread.ID) + assert.True(t, updated.IsRead) + assert.Equal(t, contact, updated.Contact) + require.NotNil(t, updated.LastMessageContent) + assert.Equal(t, "Unread inbound message", *updated.LastMessageContent) + waitForMessageThread(ctx, t, phone.PhoneNumber, contact, 10*time.Second, func(thread integrationMessageThread) bool { + return thread.IsRead + }) + + requestJSON( + ctx, + t, + http.MethodPost, + "/v1/messages/calls/missed", + phone.PhoneAPIKey, + map[string]any{ + "from": contact, + "to": phone.PhoneNumber, + "sim": "SIM1", + "timestamp": time.Now().UTC().Format(time.RFC3339Nano), + }, + http.StatusOK, + nil, + ) + + thread = waitForMessageThread(ctx, t, phone.PhoneNumber, contact, 20*time.Second, func(thread integrationMessageThread) bool { + return !thread.IsRead && + thread.LastMessageContent != nil && + *thread.LastMessageContent == "Missed phone call" + }) + assert.False(t, thread.IsRead) + + outboundContent := "Outbound activity preserves unread" + client := newAPIClient() + _, response, err := client.Messages.Send(ctx, &httpsms.MessageSendParams{ + From: phone.PhoneNumber, + To: contact, + Content: outboundContent, + }) + require.NoError(t, err) + require.Equal(t, http.StatusOK, response.HTTPResponse.StatusCode) + + thread = waitForMessageThread(ctx, t, phone.PhoneNumber, contact, 20*time.Second, func(thread integrationMessageThread) bool { + return thread.LastMessageContent != nil && + *thread.LastMessageContent == outboundContent + }) + assert.False(t, thread.IsRead, "outbound activity must not clear unread state") +} diff --git a/web/app/components/MessageThread.vue b/web/app/components/MessageThread.vue index 029599180..2a9834317 100644 --- a/web/app/components/MessageThread.vue +++ b/web/app/components/MessageThread.vue @@ -99,7 +99,11 @@ function onInstallApp() { :active="threadsStore.threadId === thread.id" > - {{ + {{ formatPhoneNumber(thread.contact) }} {{ thread.last_message_content }} diff --git a/web/app/pages/threads/[id]/index.vue b/web/app/pages/threads/[id]/index.vue index 1fd11eae1..73632a64c 100644 --- a/web/app/pages/threads/[id]/index.vue +++ b/web/app/pages/threads/[id]/index.vue @@ -112,9 +112,19 @@ function scrollToElement() { hideMessages.value = false } -function loadMessages(hide = true) { +async function markCurrentThreadRead(force = false) { + const threadId = route.params.id as string + try { + await threadsStore.markThreadRead(threadId, force) + } catch (error) { + console.error(error) + } +} + +function loadMessages(hide = true, markRead = true) { loadingMessages.value = true const threadId = route.params.id as string + if (markRead) void markCurrentThreadRead() threadsStore .loadThreadMessages(threadId) .then((response: EntitiesMessage[]) => { @@ -232,7 +242,16 @@ onMounted(async () => { if (!loadingMessages.value) loadMessages(false) }) webhookChannel.bind('message.phone.received', () => { - if (!loadingMessages.value) loadMessages(false) + if (!loadingMessages.value) { + void markCurrentThreadRead(true) + loadMessages(false, false) + } + }) + webhookChannel.bind('message.call.missed', () => { + if (!loadingMessages.value) { + void markCurrentThreadRead(true) + loadMessages(false, false) + } }) }) diff --git a/web/app/stores/threads.ts b/web/app/stores/threads.ts index aef2b2947..7e3002624 100644 --- a/web/app/stores/threads.ts +++ b/web/app/stores/threads.ts @@ -24,6 +24,13 @@ export const useThreadsStore = defineStore('threads', () => { return threads.value.find((x) => x.id === id) !== undefined } + function replaceThread(updatedThread: EntitiesMessageThread) { + const index = threads.value.findIndex( + (thread) => thread.id === updatedThread.id, + ) + if (index !== -1) threads.value[index] = updatedThread + } + async function loadThreads() { const phonesStore = usePhonesStore() if (phonesStore.owner === null && phonesStore.phones.length === 0) { @@ -93,6 +100,38 @@ export const useThreadsStore = defineStore('threads', () => { }) } + async function markThreadRead(threadId: string, force = false) { + const thread = threads.value.find((item) => item.id === threadId) + if (!thread) throw new Error(`Cannot find thread with id ${threadId}`) + if (!force && thread.is_read) return + + try { + const response = await apiFetch<{ data: EntitiesMessageThread }>( + `/v1/message-threads/${threadId}`, + { + method: 'PUT', + body: { is_read: true }, + }, + ) + replaceThread(response.data) + } catch (error) { + notificationsStore.addNotification({ + message: 'The message thread could not be marked as read', + type: 'error', + }) + try { + await loadThreads() + } catch (reloadError) { + throw new AggregateError( + [error, reloadError], + 'Could not mark the message thread as read or reload threads', + { cause: reloadError }, + ) + } + throw error + } + } + async function deleteThread(id: string) { await apiFetch(`/v1/message-threads/${id}`, { method: 'DELETE' }) threadId.value = null @@ -122,6 +161,7 @@ export const useThreadsStore = defineStore('threads', () => { setThreadId, toggleArchive, updateThread, + markThreadRead, deleteThread, resetState, } diff --git a/web/shared/types/api.ts b/web/shared/types/api.ts index 3f19b8534..27e2db111 100644 --- a/web/shared/types/api.ts +++ b/web/shared/types/api.ts @@ -205,6 +205,8 @@ export interface EntitiesMessageThread { id: string; /** @example false */ is_archived: boolean; + /** @example true */ + is_read: boolean; /** @example "This is a sample message content" */ last_message_content: string; /** @example "32343a19-da5e-4b1b-a767-3298a73703ca" */ @@ -484,7 +486,9 @@ export interface RequestsMessageSendScheduleWindow { export interface RequestsMessageThreadUpdate { /** @example true */ - is_archived: boolean; + is_archived?: boolean; + /** @example true */ + is_read?: boolean; } export interface RequestsPhoneAPIKeyStoreRequest { @@ -684,6 +688,14 @@ export interface ResponsesMessageSendSchedulesResponse { status: string; } +export interface ResponsesMessageThreadResponse { + data: EntitiesMessageThread; + /** @example "Request handled successfully" */ + message: string; + /** @example "success" */ + status: string; +} + export interface ResponsesMessageThreadsResponse { data: EntitiesMessageThread[]; /** @example "Request handled successfully" */