From df249ee5d5cf99ee624ba73b05fad5c253378544 Mon Sep 17 00:00:00 2001 From: wucm667 Date: Thu, 6 Aug 2026 02:54:16 +0800 Subject: [PATCH] fix(scheduler): diagnose legacy OpenAI exclusions --- .../service/openai_gateway_scheduling.go | 75 ++++++++++++---- ..._profit_control_legacy_diagnostics_test.go | 85 +++++++++++++++++++ 2 files changed, 144 insertions(+), 16 deletions(-) create mode 100644 backend/internal/service/openai_profit_control_legacy_diagnostics_test.go diff --git a/backend/internal/service/openai_gateway_scheduling.go b/backend/internal/service/openai_gateway_scheduling.go index c70a03802f..9780805296 100644 --- a/backend/internal/service/openai_gateway_scheduling.go +++ b/backend/internal/service/openai_gateway_scheduling.go @@ -259,6 +259,21 @@ 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 + } + // 分组利润控制:legacy 引擎的粘性/候选循环与 DB recheck 共用 + // 本判定,任何 fallback 都不能把利润不合格账号重新放回候选。 + if vetoed, _ := openAIProfitControlVetoReason(ctx, account); vetoed { + return false + } + return true +} + +// 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 { platform = normalizeOpenAICompatiblePlatform(platform) if account == nil || account.Platform != platform || !account.IsOpenAICompatible() || !account.IsSchedulableForModelWithContext(ctx, requestedModel) { return false @@ -300,11 +315,6 @@ func isOpenAICompatibleAccountEligibleForRequest(ctx context.Context, account *A if requireCompact && openAICompactSupportTier(account) == 0 { return false } - // 分组利润控制:legacy 引擎的粘性/候选循环与 DB recheck 共用 - // 本判定,任何 fallback 都不能把利润不合格账号重新放回候选。 - if vetoed, _ := openAIProfitControlVetoReason(ctx, account); vetoed { - return false - } return true } @@ -659,10 +669,10 @@ func (s *OpenAIGatewayService) selectAccountForModelWithExclusions(ctx context.C // 3. 按优先级 + LRU 选择最佳账号 // Select by priority + LRU - selected, compactBlocked := s.selectBestAccount(ctx, groupID, platform, accounts, requestedModel, excludedIDs, requireCompact, requiredCapability, preferLowUpstreamRate) + selected, compactBlocked, filterStats := s.selectBestAccount(ctx, groupID, platform, accounts, requestedModel, excludedIDs, requireCompact, requiredCapability, preferLowUpstreamRate) if selected == nil { - return nil, noAvailableOpenAISelectionError(requestedModel, compactBlocked, "") + return nil, noAvailableOpenAISelectionError(requestedModel, compactBlocked, filterStats.summary("")) } hydrated, err := s.hydrateSelectedAccount(ctx, selected) @@ -752,10 +762,12 @@ func (s *OpenAIGatewayService) tryStickySessionHit(ctx context.Context, groupID // selectBestAccount selects the best account from candidates (priority + LRU). // Returns nil if no available account. The second return reports whether at // least one candidate was filtered out solely because it lacks compact support -// (only meaningful when requireCompact=true). -func (s *OpenAIGatewayService) selectBestAccount(ctx context.Context, groupID *int64, platform string, accounts []Account, requestedModel string, excludedIDs map[int64]struct{}, requireCompact bool, requiredCapability OpenAIEndpointCapability, preferLowUpstreamRate bool) (*Account, bool) { +// (only meaningful when requireCompact=true); the third contains deterministic +// exclusion diagnostics for the evaluated snapshot. +func (s *OpenAIGatewayService) selectBestAccount(ctx context.Context, groupID *int64, platform string, accounts []Account, requestedModel string, excludedIDs map[int64]struct{}, requireCompact bool, requiredCapability OpenAIEndpointCapability, preferLowUpstreamRate bool) (*Account, bool, openAISelectionFilterStats) { platform = normalizeOpenAICompatiblePlatform(platform) compactBlocked := false + filterStats := openAISelectionFilterStats{pool: len(accounts)} needsUpstreamCheck := s.needsUpstreamChannelRestrictionCheck(ctx, groupID) eligible := make([]*Account, 0, len(accounts)) compactTiers := make(map[int64]int, len(accounts)) @@ -766,18 +778,26 @@ func (s *OpenAIGatewayService) selectBestAccount(ctx context.Context, groupID *i // 跳过被排除的账号 // Skip excluded accounts if _, excluded := excludedIDs[acc.ID]; excluded { + filterStats.exclude("excluded") continue } - fresh := s.resolveFreshSchedulableOpenAIAccount(ctx, acc, platform, requestedModel, false, requiredCapability) + fresh := s.resolveFreshSchedulableOpenAIAccountBeforeProfit(ctx, acc, platform, requestedModel, false, requiredCapability) if fresh == nil { + filterStats.exclude("ineligible") continue } - fresh = s.recheckSelectedOpenAIAccountFromDB(ctx, fresh, groupID, platform, requestedModel, false, requiredCapability) + fresh = s.recheckSelectedOpenAIAccountFromDBBeforeProfit(ctx, fresh, groupID, platform, requestedModel, false, requiredCapability) if fresh == nil { + filterStats.exclude("ineligible") continue } if needsUpstreamCheck && s.isUpstreamModelRestrictedByChannel(ctx, *groupID, fresh, requestedModel, requireCompact) { + filterStats.exclude("channel_restricted") + continue + } + if vetoed, reason := openAIProfitControlVetoReason(ctx, fresh); vetoed { + filterStats.exclude(reason) continue } compactTier := 0 @@ -785,6 +805,7 @@ func (s *OpenAIGatewayService) selectBestAccount(ctx context.Context, groupID *i compactTier = openAICompactSupportTier(fresh) if compactTier == 0 { compactBlocked = true + filterStats.exclude("compact_unsupported") continue } } @@ -794,7 +815,7 @@ func (s *OpenAIGatewayService) selectBestAccount(ctx context.Context, groupID *i } if len(eligible) == 0 { - return nil, compactBlocked + return nil, compactBlocked, filterStats } rateOrder := openAILegacyUpstreamRateOrder{} if preferLowUpstreamRate { @@ -810,7 +831,7 @@ func (s *OpenAIGatewayService) selectBestAccount(ctx context.Context, groupID *i } return s.isBetterAccount(a, b) }) - return eligible[0], compactBlocked + return eligible[0], compactBlocked, filterStats } // isBetterAccount 判断 candidate 是否比 current 更优。 @@ -1230,6 +1251,17 @@ func (s *OpenAIGatewayService) tryAcquireAccountSlot(ctx context.Context, accoun } func (s *OpenAIGatewayService) resolveFreshSchedulableOpenAIAccount(ctx context.Context, account *Account, platform string, requestedModel string, requireCompact bool, requiredCapability OpenAIEndpointCapability) *Account { + fresh := s.resolveFreshSchedulableOpenAIAccountBeforeProfit(ctx, account, platform, requestedModel, requireCompact, requiredCapability) + if fresh == nil { + return nil + } + if vetoed, _ := openAIProfitControlVetoReason(ctx, fresh); vetoed { + return nil + } + return fresh +} + +func (s *OpenAIGatewayService) resolveFreshSchedulableOpenAIAccountBeforeProfit(ctx context.Context, account *Account, platform string, requestedModel string, requireCompact bool, requiredCapability OpenAIEndpointCapability) *Account { if account == nil { return nil } @@ -1244,7 +1276,7 @@ func (s *OpenAIGatewayService) resolveFreshSchedulableOpenAIAccount(ctx context. fresh = current } - if !isOpenAICompatibleAccountEligibleForRequest(ctx, fresh, platform, requestedModel, requireCompact, requiredCapability) { + if !isOpenAICompatibleAccountEligibleForRequestBeforeProfit(ctx, fresh, platform, requestedModel, requireCompact, requiredCapability) { return nil } if !parentHealthyForShadow(fresh, s.parentAccountLookup(ctx)) { @@ -1274,12 +1306,23 @@ func (s *OpenAIGatewayService) parentAccountLookup(ctx context.Context) func(int } func (s *OpenAIGatewayService) recheckSelectedOpenAIAccountFromDB(ctx context.Context, account *Account, groupID *int64, platform string, requestedModel string, requireCompact bool, requiredCapability OpenAIEndpointCapability) *Account { + latest := s.recheckSelectedOpenAIAccountFromDBBeforeProfit(ctx, account, groupID, platform, requestedModel, requireCompact, requiredCapability) + if latest == nil { + return nil + } + if vetoed, _ := openAIProfitControlVetoReason(ctx, latest); vetoed { + return nil + } + return latest +} + +func (s *OpenAIGatewayService) recheckSelectedOpenAIAccountFromDBBeforeProfit(ctx context.Context, account *Account, groupID *int64, platform string, requestedModel string, requireCompact bool, requiredCapability OpenAIEndpointCapability) *Account { if account == nil { return nil } platform = normalizeOpenAICompatiblePlatform(platform) if s.schedulerSnapshot == nil || s.accountRepo == nil { - if !isOpenAICompatibleAccountEligibleForRequest(ctx, account, platform, requestedModel, requireCompact, requiredCapability) { + if !isOpenAICompatibleAccountEligibleForRequestBeforeProfit(ctx, account, platform, requestedModel, requireCompact, requiredCapability) { return nil } if !parentHealthyForShadow(account, s.parentAccountLookup(ctx)) { @@ -1298,7 +1341,7 @@ func (s *OpenAIGatewayService) recheckSelectedOpenAIAccountFromDB(ctx context.Co if !s.openAIAccountMatchesSchedulingGroup(latest, groupID) { return nil } - if !isOpenAICompatibleAccountEligibleForRequest(ctx, latest, platform, requestedModel, requireCompact, requiredCapability) { + if !isOpenAICompatibleAccountEligibleForRequestBeforeProfit(ctx, latest, platform, requestedModel, requireCompact, requiredCapability) { return nil } if !parentHealthyForShadow(latest, s.parentAccountLookup(ctx)) { diff --git a/backend/internal/service/openai_profit_control_legacy_diagnostics_test.go b/backend/internal/service/openai_profit_control_legacy_diagnostics_test.go new file mode 100644 index 0000000000..19e6b027e2 --- /dev/null +++ b/backend/internal/service/openai_profit_control_legacy_diagnostics_test.go @@ -0,0 +1,85 @@ +package service + +import ( + "errors" + "strings" + "testing" + "time" + + "github.com/Wei-Shaw/sub2api/internal/config" + "github.com/stretchr/testify/require" +) + +func legacyProfitDiagnosticService(accounts []Account) *OpenAIGatewayService { + return &OpenAIGatewayService{ + accountRepo: stubOpenAIAccountRepo{accounts: accounts}, + cfg: &config.Config{}, + rateLimitService: newOpenAIAdvancedSchedulerRateLimitService("false"), + concurrencyService: NewConcurrencyService(stubConcurrencyCache{}), + } +} + +func legacyProfitDiagnosticAccount(id int64) *Account { + now := time.Now() + a := upstreamCostTestAccount(id, UpstreamBillingProbeStatusOK, 0.9, now.Add(-time.Minute), 30*time.Minute) + a.Status = StatusActive + a.Schedulable = true + a.Concurrency = 1 + return a +} + +func TestSelectAccountWithScheduler_LegacyProfitDiagnostics(t *testing.T) { + groupID := int64(5313) + ctx := profitControlTestCtx(profitControlTestGroup(groupID, 0.5, 0)) + + t.Run("threshold reports deterministic pool count", func(t *testing.T) { + account := legacyProfitDiagnosticAccount(53131) + profitControlTestAccountWithRate(account, 0.9) + svc := legacyProfitDiagnosticService([]Account{*account}) + + selection, _, err := svc.SelectAccountWithScheduler(ctx, &groupID, "", "", "gpt-test", nil, OpenAIUpstreamTransportAny, false) + require.Nil(t, selection) + require.ErrorIs(t, err, ErrNoAvailableAccounts) + require.Contains(t, err.Error(), "pool=1, filtered: profit_threshold=1") + }) + + t.Run("missing account rate reports invalid rate", func(t *testing.T) { + account := upstreamCostTestOAuthAccount(53132) + account.Status = StatusActive + account.Schedulable = true + account.Concurrency = 1 + svc := legacyProfitDiagnosticService([]Account{*account}) + + selection, _, err := svc.SelectAccountWithScheduler(ctx, &groupID, "", "", "gpt-test", nil, OpenAIUpstreamTransportAny, false) + require.Nil(t, selection) + require.ErrorIs(t, err, ErrNoAvailableAccounts) + require.Contains(t, err.Error(), "profit_invalid_account_rate=1") + }) + + t.Run("model support gate does not report profit reasons", func(t *testing.T) { + account := legacyProfitDiagnosticAccount(53133) + profitControlTestAccountWithRate(account, 0.9) + account.Credentials = map[string]any{"model_mapping": map[string]any{"other-model": "other-model"}} + svc := legacyProfitDiagnosticService([]Account{*account}) + + selection, _, err := svc.SelectAccountWithScheduler(ctx, &groupID, "", "", "gpt-test", nil, OpenAIUpstreamTransportAny, false) + require.Nil(t, selection) + require.ErrorIs(t, err, ErrNoAvailableAccounts) + require.Contains(t, err.Error(), "ineligible=1") + require.NotContains(t, err.Error(), openAIProfitFilterReasonThreshold) + require.NotContains(t, err.Error(), openAIProfitFilterReasonInvalidAccountRate) + }) + + t.Run("compact preserves compact sentinel", func(t *testing.T) { + account := legacyProfitDiagnosticAccount(53134) + account.Extra = map[string]any{"openai_compact_supported": false} + profitControlTestAccountWithRate(account, 0.4) + svc := legacyProfitDiagnosticService([]Account{*account}) + + selection, _, err := svc.SelectAccountWithScheduler(ctx, &groupID, "", "", "gpt-test", nil, OpenAIUpstreamTransportAny, true) + require.Nil(t, selection) + require.ErrorIs(t, err, ErrNoAvailableCompactAccounts) + require.False(t, strings.Contains(err.Error(), "pool="), err) + require.False(t, errors.Is(err, ErrNoAvailableAccounts)) + }) +}