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"},