From bfd1257db0a4a9f1d7cd65f9e4d6bc2382622f91 Mon Sep 17 00:00:00 2001 From: wangbo Date: Fri, 31 Jul 2026 07:15:44 +0800 Subject: [PATCH] =?UTF-8?q?fix(worker):=20=E7=BB=9F=E4=B8=80=E5=BC=82?= =?UTF-8?q?=E6=AD=A5=E6=89=A7=E8=A1=8C=E5=87=86=E5=85=A5=E8=8C=83=E5=9B=B4?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Worker 执行阶段重建准入范围时遗漏 worker_capacity,导致已预留 River 槽的任务再次同步等待准入并占住执行槽。 调度、执行和候选切换统一复用 taskAdmissionScopes,确保异步任务始终携带同一全局 Worker 容量范围。 验证:Go 全量测试、runner race、go vet、迁移安全检查通过。 --- apps/api/internal/runner/admission.go | 38 ++++++++++++++++----------- apps/api/internal/runner/service.go | 6 +++-- 2 files changed, 27 insertions(+), 17 deletions(-) diff --git a/apps/api/internal/runner/admission.go b/apps/api/internal/runner/admission.go index fd91566..1697e3a 100644 --- a/apps/api/internal/runner/admission.go +++ b/apps/api/internal/runner/admission.go @@ -113,14 +113,10 @@ func (s *Service) buildTaskAdmissionPlan(ctx context.Context, task store.Gateway if !available { continue } - scopes, groupID := s.admissionScopes(ctx, user, candidate) - if task.AsyncMode { - scopes, err = s.withAsyncWorkerCapacityScope(ctx, scopes) - if err != nil { - return taskAdmissionPlan{}, err - } + scopes, groupID, scopeErr := s.taskAdmissionScopes(ctx, task, user, candidate) + if scopeErr != nil { + return taskAdmissionPlan{}, scopeErr } - scopes = acceptanceAdmissionScopes(task, scopes) hasConcurrentLimit := false for _, scope := range scopes { if scope.ConcurrentLimit > 0 { @@ -172,6 +168,23 @@ func (s *Service) withAsyncWorkerCapacityScope(ctx context.Context, scopes []sto return out, nil } +func (s *Service) taskAdmissionScopes( + ctx context.Context, + task store.GatewayTask, + user *auth.User, + candidate store.RuntimeModelCandidate, +) ([]store.AdmissionScope, string, error) { + scopes, groupID := s.admissionScopes(ctx, user, candidate) + if task.AsyncMode { + var err error + scopes, err = s.withAsyncWorkerCapacityScope(ctx, scopes) + if err != nil { + return nil, "", err + } + } + return acceptanceAdmissionScopes(task, scopes), groupID, nil +} + func acceptanceAdmissionScopes(task store.GatewayTask, scopes []store.AdmissionScope) []store.AdmissionScope { if task.RunMode != "acceptance" && task.RunMode != "acceptance_canary" { return scopes @@ -357,15 +370,10 @@ func (s *Service) ensureCandidateAdmission( body map[string]any, candidate store.RuntimeModelCandidate, ) (store.TaskAdmissionResult, bool, error) { - scopes, groupID := s.admissionScopes(ctx, user, candidate) - if task.AsyncMode { - var err error - scopes, err = s.withAsyncWorkerCapacityScope(ctx, scopes) - if err != nil { - return store.TaskAdmissionResult{}, true, err - } + scopes, groupID, err := s.taskAdmissionScopes(ctx, task, user, candidate) + if err != nil { + return store.TaskAdmissionResult{}, true, err } - scopes = acceptanceAdmissionScopes(task, scopes) hasConcurrentLimit := false for _, scope := range scopes { if scope.ConcurrentLimit > 0 { diff --git a/apps/api/internal/runner/service.go b/apps/api/internal/runner/service.go index 970698d..dc21835 100644 --- a/apps/api/internal/runner/service.go +++ b/apps/api/internal/runner/service.go @@ -544,8 +544,10 @@ func (s *Service) executeWithToken(ctx context.Context, task store.GatewayTask, var admittedLeases []store.ConcurrencyLease admittedPlatformModelID := "" if distributedAdmission && len(candidates) > 0 { - admissionScopes, groupID := s.admissionScopes(ctx, user, candidates[0]) - admissionScopes = acceptanceAdmissionScopes(task, admissionScopes) + admissionScopes, groupID, scopeErr := s.taskAdmissionScopes(ctx, task, user, candidates[0]) + if scopeErr != nil { + return Result{}, scopeErr + } hasConcurrentLimit := false for _, scope := range admissionScopes { if scope.ConcurrentLimit > 0 {