fix: outbox scheduler snapshot coalesce

This commit is contained in:
jjaw
2026-06-16 11:48:12 +08:00
committed by shaw
parent 45f3b0dd74
commit 34e66ec0a5
3 changed files with 51 additions and 3 deletions
@@ -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",
@@ -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) {
@@ -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}