fix(scheduler): 高级调度器无可用账号错误补充排除原因统计

问题背景(Refs #4599):OpenAI OAuth 账号请求 /v1/responses 全模型 503,
日志仅有 no available OpenAI accounts supporting model 且
excluded_account_count=0。selectByLoadBalance 初筛循环的每个过滤分支
(配额自动暂停、渠道上游模型限制、运行时熔断等)都是静默 continue,
最多只有 Debug 级日志,用户与维护者都无法定位候选账号被哪个隐式
过滤点剔除。

修改:
- selectByLoadBalance 初筛循环为每个 continue 分支记录原因计数
  (excluded/not_schedulable/platform_mismatch/runtime_blocked/
  privacy_not_set/transport_incompatible 及请求兼容性细分原因)
- isAccountRequestCompatible 新增 isAccountRequestCompatibleReason
  变体,细分 runtime_blocked/quota_auto_pause_<窗口>/
  shadow_parent_unhealthy/model_not_supported/
  channel_upstream_restricted/capability_mismatch;原 bool 函数改为
  包装变体,其余调用点行为不变
- noAvailableOpenAISelectionError 增加 details 参数,初筛全灭时输出
  如 (pool=3, filtered: excluded=1 model_not_supported=1
  quota_auto_pause_7d=1),原因按字典序排序保证输出确定;快照池为空
  时输出 (pool=0);选择序列耗尽路径追加 selection_order_exhausted
  标注;legacy 调用点传空保持原消息不变

兼容性:错误仍通过 Unwrap 匹配 ErrNoAvailableAccounts,分类与 ops
标记逻辑不变;客户端可见错误仍为通用分类消息,统计仅进服务端日志;
happy path 原因 map 惰性分配,无额外开销。

Refs #4599

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
superman2003
2026-07-20 21:44:42 +08:00
co-authored by Cursor
parent e625ce3b3b
commit 29bea0a758
4 changed files with 264 additions and 24 deletions
@@ -1270,6 +1270,53 @@ func (s *defaultOpenAIAccountScheduler) tryFallbackToWeightedSticky(
return nil, nil
}
// openAISelectionFilterStats counts why candidates were dropped by the
// selectByLoadBalance initial filter. Historically these exclusions were
// silent (debug logs at best), so a "no available accounts" failure with
// excluded_account_count=0 was undiagnosable from the error alone (#4599).
// The reasons map is lazily allocated: on the happy path (nothing filtered
// out, or an account is eventually selected) no extra allocation happens.
type openAISelectionFilterStats struct {
pool int
reasons map[string]int
}
func (s *openAISelectionFilterStats) exclude(reason string) {
if s.reasons == nil {
s.reasons = make(map[string]int, 4)
}
s.reasons[reason]++
}
// summary renders deterministic exclusion statistics for scheduling error
// messages, e.g. "pool=3, filtered: model_not_supported=2 quota_auto_pause_7d=1".
// Reasons are sorted lexicographically so the output is stable for tests and
// log aggregation. extra, when non-empty, is appended as a trailing marker.
func (s openAISelectionFilterStats) summary(extra string) string {
var b strings.Builder
b.WriteString("pool=")
b.WriteString(strconv.Itoa(s.pool))
if len(s.reasons) > 0 {
reasons := make([]string, 0, len(s.reasons))
for reason := range s.reasons {
reasons = append(reasons, reason)
}
sort.Strings(reasons)
b.WriteString(", filtered:")
for _, reason := range reasons {
b.WriteString(" ")
b.WriteString(reason)
b.WriteString("=")
b.WriteString(strconv.Itoa(s.reasons[reason]))
}
}
if extra != "" {
b.WriteString(", ")
b.WriteString(extra)
}
return b.String()
}
func (s *defaultOpenAIAccountScheduler) selectByLoadBalance(
ctx context.Context,
req OpenAIAccountScheduleRequest,
@@ -1280,7 +1327,7 @@ func (s *defaultOpenAIAccountScheduler) selectByLoadBalance(
return nil, 0, 0, 0, err
}
if len(accounts) == 0 {
return nil, 0, 0, 0, noAvailableOpenAISelectionError(req.RequestedModel, false)
return nil, 0, 0, 0, noAvailableOpenAISelectionError(req.RequestedModel, false, openAISelectionFilterStats{}.summary(""))
}
// require_privacy_set: 获取分组信息
@@ -1289,19 +1336,27 @@ func (s *defaultOpenAIAccountScheduler) selectByLoadBalance(
schedGroup, _ = s.service.schedulerSnapshot.GetGroupByID(ctx, *req.GroupID)
}
filterStats := openAISelectionFilterStats{pool: len(accounts)}
filtered := make([]*Account, 0, len(accounts))
loadReq := make([]AccountWithConcurrency, 0, len(accounts))
for i := range accounts {
account := &accounts[i]
if req.ExcludedIDs != nil {
if _, excluded := req.ExcludedIDs[account.ID]; excluded {
filterStats.exclude("excluded")
continue
}
}
if !account.IsSchedulable() || account.Platform != normalizeOpenAICompatiblePlatform(req.Platform) || !account.IsOpenAICompatible() {
if !account.IsSchedulable() {
filterStats.exclude("not_schedulable")
continue
}
if account.Platform != normalizeOpenAICompatiblePlatform(req.Platform) || !account.IsOpenAICompatible() {
filterStats.exclude("platform_mismatch")
continue
}
if s.service.isOpenAIAccountRequestRuntimeBlocked(account, req.RequestedModel) {
filterStats.exclude("runtime_blocked")
continue
}
// require_privacy_set: 跳过 privacy 未设置的账号并标记异常
@@ -1309,12 +1364,15 @@ func (s *defaultOpenAIAccountScheduler) selectByLoadBalance(
s.service.BlockAccountScheduling(account, time.Time{}, "privacy_not_set")
_ = s.service.accountRepo.SetError(ctx, account.ID,
fmt.Sprintf("Privacy not set, required by group [%s]", schedGroup.Name))
filterStats.exclude("privacy_not_set")
continue
}
if !s.isAccountRequestCompatible(ctx, account, req) {
if compatible, reason := s.isAccountRequestCompatibleReason(ctx, account, req); !compatible {
filterStats.exclude(reason)
continue
}
if !s.isAccountTransportCompatible(account, req.RequiredTransport) {
filterStats.exclude("transport_incompatible")
continue
}
filtered = append(filtered, account)
@@ -1324,7 +1382,7 @@ func (s *defaultOpenAIAccountScheduler) selectByLoadBalance(
})
}
if len(filtered) == 0 {
return nil, 0, 0, 0, noAvailableOpenAISelectionError(req.RequestedModel, false)
return nil, 0, 0, 0, noAvailableOpenAISelectionError(req.RequestedModel, false, filterStats.summary(""))
}
loadMap := map[int64]*AccountLoadInfo{}
@@ -1356,7 +1414,7 @@ func (s *defaultOpenAIAccountScheduler) selectByLoadBalance(
candidateCount, topK, loadSkew := regularAttempt.candidateCount, regularAttempt.topK, regularAttempt.loadSkew
fallbackErr := regularAttempt.err
if regularAttempt.err == nil {
result, candidateCount, topK, loadSkew, fallbackErr = s.finishLoadBalanceSelectionFallback(ctx, req, regularAttempt, budget)
result, candidateCount, topK, loadSkew, fallbackErr = s.finishLoadBalanceSelectionFallback(ctx, req, regularAttempt, budget, filterStats)
if fallbackErr == nil && result != nil {
return result, candidateCount, topK, loadSkew, nil
}
@@ -1364,13 +1422,13 @@ func (s *defaultOpenAIAccountScheduler) selectByLoadBalance(
// 常规池既无法获取也无法排队(含仅剩不支持 compact 的候选)时,
// 回退到订阅池的等待计划:busy-but-waitable 的订阅账号不应因常规池存在
// 而被丢弃,否则开启订阅优先反而让本可排队成功的请求硬失败。
subResult, subCandidateCount, subTopK, subLoadSkew, subErr := s.finishLoadBalanceSelectionFallback(ctx, req, attempt, budget)
subResult, subCandidateCount, subTopK, subLoadSkew, subErr := s.finishLoadBalanceSelectionFallback(ctx, req, attempt, budget, filterStats)
if subErr == nil && subResult != nil {
return subResult, subCandidateCount, subTopK, subLoadSkew, nil
}
return result, candidateCount, topK, loadSkew, fallbackErr
}
return s.finishLoadBalanceSelectionFallback(ctx, req, attempt, budget)
return s.finishLoadBalanceSelectionFallback(ctx, req, attempt, budget, filterStats)
}
}
@@ -1381,7 +1439,7 @@ func (s *defaultOpenAIAccountScheduler) selectByLoadBalance(
if attempt.result != nil {
return attempt.result, attempt.candidateCount, attempt.topK, attempt.loadSkew, nil
}
return s.finishLoadBalanceSelectionFallback(ctx, req, attempt, budget)
return s.finishLoadBalanceSelectionFallback(ctx, req, attempt, budget, filterStats)
}
func partitionOpenAIChatGPTSubscriptionAccounts(accounts []*Account) ([]*Account, []*Account) {
@@ -1511,13 +1569,14 @@ func (s *defaultOpenAIAccountScheduler) finishLoadBalanceSelectionFallback(
req OpenAIAccountScheduleRequest,
attempt openAIAccountLoadSelectionAttempt,
budget *openAISelectionProbeBudget,
filterStats openAISelectionFilterStats,
) (*AccountSelectionResult, int, int, float64, error) {
candidateCount := attempt.candidateCount
topK := attempt.topK
loadSkew := attempt.loadSkew
if len(attempt.selectionOrder) == 0 {
return nil, candidateCount, topK, loadSkew, noAvailableOpenAISelectionError(req.RequestedModel, attempt.compactBlocked)
return nil, candidateCount, topK, loadSkew, noAvailableOpenAISelectionError(req.RequestedModel, attempt.compactBlocked, filterStats.summary("selection_order_empty"))
}
if stickyFallback, stickyErr := s.tryFallbackToWeightedSticky(ctx, req); stickyErr != nil {
@@ -1552,7 +1611,7 @@ func (s *defaultOpenAIAccountScheduler) finishLoadBalanceSelectionFallback(
continue
}
if !s.consumeOpenAISelectionDBRecheck(budget) {
return nil, candidateCount, topK, loadSkew, noAvailableOpenAISelectionError(req.RequestedModel, compactBlocked)
return nil, candidateCount, topK, loadSkew, noAvailableOpenAISelectionError(req.RequestedModel, compactBlocked, filterStats.summary("selection_order_exhausted"))
}
fresh = s.service.recheckSelectedOpenAIAccountFromDB(ctx, fresh, req.GroupID, req.Platform, req.RequestedModel, false, req.RequiredCapability)
if fresh == nil || !s.isAccountTransportCompatible(fresh, req.RequiredTransport) || !s.isAccountRequestCompatible(ctx, fresh, req) {
@@ -1574,7 +1633,7 @@ func (s *defaultOpenAIAccountScheduler) finishLoadBalanceSelectionFallback(
}
}
return nil, candidateCount, topK, loadSkew, noAvailableOpenAISelectionError(req.RequestedModel, compactBlocked)
return nil, candidateCount, topK, loadSkew, noAvailableOpenAISelectionError(req.RequestedModel, compactBlocked, filterStats.summary("selection_order_exhausted"))
}
func (s *defaultOpenAIAccountScheduler) isAccountTransportCompatible(account *Account, requiredTransport OpenAIUpstreamTransport) bool {
@@ -1604,18 +1663,31 @@ func (s *defaultOpenAIAccountScheduler) lookupShadowParentAccount(ctx context.Co
}
func (s *defaultOpenAIAccountScheduler) isAccountRequestCompatible(ctx context.Context, account *Account, req OpenAIAccountScheduleRequest) bool {
compatible, _ := s.isAccountRequestCompatibleReason(ctx, account, req)
return compatible
}
// isAccountRequestCompatibleReason reports whether the account can serve the
// request, and when it cannot, names the veto point. The reason feeds
// openAISelectionFilterStats so that "no available accounts" errors state why
// each candidate was dropped instead of failing silently (#4599).
func (s *defaultOpenAIAccountScheduler) isAccountRequestCompatibleReason(ctx context.Context, account *Account, req OpenAIAccountScheduleRequest) (bool, string) {
if account == nil {
return false
return false, "account_nil"
}
if s != nil && s.service != nil && s.service.isOpenAIAccountRequestRuntimeBlocked(account, req.RequestedModel) {
return false
return false, "runtime_blocked"
}
// Quota auto-pause must be evaluated during the initial filter too. Without it the
// TopK candidate pool can be filled with paused accounts and the later fresh/DB
// rechecks won't reach healthy accounts that fell outside TopK — manifesting as
// "no available accounts" even though healthy ones exist.
if paused, _ := shouldAutoPauseOpenAIAccountByQuota(ctx, account); paused {
return false
if paused, decision := shouldAutoPauseOpenAIAccountByQuota(ctx, account); paused {
reason := "quota_auto_pause"
if decision.window != "" {
reason += "_" + decision.window
}
return false, reason
}
// 母账号健康联动:影子账号的凭据来自母账号,母账号不可调度时影子也不应被选中。
// Parent-health gate: shadow borrows the parent's credentials; an unschedulable
@@ -1623,17 +1695,20 @@ func (s *defaultOpenAIAccountScheduler) isAccountRequestCompatible(ctx context.C
if !parentHealthyForShadow(account, func(id int64) *Account {
return s.lookupShadowParentAccount(ctx, id)
}) {
return false
return false, "shadow_parent_unhealthy"
}
if req.RequestedModel != "" && !account.IsModelSupported(req.RequestedModel) {
return false
return false, "model_not_supported"
}
if req.GroupID != nil && s != nil && s.service != nil &&
s.service.needsUpstreamChannelRestrictionCheck(ctx, req.GroupID) &&
s.service.isUpstreamModelRestrictedByChannel(ctx, *req.GroupID, account, req.RequestedModel, req.RequireCompact) {
return false
return false, "channel_upstream_restricted"
}
return accountSupportsOpenAICapabilities(account, req.RequiredCapability, req.RequiredImageCapability)
if !accountSupportsOpenAICapabilities(account, req.RequiredCapability, req.RequiredImageCapability) {
return false, "capability_mismatch"
}
return true, ""
}
func (s *defaultOpenAIAccountScheduler) ReportResult(accountID int64, success bool, firstTokenMs *int) {
@@ -826,6 +826,161 @@ func TestOpenAIGatewayService_SelectAccountWithScheduler_GrokMediaCapabilityFilt
})
}
// Regression #4599: when the advanced scheduler's load-balance initial filter
// drops every candidate, the resulting "no available OpenAI accounts" error must
// carry per-reason exclusion counts. Previously all filter branches were silent
// (debug logs at best), so a 503 with excluded_account_count=0 could not be
// diagnosed from the error alone.
func TestOpenAIGatewayService_SelectAccountWithScheduler_NoAvailableErrorReportsQuotaAutoPauseExclusion(t *testing.T) {
resetOpenAIAdvancedSchedulerSettingCacheForTest()
ctx := withOpenAIQuotaAutoPauseSettings(context.Background(), OpsOpenAIAccountQuotaAutoPauseSettings{DefaultThreshold7d: 0.9})
groupID := int64(101201)
accounts := []Account{
{
ID: 38101,
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),
},
},
}
cfg := &config.Config{}
svc := &OpenAIGatewayService{
accountRepo: schedulerTestOpenAIAccountRepo{accounts: accounts},
cache: &schedulerTestGatewayCache{},
cfg: cfg,
rateLimitService: newOpenAIAdvancedSchedulerRateLimitService("true"),
concurrencyService: NewConcurrencyService(schedulerTestConcurrencyCache{}),
}
selection, decision, err := svc.SelectAccountWithScheduler(
ctx, &groupID, "", "", "gpt-5.4-mini", nil, 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=1, filtered: quota_auto_pause_7d=1)")
}
func TestOpenAIGatewayService_SelectAccountWithScheduler_NoAvailableErrorReportsModelNotSupported(t *testing.T) {
resetOpenAIAdvancedSchedulerSettingCacheForTest()
ctx := context.Background()
groupID := int64(101202)
// OAuth account with empty model_mapping: foreign-family models (grok-*)
// are not servable by Codex upstream (#3662) and are dropped by the filter.
accounts := []Account{
{
ID: 38111,
Platform: PlatformOpenAI,
Type: AccountTypeOAuth,
Status: StatusActive,
Schedulable: true,
Concurrency: 1,
},
}
svc := &OpenAIGatewayService{
accountRepo: schedulerTestOpenAIAccountRepo{accounts: accounts},
cache: &schedulerTestGatewayCache{},
cfg: &config.Config{},
rateLimitService: newOpenAIAdvancedSchedulerRateLimitService("true"),
concurrencyService: NewConcurrencyService(schedulerTestConcurrencyCache{}),
}
selection, _, err := svc.SelectAccountWithScheduler(
ctx, &groupID, "", "", "grok-4.5", nil, OpenAIUpstreamTransportAny, false,
)
require.Error(t, err)
require.ErrorIs(t, err, ErrNoAvailableAccounts)
require.Nil(t, selection)
require.EqualError(t, err, "no available OpenAI accounts supporting model: grok-4.5 (pool=1, filtered: model_not_supported=1)")
}
func TestOpenAIGatewayService_SelectAccountWithScheduler_NoAvailableErrorAggregatesReasonsDeterministically(t *testing.T) {
resetOpenAIAdvancedSchedulerSettingCacheForTest()
ctx := withOpenAIQuotaAutoPauseSettings(context.Background(), OpsOpenAIAccountQuotaAutoPauseSettings{DefaultThreshold7d: 0.9})
groupID := int64(101203)
quotaPaused := Account{
ID: 38121,
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: 38122,
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: 38123,
Platform: PlatformOpenAI,
Type: AccountTypeAPIKey,
Status: StatusActive,
Schedulable: true,
Concurrency: 1,
}
svc := &OpenAIGatewayService{
accountRepo: schedulerTestOpenAIAccountRepo{accounts: []Account{quotaPaused, mappingMiss, excluded}},
cache: &schedulerTestGatewayCache{},
cfg: &config.Config{},
rateLimitService: newOpenAIAdvancedSchedulerRateLimitService("true"),
concurrencyService: NewConcurrencyService(schedulerTestConcurrencyCache{}),
}
selection, _, err := svc.SelectAccountWithScheduler(
ctx, &groupID, "", "", "gpt-5.4-mini", map[int64]struct{}{38123: {}}, OpenAIUpstreamTransportAny, false,
)
require.Error(t, err)
require.ErrorIs(t, err, ErrNoAvailableAccounts)
require.Nil(t, selection)
// Reasons are sorted lexicographically, so the message is deterministic.
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_NoAvailableErrorReportsEmptyPool(t *testing.T) {
resetOpenAIAdvancedSchedulerSettingCacheForTest()
ctx := context.Background()
groupID := int64(101204)
svc := &OpenAIGatewayService{
accountRepo: schedulerTestOpenAIAccountRepo{},
cache: &schedulerTestGatewayCache{},
cfg: &config.Config{},
rateLimitService: newOpenAIAdvancedSchedulerRateLimitService("true"),
concurrencyService: NewConcurrencyService(schedulerTestConcurrencyCache{}),
}
selection, _, err := svc.SelectAccountWithScheduler(
ctx, &groupID, "", "", "gpt-5.1", nil, OpenAIUpstreamTransportAny, false,
)
require.Error(t, err)
require.ErrorIs(t, err, ErrNoAvailableAccounts)
require.Nil(t, selection)
require.EqualError(t, err, "no available OpenAI accounts supporting model: gpt-5.1 (pool=0)")
}
func TestOpenAIGatewayService_SelectAccountWithScheduler_EnabledUsesAdvancedPreviousResponseRouting(t *testing.T) {
resetOpenAIAdvancedSchedulerSettingCacheForTest()
@@ -244,7 +244,7 @@ func TestAdvancedSchedulerSharesProbeBudgetWithFallbackDBRechecks(t *testing.T)
require.NoError(t, err)
require.Nil(t, selection)
selection, _, _, _, err = scheduler.finishLoadBalanceSelectionFallback(
context.Background(), req, openAIAccountLoadSelectionAttempt{selectionOrder: selectionOrder}, budget,
context.Background(), req, openAIAccountLoadSelectionAttempt{selectionOrder: selectionOrder}, budget, openAISelectionFilterStats{},
)
require.Error(t, err)
@@ -170,14 +170,24 @@ func normalizeOpenAICompatiblePlatform(platform string) string {
return PlatformOpenAI
}
func noAvailableOpenAISelectionError(requestedModel string, compactBlocked bool) error {
// details carries an optional machine-parseable exclusion summary (e.g.
// "pool=2, filtered: quota_auto_pause_7d=1 runtime_blocked=1") appended in
// parentheses. It is for server-side logs / ops diagnostics only: handlers
// never forward this error text to OpenAI-platform clients (they respond with
// the generic classification message). Callers that must preserve the legacy
// message pass "".
func noAvailableOpenAISelectionError(requestedModel string, compactBlocked bool, details string) error {
if compactBlocked {
return ErrNoAvailableCompactAccounts
}
message := "no available OpenAI accounts"
if requestedModel != "" {
return openAINoAvailableSelectionError{message: fmt.Sprintf("no available OpenAI accounts supporting model: %s", requestedModel)}
message = fmt.Sprintf("no available OpenAI accounts supporting model: %s", requestedModel)
}
return openAINoAvailableSelectionError{message: "no available OpenAI accounts"}
if details != "" {
message += " (" + details + ")"
}
return openAINoAvailableSelectionError{message: message}
}
type openAINoAvailableSelectionError struct {
@@ -612,7 +622,7 @@ func (s *OpenAIGatewayService) selectAccountForModelWithExclusions(ctx context.C
selected, compactBlocked := s.selectBestAccount(ctx, groupID, platform, accounts, requestedModel, excludedIDs, requireCompact, requiredCapability, preferLowUpstreamRate)
if selected == nil {
return nil, noAvailableOpenAISelectionError(requestedModel, compactBlocked)
return nil, noAvailableOpenAISelectionError(requestedModel, compactBlocked, "")
}
hydrated, err := s.hydrateSelectedAccount(ctx, selected)