Skip to content

Commit 33fe141

Browse files
AchoArnoldCopilot
andcommitted
fix(api): preserve deleted message markers
Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: 332cfd7b-8a60-4e28-9ea8-61854aa712c3
1 parent 8e9e554 commit 33fe141

10 files changed

Lines changed: 345 additions & 101 deletions
Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,8 @@
1+
package entities
2+
3+
import "github.com/google/uuid"
4+
5+
// MessageThreadDeletedItem records a permanently deleted message activity.
6+
type MessageThreadDeletedItem struct {
7+
MessageID uuid.UUID `gorm:"primaryKey;type:uuid"`
8+
}

api/pkg/entities/message_thread_test.go

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -59,3 +59,15 @@ func TestMessageThreadUnreadItemRetainsCountedState(t *testing.T) {
5959
assert.Contains(t, counted.Tag.Get("gorm"), "not null")
6060
assert.Contains(t, counted.Tag.Get("gorm"), "default:true")
6161
}
62+
63+
func TestMessageThreadDeletedItemIsIndependentFromThreadLifecycle(t *testing.T) {
64+
itemType := reflect.TypeOf(MessageThreadDeletedItem{})
65+
66+
require.Equal(t, 1, itemType.NumField())
67+
messageID, ok := itemType.FieldByName("MessageID")
68+
require.True(t, ok)
69+
assert.Contains(t, messageID.Tag.Get("gorm"), "primaryKey")
70+
assert.Contains(t, messageID.Tag.Get("gorm"), "type:uuid")
71+
_, hasMessageThreadID := itemType.FieldByName("MessageThreadID")
72+
assert.False(t, hasMessageThreadID)
73+
}

api/pkg/listeners/message_thread_listener_test.go

Lines changed: 2 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -64,15 +64,6 @@ func TestMessageThreadListenerMarksMissedCallUnread(t *testing.T) {
6464

6565
func TestMessageThreadListenerDeletesNonLastUnreadMessage(t *testing.T) {
6666
repository, routes := newMessageThreadListenerForTest()
67-
currentLastMessageID := uuid.New()
68-
repository.thread = &entities.MessageThread{
69-
ID: uuid.New(),
70-
UserID: entities.UserID("user-id"),
71-
Owner: "+18005550199",
72-
Contact: "+18005550100",
73-
LastMessageID: &currentLastMessageID,
74-
}
75-
7667
deletedMessageID := uuid.New()
7768
previousMessageID := uuid.New()
7869
previousStatus := entities.MessageStatus(entities.MessageStatusDelivered)
@@ -95,6 +86,8 @@ func TestMessageThreadListenerDeletesNonLastUnreadMessage(t *testing.T) {
9586

9687
require.NoError(t, err)
9788
assert.Equal(t, deletedMessageID, repository.deletedUpdate.DeletedMessageID)
89+
assert.Equal(t, "+18005550199", repository.deletedUpdate.Owner)
90+
assert.Equal(t, "+18005550100", repository.deletedUpdate.Contact)
9891
require.NotNil(t, repository.deletedUpdate.LastMessageID)
9992
assert.Equal(t, previousMessageID, *repository.deletedUpdate.LastMessageID)
10093
}

api/pkg/migrations/message_thread_unread_count.go

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -20,7 +20,11 @@ func (messageThreadConversationIndex) TableName() string {
2020

2121
// MigrateMessageThreadUnreadCount migrates message thread unread count schema.
2222
func MigrateMessageThreadUnreadCount(db *gorm.DB) error {
23-
if err := db.AutoMigrate(&entities.MessageThread{}, &entities.MessageThreadUnreadItem{}); err != nil {
23+
if err := db.AutoMigrate(
24+
&entities.MessageThread{},
25+
&entities.MessageThreadUnreadItem{},
26+
&entities.MessageThreadDeletedItem{},
27+
); err != nil {
2428
return stacktrace.Propagate(err, "cannot migrate message thread unread count schema")
2529
}
2630

api/pkg/migrations/message_thread_unread_count_test.go

Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,19 @@ func TestMigrateMessageThreadUnreadCountSkipsLegacyBackfillWhenIsReadColumnMissi
2828
assert.Contains(t, strings.Join(recorder.execs, "\n"), `"counted" boolean NOT NULL DEFAULT true`)
2929
}
3030

31+
func TestMigrateMessageThreadUnreadCountCreatesIndependentDeletedItemSchema(t *testing.T) {
32+
db, recorder := newMigrationTestDB(t, migrationTestDBOptions{})
33+
34+
require.NoError(t, MigrateMessageThreadUnreadCount(db))
35+
36+
deletedItems := migrationExecIndex(recorder, `CREATE TABLE "message_thread_deleted_items"`)
37+
require.NotEqual(t, -1, deletedItems)
38+
assert.Contains(t, recorder.execs[deletedItems], `"message_id" uuid`)
39+
assert.Contains(t, recorder.execs[deletedItems], `PRIMARY KEY`)
40+
assert.NotContains(t, recorder.execs[deletedItems], `"message_thread_id"`)
41+
assert.NotContains(t, recorder.execs[deletedItems], `FOREIGN KEY`)
42+
}
43+
3144
func TestMigrateMessageThreadUnreadCountBackfillsBeforeDropAndSkipsOnSecondRun(t *testing.T) {
3245
db, recorder := newMigrationTestDB(t, migrationTestDBOptions{hasLegacyIsRead: true})
3346

api/pkg/repositories/gorm_message_thread_repository.go

Lines changed: 52 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -148,6 +148,28 @@ func insertUnreadItem(tx *gorm.DB, item entities.MessageThreadUnreadItem) (bool,
148148
return result.RowsAffected == 1, nil
149149
}
150150

151+
func insertDeletedItem(tx *gorm.DB, messageID uuid.UUID) error {
152+
if err := tx.
153+
Clauses(clause.OnConflict{DoNothing: true}).
154+
Create(&entities.MessageThreadDeletedItem{MessageID: messageID}).
155+
Error; err != nil {
156+
return stacktrace.Propagatef(err, "cannot insert deleted message marker for message [%s]", messageID)
157+
}
158+
return nil
159+
}
160+
161+
func isDeletedItem(tx *gorm.DB, messageID uuid.UUID) (bool, error) {
162+
item := new(entities.MessageThreadDeletedItem)
163+
result := tx.
164+
Where("message_id = ?", messageID).
165+
Limit(1).
166+
Find(item)
167+
if result.Error != nil {
168+
return false, stacktrace.Propagatef(result.Error, "cannot load deleted message marker for message [%s]", messageID)
169+
}
170+
return result.RowsAffected != 0, nil
171+
}
172+
151173
func markUnreadItemDeleted(tx *gorm.DB, messageID uuid.UUID, threadID uuid.UUID) (bool, error) {
152174
result := tx.
153175
Model(&entities.MessageThreadUnreadItem{}).
@@ -251,26 +273,33 @@ func (repository *gormMessageThreadRepository) UpdateAfterDeletedMessage(ctx con
251273

252274
err := crdbgorm.ExecuteTx(ctx, repository.db, nil, func(tx *gorm.DB) error {
253275
tx = tx.WithContext(ctx)
254-
thread, err := lockMessageThread(tx, params.UserID, params.MessageThreadID)
276+
if err := insertDeletedItem(tx, params.DeletedMessageID); err != nil {
277+
return err
278+
}
279+
280+
thread, err := lockMessageThreadByConversation(tx, params.UserID, params.Owner, params.Contact)
281+
if stacktrace.GetCode(err) == ErrCodeNotFound {
282+
return nil
283+
}
255284
if err != nil {
256285
return err
257286
}
258287

259-
deleted, err := markUnreadItemDeleted(tx, params.DeletedMessageID, params.MessageThreadID)
288+
deleted, err := markUnreadItemDeleted(tx, params.DeletedMessageID, thread.ID)
260289
if err != nil {
261290
return err
262291
}
263292
if deleted {
264293
if err := tx.
265294
Model(thread).
266295
Where("user_id = ?", params.UserID).
267-
Where("id = ?", params.MessageThreadID).
296+
Where("id = ?", thread.ID).
268297
UpdateColumn("unread_count", gorm.Expr("GREATEST(unread_count - 1, 0)")).
269298
Error; err != nil {
270299
return stacktrace.Propagatef(
271300
err,
272301
"cannot decrement unread count for thread [%s] and user [%s]",
273-
params.MessageThreadID,
302+
thread.ID,
274303
params.UserID,
275304
)
276305
}
@@ -285,13 +314,13 @@ func (repository *gormMessageThreadRepository) UpdateAfterDeletedMessage(ctx con
285314
if params.LastMessageID == nil {
286315
if err := tx.
287316
Where("user_id = ?", params.UserID).
288-
Where("id = ?", params.MessageThreadID).
317+
Where("id = ?", thread.ID).
289318
Delete(&entities.MessageThread{}).
290319
Error; err != nil {
291320
return stacktrace.Propagatef(
292321
err,
293322
"cannot delete message thread [%s] for user [%s] after deleting final message [%s]",
294-
params.MessageThreadID,
323+
thread.ID,
295324
params.UserID,
296325
params.DeletedMessageID,
297326
)
@@ -306,13 +335,13 @@ func (repository *gormMessageThreadRepository) UpdateAfterDeletedMessage(ctx con
306335
if err := tx.
307336
Model(thread).
308337
Where("user_id = ?", params.UserID).
309-
Where("id = ?", params.MessageThreadID).
338+
Where("id = ?", thread.ID).
310339
Updates(updates).
311340
Error; err != nil {
312341
return stacktrace.Propagatef(
313342
err,
314343
"cannot update deleted-message metadata for thread [%s] and user [%s]",
315-
params.MessageThreadID,
344+
thread.ID,
316345
params.UserID,
317346
)
318347
}
@@ -323,9 +352,10 @@ func (repository *gormMessageThreadRepository) UpdateAfterDeletedMessage(ctx con
323352
span,
324353
stacktrace.Propagatef(
325354
err,
326-
"cannot apply deleted message [%s] to thread [%s] for user [%s]",
355+
"cannot apply deleted message [%s] to conversation [%s/%s] for user [%s]",
327356
params.DeletedMessageID,
328-
params.MessageThreadID,
357+
params.Owner,
358+
params.Contact,
329359
params.UserID,
330360
),
331361
)
@@ -341,6 +371,13 @@ func (repository *gormMessageThreadRepository) Store(ctx context.Context, params
341371

342372
err := crdbgorm.ExecuteTx(ctx, repository.db, nil, func(tx *gorm.DB) error {
343373
tx = tx.WithContext(ctx)
374+
if params.Thread.LastMessageID != nil {
375+
deleted, err := isDeletedItem(tx, *params.Thread.LastMessageID)
376+
if err != nil || deleted {
377+
return err
378+
}
379+
}
380+
344381
candidate := *params.Thread
345382
result := tx.Clauses(clause.OnConflict{DoNothing: true}).Create(&candidate)
346383
if result.Error != nil {
@@ -417,6 +454,11 @@ func (repository *gormMessageThreadRepository) UpdateActivity(ctx context.Contex
417454

418455
err := crdbgorm.ExecuteTx(ctx, repository.db, nil, func(tx *gorm.DB) error {
419456
tx = tx.WithContext(ctx)
457+
deleted, err := isDeletedItem(tx, params.MessageID)
458+
if err != nil || deleted {
459+
return err
460+
}
461+
420462
thread, err := lockMessageThread(tx, params.UserID, params.MessageThreadID)
421463
if err != nil {
422464
return err

0 commit comments

Comments
 (0)