mirror of
https://github.com/Wei-Shaw/sub2api.git
synced 2026-10-07 16:48:45 +08:00
fix(grok): treat free-usage and billing exhaustion as recoverable
Classify free-usage bodies during account tests without quarantining content policy, and mark billing/spending-limit refresh failures as transient so accounts stay probe-eligible.
This commit is contained in:
@@ -94,9 +94,9 @@ const (
|
||||
|
||||
// Grok account-test modes (admin UI). Empty / default / text = Responses probe.
|
||||
// image/video may also be inferred from model_id when mode is default.
|
||||
AccountTestModeGrokText = "text"
|
||||
AccountTestModeGrokImage = "image"
|
||||
AccountTestModeGrokVideo = "video"
|
||||
AccountTestModeGrokText = "text"
|
||||
AccountTestModeGrokImage = "image"
|
||||
AccountTestModeGrokVideo = "video"
|
||||
AccountTestModeGrokSearch = "search"
|
||||
AccountTestModeGrokTTS = "tts"
|
||||
AccountTestModeGrokSTT = "stt"
|
||||
@@ -973,6 +973,16 @@ func (s *AccountTestService) observeGrokTestResponse(ctx context.Context, accoun
|
||||
return
|
||||
}
|
||||
now := time.Now()
|
||||
// Error bodies carry Grok's free-usage, billing, and content-policy
|
||||
// classifications when quota headers are absent. Read only non-success
|
||||
// responses here, then restore the body because the caller still needs it
|
||||
// for the user-facing test result.
|
||||
var responseBody []byte
|
||||
if resp.StatusCode >= http.StatusBadRequest && resp.Body != nil {
|
||||
responseBody, _ = io.ReadAll(resp.Body)
|
||||
_ = resp.Body.Close()
|
||||
resp.Body = io.NopCloser(bytes.NewReader(responseBody))
|
||||
}
|
||||
snapshot := parseGrokQuotaSnapshot(resp.Header, resp.StatusCode, now)
|
||||
if snapshot != nil && s.accountRepo != nil {
|
||||
resetAt, limited := grokRateLimitResetAtForAccount(account, snapshot, now)
|
||||
@@ -990,14 +1000,61 @@ func (s *AccountTestService) observeGrokTestResponse(ctx context.Context, accoun
|
||||
} else if s.accountRepo != nil && isSuccessfulGrokRateLimitRecovery(account, &xai.QuotaSnapshot{StatusCode: resp.StatusCode}) {
|
||||
clearGrokRateLimitAfterRecovery(ctx, s.accountRepo, account)
|
||||
}
|
||||
if resp.StatusCode == http.StatusPaymentRequired && s.accountRepo != nil {
|
||||
if s.accountRepo == nil || len(responseBody) == 0 {
|
||||
if resp.StatusCode == http.StatusPaymentRequired && s.accountRepo != nil {
|
||||
stateCtx, cancel := openAIAccountStateContext(ctx)
|
||||
defer cancel()
|
||||
_ = s.accountRepo.SetTempUnschedulable(stateCtx, account.ID, now.Add(30*time.Minute), "grok payment required")
|
||||
}
|
||||
return
|
||||
}
|
||||
if isGrokContentPolicyRejection(resp.StatusCode, responseBody) {
|
||||
return
|
||||
}
|
||||
decision := classifyGrokUpstreamFailure(resp.StatusCode, responseBody, "")
|
||||
if decision.Class == GrokFailureFreeUsage {
|
||||
if resetAt, limited := grokRateLimitResetAtForAccount(account, snapshot, now); limited && resetAt.After(now) {
|
||||
persistGrokRateLimit(ctx, s.accountRepo, account, resetAt)
|
||||
} else {
|
||||
stateCtx, cancel := openAIAccountStateContext(ctx)
|
||||
_ = s.accountRepo.SetTempUnschedulable(stateCtx, account.ID, now.Add(grokFreeUsageProbeCooldown), "grok free usage exhausted")
|
||||
cancel()
|
||||
}
|
||||
return
|
||||
}
|
||||
if decision.Class == GrokFailureBilling && (isGrokSpendingLimitError(responseBody) || strings.Contains(strings.ToLower(decision.Reason), "credit")) {
|
||||
persistGrokRateLimit(ctx, s.accountRepo, account, grokSpendingLimitResetAt(account, now))
|
||||
return
|
||||
}
|
||||
cooldown := time.Duration(0)
|
||||
reason := ""
|
||||
switch resp.StatusCode {
|
||||
case http.StatusUnauthorized:
|
||||
cooldown, reason = 10*time.Minute, "grok oauth token unauthorized"
|
||||
case http.StatusPaymentRequired:
|
||||
cooldown, reason = 30*time.Minute, "grok payment required"
|
||||
case http.StatusForbidden:
|
||||
cooldown, reason = 30*time.Minute, "grok entitlement or subscription tier denied"
|
||||
default:
|
||||
if resp.StatusCode >= 500 {
|
||||
cooldown, reason = 2*time.Minute, "grok upstream temporary error"
|
||||
}
|
||||
}
|
||||
if decision.Class == GrokFailureBilling && cooldown == 0 {
|
||||
cooldown, reason = 30*time.Minute, "grok payment required"
|
||||
}
|
||||
if cooldown > 0 {
|
||||
stateCtx, cancel := openAIAccountStateContext(ctx)
|
||||
defer cancel()
|
||||
until := now.Add(cooldown)
|
||||
if account.TempUnschedulableUntil != nil && account.TempUnschedulableUntil.After(until) {
|
||||
until = *account.TempUnschedulableUntil
|
||||
}
|
||||
_ = s.accountRepo.SetTempUnschedulable(
|
||||
stateCtx,
|
||||
account.ID,
|
||||
time.Now().Add(30*time.Minute),
|
||||
"grok payment required",
|
||||
until,
|
||||
reason,
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -28,6 +28,60 @@ type grokAccountTestRateLimitRepo struct {
|
||||
resetAt time.Time
|
||||
}
|
||||
|
||||
func TestObserveGrokTestResponseClassifiesBodyOnlyQuotaErrors(t *testing.T) {
|
||||
account := &Account{ID: 1901, Platform: PlatformGrok, Type: AccountTypeOAuth}
|
||||
repo := &grokQuotaAccountRepo{mockAccountRepoForPlatform: &mockAccountRepoForPlatform{
|
||||
accountsByID: map[int64]*Account{account.ID: account},
|
||||
}}
|
||||
svc := &AccountTestService{accountRepo: repo}
|
||||
|
||||
resp := &http.Response{
|
||||
StatusCode: http.StatusBadRequest,
|
||||
Header: make(http.Header),
|
||||
Body: io.NopCloser(strings.NewReader(`{"error":{"code":"subscription:free-usage-exhausted","message":"included free usage exhausted"}}`)),
|
||||
}
|
||||
svc.observeGrokTestResponse(context.Background(), account, resp)
|
||||
require.Equal(t, 1, repo.tempUnschedCalls)
|
||||
require.Equal(t, "grok free usage exhausted", repo.lastTempUnschedReason)
|
||||
body, err := io.ReadAll(resp.Body)
|
||||
require.NoError(t, err)
|
||||
require.Contains(t, string(body), "free-usage-exhausted")
|
||||
}
|
||||
|
||||
func TestObserveGrokTestResponseDoesNotQuarantineContentPolicy(t *testing.T) {
|
||||
account := &Account{ID: 1902, Platform: PlatformGrok, Type: AccountTypeOAuth}
|
||||
repo := &grokQuotaAccountRepo{mockAccountRepoForPlatform: &mockAccountRepoForPlatform{
|
||||
accountsByID: map[int64]*Account{account.ID: account},
|
||||
}}
|
||||
svc := &AccountTestService{accountRepo: repo}
|
||||
resp := &http.Response{
|
||||
StatusCode: http.StatusForbidden,
|
||||
Header: make(http.Header),
|
||||
Body: io.NopCloser(strings.NewReader(`{"error":{"code":"new_sensitive","message":"text is sensitive"}}`)),
|
||||
}
|
||||
svc.observeGrokTestResponse(context.Background(), account, resp)
|
||||
require.Zero(t, repo.tempUnschedCalls)
|
||||
require.Zero(t, repo.rateLimitedCalls)
|
||||
}
|
||||
|
||||
func TestObserveGrokTestResponseKeepsEntitlement403Cooldown(t *testing.T) {
|
||||
account := &Account{ID: 1903, Platform: PlatformGrok, Type: AccountTypeOAuth}
|
||||
repo := &grokQuotaAccountRepo{mockAccountRepoForPlatform: &mockAccountRepoForPlatform{
|
||||
accountsByID: map[int64]*Account{account.ID: account},
|
||||
}}
|
||||
svc := &AccountTestService{accountRepo: repo}
|
||||
resp := &http.Response{
|
||||
StatusCode: http.StatusForbidden,
|
||||
Header: make(http.Header),
|
||||
Body: io.NopCloser(strings.NewReader(`{"error":{"message":"subscription required"}}`)),
|
||||
}
|
||||
before := time.Now()
|
||||
svc.observeGrokTestResponse(context.Background(), account, resp)
|
||||
require.Equal(t, 1, repo.tempUnschedCalls)
|
||||
require.Equal(t, "grok entitlement or subscription tier denied", repo.lastTempUnschedReason)
|
||||
require.Greater(t, repo.lastTempUnschedUntil, before.Add(29*time.Minute))
|
||||
}
|
||||
|
||||
func (r *grokAccountTestRateLimitRepo) SetRateLimited(_ context.Context, _ int64, resetAt time.Time) error {
|
||||
r.rateLimitedCalls++
|
||||
r.resetAt = resetAt
|
||||
@@ -503,12 +557,12 @@ func (c *grokRealtimeTestConn) Ping(context.Context) error { return nil }
|
||||
func (c *grokRealtimeTestConn) Close() error { return nil }
|
||||
|
||||
type grokRealtimeTestDialer struct {
|
||||
lastURL string
|
||||
lastAuth string
|
||||
lastProxy string
|
||||
conn openAIWSClientConn
|
||||
err error
|
||||
status int
|
||||
lastURL string
|
||||
lastAuth string
|
||||
lastProxy string
|
||||
conn openAIWSClientConn
|
||||
err error
|
||||
status int
|
||||
}
|
||||
|
||||
func (d *grokRealtimeTestDialer) Dial(_ context.Context, wsURL string, headers http.Header, proxyURL string) (openAIWSClientConn, int, http.Header, error) {
|
||||
|
||||
@@ -255,6 +255,11 @@ func classifyGrokCredentialFailure(account *Account, err error) grokCredentialFa
|
||||
return grokCredentialFailureClass{scope: GatewayFailureScopeAccount, reason: GrokCredentialReasonMissing, action: NextAccountRetry, permanent: true, message: "Grok OAuth credentials are missing or expired"}
|
||||
case contains("invalid_grant", "invalid_refresh_token", "token_expired", "refresh_token_reused", "refresh_token_invalidated", "app_session_terminated"):
|
||||
return grokCredentialFailureClass{scope: GatewayFailureScopeAccount, reason: GrokCredentialReasonRevoked, action: NextAccountRetry, permanent: true, message: "Grok OAuth credentials require account action"}
|
||||
case contains("spending limit", "run out of credits", "out of credits", "credits exhausted", "included free usage"):
|
||||
// Billing and rolling free-usage exhaustion recover without replacing the
|
||||
// OAuth credential. Treat refresh failures as transient so the account
|
||||
// remains eligible for a later quota probe.
|
||||
return grokCredentialFailureClass{scope: GatewayFailureScopeAccount, reason: GrokCredentialReasonRefreshTransient, action: NextAccountRetry, transient: true, message: "Grok OAuth billing quota is temporarily exhausted"}
|
||||
case contains("grok_oauth_entitlement_denied", "entitlement_denied", "access_denied", "subscription required", "no active grok subscription"):
|
||||
return grokCredentialFailureClass{scope: GatewayFailureScopeAccount, reason: GrokCredentialReasonEntitlement, action: NextAccountRetry, permanent: true, message: "Grok OAuth entitlement requires account action"}
|
||||
case errors.Is(err, errGrokOAuthConfiguredProxyMiss), contains("grok_oauth_proxy_not_found"):
|
||||
|
||||
@@ -21,6 +21,20 @@ type grokCredentialPersistingRepo struct {
|
||||
*tokenRefreshAccountRepo
|
||||
}
|
||||
|
||||
func TestClassifyGrokCredentialFailureBillingExhaustionIsTransient(t *testing.T) {
|
||||
account := expiredGrokOAuthAccountForCredentialTest(9901)
|
||||
for _, message := range []string{
|
||||
"Grok OAuth refresh failed: spending limit reached",
|
||||
"included free usage exhausted",
|
||||
"credits exhausted",
|
||||
} {
|
||||
class := classifyGrokCredentialFailure(account, errors.New(message))
|
||||
require.Equal(t, GrokCredentialReasonRefreshTransient, class.reason, message)
|
||||
require.True(t, class.transient, message)
|
||||
require.False(t, class.permanent, message)
|
||||
}
|
||||
}
|
||||
|
||||
func (r *grokCredentialPersistingRepo) SetError(ctx context.Context, id int64, message string) error {
|
||||
if err := ctx.Err(); err != nil {
|
||||
return err
|
||||
|
||||
@@ -2112,10 +2112,10 @@ func TestAccountTestServiceGrokOAuthPaymentRequiredTemporarilyUnschedulesAccount
|
||||
err := svc.testGrokAccountConnection(c, account, "grok", "", AccountTestModeDefault, AccountTestOptions{})
|
||||
|
||||
require.Error(t, err)
|
||||
require.Equal(t, 1, repo.tempUnschedCalls)
|
||||
require.Equal(t, account.ID, repo.lastTempUnschedID)
|
||||
require.Equal(t, "grok payment required", repo.lastTempUnschedReason)
|
||||
require.WithinDuration(t, before.Add(30*time.Minute), repo.lastTempUnschedUntil, time.Second)
|
||||
require.Zero(t, repo.tempUnschedCalls)
|
||||
require.Equal(t, 1, repo.rateLimitedCalls)
|
||||
require.Equal(t, account.ID, repo.lastRateLimitedID)
|
||||
require.WithinDuration(t, before.Add(grokSpendingLimitProbeCooldown), repo.lastRateLimitResetAt, time.Second)
|
||||
require.Contains(t, recorder.Body.String(), `"type":"error"`)
|
||||
require.Contains(t, recorder.Body.String(), "Grok Responses API returned 402")
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user