mirror of
https://github.com/Wei-Shaw/sub2api.git
synced 2026-10-07 16:48:45 +08:00
fix(scheduler): diagnose legacy OpenAI exclusions
This commit is contained in:
@@ -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)) {
|
||||
|
||||
@@ -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))
|
||||
})
|
||||
}
|
||||
Reference in New Issue
Block a user