mirror of
https://github.com/Wei-Shaw/sub2api.git
synced 2026-10-07 17:08:33 +08:00
fix(grok): unify pool mode bypass for all default cooldown paths
- Move pool mode check before the status switch so 401/402/403/5xx all skip tempUnschedule consistently - Skip rateLimitGrok in updateGrokUsageSnapshot for pool mode - Expand test coverage to all affected status codes
This commit is contained in:
@@ -259,32 +259,40 @@ func TestHandleGrokAccountUpstreamErrorEntitlement403KeepsDefaultCooldown(t *tes
|
||||
require.Less(t, repo.lastTempUnschedUntil, before.Add(31*time.Minute))
|
||||
}
|
||||
|
||||
func TestHandleGrokAccountUpstreamErrorEntitlement403RespectsPoolMode(t *testing.T) {
|
||||
t.Run("pool mode keeps scheduling state", func(t *testing.T) {
|
||||
repo := &grokQuotaAccountRepo{}
|
||||
svc := &OpenAIGatewayService{accountRepo: repo}
|
||||
account := &Account{
|
||||
ID: 4722,
|
||||
Platform: PlatformGrok,
|
||||
Type: AccountTypeAPIKey,
|
||||
Credentials: map[string]any{
|
||||
"pool_mode": true,
|
||||
},
|
||||
}
|
||||
body := []byte(`{"error":{"message":"grok access or entitlement denied"}}`)
|
||||
func TestHandleGrokAccountUpstreamErrorDefaultCooldownsRespectPoolMode(t *testing.T) {
|
||||
for _, statusCode := range []int{
|
||||
http.StatusUnauthorized,
|
||||
http.StatusPaymentRequired,
|
||||
http.StatusForbidden,
|
||||
http.StatusInternalServerError,
|
||||
} {
|
||||
t.Run(http.StatusText(statusCode), func(t *testing.T) {
|
||||
repo := &grokQuotaAccountRepo{}
|
||||
svc := &OpenAIGatewayService{accountRepo: repo}
|
||||
account := &Account{
|
||||
ID: int64(4800 + statusCode),
|
||||
Platform: PlatformGrok,
|
||||
Type: AccountTypeAPIKey,
|
||||
Credentials: map[string]any{
|
||||
"pool_mode": true,
|
||||
},
|
||||
}
|
||||
body := []byte(`{"error":{"message":"grok access or entitlement denied"}}`)
|
||||
|
||||
svc.handleGrokAccountUpstreamError(
|
||||
context.Background(), account, http.StatusForbidden, nil,
|
||||
body,
|
||||
)
|
||||
svc.handleGrokAccountUpstreamError(
|
||||
context.Background(), account, statusCode, nil, body,
|
||||
)
|
||||
|
||||
require.Zero(t, repo.tempUnschedCalls)
|
||||
require.False(t, svc.isOpenAIAccountRuntimeBlocked(account))
|
||||
require.Nil(t, account.TempUnschedulableUntil)
|
||||
require.Empty(t, account.TempUnschedulableReason)
|
||||
require.True(t, svc.shouldFailoverGrokUpstreamError(http.StatusForbidden, body))
|
||||
require.True(t, account.IsPoolModeRetryableStatus(http.StatusForbidden))
|
||||
})
|
||||
require.Zero(t, repo.tempUnschedCalls)
|
||||
require.False(t, svc.isOpenAIAccountRuntimeBlocked(account))
|
||||
require.Nil(t, account.TempUnschedulableUntil)
|
||||
require.Empty(t, account.TempUnschedulableReason)
|
||||
require.True(t, svc.shouldFailoverGrokUpstreamError(statusCode, body))
|
||||
})
|
||||
}
|
||||
|
||||
account := &Account{Type: AccountTypeAPIKey, Credentials: map[string]any{"pool_mode": true}}
|
||||
require.True(t, account.IsPoolModeRetryableStatus(http.StatusForbidden))
|
||||
|
||||
t.Run("explicit temporary rule still applies", func(t *testing.T) {
|
||||
repo := &grokQuotaAccountRepo{}
|
||||
|
||||
@@ -1129,12 +1129,11 @@ func (s *OpenAIGatewayService) updateGrokUsageSnapshot(ctx context.Context, acco
|
||||
grokQuotaSnapshotExtraKey: snapshot,
|
||||
})
|
||||
}
|
||||
// Error responses are reconciled by handleGrokAccountUpstreamError, which
|
||||
// also installs the immediate in-memory scheduling block. Successful
|
||||
// responses can still consume the last available request/token, so persist
|
||||
// that exhausted window here as a real rate limit rather than relying only
|
||||
// on the passive snapshot scheduler check.
|
||||
if hasActiveLimit {
|
||||
// Error responses are reconciled by handleGrokAccountUpstreamError. Pool-mode
|
||||
// API keys retain the snapshot for observability but leave account health to
|
||||
// the upstream pool. Other accounts install the immediate runtime and durable
|
||||
// rate-limit state when the observed window is exhausted.
|
||||
if hasActiveLimit && !account.IsPoolMode() {
|
||||
s.rateLimitGrok(stateCtx, account, resetAt)
|
||||
} else if recovery {
|
||||
clearGrokRateLimitAfterRecovery(stateCtx, s.accountRepo, account)
|
||||
@@ -1356,23 +1355,24 @@ func (s *OpenAIGatewayService) handleGrokAccountUpstreamError(ctx context.Contex
|
||||
}
|
||||
now := time.Now()
|
||||
s.updateGrokUsageSnapshot(ctx, account, parseGrokQuotaSnapshot(headers, statusCode, now))
|
||||
if statusCode == http.StatusForbidden && s.applyGrokForbiddenPolicy(ctx, account, responseBody) {
|
||||
return
|
||||
}
|
||||
if account.IsPoolMode() {
|
||||
slog.Info("grok_pool_mode_error_state_skipped", "account_id", account.ID, "status_code", statusCode)
|
||||
return
|
||||
}
|
||||
switch statusCode {
|
||||
case http.StatusUnauthorized:
|
||||
s.tempUnscheduleGrok(ctx, account, 10*time.Minute, "grok credentials unauthorized")
|
||||
case http.StatusPaymentRequired:
|
||||
s.tempUnscheduleGrok(ctx, account, 30*time.Minute, "grok payment required")
|
||||
case http.StatusForbidden:
|
||||
if s.applyGrokForbiddenPolicy(ctx, account, responseBody) {
|
||||
return
|
||||
}
|
||||
if account.IsPoolMode() {
|
||||
return
|
||||
}
|
||||
s.tempUnscheduleGrok(ctx, account, 30*time.Minute, "grok access or entitlement denied")
|
||||
case http.StatusTooManyRequests:
|
||||
// updateGrokUsageSnapshot installs both runtime and durable rate-limit state.
|
||||
// updateGrokUsageSnapshot installs rate-limit state for non-pool accounts.
|
||||
default:
|
||||
if statusCode >= 500 && !account.IsPoolMode() {
|
||||
if statusCode >= 500 {
|
||||
s.tempUnscheduleGrok(ctx, account, 2*time.Minute, "grok upstream temporary error")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -2612,6 +2612,30 @@ func TestHandleGrokAccountUpstreamError429SetsRateLimitedFromRetryAfter(t *testi
|
||||
require.Zero(t, repo.tempUnschedCalls)
|
||||
}
|
||||
|
||||
func TestHandleGrokAccountUpstreamError429PoolModeKeepsSchedulingState(t *testing.T) {
|
||||
account := &Account{
|
||||
ID: 613,
|
||||
Platform: PlatformGrok,
|
||||
Type: AccountTypeAPIKey,
|
||||
Credentials: map[string]any{
|
||||
"pool_mode": true,
|
||||
},
|
||||
}
|
||||
repo := &grokQuotaAccountRepo{}
|
||||
svc := &OpenAIGatewayService{accountRepo: repo}
|
||||
|
||||
svc.handleGrokAccountUpstreamError(
|
||||
context.Background(), account, http.StatusTooManyRequests,
|
||||
http.Header{"Retry-After": []string{"45"}}, nil,
|
||||
)
|
||||
|
||||
require.Equal(t, 1, repo.updateCalls, "pool mode should retain the quota snapshot for observability")
|
||||
require.Zero(t, repo.rateLimitedCalls)
|
||||
require.Zero(t, repo.tempUnschedCalls)
|
||||
require.False(t, svc.isOpenAIAccountRuntimeBlocked(account))
|
||||
require.Nil(t, account.RateLimitResetAt)
|
||||
}
|
||||
|
||||
func TestHandleGrokAccountUpstreamError402RecoversAfterCooldownExpiry(t *testing.T) {
|
||||
account := &Account{
|
||||
ID: 610, Platform: PlatformGrok, Type: AccountTypeOAuth,
|
||||
|
||||
Reference in New Issue
Block a user