From 34e66ec0a5453b54739d21fb3fd4196a183cc7cb Mon Sep 17 00:00:00 2001 From: jjaw Date: Fri, 12 Jun 2026 02:07:28 +0800 Subject: [PATCH] fix: outbox scheduler snapshot coalesce --- .../account_repo_integration_test.go | 31 +++++++++++++++++++ .../ops_write_pressure_integration_test.go | 21 +++++++++++-- .../repository/scheduler_outbox_repo.go | 2 +- 3 files changed, 51 insertions(+), 3 deletions(-) diff --git a/backend/internal/repository/account_repo_integration_test.go b/backend/internal/repository/account_repo_integration_test.go index c4d65a665c..b216ecd04d 100644 --- a/backend/internal/repository/account_repo_integration_test.go +++ b/backend/internal/repository/account_repo_integration_test.go @@ -752,6 +752,37 @@ func (s *AccountRepoSuite) TestTempUnschedulableFieldsLoadedByGetByIDAndGetByIDs s.Require().Equal("", cacheRecorder.setAccounts[0].TempUnschedulableReason) } +func (s *AccountRepoSuite) TestSetTempUnschedulableSkipsOutboxWhenWindowDoesNotExtend() { + account := mustCreateAccount(s.T(), s.client, &service.Account{Name: "acc-temp-noop"}) + cacheRecorder := &schedulerCacheRecorder{} + s.repo.schedulerCache = cacheRecorder + + _, err := s.repo.sql.ExecContext(s.ctx, "TRUNCATE scheduler_outbox") + s.Require().NoError(err) + + until := time.Now().Add(30 * time.Minute).UTC().Truncate(time.Second) + s.Require().NoError(s.repo.SetTempUnschedulable(s.ctx, account.ID, until, "first")) + + var count int + err = scanSingleRow(s.ctx, s.repo.sql, "SELECT COUNT(*) FROM scheduler_outbox", nil, &count) + s.Require().NoError(err) + s.Require().Equal(1, count) + s.Require().Len(cacheRecorder.setAccounts, 1) + + s.Require().NoError(s.repo.SetTempUnschedulable(s.ctx, account.ID, until.Add(-5*time.Minute), "older")) + + err = scanSingleRow(s.ctx, s.repo.sql, "SELECT COUNT(*) FROM scheduler_outbox", nil, &count) + s.Require().NoError(err) + s.Require().Equal(1, count) + s.Require().Len(cacheRecorder.setAccounts, 1) + + got, err := s.repo.GetByID(s.ctx, account.ID) + s.Require().NoError(err) + s.Require().Equal("first", got.TempUnschedulableReason) + s.Require().NotNil(got.TempUnschedulableUntil) + s.Require().WithinDuration(until, *got.TempUnschedulableUntil, time.Second) +} + func (s *AccountRepoSuite) TestClearModelRateLimits_SyncsSchedulerSnapshot() { account := mustCreateAccount(s.T(), s.client, &service.Account{ Name: "acc-clear-model-rate", diff --git a/backend/internal/repository/ops_write_pressure_integration_test.go b/backend/internal/repository/ops_write_pressure_integration_test.go index ebb7a84226..4520891ec6 100644 --- a/backend/internal/repository/ops_write_pressure_integration_test.go +++ b/backend/internal/repository/ops_write_pressure_integration_test.go @@ -57,10 +57,27 @@ 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) - time.Sleep(schedulerOutboxDedupWindow + 150*time.Millisecond) + require.GreaterOrEqual(t, schedulerOutboxDedupWindow, 10*time.Second) + + 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, 2, count) + require.Equal(t, 1, count) +} + +func TestEnqueueSchedulerOutbox_CoalescesAccountStateBurst(t *testing.T) { + ctx := context.Background() + _, _ = integrationDB.ExecContext(ctx, "TRUNCATE scheduler_outbox RESTART IDENTITY") + + accountID := int64(22345) + for range 50 { + 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)) + t.Logf("same-account account_changed burst: calls=50 inserted=%d dedup_window=%s", count, schedulerOutboxDedupWindow) + require.Equal(t, 1, count) } func TestEnqueueSchedulerOutbox_DoesNotDeduplicateLastUsed(t *testing.T) { diff --git a/backend/internal/repository/scheduler_outbox_repo.go b/backend/internal/repository/scheduler_outbox_repo.go index 4b9a9f58b1..586efd0582 100644 --- a/backend/internal/repository/scheduler_outbox_repo.go +++ b/backend/internal/repository/scheduler_outbox_repo.go @@ -13,7 +13,7 @@ type schedulerOutboxRepository struct { db *sql.DB } -const schedulerOutboxDedupWindow = time.Second +const schedulerOutboxDedupWindow = 10 * time.Second func NewSchedulerOutboxRepository(db *sql.DB) service.SchedulerOutboxRepository { return &schedulerOutboxRepository{db: db}