From d1b482c2f322b59d34c0eb221c0cbab969123d66 Mon Sep 17 00:00:00 2001 From: wangbo Date: Fri, 31 Jul 2026 07:23:46 +0800 Subject: [PATCH] =?UTF-8?q?fix(worker):=20=E5=AE=B9=E9=87=8F=E9=A5=B1?= =?UTF-8?q?=E5=92=8C=E5=90=8E=E5=81=9C=E6=AD=A2=E9=87=8D=E5=A4=8D=E5=87=86?= =?UTF-8?q?=E5=85=A5=E6=89=AB=E6=8F=8F?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 双 Worker 每秒遍历全部 waiting 任务,在全局执行槽已满时仍重复探测,导致 River 执行协程长时间等待同一 worker_capacity 进程锁。 异步 FIFO 头部未准入时立即结束本轮扫描,等待租约释放通知后再继续;任务级错误仍跳过并处理后续任务。 验证:Go 全量测试、runner race、go vet、迁移安全检查通过。 --- apps/api/internal/runner/admission.go | 27 +++++++++++++++++++++------ 1 file changed, 21 insertions(+), 6 deletions(-) diff --git a/apps/api/internal/runner/admission.go b/apps/api/internal/runner/admission.go index 1697e3a..1d07937 100644 --- a/apps/api/internal/runner/admission.go +++ b/apps/api/internal/runner/admission.go @@ -603,19 +603,25 @@ func (s *Service) SubmitAsyncTask(ctx context.Context, task store.GatewayTask) e return nil } -func (s *Service) dispatchWaitingAsyncTask(ctx context.Context, task store.GatewayTask) error { +func (s *Service) dispatchWaitingAsyncTask(ctx context.Context, task store.GatewayTask) (bool, error) { user := authUserFromTask(task) plan, err := s.buildTaskAdmissionPlan(ctx, task, user) if err != nil { - return err + return false, err } if !plan.Eligible { - return s.EnqueueAsyncTask(ctx, task) + if err := s.EnqueueAsyncTask(ctx, task); err != nil { + return false, err + } + return true, nil } - _, err = s.tryTaskAdmissionWithAdmittedHook(ctx, task, plan, "", func(tx pgx.Tx) error { + result, err := s.tryTaskAdmissionWithAdmittedHook(ctx, task, plan, "", func(tx pgx.Tx) error { return s.enqueueAsyncTaskTx(ctx, tx, task.ID, asyncTaskInsertOpts(task)) }) - return err + if err != nil { + return false, err + } + return result.Admitted, nil } func (s *Service) cancelAsyncSubmissionIfDisconnected(ctx context.Context, taskID string) bool { @@ -660,8 +666,17 @@ func (s *Service) dispatchWaitingAsyncAdmissions(ctx context.Context) { if task.Status != "queued" { continue } - if err := s.dispatchWaitingAsyncTask(ctx, task); err != nil { + admitted, err := s.dispatchWaitingAsyncTask(ctx, task) + if err != nil { s.logger.Warn("dispatch waiting async admission failed", "taskID", taskID, "error", err) + continue + } + if !admitted { + // Every asynchronous task shares the global worker-capacity FIFO + // scope. Once its current head cannot be admitted, later tasks + // cannot make progress until a lease is released. + s.observeTaskAdmission("dispatch_saturated") + break } } }