Merge pull request #6124 from anguobao123/codex/diagnose-openai-load-batch-exclusions

fix(scheduler): diagnose load-batch OpenAI exclusions
This commit is contained in:
Wesley Liddick
2026-08-24 11:21:30 +08:00
committed by GitHub
2 changed files with 118 additions and 16 deletions
@@ -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()
@@ -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 {