Merge pull request #5316 from wucm667/fix/issue-5313-legacy-scheduler-diagnostics

fix(scheduler): diagnose legacy OpenAI exclusions
This commit is contained in:
Wesley Liddick
2026-08-10 10:52:14 +08:00
committed by GitHub
2 changed files with 144 additions and 16 deletions
@@ -312,6 +312,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
@@ -353,11 +368,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
}
@@ -712,10 +722,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)
@@ -805,10 +815,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))
@@ -819,18 +831,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
@@ -838,6 +858,7 @@ func (s *OpenAIGatewayService) selectBestAccount(ctx context.Context, groupID *i
compactTier = openAICompactSupportTier(fresh)
if compactTier == 0 {
compactBlocked = true
filterStats.exclude("compact_unsupported")
continue
}
}
@@ -847,7 +868,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 {
@@ -863,7 +884,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 更优。
@@ -1294,6 +1315,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
}
@@ -1308,7 +1340,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)) {
@@ -1341,12 +1373,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 s.isOpenAIAccountBlockedBySchedulingThreshold(ctx, account) {
@@ -1368,7 +1411,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)) {
@@ -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))
})
}