diff --git a/backend/internal/service/widget_service_test.go b/backend/internal/service/widget_service_test.go index e80700b6..1591beff 100644 --- a/backend/internal/service/widget_service_test.go +++ b/backend/internal/service/widget_service_test.go @@ -202,7 +202,7 @@ func TestWidgetService_RealtimeFailureDoesNotBlockCaptain(t *testing.T) { require.Equal(t, model.BackgroundJobStatusRetrying, realtimeJob.Status) } -func TestWidgetService_RealtimeJobFailureRollsBackMessage(t *testing.T) { +func TestWidgetService_SecondRealtimeJobFailureRollsBackMessage(t *testing.T) { db, svc := setupWidgetServiceTest(t) seedWidgetInbox(t, db) pool := worker.NewWorkerPool(db, nil) @@ -217,9 +217,11 @@ func TestWidgetService_RealtimeJobFailureRollsBackMessage(t *testing.T) { initResp, err := svc.Init(context.Background(), WidgetInitRequest{WebsiteToken: "test_ws_token_123"}) require.NoError(t, err) - failBackgroundJobCreate(t, db) + require.NoError(t, db.Exec(`CREATE TRIGGER fail_token_realtime_job BEFORE INSERT ON background_jobs + WHEN CAST(NEW.payload AS TEXT) LIKE '%"target":"pubsub_token"%' + BEGIN SELECT RAISE(FAIL, 'forced second realtime job insert failure'); END`).Error) _, err = svc.SendMessage(context.Background(), WidgetSendMessageRequest{WidgetToken: initResp.WidgetToken, Content: "must roll back"}) - require.ErrorContains(t, err, "forced background job insert failure") + require.ErrorContains(t, err, "forced second realtime job insert failure") var messages, jobs int64 require.NoError(t, db.Model(&model.Message{}).Where("content = ?", "must roll back").Count(&messages).Error) diff --git a/backend/internal/ws/event_publisher_test.go b/backend/internal/ws/event_publisher_test.go index 02db3066..7923149f 100644 --- a/backend/internal/ws/event_publisher_test.go +++ b/backend/internal/ws/event_publisher_test.go @@ -758,6 +758,32 @@ func TestEventPublisher_WidgetEvent_PubsubTokenRoomDelivery(t *testing.T) { assert.Equal(t, uint(1), roomMsg.AccountID) } +func TestEventPublisher_DurableEnqueueKeepsEarlierTargetWhenSecondInsertFails(t *testing.T) { + db, err := gorm.Open(sqlite.Open("file:"+t.Name()+"?mode=memory&cache=shared"), &gorm.Config{}) + require.NoError(t, err) + require.NoError(t, db.AutoMigrate(&model.BackgroundJob{})) + + publisher := NewEventPublisherLocal(nil, nil) + publisher.SetWorkerPool(worker.NewWorkerPool(db)) + require.NoError(t, publisher.PublishEvent(1, EventConversationUpdated, map[string]any{"id": 7})) + + var jobs []model.BackgroundJob + require.NoError(t, db.Where("job_type = ?", taskTypeRealtimeEventPublish).Find(&jobs).Error) + require.Len(t, jobs, 1, "PublishEvent must enqueue its account target") + require.NoError(t, db.Exec("DELETE FROM background_jobs").Error) + require.NoError(t, db.Exec(`CREATE TRIGGER fail_token_realtime_job BEFORE INSERT ON background_jobs + WHEN CAST(NEW.payload AS TEXT) LIKE '%"target":"pubsub_token"%' + BEGIN SELECT RAISE(FAIL, 'forced second realtime job insert failure'); END`).Error) + + err = publisher.PublishWidgetEvent(1, "visitor", EventConversationUpdated, map[string]any{"id": 8}) + require.ErrorContains(t, err, "forced second realtime job insert failure") + require.NoError(t, db.Where("job_type = ?", taskTypeRealtimeEventPublish).Find(&jobs).Error) + require.Len(t, jobs, 1, "the successful account target must remain durable") + var persisted realtimeEventPublishJob + require.NoError(t, json.Unmarshal(jobs[0].Payload, &persisted)) + require.Equal(t, realtimeTargetAccount, persisted.Target) +} + func TestEventPublisher_DurableWidgetPublishRetriesPartialFailure(t *testing.T) { db, err := gorm.Open(sqlite.Open("file:"+t.Name()+"?mode=memory&cache=shared"), &gorm.Config{}) require.NoError(t, err)