fix: safely coalesce scheduler outbox events

This commit is contained in:
jjaw
2026-06-16 11:48:12 +08:00
committed by shaw
parent 34e66ec0a5
commit 3ef70b045d
7 changed files with 81 additions and 18 deletions
@@ -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))
@@ -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")
@@ -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,
@@ -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
}
@@ -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
@@ -0,0 +1,2 @@
ALTER TABLE scheduler_outbox
ADD COLUMN IF NOT EXISTS dedup_key TEXT;
@@ -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;