fix(worker): 固定异步准入候选避免执行槽空转

异步任务在调度器准入后,Worker 会再次按实时负载排序候选;排序变化会让已准入任务退回 waiting,同时继续占用 River 执行槽,导致高并发吞吐塌陷。

执行前复用已持久化的 admitted 候选并稳定置顶,候选失效时仍保留原有重选路径;增加候选固定与非准入场景单测。

验证:Go 全量测试、runner race、go vet、OpenAPI 生成一致性、迁移安全检查通过。
This commit is contained in:
2026-07-31 07:06:06 +08:00
parent 3d9ae74b87
commit 238798e47c
3 changed files with 109 additions and 10 deletions
+42 -9
View File
@@ -187,6 +187,45 @@ func acceptanceAdmissionScopes(task store.GatewayTask, scopes []store.AdmissionS
return out
}
func (s *Service) loadAsyncTaskAdmission(ctx context.Context, task store.GatewayTask) (*store.TaskAdmission, error) {
if !task.AsyncMode {
return nil, nil
}
admission, err := s.store.GetTaskAdmission(ctx, task.ID)
if errors.Is(err, pgx.ErrNoRows) {
return nil, nil
}
if err != nil {
return nil, err
}
return &admission, nil
}
func pinCandidatesToTaskAdmission(
candidates []store.RuntimeModelCandidate,
admission *store.TaskAdmission,
) ([]store.RuntimeModelCandidate, bool) {
if admission == nil || admission.Status != "admitted" || len(candidates) < 2 {
return candidates, false
}
pinnedIndex := -1
for index := range candidates {
if candidates[index].PlatformID == admission.PlatformID &&
candidates[index].PlatformModelID == admission.PlatformModelID {
pinnedIndex = index
break
}
}
if pinnedIndex <= 0 {
return candidates, false
}
out := append([]store.RuntimeModelCandidate(nil), candidates...)
pinned := out[pinnedIndex]
copy(out[1:pinnedIndex+1], out[:pinnedIndex])
out[0] = pinned
return out, true
}
func (s *Service) tryTaskAdmission(ctx context.Context, task store.GatewayTask, plan taskAdmissionPlan, waiterID string) (store.TaskAdmissionResult, error) {
return s.tryTaskAdmissionWithAdmittedHook(ctx, task, plan, waiterID, nil)
}
@@ -195,17 +234,11 @@ func (s *Service) activeAsyncTaskAdmission(
ctx context.Context,
task store.GatewayTask,
plan taskAdmissionPlan,
admission *store.TaskAdmission,
) (store.TaskAdmissionResult, bool, error) {
if !task.AsyncMode {
if !task.AsyncMode || admission == nil {
return store.TaskAdmissionResult{}, false, nil
}
admission, err := s.store.GetTaskAdmission(ctx, task.ID)
if errors.Is(err, pgx.ErrNoRows) {
return store.TaskAdmissionResult{}, false, nil
}
if err != nil {
return store.TaskAdmissionResult{}, false, err
}
if admission.Status != "admitted" ||
admission.PlatformID != plan.Candidate.PlatformID ||
admission.PlatformModelID != plan.Candidate.PlatformModelID ||
@@ -220,7 +253,7 @@ func (s *Service) activeAsyncTaskAdmission(
return store.TaskAdmissionResult{}, false, nil
}
return store.TaskAdmissionResult{
Admission: admission,
Admission: *admission,
Admitted: true,
Leases: leases,
}, true, nil
@@ -61,3 +61,57 @@ func TestAcceptanceAdmissionScopesLeaveProductionPolicyUnchanged(t *testing.T) {
t.Fatalf("production admission policy changed: %+v", got[0])
}
}
func TestPinCandidatesToTaskAdmissionPreservesAdmittedCandidate(t *testing.T) {
input := []store.RuntimeModelCandidate{
{PlatformID: "platform-a", PlatformModelID: "model-a"},
{PlatformID: "platform-b", PlatformModelID: "model-b"},
{PlatformID: "platform-c", PlatformModelID: "model-c"},
}
admission := &store.TaskAdmission{
Status: "admitted",
PlatformID: "platform-b",
PlatformModelID: "model-b",
}
got, pinned := pinCandidatesToTaskAdmission(input, admission)
if !pinned {
t.Fatal("expected admitted candidate to be pinned")
}
if got[0].PlatformModelID != "model-b" || got[1].PlatformModelID != "model-a" || got[2].PlatformModelID != "model-c" {
t.Fatalf("unexpected candidate order: %+v", got)
}
if input[0].PlatformModelID != "model-a" {
t.Fatalf("candidate pinning mutated the caller slice: %+v", input)
}
}
func TestPinCandidatesToTaskAdmissionIgnoresWaitingOrMissingCandidate(t *testing.T) {
input := []store.RuntimeModelCandidate{
{PlatformID: "platform-a", PlatformModelID: "model-a"},
{PlatformID: "platform-b", PlatformModelID: "model-b"},
}
for name, admission := range map[string]*store.TaskAdmission{
"waiting": {
Status: "waiting",
PlatformID: "platform-b",
PlatformModelID: "model-b",
},
"missing": {
Status: "admitted",
PlatformID: "platform-c",
PlatformModelID: "model-c",
},
} {
t.Run(name, func(t *testing.T) {
got, pinned := pinCandidatesToTaskAdmission(input, admission)
if pinned {
t.Fatalf("unexpected candidate pin for %s: %+v", name, got)
}
if got[0].PlatformModelID != "model-a" {
t.Fatalf("candidate order changed for %s: %+v", name, got)
}
})
}
}
+13 -1
View File
@@ -422,6 +422,18 @@ func (s *Service) executeWithToken(ctx context.Context, task store.GatewayTask,
return Result{Task: failed, Output: failed.Result}, err
}
}
var asyncAdmission *store.TaskAdmission
if distributedAdmission && task.AsyncMode {
asyncAdmission, err = s.loadAsyncTaskAdmission(ctx, task)
if err != nil {
return Result{}, err
}
var pinned bool
candidates, pinned = pinCandidatesToTaskAdmission(candidates, asyncAdmission)
if pinned {
s.observeTaskAdmission("candidate_pinned")
}
}
pricingByCandidate := map[string]resolvedPricing{}
preprocessingByCandidate := map[string]parameterPreprocessResult{}
reservationBillings := []any(nil)
@@ -561,7 +573,7 @@ func (s *Service) executeWithToken(ctx context.Context, task store.GatewayTask,
var admissionErr error
if task.AsyncMode {
var alreadyAdmitted bool
admissionResult, alreadyAdmitted, admissionErr = s.activeAsyncTaskAdmission(ctx, task, plan)
admissionResult, alreadyAdmitted, admissionErr = s.activeAsyncTaskAdmission(ctx, task, plan, asyncAdmission)
if admissionErr == nil && !alreadyAdmitted {
admissionResult, admissionErr = s.tryTaskAdmission(ctx, task, plan, "")
}