diff --git a/backend/internal/repository/ops_write_pressure_integration_test.go b/backend/internal/repository/ops_write_pressure_integration_test.go index 91987c01eb..0eb0385848 100644 --- a/backend/internal/repository/ops_write_pressure_integration_test.go +++ b/backend/internal/repository/ops_write_pressure_integration_test.go @@ -59,13 +59,38 @@ func TestEnqueueSchedulerOutbox_DeduplicatesIdempotentEvents(t *testing.T) { var firstID int64 require.NoError(t, integrationDB.QueryRowContext(ctx, "SELECT id FROM scheduler_outbox WHERE event_type = $1", service.SchedulerOutboxEventAccountChanged).Scan(&firstID)) - require.NoError(t, NewSchedulerOutboxRepository(integrationDB).MarkProcessed(ctx, []int64{firstID})) + events, err := NewSchedulerOutboxRepository(integrationDB).ListAfterAndReleaseDedup(ctx, 0, 100) + require.NoError(t, err) + require.Len(t, events, 1) + require.Equal(t, firstID, events[0].ID) require.NoError(t, enqueueSchedulerOutbox(ctx, integrationDB, service.SchedulerOutboxEventAccountChanged, &accountID, nil, nil)) require.NoError(t, integrationDB.QueryRowContext(ctx, "SELECT COUNT(*) FROM scheduler_outbox WHERE event_type = $1", service.SchedulerOutboxEventAccountChanged).Scan(&count)) require.Equal(t, 2, count) } +func TestSchedulerOutbox_ListAfterAndReleaseDedup_AllowsSameKeyWhileEventInFlight(t *testing.T) { + ctx := context.Background() + _, _ = integrationDB.ExecContext(ctx, "TRUNCATE scheduler_outbox RESTART IDENTITY") + + accountID := int64(17345) + require.NoError(t, enqueueSchedulerOutbox(ctx, integrationDB, service.SchedulerOutboxEventAccountChanged, &accountID, nil, nil)) + + events, err := NewSchedulerOutboxRepository(integrationDB).ListAfterAndReleaseDedup(ctx, 0, 100) + require.NoError(t, err) + require.Len(t, events, 1) + + require.NoError(t, enqueueSchedulerOutbox(ctx, integrationDB, service.SchedulerOutboxEventAccountChanged, &accountID, nil, nil)) + + var count int + require.NoError(t, integrationDB.QueryRowContext(ctx, "SELECT COUNT(*) FROM scheduler_outbox WHERE event_type = $1", service.SchedulerOutboxEventAccountChanged).Scan(&count)) + require.Equal(t, 2, count) + + var pendingKeys int + require.NoError(t, integrationDB.QueryRowContext(ctx, "SELECT COUNT(*) FROM scheduler_outbox WHERE dedup_key IS NOT NULL").Scan(&pendingKeys)) + require.Equal(t, 1, pendingKeys) +} + func TestEnqueueSchedulerOutbox_CoalescesAccountStateBurst(t *testing.T) { ctx := context.Background() _, _ = integrationDB.ExecContext(ctx, "TRUNCATE scheduler_outbox RESTART IDENTITY") diff --git a/backend/internal/repository/scheduler_outbox_repo.go b/backend/internal/repository/scheduler_outbox_repo.go index b60d1648f5..9edde6945c 100644 --- a/backend/internal/repository/scheduler_outbox_repo.go +++ b/backend/internal/repository/scheduler_outbox_repo.go @@ -10,7 +10,6 @@ import ( "strconv" "github.com/Wei-Shaw/sub2api/internal/service" - "github.com/lib/pq" ) type schedulerOutboxRepository struct { @@ -21,16 +20,30 @@ func NewSchedulerOutboxRepository(db *sql.DB) service.SchedulerOutboxRepository return &schedulerOutboxRepository{db: db} } -func (r *schedulerOutboxRepository) ListAfter(ctx context.Context, afterID int64, limit int) ([]service.SchedulerOutboxEvent, error) { +func (r *schedulerOutboxRepository) ListAfterAndReleaseDedup(ctx context.Context, afterID int64, limit int) ([]service.SchedulerOutboxEvent, error) { if limit <= 0 { limit = 100 } rows, err := r.db.QueryContext(ctx, ` - SELECT id, event_type, account_id, group_id, payload, created_at - FROM scheduler_outbox - WHERE id > $1 - ORDER BY id ASC - LIMIT $2 + WITH selected AS MATERIALIZED ( + SELECT id, event_type, account_id, group_id, payload, created_at + FROM scheduler_outbox + WHERE id > $1 + ORDER BY id ASC + LIMIT $2 + FOR UPDATE + ), released AS ( + UPDATE scheduler_outbox AS o + SET dedup_key = NULL + FROM selected AS s + WHERE o.id = s.id + AND o.dedup_key IS NOT NULL + RETURNING o.id + ) + SELECT s.id, s.event_type, s.account_id, s.group_id, s.payload, s.created_at + FROM selected AS s + CROSS JOIN (SELECT COUNT(*) FROM released) AS release_barrier + ORDER BY s.id ASC `, afterID, limit) if err != nil { return nil, err @@ -81,19 +94,6 @@ func (r *schedulerOutboxRepository) MaxID(ctx context.Context) (int64, error) { return maxID, nil } -func (r *schedulerOutboxRepository) MarkProcessed(ctx context.Context, eventIDs []int64) error { - if len(eventIDs) == 0 { - return nil - } - _, err := r.db.ExecContext(ctx, ` - UPDATE scheduler_outbox - SET dedup_key = NULL - WHERE id = ANY($1) - AND dedup_key IS NOT NULL - `, pq.Array(eventIDs)) - return err -} - func enqueueSchedulerOutbox(ctx context.Context, exec sqlExecutor, eventType string, accountID *int64, groupID *int64, payload any) error { if exec == nil { return nil diff --git a/backend/internal/service/scheduler_outbox.go b/backend/internal/service/scheduler_outbox.go index 0fe977d25e..c138b7e5a5 100644 --- a/backend/internal/service/scheduler_outbox.go +++ b/backend/internal/service/scheduler_outbox.go @@ -16,7 +16,6 @@ type SchedulerOutboxEvent struct { // SchedulerOutboxRepository 提供调度 outbox 的读取接口。 type SchedulerOutboxRepository interface { - ListAfter(ctx context.Context, afterID int64, limit int) ([]SchedulerOutboxEvent, error) + ListAfterAndReleaseDedup(ctx context.Context, afterID int64, limit int) ([]SchedulerOutboxEvent, error) MaxID(ctx context.Context) (int64, error) - MarkProcessed(ctx context.Context, eventIDs []int64) error } diff --git a/backend/internal/service/scheduler_snapshot_service.go b/backend/internal/service/scheduler_snapshot_service.go index 23025ad719..6b25c0430f 100644 --- a/backend/internal/service/scheduler_snapshot_service.go +++ b/backend/internal/service/scheduler_snapshot_service.go @@ -242,7 +242,7 @@ func (s *SchedulerSnapshotService) pollOutbox() { return } - events, err := s.outboxRepo.ListAfter(ctx, watermark, 200) + events, err := s.outboxRepo.ListAfterAndReleaseDedup(ctx, watermark, 200) if err != nil { logger.LegacyPrintf("service.scheduler_snapshot", "[Scheduler] outbox poll failed: %v", err) return @@ -253,7 +253,6 @@ func (s *SchedulerSnapshotService) pollOutbox() { watermarkForCheck := watermark seen := make(map[batchSeenKey]struct{}) - processedIDs := make([]int64, 0, len(events)) for _, event := range events { eventCtx, cancel := context.WithTimeout(context.Background(), outboxEventTimeout) err := s.handleOutboxEvent(eventCtx, event, seen) @@ -262,15 +261,6 @@ func (s *SchedulerSnapshotService) pollOutbox() { logger.LegacyPrintf("service.scheduler_snapshot", "[Scheduler] outbox handle failed: id=%d type=%s err=%v", event.ID, event.EventType, err) return } - processedIDs = append(processedIDs, event.ID) - } - - markCtx, markCancel := context.WithTimeout(context.Background(), 5*time.Second) - markErr := s.outboxRepo.MarkProcessed(markCtx, processedIDs) - markCancel() - if markErr != nil { - logger.LegacyPrintf("service.scheduler_snapshot", "[Scheduler] outbox mark processed failed: %v", markErr) - return } lastID := events[len(events)-1].ID