From 4455f8938d6580d5b8cf757c58a88ff0e6a62572 Mon Sep 17 00:00:00 2001 From: Rogee Date: Mon, 24 Aug 2026 16:46:35 +0800 Subject: [PATCH] HH-574: enforce stale job claim limit (#168) * fix(HH-574): enforce stale job claim limit * fix(HH-574): guard exhausted retry claims --------- Co-authored-by: Rogee --- backend/internal/worker/worker.go | 34 +++++++++-- backend/internal/worker/worker_test.go | 60 +++++++++++++++++++ .../message-lifecycle-and-ai-flow.md | 7 ++- .../worker-redis-queue-migration.md | 4 +- 4 files changed, 95 insertions(+), 10 deletions(-) diff --git a/backend/internal/worker/worker.go b/backend/internal/worker/worker.go index f192dea6..96cc30c8 100644 --- a/backend/internal/worker/worker.go +++ b/backend/internal/worker/worker.go @@ -235,6 +235,8 @@ func WithScheduledAt(at time.Time) EnqueueOption { } } +// WithMaxAttempts sets the absolute claim limit. A claim interrupted by a +// process crash still consumes an attempt. func WithMaxAttempts(max int) EnqueueOption { return func(job *model.BackgroundJob) { if max > 0 { @@ -501,9 +503,15 @@ func (wp *WorkerPool) RequeueStaleJobs(ctx context.Context) (int64, error) { if wp.db == nil { return 0, ErrWorkerDatabaseRequired } - cutoff := wp.now().Add(-wp.staleLockTimeout) + now := wp.now() + cutoff := now.Add(-wp.staleLockTimeout) updates := map[string]any{ - "status": model.BackgroundJobStatusRetrying, + "status": gorm.Expr( + "CASE WHEN attempts >= max_attempts THEN ? ELSE ? END", + model.BackgroundJobStatusDead, + model.BackgroundJobStatusRetrying, + ), + "failed_at": gorm.Expr("CASE WHEN attempts >= max_attempts THEN ? ELSE failed_at END", now), "locked_at": nil, "locked_by": "", } @@ -513,6 +521,15 @@ func (wp *WorkerPool) RequeueStaleJobs(ctx context.Context) (int64, error) { return result.RowsAffected, result.Error } +func archiveExhaustedRetries(db *gorm.DB, now time.Time) error { + return db.Model(&model.BackgroundJob{}). + Where("status = ? AND attempts >= max_attempts", model.BackgroundJobStatusRetrying). + Updates(map[string]any{ + "status": model.BackgroundJobStatusDead, + "failed_at": now, + }).Error +} + // run is the per-goroutine consume loop. When Redis is configured it uses // XREADGROUP BLOCK; otherwise it falls back to DB polling. func (wp *WorkerPool) run(ctx, jobCtx context.Context, index int) { @@ -620,6 +637,11 @@ func (wp *WorkerPool) processRedisMessage(lifecycleCtx, jobCtx context.Context, applogger.L().Errorf("load job %d from DB failed: %v", jobID, err) return } + now := wp.now() + if err := archiveExhaustedRetries(wp.db.WithContext(jobCtx).Where("id = ?", jobID), now); err != nil { + applogger.L().Errorf("archive exhausted job %d failed: %v", jobID, err) + return + } // Claim check: atomically transition queued/retrying → running. // RowsAffected == 0 means another consumer already claimed it or the @@ -629,9 +651,8 @@ func (wp *WorkerPool) processRedisMessage(lifecycleCtx, jobCtx context.Context, wp.claimMu.RUnlock() return } - now := wp.now() result := wp.db.WithContext(jobCtx).Model(&model.BackgroundJob{}). - Where("id = ? AND status IN ? AND scheduled_at <= ?", + Where("id = ? AND status IN ? AND scheduled_at <= ? AND attempts < max_attempts", jobID, []string{model.BackgroundJobStatusQueued, model.BackgroundJobStatusRetrying}, now, @@ -780,9 +801,12 @@ func (wp *WorkerPool) claimNext(lifecycleCtx, jobCtx context.Context) (*model.Ba return nil, err } } + if err := archiveExhaustedRetries(wp.db.WithContext(jobCtx), wp.now()); err != nil { + return nil, err + } var job model.BackgroundJob err := wp.db.WithContext(jobCtx).Transaction(func(tx *gorm.DB) error { - query := tx.Where("status IN ? AND scheduled_at <= ?", []string{model.BackgroundJobStatusQueued, model.BackgroundJobStatusRetrying}, wp.now()) + query := tx.Where("status IN ? AND scheduled_at <= ? AND attempts < max_attempts", []string{model.BackgroundJobStatusQueued, model.BackgroundJobStatusRetrying}, wp.now()) if len(wp.queues) > 0 { query = query.Where("queue IN ?", wp.queues) } diff --git a/backend/internal/worker/worker_test.go b/backend/internal/worker/worker_test.go index e21dd3f1..692d15b1 100644 --- a/backend/internal/worker/worker_test.go +++ b/backend/internal/worker/worker_test.go @@ -207,6 +207,66 @@ func TestWorkerPoolRequeuesStaleRunningJobs(t *testing.T) { } } +func TestWorkerPoolDeadLettersJobWhenFinalClaimGoesStale(t *testing.T) { + db := newWorkerTestDB(t) + now := time.Date(2026, 6, 5, 10, 0, 0, 0, time.UTC) + wp := NewWorkerPoolWithOptions(db, WithNow(func() time.Time { return now }), WithStaleLockTimeout(time.Minute)) + job, err := wp.Enqueue(context.Background(), "crashing_job", nil, WithMaxAttempts(3)) + require.NoError(t, err) + + for attempt := 1; attempt <= job.MaxAttempts; attempt++ { + claimed, claimErr := wp.claimNext(nil, context.Background()) + require.NoError(t, claimErr) + require.Equal(t, attempt, claimed.Attempts) + + now = now.Add(2 * time.Minute) + count, requeueErr := wp.RequeueStaleJobs(context.Background()) + require.NoError(t, requeueErr) + require.EqualValues(t, 1, count) + } + + reloaded := loadJob(t, db, job.ID) + require.Equal(t, model.BackgroundJobStatusDead, reloaded.Status) + require.Equal(t, job.MaxAttempts, reloaded.Attempts) + require.NotNil(t, reloaded.FailedAt) + processed, err := wp.ProcessOne(context.Background()) + require.NoError(t, err) + require.False(t, processed, "an exhausted stale job must not receive a fourth claim") +} + +func TestWorkerPoolArchivesPersistedExhaustedRetry(t *testing.T) { + db := newWorkerTestDB(t) + now := time.Date(2026, 6, 5, 10, 0, 0, 0, time.UTC) + wp := NewWorkerPoolWithOptions(db, WithNow(func() time.Time { return now })) + var handled atomic.Int32 + wp.Register("legacy_retry", func(context.Context, *model.BackgroundJob) error { + handled.Add(1) + return nil + }) + job, err := wp.Enqueue(context.Background(), "legacy_retry", nil, WithMaxAttempts(3)) + require.NoError(t, err) + require.NoError(t, db.Model(job).Updates(map[string]any{ + "status": model.BackgroundJobStatusRetrying, + "attempts": job.MaxAttempts, + }).Error) + + processed, err := wp.ProcessOne(context.Background()) + require.NoError(t, err) + require.False(t, processed) + require.Zero(t, handled.Load()) + reloaded := loadJob(t, db, job.ID) + require.Equal(t, model.BackgroundJobStatusDead, reloaded.Status) + require.Equal(t, job.MaxAttempts, reloaded.Attempts) + require.NotNil(t, reloaded.FailedAt) + require.Equal(t, now, *reloaded.FailedAt) + failedAt := *reloaded.FailedAt + now = now.Add(time.Hour) + processed, err = wp.ProcessOne(context.Background()) + require.NoError(t, err) + require.False(t, processed) + require.Equal(t, failedAt, *loadJob(t, db, job.ID).FailedAt) +} + func TestWorkerPoolStartAndStopProcessJobs(t *testing.T) { db := newWorkerTestDB(t) now := time.Date(2026, 6, 5, 10, 0, 0, 0, time.UTC) diff --git a/docs/architecture/message-lifecycle-and-ai-flow.md b/docs/architecture/message-lifecycle-and-ai-flow.md index 6c479f62..024c747d 100644 --- a/docs/architecture/message-lifecycle-and-ai-flow.md +++ b/docs/architecture/message-lifecycle-and-ai-flow.md @@ -159,9 +159,10 @@ SSE 不保存离线或慢消费者的确认状态,所以该契约保证发布 该保证从 job 成功写入 `background_jobs` 开始。事务内入队失败会使业务事务回滚; 非事务入队按 account、pubsub token 顺序逐个执行,失败 target 没有 durable job, -同一事件的其他 target 可能已成功入队。每个 target job 的 `MaxAttempts=3`;达到上限 -(或遇到 permanent error)进入 `dead` 后不会再被自动 -重试,也不再保证发布,消费者仍需通过 REST reconciliation 恢复持久化状态。 +同一事件的其他 target 可能已成功入队。每个 target job 的 `MaxAttempts=3` 是绝对 +claim 次数上限(进程在 claim 后崩溃也消耗一次);达到上限(或遇到 +permanent error)进入 `dead` 后不会再被自动重试,也不再保证发布,消费者仍需 +通过 REST reconciliation 恢复持久化状态。 消费者必须把 realtime 事件作为可重放通知处理: diff --git a/docs/architecture/worker-redis-queue-migration.md b/docs/architecture/worker-redis-queue-migration.md index fa7e066e..0130c877 100644 --- a/docs/architecture/worker-redis-queue-migration.md +++ b/docs/architecture/worker-redis-queue-migration.md @@ -429,7 +429,7 @@ func (wp *WorkerPool) sweepDueJobs(ctx context.Context) error { **当前**:扫描 `status=running AND locked_at < cutoff`,重置为 `retrying`。 -**迁移后**:逻辑不变,仍由 `Start()` 在启动时调用一次(回收上次崩溃时正在执行的 job)。但新增:在 sweep 循环中也定期调用,回收当前实例运行期间卡住的 job。 +**迁移后**:仍由 `Start()` 在启动时调用一次(回收上次崩溃时正在执行的 job),并在 sweep 循环中定期调用。stale job 未耗尽 `max_attempts` 时转为 `retrying`;已耗尽时原子转为 `dead` 并写入 `failed_at`,不再发生下一次 claim。 ```go // Start() 中调用(不变) @@ -447,7 +447,7 @@ func (wp *WorkerPool) sweepLoop(ctx context.Context) { } ``` -**影响分析**:`RequeueStaleJobs` 方法签名和 DB 操作不变。stale lock 回收后 job 状态变为 `retrying`,会被下一轮 sweep 扫描到并重新 XADD 到 Redis。 +**影响分析**:`RequeueStaleJobs` 方法签名不变。只有未耗尽的 stale job 会被下一轮 sweep 重新 XADD 到 Redis;耗尽的 job 保持 `dead` 终态。 ### 3.6 Start/Stop 变更