From 14d13ad3a7efcd3c959a9228abad7d60162a3f5f Mon Sep 17 00:00:00 2001 From: wangbo Date: Fri, 31 Jul 2026 07:46:47 +0800 Subject: [PATCH] =?UTF-8?q?fix(worker):=20=E9=9A=94=E7=A6=BB=E4=BB=BB?= =?UTF-8?q?=E5=8A=A1=E5=88=9B=E5=BB=BA=E4=B8=8E=E9=98=9F=E5=88=97=E6=81=A2?= =?UTF-8?q?=E5=A4=8D=E7=AA=97=E5=8F=A3?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 高媒体并发下,任务创建后的请求恢复和候选准入可能持续数十秒;通用 River 恢复扫描会在准入记录落库前误把新任务当成崩溃遗留任务,提前创建 Job 并占满 Worker 执行槽。\n\n为无活动 River Job 的通用兜底恢复增加 5 分钟新任务保护窗。正常提交、准入调度以及已有 River Job 的 Worker 强杀恢复保持原路径;新增 PostgreSQL 集成回归,验证新任务不被抢占且超过保护窗后仍可恢复。\n\n验证:Go 全量测试、go vet、runner/store race、gofmt、迁移安全检查通过。 --- .../store/admission_queue_integration_test.go | 30 +++++++++++++++++++ apps/api/internal/store/tasks_runtime.go | 8 ++++- 2 files changed, 37 insertions(+), 1 deletion(-) diff --git a/apps/api/internal/store/admission_queue_integration_test.go b/apps/api/internal/store/admission_queue_integration_test.go index a852b47..b1c90ac 100644 --- a/apps/api/internal/store/admission_queue_integration_test.go +++ b/apps/api/internal/store/admission_queue_integration_test.go @@ -471,6 +471,36 @@ WHERE id = $1::uuid`, queuedAtomicTask.ID).Scan(&riverJobID, &queuedLeases); err t.Fatal("generic River recovery claimed a task owned by the async admission dispatcher") } } + recoveryGraceTask := createTask(true) + recoverableTasks, err = first.ListRecoverableAsyncTasks(ctx, 1000) + if err != nil { + t.Fatalf("list fresh recoverable tasks: %v", err) + } + for _, item := range recoverableTasks { + if item.ID == recoveryGraceTask.ID { + t.Fatal("generic River recovery claimed a task still preparing admission") + } + } + if _, err := first.pool.Exec(ctx, ` +UPDATE gateway_tasks +SET created_at = now() - interval '6 minutes' +WHERE id = $1::uuid`, recoveryGraceTask.ID); err != nil { + t.Fatalf("age task beyond River recovery grace period: %v", err) + } + recoverableTasks, err = first.ListRecoverableAsyncTasks(ctx, 1000) + if err != nil { + t.Fatalf("list aged recoverable tasks: %v", err) + } + recoveryGraceFound := false + for _, item := range recoverableTasks { + if item.ID == recoveryGraceTask.ID { + recoveryGraceFound = true + break + } + } + if !recoveryGraceFound { + t.Fatal("generic River recovery did not claim a task beyond the recovery grace period") + } result, err = second.TryTaskAdmission(ctx, inputFor(queuedAtomicTask, 100, "")) if err != nil || !result.Admitted || len(result.Leases) != 1 { t.Fatalf("worker-time queued admission result=%+v err=%v", result, err) diff --git a/apps/api/internal/store/tasks_runtime.go b/apps/api/internal/store/tasks_runtime.go index ae2940c..fd594c0 100644 --- a/apps/api/internal/store/tasks_runtime.go +++ b/apps/api/internal/store/tasks_runtime.go @@ -1043,11 +1043,17 @@ func (s *Store) ListRecoverableAsyncTasks(ctx context.Context, limit int) ([]Asy if limit <= 0 { limit = 500 } + // API handlers create the durable task before they finish media restoration + // and register distributed admission. Give that normal submission path + // ownership of fresh tasks so the generic crash recovery scan cannot create + // a River job inside the create-to-admission window. + const minimumRecoveryAge = 5 * time.Minute rows, err := s.pool.Query(ctx, ` SELECT task.id::text, task.priority, task.next_run_at FROM gateway_tasks task LEFT JOIN river_job job ON job.id = task.river_job_id WHERE task.async_mode = true + AND task.created_at <= now() - $2::interval AND ( task.status = 'queued' OR ( @@ -1067,7 +1073,7 @@ WHERE task.async_mode = true AND admission.status = 'waiting' ) ORDER BY task.priority ASC, task.created_at ASC -LIMIT $1`, limit) +LIMIT $1`, limit, minimumRecoveryAge.String()) if err != nil { return nil, err }