From b933d597839534d9cfa23768ea7cbaa75058f969 Mon Sep 17 00:00:00 2001 From: wangbo Date: Mon, 3 Aug 2026 12:19:52 +0800 Subject: [PATCH] =?UTF-8?q?fix(routing):=20=E9=A2=9D=E5=BA=A6=E9=A5=B1?= =?UTF-8?q?=E5=92=8C=E6=97=B6=E8=A7=A3=E9=99=A4=E4=BB=BB=E5=8A=A1=E5=80=99?= =?UTF-8?q?=E9=80=89=E5=9B=BA=E5=AE=9A?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 已取得并发租约的任务继续保持候选绑定;仅当已绑定平台的 RPM 或 TPM 已满且存在可用候选时,允许任务在执行前原子迁移到下一平台,避免等待窗口重置。\n\n验证:\n- go test ./internal/runner -count=1\n- go vet ./internal/runner --- apps/api/internal/runner/admission.go | 10 ++++ apps/api/internal/runner/admission_test.go | 64 ++++++++++++++++++++++ 2 files changed, 74 insertions(+) diff --git a/apps/api/internal/runner/admission.go b/apps/api/internal/runner/admission.go index 1e7bfae..b876e19 100644 --- a/apps/api/internal/runner/admission.go +++ b/apps/api/internal/runner/admission.go @@ -293,6 +293,10 @@ func pinCandidatesToTaskAdmission( if pinnedIndex <= 0 { return candidates, false } + pinnedCandidate := candidates[pinnedIndex] + if providerRateQuotaFull(pinnedCandidate) && !providerRateQuotaFull(candidates[0]) { + return candidates, false + } out := append([]store.RuntimeModelCandidate(nil), candidates...) pinned := out[pinnedIndex] copy(out[1:pinnedIndex+1], out[:pinnedIndex]) @@ -300,6 +304,12 @@ func pinCandidatesToTaskAdmission( return out, true } +func providerRateQuotaFull(candidate store.RuntimeModelCandidate) bool { + metrics := candidate.LoadMetrics + return (metrics.RPMLimit > 0 && metrics.RPMCurrent >= metrics.RPMLimit) || + (metrics.TPMLimit > 0 && metrics.TPMCurrent >= metrics.TPMLimit) +} + 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) } diff --git a/apps/api/internal/runner/admission_test.go b/apps/api/internal/runner/admission_test.go index 683bf89..edef8c1 100644 --- a/apps/api/internal/runner/admission_test.go +++ b/apps/api/internal/runner/admission_test.go @@ -227,6 +227,70 @@ func TestPinCandidatesToTaskAdmissionAllowsRequestedReselection(t *testing.T) { } } +func TestPinCandidatesToTaskAdmissionAllowsProviderRateQuotaMigration(t *testing.T) { + input := []store.RuntimeModelCandidate{ + { + PlatformID: "platform-b", + PlatformModelID: "model-b", + PlatformPriority: 200, + LoadMetrics: store.RuntimeCandidateLoadMetrics{ + RPMLimit: 12, + RPMCurrent: 0, + }, + }, + { + PlatformID: "platform-a", + PlatformModelID: "model-a", + PlatformPriority: 100, + LoadAvoided: true, + LoadMetrics: store.RuntimeCandidateLoadMetrics{ + RPMLimit: 4, + RPMCurrent: 4, + }, + }, + } + admission := &store.TaskAdmission{ + Status: "admitted", + PlatformID: "platform-a", + PlatformModelID: "model-a", + } + + got, pinned := pinCandidatesToTaskAdmission(input, admission) + + if pinned { + t.Fatalf("rate-full provider remained pinned: %+v", got) + } + if got[0].PlatformModelID != "model-b" { + t.Fatalf("available fallback was not preserved: %+v", got) + } +} + +func TestPinCandidatesToTaskAdmissionKeepsConcurrencyLeaseBinding(t *testing.T) { + input := []store.RuntimeModelCandidate{ + {PlatformID: "platform-b", PlatformModelID: "model-b"}, + { + PlatformID: "platform-a", + PlatformModelID: "model-a", + LoadAvoided: true, + LoadMetrics: store.RuntimeCandidateLoadMetrics{ + ConcurrentLimit: 2, + ConcurrentCurrent: 2, + }, + }, + } + admission := &store.TaskAdmission{ + Status: "admitted", + PlatformID: "platform-a", + PlatformModelID: "model-a", + } + + got, pinned := pinCandidatesToTaskAdmission(input, admission) + + if !pinned || got[0].PlatformModelID != "model-a" { + t.Fatalf("valid concurrency lease binding was not preserved: pinned=%v candidates=%+v", pinned, got) + } +} + func TestPinCandidatesToTaskAdmissionIgnoresMissingCandidate(t *testing.T) { input := []store.RuntimeModelCandidate{ {PlatformID: "platform-a", PlatformModelID: "model-a"},