test(HH-623): cover realtime enqueue boundaries (#170)
Co-authored-by: Rogee <rogee@ipao.vip>
This commit is contained in:
@@ -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)
|
||||
|
||||
@@ -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)
|
||||
|
||||
Reference in New Issue
Block a user