修正 Grok 429 边界与 stream idle 重试上限

This commit is contained in:
IanShaw
2026-08-20 06:43:31 -07:00
parent ab9cb69e7e
commit 2ab24a1e77
7 changed files with 15 additions and 16 deletions
@@ -171,6 +171,9 @@ func (s *FailoverState) HandleFailoverError(
// 同账号重试不算切换账号,粘性会话仅在实际切换时强制缓存计费。
retryCount := s.SameAccountRetryCount[accountID]
if failoverErr.SameAccountRetryMax > 0 && (retryLimit <= 0 || failoverErr.SameAccountRetryMax < retryLimit) {
retryLimit = failoverErr.SameAccountRetryMax
}
sameAccountRetryAllowed := failoverErr.RetryableOnSameAccount && retryLimit > 0 && retryCount < retryLimit
if sameAccountRetryAllowed && !failoverErr.SameAccountRetryDeadline.IsZero() {
sameAccountRetryAllowed = time.Now().Before(failoverErr.SameAccountRetryDeadline)
@@ -677,6 +677,7 @@ type UpstreamFailoverError struct {
RetryableOnSameAccount bool // 临时性错误(如 Google 间歇性 400、空响应),应在同一账号上重试 N 次再切换
SameAccountRetryDelay time.Duration // 同账号重试的最小间隔;零值使用 handler 默认值
SameAccountRetryDeadline time.Time // 同账号重试截止时间;零值表示仅受 retryLimit 限制
SameAccountRetryMax int // 可选的错误级同账号重试上限,低于 handler 默认预算时优先采用
RequestScopedTransient bool // 故障因素与账号无关(如上游按客户端身份/模型容量降载):可同账号重试,但不得据此对账号做临时封禁
SafeToFailoverAfterWrite bool // 仅写出 SSE 注释等非语义字节时,仍可在同一客户端流中切换账号
Stage GatewayFailureStage
@@ -38,6 +38,7 @@ func grokStreamIdleFailoverError(account *Account, idle time.Duration) *Upstream
// the request's retry limit.
RetryableOnSameAccount: account != nil && account.Platform == PlatformGrok,
RequestScopedTransient: true,
SameAccountRetryMax: 1,
// Permit at most one same-account replay after the idle failure. The
// deadline is anchored at failure time, so a hung stream cannot consume
// the normal three-attempt budget before failover.
@@ -23,6 +23,7 @@ func TestGrokStreamIdleFailoverError(t *testing.T) {
require.True(t, err.SafeToFailoverAfterWrite)
require.True(t, err.RetryableOnSameAccount)
require.True(t, err.RequestScopedTransient)
require.Equal(t, 1, err.SameAccountRetryMax)
require.Contains(t, string(err.ResponseBody), "empty_upstream")
require.WithinDuration(t, time.Now().Add(180*time.Second), err.SameAccountRetryDeadline, 2*time.Second)
}
@@ -402,13 +402,6 @@ func grokRetryableOnSameAccount(account *Account, statusCode int, responseBody [
if statusCode == http.StatusTooManyRequests {
return true
}
case GrokFailureRateLimit:
// A transient 429 does not identify a bad credential. Give every Grok
// account a bounded same-account retry window before failover; the
// failover loop still caps attempts and the client receives 429 after it.
if statusCode == http.StatusTooManyRequests {
return true
}
}
return account.IsPoolMode() && account.IsPoolModeRetryableStatus(statusCode)
}
@@ -418,7 +411,7 @@ func grokSameAccountRetryMetadata(account *Account, statusCode int, responseBody
return false, 0, time.Time{}
}
decision := classifyGrokUpstreamFailure(statusCode, responseBody, "")
if decision.Class != GrokFailureModelCapacity && decision.Class != GrokFailureRateLimit {
if decision.Class != GrokFailureModelCapacity {
return true, 0, time.Time{}
}
return true, 500 * time.Millisecond, time.Now().Add(30 * time.Second)
@@ -81,7 +81,7 @@ func TestGrokRetryableOnSameAccount_CapacityAndRateLimit(t *testing.T) {
account := &Account{ID: 9105, Platform: PlatformGrok, Type: AccountTypeOAuth}
require.True(t, grokRetryableOnSameAccount(account, http.StatusTooManyRequests,
[]byte(`{"error":{"message":"The model is currently at capacity due to high demand"}}`)))
require.True(t, grokRetryableOnSameAccount(account, http.StatusTooManyRequests,
require.False(t, grokRetryableOnSameAccount(account, http.StatusTooManyRequests,
[]byte(`{"error":{"message":"rate limit exceeded"}}`)))
require.False(t, grokRetryableOnSameAccount(account, http.StatusPaymentRequired,
[]byte(`{"error":{"message":"You have run out of credits or need a Grok subscription"}}`)))
@@ -118,9 +118,9 @@ func TestGrokSameAccountRetryMetadata_CapacityDeadline(t *testing.T) {
retryable, delay, deadline = grokSameAccountRetryMetadata(account, http.StatusTooManyRequests,
[]byte(`{"error":{"message":"rate limit exceeded"}}`))
require.True(t, retryable)
require.Equal(t, 500*time.Millisecond, delay)
require.WithinDuration(t, time.Now().Add(30*time.Second), deadline, 2*time.Second)
require.False(t, retryable)
require.Zero(t, delay)
require.True(t, deadline.IsZero())
}
func TestClassifyGrokUpstreamFailure_ValidationNoCool(t *testing.T) {
@@ -1794,7 +1794,7 @@ func TestForwardAsChatCompletionsForGrokStopFallsBackToXAIChatCompletions(t *tes
require.Equal(t, "grok-4.6", gjson.GetBytes(upstream.lastBody, "model").String())
require.False(t, gjson.GetBytes(upstream.lastBody, "prompt_cache_key").Exists())
require.Equal(t, "grok", result.Model)
require.Equal(t, "grok-4.5", result.UpstreamModel)
require.Equal(t, "grok-4.6", result.UpstreamModel)
require.Equal(t, 1, result.Usage.InputTokens)
require.Equal(t, 2, result.Usage.OutputTokens)
require.Equal(t, 1, result.Usage.CacheReadInputTokens)
@@ -1848,7 +1848,7 @@ func TestForwardGrokResponsesStreamingDefaultsEmptyModelTo45AndSnapshots(t *test
require.Equal(t, xai.DefaultCLIBaseURL+"/responses", upstream.lastReq.URL.String())
require.Equal(t, "Bearer access-token", upstream.lastReq.Header.Get("Authorization"))
require.Equal(t, "responses=experimental", upstream.lastReq.Header.Get("OpenAI-Beta"))
require.Equal(t, "grok-4.6", gjson.GetBytes(upstream.lastBody, "model").String())
require.Equal(t, "grok-4.5", gjson.GetBytes(upstream.lastBody, "model").String())
require.NotEmpty(t, gjson.GetBytes(upstream.lastBody, "prompt_cache_key").String())
require.Equal(t, gjson.GetBytes(upstream.lastBody, "prompt_cache_key").String(), upstream.lastReq.Header.Get(grokConversationIDHeader))
require.Equal(t, "web_search", gjson.GetBytes(upstream.lastBody, "tools.0.type").String())
@@ -2367,7 +2367,7 @@ func TestForwardAsChatCompletionsForGrokStreamingUsesRawXAIChatCompletions(t *te
require.Equal(t, "Bearer access-token", upstream.lastReq.Header.Get("Authorization"))
require.Equal(t, "text/event-stream", upstream.lastReq.Header.Get("Accept"))
require.Equal(t, xai.CLIUserAgent(xai.CLIClientVersion), upstream.lastReq.Header.Get("User-Agent"))
require.Equal(t, "grok-4.5", gjson.GetBytes(upstream.lastBody, "model").String())
require.Equal(t, "grok-4.6", gjson.GetBytes(upstream.lastBody, "model").String())
require.True(t, gjson.GetBytes(upstream.lastBody, "stream_options.include_usage").Bool())
require.True(t, result.Stream)
require.Equal(t, 6, result.Usage.InputTokens)
@@ -2648,7 +2648,7 @@ func TestForwardAsAnthropicForGrokUsesXAIResponses(t *testing.T) {
require.Equal(t, "grok-experimental", upstream.lastReq.Header.Get("OpenAI-Beta"))
require.Empty(t, upstream.lastReq.Header.Get("originator"))
require.Empty(t, upstream.lastReq.Header.Get("version"))
require.Equal(t, "grok-4.5", gjson.GetBytes(upstream.lastBody, "model").String())
require.Equal(t, "grok-4.6", gjson.GetBytes(upstream.lastBody, "model").String())
require.NotEmpty(t, gjson.GetBytes(upstream.lastBody, "prompt_cache_key").String())
require.Equal(t, gjson.GetBytes(upstream.lastBody, "prompt_cache_key").String(), upstream.lastReq.Header.Get(grokConversationIDHeader))
require.Equal(t, "web_search", gjson.GetBytes(upstream.lastBody, "tools.0.type").String())