fix(openai): pause OAuth accounts on exhausted quota 429

This commit is contained in:
shaw
2026-08-25 20:25:22 +08:00
parent aa2c4e8d13
commit f1aadd48d5
4 changed files with 118 additions and 14 deletions
@@ -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)
}
@@ -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) {
@@ -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
}
@@ -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,