diff --git a/backend/internal/repository/migrations_schema_integration_test.go b/backend/internal/repository/migrations_schema_integration_test.go index 7d5f66c7f5..b1ea0f990f 100644 --- a/backend/internal/repository/migrations_schema_integration_test.go +++ b/backend/internal/repository/migrations_schema_integration_test.go @@ -97,6 +97,10 @@ func TestMigrationsRunner_IsIdempotent_AndSchemaIsUpToDate(t *testing.T) { require.NoError(t, tx.QueryRowContext(context.Background(), "SELECT to_regclass('public.security_secrets')").Scan(&securitySecretsRegclass)) require.True(t, securitySecretsRegclass.Valid, "expected security_secrets table to exist") + // scheduler_outbox pending dedup support + requireColumn(t, tx, "scheduler_outbox", "dedup_key", "text", 0, true) + requireIndex(t, tx, "scheduler_outbox", "idx_scheduler_outbox_pending_dedup_key") + // user_allowed_groups table should exist var uagRegclass sql.NullString require.NoError(t, tx.QueryRowContext(context.Background(), "SELECT to_regclass('public.user_allowed_groups')").Scan(&uagRegclass)) diff --git a/backend/internal/repository/ops_write_pressure_integration_test.go b/backend/internal/repository/ops_write_pressure_integration_test.go index 4520891ec6..91987c01eb 100644 --- a/backend/internal/repository/ops_write_pressure_integration_test.go +++ b/backend/internal/repository/ops_write_pressure_integration_test.go @@ -57,12 +57,13 @@ func TestEnqueueSchedulerOutbox_DeduplicatesIdempotentEvents(t *testing.T) { require.NoError(t, integrationDB.QueryRowContext(ctx, "SELECT COUNT(*) FROM scheduler_outbox WHERE event_type = $1", service.SchedulerOutboxEventAccountChanged).Scan(&count)) require.Equal(t, 1, count) - require.GreaterOrEqual(t, schedulerOutboxDedupWindow, 10*time.Second) + 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})) - time.Sleep(1200 * time.Millisecond) 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, 1, count) + require.Equal(t, 2, count) } func TestEnqueueSchedulerOutbox_CoalescesAccountStateBurst(t *testing.T) { @@ -76,10 +77,25 @@ func TestEnqueueSchedulerOutbox_CoalescesAccountStateBurst(t *testing.T) { var count int require.NoError(t, integrationDB.QueryRowContext(ctx, "SELECT COUNT(*) FROM scheduler_outbox WHERE event_type = $1", service.SchedulerOutboxEventAccountChanged).Scan(&count)) - t.Logf("same-account account_changed burst: calls=50 inserted=%d dedup_window=%s", count, schedulerOutboxDedupWindow) + t.Logf("same-account account_changed burst: calls=50 inserted=%d", count) require.Equal(t, 1, count) } +func TestEnqueueSchedulerOutbox_DoesNotDeduplicateDifferentPayload(t *testing.T) { + ctx := context.Background() + _, _ = integrationDB.ExecContext(ctx, "TRUNCATE scheduler_outbox RESTART IDENTITY") + + accountID := int64(32345) + payload1 := map[string]any{"group_ids": []int64{1}} + payload2 := map[string]any{"group_ids": []int64{2}} + require.NoError(t, enqueueSchedulerOutbox(ctx, integrationDB, service.SchedulerOutboxEventAccountChanged, &accountID, nil, payload1)) + require.NoError(t, enqueueSchedulerOutbox(ctx, integrationDB, service.SchedulerOutboxEventAccountChanged, &accountID, nil, payload2)) + + 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) +} + func TestEnqueueSchedulerOutbox_DoesNotDeduplicateLastUsed(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 586efd0582..b60d1648f5 100644 --- a/backend/internal/repository/scheduler_outbox_repo.go +++ b/backend/internal/repository/scheduler_outbox_repo.go @@ -2,19 +2,21 @@ package repository import ( "context" + "crypto/sha256" "database/sql" + "encoding/hex" "encoding/json" - "time" + "fmt" + "strconv" "github.com/Wei-Shaw/sub2api/internal/service" + "github.com/lib/pq" ) type schedulerOutboxRepository struct { db *sql.DB } -const schedulerOutboxDedupWindow = 10 * time.Second - func NewSchedulerOutboxRepository(db *sql.DB) service.SchedulerOutboxRepository { return &schedulerOutboxRepository{db: db} } @@ -79,17 +81,32 @@ 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 } var payloadArg any + var payloadJSON []byte if payload != nil { encoded, err := json.Marshal(payload) if err != nil { return err } payloadArg = encoded + payloadJSON = encoded } query := ` INSERT INTO scheduler_outbox (event_type, account_id, group_id, payload) @@ -97,24 +114,34 @@ func enqueueSchedulerOutbox(ctx context.Context, exec sqlExecutor, eventType str ` args := []any{eventType, accountID, groupID, payloadArg} if schedulerOutboxEventSupportsDedup(eventType) { + dedupKey := schedulerOutboxDedupKey(eventType, accountID, groupID, payloadJSON) query = ` - INSERT INTO scheduler_outbox (event_type, account_id, group_id, payload) - SELECT $1, $2, $3, $4 - WHERE NOT EXISTS ( - SELECT 1 - FROM scheduler_outbox - WHERE event_type = $1 - AND account_id IS NOT DISTINCT FROM $2 - AND group_id IS NOT DISTINCT FROM $3 - AND created_at >= NOW() - make_interval(secs => $5) - ) + INSERT INTO scheduler_outbox (event_type, account_id, group_id, payload, dedup_key) + VALUES ($1, $2, $3, $4, $5) + ON CONFLICT (dedup_key) WHERE dedup_key IS NOT NULL DO NOTHING ` - args = append(args, schedulerOutboxDedupWindow.Seconds()) + args = append(args, dedupKey) } _, err := exec.ExecContext(ctx, query, args...) return err } +func schedulerOutboxDedupKey(eventType string, accountID *int64, groupID *int64, payloadJSON []byte) string { + h := sha256.New() + _, _ = h.Write([]byte(eventType)) + _, _ = h.Write([]byte{0}) + if accountID != nil { + _, _ = h.Write([]byte(strconv.FormatInt(*accountID, 10))) + } + _, _ = h.Write([]byte{0}) + if groupID != nil { + _, _ = h.Write([]byte(strconv.FormatInt(*groupID, 10))) + } + _, _ = h.Write([]byte{0}) + _, _ = h.Write(payloadJSON) + return fmt.Sprintf("scheduler_outbox:%s", hex.EncodeToString(h.Sum(nil))) +} + func schedulerOutboxEventSupportsDedup(eventType string) bool { switch eventType { case service.SchedulerOutboxEventAccountChanged, diff --git a/backend/internal/service/scheduler_outbox.go b/backend/internal/service/scheduler_outbox.go index 32bfcfaaa1..0fe977d25e 100644 --- a/backend/internal/service/scheduler_outbox.go +++ b/backend/internal/service/scheduler_outbox.go @@ -18,4 +18,5 @@ type SchedulerOutboxEvent struct { type SchedulerOutboxRepository interface { ListAfter(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 a68cdf0c77..23025ad719 100644 --- a/backend/internal/service/scheduler_snapshot_service.go +++ b/backend/internal/service/scheduler_snapshot_service.go @@ -253,6 +253,7 @@ 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) @@ -261,6 +262,15 @@ 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 diff --git a/backend/migrations/151_scheduler_outbox_dedup_key.sql b/backend/migrations/151_scheduler_outbox_dedup_key.sql new file mode 100644 index 0000000000..c1d320ef67 --- /dev/null +++ b/backend/migrations/151_scheduler_outbox_dedup_key.sql @@ -0,0 +1,2 @@ +ALTER TABLE scheduler_outbox + ADD COLUMN IF NOT EXISTS dedup_key TEXT; diff --git a/backend/migrations/152_scheduler_outbox_pending_dedup_key_index_notx.sql b/backend/migrations/152_scheduler_outbox_pending_dedup_key_index_notx.sql new file mode 100644 index 0000000000..4be7f610c9 --- /dev/null +++ b/backend/migrations/152_scheduler_outbox_pending_dedup_key_index_notx.sql @@ -0,0 +1,3 @@ +CREATE UNIQUE INDEX CONCURRENTLY IF NOT EXISTS idx_scheduler_outbox_pending_dedup_key + ON scheduler_outbox (dedup_key) + WHERE dedup_key IS NOT NULL;