From f1aadd48d5ca295cb1babc5e65412af2f46d1247 Mon Sep 17 00:00:00 2001 From: shaw Date: Tue, 25 Aug 2026 20:25:22 +0800 Subject: [PATCH] fix(openai): pause OAuth accounts on exhausted quota 429 --- .../openai_account_runtime_block_fastpath.go | 72 ++++++++++++++++--- ...nai_account_runtime_block_fastpath_test.go | 36 +++++++++- .../service/openai_gateway_passthrough.go | 9 ++- .../service/openai_gateway_upstream_errors.go | 15 +++- 4 files changed, 118 insertions(+), 14 deletions(-) diff --git a/backend/internal/service/openai_account_runtime_block_fastpath.go b/backend/internal/service/openai_account_runtime_block_fastpath.go index 1d8da95913..1f402202ca 100644 --- a/backend/internal/service/openai_account_runtime_block_fastpath.go +++ b/backend/internal/service/openai_account_runtime_block_fastpath.go @@ -28,6 +28,49 @@ type OpenAIOAuth429FailoverState struct { grokOAuth429FollowupPending bool } +type openAIOAuth429Disposition uint8 + +const ( + openAIOAuth429Transient openAIOAuth429Disposition = iota + openAIOAuth429Quota5h + openAIOAuth429Quota7d + openAIOAuth429QuotaReset +) + +// classifyOpenAIOAuth429 区分账号配额耗尽信号与普通瞬时 429。明确窗口达到 +// 100% 时以该窗口为准;没有 100% 标记但包含重置头时,沿用 v179 的兼容语义, +// 仍视为配额限流信号。 +func classifyOpenAIOAuth429(headers http.Header, responseBody []byte) (openAIOAuth429Disposition, *time.Time) { + if snapshot := ParseCodexRateLimitHeaders(headers); snapshot != nil { + if normalized := snapshot.Normalize(); normalized != nil { + if normalized.Used7dPercent != nil && *normalized.Used7dPercent >= 100 { + if normalized.Reset7dSeconds != nil { + now := time.Now() + resetAt := now.Add(time.Duration(*normalized.Reset7dSeconds) * time.Second) + return openAIOAuth429Quota7d, &resetAt + } + return openAIOAuth429Quota7d, nil + } + if normalized.Used5hPercent != nil && *normalized.Used5hPercent >= 100 { + if normalized.Reset5hSeconds != nil { + now := time.Now() + resetAt := now.Add(time.Duration(*normalized.Reset5hSeconds) * time.Second) + return openAIOAuth429Quota5h, &resetAt + } + return openAIOAuth429Quota5h, nil + } + } + } + if resetAt := calculateOpenAI429ResetTime(headers); resetAt != nil { + return openAIOAuth429QuotaReset, resetAt + } + if resetUnix := parseOpenAIRateLimitResetTime(responseBody); resetUnix != nil { + resetAt := time.Unix(*resetUnix, 0) + return openAIOAuth429QuotaReset, &resetAt + } + return openAIOAuth429Transient, nil +} + func openAIAccountStateContext(ctx context.Context) (context.Context, context.CancelFunc) { base := context.Background() if ctx != nil { @@ -165,19 +208,16 @@ func (s *OpenAIGatewayService) markOpenAIOAuth429RateLimited(ctx context.Context return } s.recordOpenAIOAuth429() - if s.openAIOAuth429RetryWindowActive(account) { + disposition, resetAt := classifyOpenAIOAuth429(headers, responseBody) + if disposition == openAIOAuth429Transient && s.openAIOAuth429RetryWindowActive(account) { return } cooldownUntil := time.Now().Add(openAIOAuth429FallbackCooldown) - if s.rateLimitService != nil { - if resetAt := s.rateLimitService.calculateOpenAI429ResetTime(headers); resetAt != nil && resetAt.After(time.Now()) { - cooldownUntil = *resetAt - } else if resetUnix := parseOpenAIRateLimitResetTime(responseBody); resetUnix != nil { - if resetAt := time.Unix(*resetUnix, 0); resetAt.After(time.Now()) { - cooldownUntil = resetAt - } - } else if cooldown, ok := s.rateLimitService.get429FallbackCooldown(ctx, account); ok && cooldown > 0 { + if resetAt != nil && resetAt.After(time.Now()) { + cooldownUntil = *resetAt + } else if s.rateLimitService != nil { + if cooldown, ok := s.rateLimitService.get429FallbackCooldown(ctx, account); ok && cooldown > 0 { cooldownUntil = time.Now().Add(cooldown) } } @@ -186,9 +226,17 @@ func (s *OpenAIGatewayService) markOpenAIOAuth429RateLimited(ctx context.Context } func (s *OpenAIGatewayService) shouldRetryOpenAIOAuth429OnSameAccount(account *Account, statusCode int, shouldDisable bool) bool { + return s.shouldRetryOpenAIOAuth429OnSameAccountWithResponse(account, statusCode, shouldDisable, nil, nil) +} + +func (s *OpenAIGatewayService) shouldRetryOpenAIOAuth429OnSameAccountWithResponse(account *Account, statusCode int, shouldDisable bool, headers http.Header, responseBody []byte) bool { if shouldDisable || statusCode != http.StatusTooManyRequests || !isOpenAIOAuthAccount(account) || account.IsShadow() { return false } + disposition, _ := classifyOpenAIOAuth429(headers, responseBody) + if disposition != openAIOAuth429Transient { + return false + } // markOpenAIOAuth429RateLimited parks the account once the window expires. // Do not accidentally create a fresh window after that transition. if s.isOpenAIAccountRuntimeBlocked(account) { @@ -199,10 +247,14 @@ func (s *OpenAIGatewayService) shouldRetryOpenAIOAuth429OnSameAccount(account *A // ShouldRetryOpenAIOAuth429 lets RateLimitService defer persistent account // cooldown until the gateway's same-account retry window is exhausted. -func (s *OpenAIGatewayService) ShouldRetryOpenAIOAuth429(account *Account, _ http.Header, _ []byte) bool { +func (s *OpenAIGatewayService) ShouldRetryOpenAIOAuth429(account *Account, headers http.Header, responseBody []byte) bool { if s == nil || !isOpenAIOAuthAccount(account) || account.IsShadow() || s.isOpenAIAccountRuntimeBlocked(account) { return false } + disposition, _ := classifyOpenAIOAuth429(headers, responseBody) + if disposition != openAIOAuth429Transient { + return false + } return s.openAIOAuth429RetryWindowActive(account) } diff --git a/backend/internal/service/openai_account_runtime_block_fastpath_test.go b/backend/internal/service/openai_account_runtime_block_fastpath_test.go index fcda6d6a05..89d8c5b5ec 100644 --- a/backend/internal/service/openai_account_runtime_block_fastpath_test.go +++ b/backend/internal/service/openai_account_runtime_block_fastpath_test.go @@ -14,7 +14,7 @@ import ( ) type oauth429RateLimitRepo struct { - AccountRepository + mockAccountRepoForGemini setRateLimitedCalls int lastRateLimitedUntil time.Time } @@ -62,6 +62,38 @@ func TestOpenAI429FastPath_BlocksOAuthOnlyAfterRetryWindow(t *testing.T) { require.False(t, svc.shouldRetryOpenAIOAuth429OnSameAccount(account, http.StatusTooManyRequests, false)) } +func TestOpenAI429FastPath_BlocksOAuthImmediatelyWhenSevenDayQuotaIsExhausted(t *testing.T) { + repo := &oauth429RateLimitRepo{} + rateLimits := NewRateLimitService(repo, nil, &config.Config{}, nil, nil) + svc := &OpenAIGatewayService{rateLimitService: rateLimits} + rateLimits.SetAccountRuntimeBlocker(svc) + account := &Account{ID: 423, Platform: PlatformOpenAI, Type: AccountTypeOAuth} + headers := http.Header{} + headers.Set("x-codex-primary-used-percent", "100") + headers.Set("x-codex-primary-reset-after-seconds", "604800") + headers.Set("x-codex-primary-window-minutes", "10080") + headers.Set("x-codex-secondary-used-percent", "20") + headers.Set("x-codex-secondary-reset-after-seconds", "3600") + headers.Set("x-codex-secondary-window-minutes", "300") + + shouldDisable := svc.handleOpenAIAccountUpstreamError(context.Background(), account, http.StatusTooManyRequests, headers, []byte(`{"error":{"type":"rate_limit_error","code":"rate_limit_exceeded"}}`)) + + require.False(t, shouldDisable) + require.True(t, svc.isOpenAIAccountRuntimeBlocked(account)) + require.Equal(t, 1, repo.setRateLimitedCalls) + require.Greater(t, time.Until(repo.lastRateLimitedUntil), 6*24*time.Hour) + require.False(t, svc.ShouldRetryOpenAIOAuth429(account, headers, nil)) +} + +func TestOpenAI429FastPath_RetriesOAuthWhenNoQuotaSignalExists(t *testing.T) { + svc := &OpenAIGatewayService{} + account := &Account{ID: 424, Platform: PlatformOpenAI, Type: AccountTypeOAuth} + headers := http.Header{"Retry-After": []string{"1"}} + + require.True(t, svc.ShouldRetryOpenAIOAuth429(account, headers, []byte(`{"error":{"type":"rate_limit_error","message":"try again"}}`))) + require.False(t, svc.isOpenAIAccountRuntimeBlocked(account)) +} + func TestOpenAIStream429IgnoresSuccessfulQuotaSnapshotHeaders(t *testing.T) { repo := &oauth429RateLimitRepo{} rateLimits := NewRateLimitService(repo, nil, &config.Config{}, nil, nil) @@ -164,7 +196,7 @@ func TestOpenAI429FastPath_SkipsSparkShadow(t *testing.T) { svc.markOpenAIOAuth429RateLimited(context.Background(), normal, headers, nil) require.False(t, svc.isOpenAIAccountRuntimeBlocked(shadow), "spark shadow must not be runtime-blocked by /responses global 429") - require.False(t, svc.isOpenAIAccountRuntimeBlocked(normal), "normal OpenAI OAuth account stays schedulable during its retry window") + require.True(t, svc.isOpenAIAccountRuntimeBlocked(normal), "normal OpenAI OAuth account with an exhausted 5h window must be paused") } func TestOpenAIRuntimeBlock_AppliesToOpenAIAPIKeyWhenRateLimitServiceStopsScheduling(t *testing.T) { diff --git a/backend/internal/service/openai_gateway_passthrough.go b/backend/internal/service/openai_gateway_passthrough.go index b87ce9df87..e3f19f335e 100644 --- a/backend/internal/service/openai_gateway_passthrough.go +++ b/backend/internal/service/openai_gateway_passthrough.go @@ -1678,7 +1678,14 @@ func (s *OpenAIGatewayService) newOpenAIStreamFailoverError( }, }) retryableOnSameAccount := openAIStreamFailedEventRetryableOnSameAccount(account, payload, message) - failoverErr := s.newOpenAIAccountFailoverError(account, statusCode, headers, payload, message, shouldDisable, retryableOnSameAccount) + // 流终止事件承载在 HTTP 200 内,外层响应头描述的是成功流状态,而不是语义上的 + // 429 事件。仅在配额分类时忽略这些头;故障转移错误仍保留它们,使 Retry-After + // 和请求 ID 能继续传递给后续处理。 + classificationHeaders := headers + if statusCode == http.StatusTooManyRequests { + classificationHeaders = nil + } + failoverErr := s.newOpenAIAccountFailoverErrorWithClassificationHeaders(account, statusCode, headers, classificationHeaders, payload, message, shouldDisable, retryableOnSameAccount) if failoverErr.IsCredentialFailure() || failoverErr.RequestScopedTransient { return failoverErr } diff --git a/backend/internal/service/openai_gateway_upstream_errors.go b/backend/internal/service/openai_gateway_upstream_errors.go index 4a4e9ab7fc..70c000e321 100644 --- a/backend/internal/service/openai_gateway_upstream_errors.go +++ b/backend/internal/service/openai_gateway_upstream_errors.go @@ -333,7 +333,20 @@ func (s *OpenAIGatewayService) newOpenAIAccountFailoverError( shouldDisable bool, retryableOnSameAccount bool, ) *UpstreamFailoverError { - oauth429Retry := s.shouldRetryOpenAIOAuth429OnSameAccount(account, statusCode, shouldDisable) + return s.newOpenAIAccountFailoverErrorWithClassificationHeaders(account, statusCode, responseHeaders, responseHeaders, responseBody, upstreamMsg, shouldDisable, retryableOnSameAccount) +} + +func (s *OpenAIGatewayService) newOpenAIAccountFailoverErrorWithClassificationHeaders( + account *Account, + statusCode int, + responseHeaders http.Header, + classificationHeaders http.Header, + responseBody []byte, + upstreamMsg string, + shouldDisable bool, + retryableOnSameAccount bool, +) *UpstreamFailoverError { + oauth429Retry := s.shouldRetryOpenAIOAuth429OnSameAccountWithResponse(account, statusCode, shouldDisable, classificationHeaders, responseBody) failoverErr := newOpenAIUpstreamFailoverError( statusCode, responseHeaders,