mirror of
https://github.com/Wei-Shaw/sub2api.git
synced 2026-09-24 16:05:44 +08:00
fix: release scheduler outbox dedup on claim
This commit is contained in:
@@ -59,13 +59,38 @@ func TestEnqueueSchedulerOutbox_DeduplicatesIdempotentEvents(t *testing.T) {
|
||||
|
||||
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}))
|
||||
events, err := NewSchedulerOutboxRepository(integrationDB).ListAfterAndReleaseDedup(ctx, 0, 100)
|
||||
require.NoError(t, err)
|
||||
require.Len(t, events, 1)
|
||||
require.Equal(t, firstID, events[0].ID)
|
||||
|
||||
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)
|
||||
}
|
||||
|
||||
func TestSchedulerOutbox_ListAfterAndReleaseDedup_AllowsSameKeyWhileEventInFlight(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
_, _ = integrationDB.ExecContext(ctx, "TRUNCATE scheduler_outbox RESTART IDENTITY")
|
||||
|
||||
accountID := int64(17345)
|
||||
require.NoError(t, enqueueSchedulerOutbox(ctx, integrationDB, service.SchedulerOutboxEventAccountChanged, &accountID, nil, nil))
|
||||
|
||||
events, err := NewSchedulerOutboxRepository(integrationDB).ListAfterAndReleaseDedup(ctx, 0, 100)
|
||||
require.NoError(t, err)
|
||||
require.Len(t, events, 1)
|
||||
|
||||
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))
|
||||
require.Equal(t, 2, count)
|
||||
|
||||
var pendingKeys int
|
||||
require.NoError(t, integrationDB.QueryRowContext(ctx, "SELECT COUNT(*) FROM scheduler_outbox WHERE dedup_key IS NOT NULL").Scan(&pendingKeys))
|
||||
require.Equal(t, 1, pendingKeys)
|
||||
}
|
||||
|
||||
func TestEnqueueSchedulerOutbox_CoalescesAccountStateBurst(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
_, _ = integrationDB.ExecContext(ctx, "TRUNCATE scheduler_outbox RESTART IDENTITY")
|
||||
|
||||
@@ -10,7 +10,6 @@ import (
|
||||
"strconv"
|
||||
|
||||
"github.com/Wei-Shaw/sub2api/internal/service"
|
||||
"github.com/lib/pq"
|
||||
)
|
||||
|
||||
type schedulerOutboxRepository struct {
|
||||
@@ -21,16 +20,30 @@ func NewSchedulerOutboxRepository(db *sql.DB) service.SchedulerOutboxRepository
|
||||
return &schedulerOutboxRepository{db: db}
|
||||
}
|
||||
|
||||
func (r *schedulerOutboxRepository) ListAfter(ctx context.Context, afterID int64, limit int) ([]service.SchedulerOutboxEvent, error) {
|
||||
func (r *schedulerOutboxRepository) ListAfterAndReleaseDedup(ctx context.Context, afterID int64, limit int) ([]service.SchedulerOutboxEvent, error) {
|
||||
if limit <= 0 {
|
||||
limit = 100
|
||||
}
|
||||
rows, err := r.db.QueryContext(ctx, `
|
||||
SELECT id, event_type, account_id, group_id, payload, created_at
|
||||
FROM scheduler_outbox
|
||||
WHERE id > $1
|
||||
ORDER BY id ASC
|
||||
LIMIT $2
|
||||
WITH selected AS MATERIALIZED (
|
||||
SELECT id, event_type, account_id, group_id, payload, created_at
|
||||
FROM scheduler_outbox
|
||||
WHERE id > $1
|
||||
ORDER BY id ASC
|
||||
LIMIT $2
|
||||
FOR UPDATE
|
||||
), released AS (
|
||||
UPDATE scheduler_outbox AS o
|
||||
SET dedup_key = NULL
|
||||
FROM selected AS s
|
||||
WHERE o.id = s.id
|
||||
AND o.dedup_key IS NOT NULL
|
||||
RETURNING o.id
|
||||
)
|
||||
SELECT s.id, s.event_type, s.account_id, s.group_id, s.payload, s.created_at
|
||||
FROM selected AS s
|
||||
CROSS JOIN (SELECT COUNT(*) FROM released) AS release_barrier
|
||||
ORDER BY s.id ASC
|
||||
`, afterID, limit)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
@@ -81,19 +94,6 @@ 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
|
||||
|
||||
@@ -16,7 +16,6 @@ type SchedulerOutboxEvent struct {
|
||||
|
||||
// SchedulerOutboxRepository 提供调度 outbox 的读取接口。
|
||||
type SchedulerOutboxRepository interface {
|
||||
ListAfter(ctx context.Context, afterID int64, limit int) ([]SchedulerOutboxEvent, error)
|
||||
ListAfterAndReleaseDedup(ctx context.Context, afterID int64, limit int) ([]SchedulerOutboxEvent, error)
|
||||
MaxID(ctx context.Context) (int64, error)
|
||||
MarkProcessed(ctx context.Context, eventIDs []int64) error
|
||||
}
|
||||
|
||||
@@ -242,7 +242,7 @@ func (s *SchedulerSnapshotService) pollOutbox() {
|
||||
return
|
||||
}
|
||||
|
||||
events, err := s.outboxRepo.ListAfter(ctx, watermark, 200)
|
||||
events, err := s.outboxRepo.ListAfterAndReleaseDedup(ctx, watermark, 200)
|
||||
if err != nil {
|
||||
logger.LegacyPrintf("service.scheduler_snapshot", "[Scheduler] outbox poll failed: %v", err)
|
||||
return
|
||||
@@ -253,7 +253,6 @@ 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)
|
||||
@@ -262,15 +261,6 @@ 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
|
||||
|
||||
Reference in New Issue
Block a user