From cb14935e9acb5371606218cd3da338f937d93ebf Mon Sep 17 00:00:00 2001 From: jjaw Date: Tue, 16 Jun 2026 11:56:40 +0800 Subject: [PATCH 1/2] fix: cleanup consumed scheduler outbox rows --- .../repository/scheduler_outbox_repo.go | 61 ++++ .../repository/scheduler_outbox_repo_test.go | 91 ++++++ backend/internal/service/scheduler_outbox.go | 8 + .../scheduler_snapshot_outbox_cleanup_test.go | 269 ++++++++++++++++++ .../service/scheduler_snapshot_service.go | 36 ++- 5 files changed, 464 insertions(+), 1 deletion(-) create mode 100644 backend/internal/repository/scheduler_outbox_repo_test.go create mode 100644 backend/internal/service/scheduler_snapshot_outbox_cleanup_test.go diff --git a/backend/internal/repository/scheduler_outbox_repo.go b/backend/internal/repository/scheduler_outbox_repo.go index 9edde6945c..cea5d451f6 100644 --- a/backend/internal/repository/scheduler_outbox_repo.go +++ b/backend/internal/repository/scheduler_outbox_repo.go @@ -8,6 +8,7 @@ import ( "encoding/json" "fmt" "strconv" + "time" "github.com/Wei-Shaw/sub2api/internal/service" ) @@ -16,6 +17,12 @@ type schedulerOutboxRepository struct { db *sql.DB } +type schedulerOutboxCleanupLease struct { + conn *sql.Conn +} + +const schedulerOutboxDefaultCleanSize = 5000 + func NewSchedulerOutboxRepository(db *sql.DB) service.SchedulerOutboxRepository { return &schedulerOutboxRepository{db: db} } @@ -94,6 +101,60 @@ func (r *schedulerOutboxRepository) MaxID(ctx context.Context) (int64, error) { return maxID, nil } +func (r *schedulerOutboxRepository) DeleteConsumedUpTo(ctx context.Context, watermark int64, limit int) (int64, error) { + if watermark <= 0 { + return 0, nil + } + if limit <= 0 { + limit = schedulerOutboxDefaultCleanSize + } + result, err := r.db.ExecContext(ctx, ` + WITH doomed AS ( + SELECT id + FROM scheduler_outbox + WHERE id <= $1 + ORDER BY id ASC + LIMIT $2 + ) + DELETE FROM scheduler_outbox o + USING doomed d + WHERE o.id = d.id + `, watermark, limit) + if err != nil { + return 0, err + } + return result.RowsAffected() +} + +func (r *schedulerOutboxRepository) TryAcquireCleanupLock(ctx context.Context) (service.SchedulerOutboxCleanupLease, bool, error) { + conn, err := r.db.Conn(ctx) + if err != nil { + return nil, false, err + } + + var acquired bool + if err := conn.QueryRowContext(ctx, "SELECT pg_try_advisory_lock(hashtext('scheduler_outbox_cleanup'))").Scan(&acquired); err != nil { + _ = conn.Close() + return nil, false, err + } + if !acquired { + _ = conn.Close() + return nil, false, nil + } + return &schedulerOutboxCleanupLease{conn: conn}, true, nil +} + +func (l *schedulerOutboxCleanupLease) Release() { + if l == nil || l.conn == nil { + return + } + ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) + defer cancel() + _, _ = l.conn.ExecContext(ctx, "SELECT pg_advisory_unlock(hashtext('scheduler_outbox_cleanup'))") + _ = l.conn.Close() + l.conn = nil +} + 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/repository/scheduler_outbox_repo_test.go b/backend/internal/repository/scheduler_outbox_repo_test.go new file mode 100644 index 0000000000..ea91f6c73e --- /dev/null +++ b/backend/internal/repository/scheduler_outbox_repo_test.go @@ -0,0 +1,91 @@ +package repository + +import ( + "context" + "regexp" + "testing" + + sqlmock "github.com/DATA-DOG/go-sqlmock" + "github.com/stretchr/testify/require" +) + +func TestSchedulerOutboxRepositoryDeleteConsumedUpToUsesBoundedCTE(t *testing.T) { + db, mock, err := sqlmock.New() + require.NoError(t, err) + defer func() { _ = db.Close() }() + + repo := &schedulerOutboxRepository{db: db} + const expectedSQL = ` + WITH doomed AS ( + SELECT id + FROM scheduler_outbox + WHERE id <= $1 + ORDER BY id ASC + LIMIT $2 + ) + DELETE FROM scheduler_outbox o + USING doomed d + WHERE o.id = d.id + ` + mock.ExpectExec(regexp.QuoteMeta(expectedSQL)). + WithArgs(int64(42), 5000). + WillReturnResult(sqlmock.NewResult(0, 17)) + + deleted, err := repo.DeleteConsumedUpTo(context.Background(), 42, 5000) + + require.NoError(t, err) + require.EqualValues(t, 17, deleted) + require.NoError(t, mock.ExpectationsWereMet()) +} + +func TestSchedulerOutboxRepositoryDeleteConsumedUpToSkipsNonPositiveWatermark(t *testing.T) { + db, mock, err := sqlmock.New() + require.NoError(t, err) + defer func() { _ = db.Close() }() + + repo := &schedulerOutboxRepository{db: db} + + deleted, err := repo.DeleteConsumedUpTo(context.Background(), 0, 5000) + + require.NoError(t, err) + require.EqualValues(t, 0, deleted) + require.NoError(t, mock.ExpectationsWereMet()) +} + +func TestSchedulerOutboxRepositoryTryAcquireCleanupLock(t *testing.T) { + db, mock, err := sqlmock.New() + require.NoError(t, err) + defer func() { _ = db.Close() }() + + repo := &schedulerOutboxRepository{db: db} + mock.ExpectQuery(regexp.QuoteMeta("SELECT pg_try_advisory_lock(hashtext('scheduler_outbox_cleanup'))")). + WillReturnRows(sqlmock.NewRows([]string{"pg_try_advisory_lock"}).AddRow(true)) + mock.ExpectExec(regexp.QuoteMeta("SELECT pg_advisory_unlock(hashtext('scheduler_outbox_cleanup'))")). + WillReturnResult(sqlmock.NewResult(0, 1)) + + lease, acquired, err := repo.TryAcquireCleanupLock(context.Background()) + require.NoError(t, err) + require.True(t, acquired) + require.NotNil(t, lease) + + lease.Release() + + require.NoError(t, mock.ExpectationsWereMet()) +} + +func TestSchedulerOutboxRepositoryTryAcquireCleanupLockUnavailable(t *testing.T) { + db, mock, err := sqlmock.New() + require.NoError(t, err) + defer func() { _ = db.Close() }() + + repo := &schedulerOutboxRepository{db: db} + mock.ExpectQuery(regexp.QuoteMeta("SELECT pg_try_advisory_lock(hashtext('scheduler_outbox_cleanup'))")). + WillReturnRows(sqlmock.NewRows([]string{"pg_try_advisory_lock"}).AddRow(false)) + + lease, acquired, err := repo.TryAcquireCleanupLock(context.Background()) + require.NoError(t, err) + require.False(t, acquired) + require.Nil(t, lease) + + require.NoError(t, mock.ExpectationsWereMet()) +} diff --git a/backend/internal/service/scheduler_outbox.go b/backend/internal/service/scheduler_outbox.go index c138b7e5a5..2b7665ad78 100644 --- a/backend/internal/service/scheduler_outbox.go +++ b/backend/internal/service/scheduler_outbox.go @@ -18,4 +18,12 @@ type SchedulerOutboxEvent struct { type SchedulerOutboxRepository interface { ListAfterAndReleaseDedup(ctx context.Context, afterID int64, limit int) ([]SchedulerOutboxEvent, error) MaxID(ctx context.Context) (int64, error) + DeleteConsumedUpTo(ctx context.Context, watermark int64, limit int) (int64, error) + TryAcquireCleanupLock(ctx context.Context) (SchedulerOutboxCleanupLease, bool, error) +} + +// SchedulerOutboxCleanupLease holds the PostgreSQL advisory lock used by +// scheduler outbox cleanup. +type SchedulerOutboxCleanupLease interface { + Release() } diff --git a/backend/internal/service/scheduler_snapshot_outbox_cleanup_test.go b/backend/internal/service/scheduler_snapshot_outbox_cleanup_test.go new file mode 100644 index 0000000000..535f8d54e2 --- /dev/null +++ b/backend/internal/service/scheduler_snapshot_outbox_cleanup_test.go @@ -0,0 +1,269 @@ +package service + +import ( + "context" + "errors" + "reflect" + "testing" + "time" +) + +type outboxCleanupCache struct { + watermark int64 + setWatermarks []int64 + updateErr error +} + +func (c *outboxCleanupCache) GetSnapshot(ctx context.Context, bucket SchedulerBucket) ([]*Account, bool, error) { + return nil, false, nil +} + +func (c *outboxCleanupCache) SetSnapshot(ctx context.Context, bucket SchedulerBucket, accounts []Account) error { + return nil +} + +func (c *outboxCleanupCache) GetAccount(ctx context.Context, accountID int64) (*Account, error) { + return nil, nil +} + +func (c *outboxCleanupCache) SetAccount(ctx context.Context, account *Account) error { + return nil +} + +func (c *outboxCleanupCache) DeleteAccount(ctx context.Context, accountID int64) error { + return nil +} + +func (c *outboxCleanupCache) UpdateLastUsed(ctx context.Context, updates map[int64]time.Time) error { + return c.updateErr +} + +func (c *outboxCleanupCache) TryLockBucket(ctx context.Context, bucket SchedulerBucket, ttl time.Duration) (bool, error) { + return true, nil +} + +func (c *outboxCleanupCache) UnlockBucket(ctx context.Context, bucket SchedulerBucket) error { + return nil +} + +func (c *outboxCleanupCache) ListBuckets(ctx context.Context) ([]SchedulerBucket, error) { + return nil, nil +} + +func (c *outboxCleanupCache) GetOutboxWatermark(ctx context.Context) (int64, error) { + return c.watermark, nil +} + +func (c *outboxCleanupCache) SetOutboxWatermark(ctx context.Context, id int64) error { + c.watermark = id + c.setWatermarks = append(c.setWatermarks, id) + return nil +} + +type outboxCleanupDeleteCall struct { + watermark int64 + limit int +} + +type outboxCleanupRepo struct { + events []SchedulerOutboxEvent + rows []int64 + lockAcquired bool + lockAttempts int + releaseCount int + deleteCalls []outboxCleanupDeleteCall +} + +func (r *outboxCleanupRepo) ListAfterAndReleaseDedup(ctx context.Context, afterID int64, limit int) ([]SchedulerOutboxEvent, error) { + events := make([]SchedulerOutboxEvent, 0, len(r.events)) + for _, event := range r.events { + if event.ID <= afterID { + continue + } + events = append(events, event) + if limit > 0 && len(events) >= limit { + break + } + } + return events, nil +} + +func (r *outboxCleanupRepo) MaxID(ctx context.Context) (int64, error) { + var maxID int64 + for _, id := range r.rows { + if id > maxID { + maxID = id + } + } + return maxID, nil +} + +func (r *outboxCleanupRepo) DeleteConsumedUpTo(ctx context.Context, watermark int64, limit int) (int64, error) { + r.deleteCalls = append(r.deleteCalls, outboxCleanupDeleteCall{ + watermark: watermark, + limit: limit, + }) + if watermark <= 0 || limit <= 0 { + return 0, nil + } + + deleted := int64(0) + kept := make([]int64, 0, len(r.rows)) + for _, id := range r.rows { + if id <= watermark && deleted < int64(limit) { + deleted++ + continue + } + kept = append(kept, id) + } + r.rows = kept + return deleted, nil +} + +func (r *outboxCleanupRepo) TryAcquireCleanupLock(ctx context.Context) (SchedulerOutboxCleanupLease, bool, error) { + r.lockAttempts++ + if !r.lockAcquired { + return nil, false, nil + } + return outboxCleanupLease{release: func() { + r.releaseCount++ + }}, true, nil +} + +type outboxCleanupLease struct { + release func() +} + +func (l outboxCleanupLease) Release() { + if l.release != nil { + l.release() + } +} + +func TestSchedulerSnapshotServicePollOutboxCleansConsumedRowsAfterWatermark(t *testing.T) { + cache := &outboxCleanupCache{} + repo := &outboxCleanupRepo{ + events: []SchedulerOutboxEvent{ + {ID: 10000, EventType: SchedulerOutboxEventAccountLastUsed}, + }, + rows: int64Range(1, 10003), + lockAcquired: true, + } + svc := NewSchedulerSnapshotService(cache, repo, nil, nil, nil) + + svc.pollOutbox() + + if cache.watermark != 10000 { + t.Fatalf("expected watermark 10000, got %d", cache.watermark) + } + if !reflect.DeepEqual(cache.setWatermarks, []int64{10000}) { + t.Fatalf("unexpected watermark writes: %#v", cache.setWatermarks) + } + if !reflect.DeepEqual(repo.rows, []int64{10001, 10002, 10003}) { + t.Fatalf("expected rows above watermark to remain, got %#v", repo.rows) + } + if repo.lockAttempts != 1 || repo.releaseCount != 1 { + t.Fatalf("expected one lock acquire/release, got acquire=%d release=%d", repo.lockAttempts, repo.releaseCount) + } + if len(repo.deleteCalls) != 3 { + t.Fatalf("expected cleanup to loop until a short batch, got %d calls", len(repo.deleteCalls)) + } + for _, call := range repo.deleteCalls { + if call.watermark != 10000 || call.limit != schedulerOutboxCleanupBatch { + t.Fatalf("unexpected cleanup call: %#v", call) + } + } +} + +func TestSchedulerSnapshotServicePollOutboxSkipsCleanupWhenLockUnavailable(t *testing.T) { + cache := &outboxCleanupCache{} + repo := &outboxCleanupRepo{ + events: []SchedulerOutboxEvent{ + {ID: 3, EventType: SchedulerOutboxEventAccountLastUsed}, + }, + rows: []int64{1, 2, 3, 4}, + lockAcquired: false, + } + svc := NewSchedulerSnapshotService(cache, repo, nil, nil, nil) + + svc.pollOutbox() + + if cache.watermark != 3 { + t.Fatalf("expected watermark 3, got %d", cache.watermark) + } + if !reflect.DeepEqual(repo.rows, []int64{1, 2, 3, 4}) { + t.Fatalf("expected cleanup to skip all rows, got %#v", repo.rows) + } + if repo.lockAttempts != 1 { + t.Fatalf("expected one lock attempt, got %d", repo.lockAttempts) + } + if len(repo.deleteCalls) != 0 { + t.Fatalf("expected no delete calls, got %#v", repo.deleteCalls) + } + if repo.releaseCount != 0 { + t.Fatalf("expected no release without lock, got %d", repo.releaseCount) + } +} + +func TestSchedulerSnapshotServicePollOutboxDoesNotCleanupOnHandleFailure(t *testing.T) { + cache := &outboxCleanupCache{ + updateErr: errors.New("cache update failed"), + } + repo := &outboxCleanupRepo{ + events: []SchedulerOutboxEvent{ + { + ID: 5, + EventType: SchedulerOutboxEventAccountLastUsed, + Payload: map[string]any{ + "last_used": map[string]any{"101": float64(123)}, + }, + }, + }, + rows: []int64{1, 2, 3, 4, 5, 6}, + lockAcquired: true, + } + svc := NewSchedulerSnapshotService(cache, repo, nil, nil, nil) + + svc.pollOutbox() + + if len(cache.setWatermarks) != 0 { + t.Fatalf("expected no watermark write on handle failure, got %#v", cache.setWatermarks) + } + if repo.lockAttempts != 0 { + t.Fatalf("expected cleanup lock not to be attempted, got %d", repo.lockAttempts) + } + if len(repo.deleteCalls) != 0 { + t.Fatalf("expected no delete calls, got %#v", repo.deleteCalls) + } + if !reflect.DeepEqual(repo.rows, []int64{1, 2, 3, 4, 5, 6}) { + t.Fatalf("expected rows unchanged, got %#v", repo.rows) + } +} + +func TestSchedulerSnapshotServiceCleanupSkipsNonPositiveWatermark(t *testing.T) { + repo := &outboxCleanupRepo{ + rows: []int64{1, 2, 3}, + lockAcquired: true, + } + svc := NewSchedulerSnapshotService(&outboxCleanupCache{}, repo, nil, nil, nil) + + svc.cleanupConsumedOutbox(0) + + if repo.lockAttempts != 0 { + t.Fatalf("expected no lock attempt for non-positive watermark, got %d", repo.lockAttempts) + } + if len(repo.deleteCalls) != 0 { + t.Fatalf("expected no delete calls, got %#v", repo.deleteCalls) + } + if !reflect.DeepEqual(repo.rows, []int64{1, 2, 3}) { + t.Fatalf("expected rows unchanged, got %#v", repo.rows) + } +} + +func int64Range(start, end int64) []int64 { + values := make([]int64, 0, end-start+1) + for id := start; id <= end; id++ { + values = append(values, id) + } + return values +} diff --git a/backend/internal/service/scheduler_snapshot_service.go b/backend/internal/service/scheduler_snapshot_service.go index 6b25c0430f..fd48da27af 100644 --- a/backend/internal/service/scheduler_snapshot_service.go +++ b/backend/internal/service/scheduler_snapshot_service.go @@ -18,7 +18,10 @@ var ( ErrSchedulerFallbackLimited = errors.New("scheduler db fallback limited") ) -const outboxEventTimeout = 2 * time.Minute +const ( + outboxEventTimeout = 2 * time.Minute + schedulerOutboxCleanupBatch = 5000 +) // batchSeenKey tracks which (groupID, platform) bucket sets have already been // rebuilt within a single pollOutbox call, to avoid redundant work when multiple @@ -280,11 +283,42 @@ func (s *SchedulerSnapshotService) pollOutbox() { logger.LegacyPrintf("service.scheduler_snapshot", "[Scheduler] outbox watermark write failed: %v", wmErr) } else { watermarkForCheck = lastID + s.cleanupConsumedOutbox(lastID) } s.checkOutboxLag(ctx, events[0], watermarkForCheck) } +func (s *SchedulerSnapshotService) cleanupConsumedOutbox(watermark int64) { + if s == nil || s.outboxRepo == nil || watermark <= 0 { + return + } + + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + + lease, acquired, err := s.outboxRepo.TryAcquireCleanupLock(ctx) + if err != nil { + logger.LegacyPrintf("service.scheduler_snapshot", "[Scheduler] outbox cleanup lock failed: %v", err) + return + } + if !acquired { + return + } + defer lease.Release() + + for { + deleted, err := s.outboxRepo.DeleteConsumedUpTo(ctx, watermark, schedulerOutboxCleanupBatch) + if err != nil { + logger.LegacyPrintf("service.scheduler_snapshot", "[Scheduler] outbox cleanup failed: watermark=%d err=%v", watermark, err) + return + } + if deleted == 0 || deleted < schedulerOutboxCleanupBatch { + return + } + } +} + func (s *SchedulerSnapshotService) handleOutboxEvent(ctx context.Context, event SchedulerOutboxEvent, seen map[batchSeenKey]struct{}) error { switch event.EventType { case SchedulerOutboxEventAccountLastUsed: From 31dc8913acee0b0788cdda61be4c2a10f209f11d Mon Sep 17 00:00:00 2001 From: shaw Date: Tue, 16 Jun 2026 11:57:32 +0800 Subject: [PATCH 2/2] chore(outbox-cleanup): add 10s grace to defend against id-vs-commit race PG sequences advance outside transactions, so a slow committer can hold an id that gets surpassed by the watermark before its commit becomes visible. Without the grace period, cleanup would delete such rows before the snapshot poller ever sees them. 10s is a comfortable upper bound on realistic enqueueSchedulerOutbox commit latency. --- backend/internal/repository/scheduler_outbox_repo.go | 5 +++++ backend/internal/repository/scheduler_outbox_repo_test.go | 1 + 2 files changed, 6 insertions(+) diff --git a/backend/internal/repository/scheduler_outbox_repo.go b/backend/internal/repository/scheduler_outbox_repo.go index cea5d451f6..59772fbb5a 100644 --- a/backend/internal/repository/scheduler_outbox_repo.go +++ b/backend/internal/repository/scheduler_outbox_repo.go @@ -108,11 +108,16 @@ func (r *schedulerOutboxRepository) DeleteConsumedUpTo(ctx context.Context, wate if limit <= 0 { limit = schedulerOutboxDefaultCleanSize } + // created_at < NOW() - INTERVAL '10 seconds' 防御 PG 序列号在事务内提前分配但 + // 提交延迟的竞争:若某 Tx 在 watermark 推进前持有 id=N(未提交),watermark + // 跨过 N 后该 Tx 才提交,此时 row N 已经"低于 watermark"但从未被 poll;10s + // 宽限期让此类慢事务有机会提交后被消费,再被 cleanup 删除。 result, err := r.db.ExecContext(ctx, ` WITH doomed AS ( SELECT id FROM scheduler_outbox WHERE id <= $1 + AND created_at < NOW() - INTERVAL '10 seconds' ORDER BY id ASC LIMIT $2 ) diff --git a/backend/internal/repository/scheduler_outbox_repo_test.go b/backend/internal/repository/scheduler_outbox_repo_test.go index ea91f6c73e..6b4a1af9e9 100644 --- a/backend/internal/repository/scheduler_outbox_repo_test.go +++ b/backend/internal/repository/scheduler_outbox_repo_test.go @@ -20,6 +20,7 @@ func TestSchedulerOutboxRepositoryDeleteConsumedUpToUsesBoundedCTE(t *testing.T) SELECT id FROM scheduler_outbox WHERE id <= $1 + AND created_at < NOW() - INTERVAL '10 seconds' ORDER BY id ASC LIMIT $2 )