From 3fd66a33be9360f1852171586907a5850b289507 Mon Sep 17 00:00:00 2001 From: anguobao123 <205642466+anguobao123@users.noreply.github.com> Date: Mon, 24 Aug 2026 01:25:54 +0800 Subject: [PATCH] fix(scheduler): diagnose load-batch OpenAI exclusions --- .../service/openai_account_scheduler_test.go | 69 +++++++++++++++++++ .../service/openai_gateway_scheduling.go | 65 ++++++++++++----- 2 files changed, 118 insertions(+), 16 deletions(-) diff --git a/backend/internal/service/openai_account_scheduler_test.go b/backend/internal/service/openai_account_scheduler_test.go index 3e8c241b2a..081cf30f9a 100644 --- a/backend/internal/service/openai_account_scheduler_test.go +++ b/backend/internal/service/openai_account_scheduler_test.go @@ -466,6 +466,75 @@ func TestOpenAIGatewayService_SelectAccountWithScheduler_DefaultDisabledUsesLega require.False(t, decision.StickyPreviousHit) } +// Regression: the legacy load-batch path had two bare ErrNoAvailableAccounts +// exits that bypassed the diagnostics added for both the advanced scheduler and +// the non-batched legacy selector. This is the default path when load batching +// is enabled, so quota auto-pause could still surface as an opaque 503. +func TestOpenAIGatewayService_SelectAccountWithScheduler_DefaultDisabled_LoadBatchReportsFilterReasons(t *testing.T) { + resetOpenAIAdvancedSchedulerSettingCacheForTest() + + ctx := withOpenAIQuotaAutoPauseSettings(context.Background(), OpsOpenAIAccountQuotaAutoPauseSettings{DefaultThreshold7d: 0.9}) + groupID := int64(10107) + quotaPaused := Account{ + ID: 36003, + Platform: PlatformOpenAI, + Type: AccountTypeOAuth, + Status: StatusActive, + Schedulable: true, + Concurrency: 1, + Extra: map[string]any{ + "codex_7d_used_percent": 95.0, + "codex_7d_reset_at": time.Now().Add(24 * time.Hour).Format(time.RFC3339), + "codex_usage_updated_at": time.Now().Add(-time.Minute).Format(time.RFC3339), + }, + } + mappingMiss := Account{ + ID: 36004, + Platform: PlatformOpenAI, + Type: AccountTypeAPIKey, + Status: StatusActive, + Schedulable: true, + Concurrency: 1, + Credentials: map[string]any{ + "model_mapping": map[string]any{"gpt-4o": "gpt-4o"}, + }, + } + excluded := Account{ + ID: 36005, + Platform: PlatformOpenAI, + Type: AccountTypeAPIKey, + Status: StatusActive, + Schedulable: true, + Concurrency: 1, + } + cfg := &config.Config{} + cfg.Gateway.Scheduling.LoadBatchEnabled = true + svc := &OpenAIGatewayService{ + accountRepo: schedulerTestOpenAIAccountRepo{accounts: []Account{quotaPaused, mappingMiss, excluded}}, + cache: &schedulerTestGatewayCache{}, + cfg: cfg, + concurrencyService: NewConcurrencyService(schedulerTestConcurrencyCache{}), + } + + require.False(t, svc.isOpenAIAdvancedSchedulerEnabled(ctx)) + selection, decision, err := svc.SelectAccountWithScheduler( + ctx, + &groupID, + "", + "", + "gpt-5.4-mini", + map[int64]struct{}{excluded.ID: {}}, + OpenAIUpstreamTransportAny, + false, + ) + + require.Error(t, err) + require.ErrorIs(t, err, ErrNoAvailableAccounts) + require.Nil(t, selection) + require.Equal(t, openAIAccountScheduleLayerLoadBalance, decision.Layer) + require.EqualError(t, err, "no available OpenAI accounts supporting model: gpt-5.4-mini (pool=3, filtered: excluded=1 model_not_supported=1 quota_auto_pause_7d=1)") +} + func TestOpenAIGatewayService_SelectAccountWithScheduler_DefaultDisabled_RequiredWSV2_SkipsHTTPOnlyAccount(t *testing.T) { resetOpenAIAdvancedSchedulerSettingCacheForTest() diff --git a/backend/internal/service/openai_gateway_scheduling.go b/backend/internal/service/openai_gateway_scheduling.go index 42e1c3fac8..5a0de5cd65 100644 --- a/backend/internal/service/openai_gateway_scheduling.go +++ b/backend/internal/service/openai_gateway_scheduling.go @@ -360,24 +360,45 @@ func openAICompactSupportTier(account *Account) int { // 注意:对 spark 影子账号,调用方还须额外调用 parentHealthyForShadow(account, lookup) // 检查母账号凭据可用性;该检查未内置于本函数,以避免注入 DB 依赖。 func isOpenAICompatibleAccountEligibleForRequest(ctx context.Context, account *Account, platform string, requestedModel string, requireCompact bool, requiredCapability OpenAIEndpointCapability) bool { - if !isOpenAICompatibleAccountEligibleForRequestBeforeProfit(ctx, account, platform, requestedModel, requireCompact, requiredCapability) { - return false + return openAICompatibleAccountEligibilityFailureReason(ctx, account, platform, requestedModel, requireCompact, requiredCapability) == "" +} + +// openAICompatibleAccountEligibilityFailureReason mirrors the legacy boolean +// eligibility check while naming its first veto point. Load-batch selection uses +// the reason only for server-side no-account diagnostics; the admission behavior +// remains unchanged. +func openAICompatibleAccountEligibilityFailureReason(ctx context.Context, account *Account, platform string, requestedModel string, requireCompact bool, requiredCapability OpenAIEndpointCapability) string { + if reason := openAICompatibleAccountEligibilityFailureReasonBeforeProfit(ctx, account, platform, requestedModel, requireCompact, requiredCapability); reason != "" { + return reason } // 分组利润控制:legacy 引擎的粘性/候选循环与 DB recheck 共用 // 本判定,任何 fallback 都不能把利润不合格账号重新放回候选。 - if vetoed, _ := openAIProfitControlVetoReason(ctx, account); vetoed { - return false + if vetoed, reason := openAIProfitControlVetoReason(ctx, account); vetoed { + return reason } - return true + return "" } // isOpenAICompatibleAccountEligibleForRequestBeforeProfit applies every // ordinary scheduling gate. Legacy selection uses it before classifying the // profit veto so earlier failures retain their actual reason. func isOpenAICompatibleAccountEligibleForRequestBeforeProfit(ctx context.Context, account *Account, platform string, requestedModel string, requireCompact bool, requiredCapability OpenAIEndpointCapability) bool { + return openAICompatibleAccountEligibilityFailureReasonBeforeProfit(ctx, account, platform, requestedModel, requireCompact, requiredCapability) == "" +} + +func openAICompatibleAccountEligibilityFailureReasonBeforeProfit(ctx context.Context, account *Account, platform string, requestedModel string, requireCompact bool, requiredCapability OpenAIEndpointCapability) string { platform = NormalizeOpenAICompatiblePlatform(platform) - if account == nil || account.Platform != platform || !account.IsOpenAICompatible() || !account.IsSchedulableForModelWithContext(ctx, requestedModel) { - return false + if account == nil { + return "account_nil" + } + if account.Platform != platform || !account.IsOpenAICompatible() { + return "platform_mismatch" + } + if !account.IsSchedulableForModelWithContext(ctx, requestedModel) { + if account.IsSchedulable() { + return "model_rate_limited" + } + return "not_schedulable" } if account.IsOpenAI() { if paused, reason := shouldAutoPauseOpenAIAccountByQuota(ctx, account); paused { @@ -389,7 +410,10 @@ func isOpenAICompatibleAccountEligibleForRequestBeforeProfit(ctx context.Context "threshold", reason.threshold, "utilization", reason.utilization, ) - return false + if reason.window != "" { + return "quota_auto_pause_" + reason.window + } + return "quota_auto_pause" } } if account.IsGrok() { @@ -400,23 +424,26 @@ func isOpenAICompatibleAccountEligibleForRequestBeforeProfit(ctx context.Context "threshold", reason.threshold, "utilization", reason.utilization, ) - return false + if reason.window != "" { + return "quota_auto_pause_" + reason.window + } + return "quota_auto_pause" } } if requestedModel != "" && !account.IsModelSupported(requestedModel) { - return false + return "model_not_supported" } if !account.SupportsOpenAIEndpointCapability(requiredCapability) { if account.IsGrok() && requiredCapability == OpenAIEndpointCapabilityGrokMediaGeneration { _, reason := account.GrokMediaGenerationEligibility() slog.Debug("grok_media_account_ineligible", "account_id", account.ID, "reason", reason) } - return false + return "capability_mismatch" } if requireCompact && openAICompactSupportTier(account) == 0 { - return false + return "compact_unsupported" } - return true + return "" } type openAIQuotaAutoPauseDecision struct { @@ -1101,7 +1128,7 @@ func (s *OpenAIGatewayService) selectAccountWithLoadAwareness(ctx context.Contex return nil, err } if len(accounts) == 0 { - return nil, ErrNoAvailableAccounts + return nil, noAvailableOpenAISelectionError(requestedModel, false, openAISelectionFilterStats{}.summary("")) } isExcluded := func(accountID int64) bool { @@ -1176,25 +1203,31 @@ func (s *OpenAIGatewayService) selectAccountWithLoadAwareness(ctx context.Contex return a } baseCandidateCount := 0 + filterStats := openAISelectionFilterStats{pool: len(accounts)} candidates := make([]*Account, 0, len(accounts)) for i := range accounts { acc := &accounts[i] if isExcluded(acc.ID) { + filterStats.exclude("excluded") continue } // Scheduler snapshots can be temporarily stale (bucket rebuild is throttled); // re-check schedulability here so recently rate-limited/overloaded accounts // are not selected again before the bucket is rebuilt. - if !isOpenAICompatibleAccountEligibleForRequest(ctx, acc, platform, requestedModel, false, requiredCapability) { + if reason := openAICompatibleAccountEligibilityFailureReason(ctx, acc, platform, requestedModel, false, requiredCapability); reason != "" { + filterStats.exclude(reason) continue } if !parentHealthyForShadow(acc, parentLookupL2) { + filterStats.exclude("shadow_parent_unhealthy") continue } if s.isOpenAIAccountRequestRuntimeBlocked(acc, requestedModel) { + filterStats.exclude("runtime_blocked") continue } if needsUpstreamCheck && s.isUpstreamModelRestrictedByChannel(ctx, *groupID, acc, requestedModel, requireCompact) { + filterStats.exclude("channel_upstream_restricted") continue } baseCandidateCount++ @@ -1202,7 +1235,7 @@ func (s *OpenAIGatewayService) selectAccountWithLoadAwareness(ctx context.Contex } if len(candidates) == 0 { - return nil, ErrNoAvailableAccounts + return nil, noAvailableOpenAISelectionError(requestedModel, false, filterStats.summary("")) } rateOrder := openAILegacyUpstreamRateOrder{} if preferLowUpstreamRate {