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 <rogee@ipao.vip>
This commit is contained in:
Rogee
2026-08-24 16:46:35 +08:00
committed by GitHub
co-authored by rogee
parent 84d701dce5
commit 4455f8938d
4 changed files with 95 additions and 10 deletions
+29 -5
View File
@@ -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)
}
+60
View File
@@ -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)
@@ -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 事件作为可重放通知处理:
@@ -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 变更