fix(claude): add Anthropic unified rate-limit parsing helpers

- Add new Claude rate-limit helper utilities to detect unified header-based quota rejections (5h/7d), handling case-insensitive headers and missing canonical forms.
- Compute deterministic retry duration from `Retry-After` and header reset timestamps, preferring the longest applicable unified window so cooldowns can be applied consistently.
- Introduce reusable rate-limit classification logic to distinguish header-authoritative rejections from request-scoped errors (e.g. fast-mode entitlement failures) for correct credential/model cooldown behavior.

Closes: #4874
This commit is contained in:
Luis Pater
2026-08-15 14:08:00 +08:00
parent 297139cc8d
commit ac82bedfaf
18 changed files with 1724 additions and 29 deletions

View File

@@ -256,7 +256,7 @@ func TestClassifyClaudeUpstreamError_FastModeCreditsIsRequestScoped(t *testing.T
[]byte(`{"type":"error","error":{"type":"rate_limit_error","message":"Fast mode requires usage credits"}}`),
}
for _, body := range bodies {
err := classifyClaudeUpstreamError(http.StatusTooManyRequests, body)
err := classifyClaudeUpstreamError(http.StatusTooManyRequests, nil, body)
scoped, ok := err.(cliproxyexecutor.RequestScopedError)
if !ok || !scoped.IsRequestScoped() {
@@ -281,7 +281,7 @@ func TestClassifyClaudeUpstreamError_RealRateLimitStaysCredentialScoped(t *testi
[]byte(`{"type":"error","error":{"type":"rate_limit_error","message":"This organization has exceeded its usage limit."}}`),
}
for _, body := range cases {
err := classifyClaudeUpstreamError(http.StatusTooManyRequests, body)
err := classifyClaudeUpstreamError(http.StatusTooManyRequests, nil, body)
if scoped, ok := err.(cliproxyexecutor.RequestScopedError); ok && scoped.IsRequestScoped() {
t.Fatalf("genuine rate limit was misclassified as request-scoped: %s", body)
}
@@ -292,7 +292,7 @@ func TestClassifyClaudeUpstreamError_OtherStatusesUnaffected(t *testing.T) {
body := []byte(`{"error":{"message":"Usage credits are required for fast mode."}}`)
// Only 429 carries the entitlement refusal; a 500 mentioning it is still a
// credential-scoped failure worth rotating away from.
err := classifyClaudeUpstreamError(http.StatusInternalServerError, body)
err := classifyClaudeUpstreamError(http.StatusInternalServerError, nil, body)
if scoped, ok := err.(cliproxyexecutor.RequestScopedError); ok && scoped.IsRequestScoped() {
t.Fatal("non-429 status was misclassified as request-scoped")
}

View File

@@ -239,7 +239,11 @@ func (e *ClaudeExecutor) Execute(ctx context.Context, auth *cliproxyauth.Auth, r
helps.RecordAPIResponseError(ctx, e.cfg, decErr)
msg := fmt.Sprintf("failed to decode error response body: %v", decErr)
helps.LogWithRequestID(ctx).Warn(msg)
return resp, wrapClaudeFastRequestError(fastRequest, httpResp.StatusCode, statusErr{code: httpResp.StatusCode, msg: msg})
errClassified := classifyClaudeUpstreamError(httpResp.StatusCode, httpResp.Header, []byte(msg))
if fastRequest {
return resp, wrapClaudeFastRequestError(fastRequest, httpResp.StatusCode, errClassified)
}
return resp, errClassified
}
b, readErr := io.ReadAll(errBody)
if readErr != nil {
@@ -256,7 +260,7 @@ func (e *ClaudeExecutor) Execute(ctx context.Context, auth *cliproxyauth.Auth, r
if fastRequest {
return resp, newClaudeFastDirectResponseError(httpResp, b)
}
return resp, classifyClaudeUpstreamError(httpResp.StatusCode, b)
return resp, classifyClaudeUpstreamError(httpResp.StatusCode, httpResp.Header, b)
}
decodedBody, err := decodeResponseBody(httpResp.Body, claudeResponseContentEncoding(httpResp.Header))
if err != nil {

View File

@@ -2,20 +2,25 @@ package executor
import (
"bytes"
"errors"
"fmt"
"net/http"
"strings"
"time"
"github.com/router-for-me/CLIProxyAPI/v7/internal/runtime/executor/helps"
cliproxyauth "github.com/router-for-me/CLIProxyAPI/v7/sdk/cliproxy/auth"
cliproxyexecutor "github.com/router-for-me/CLIProxyAPI/v7/sdk/cliproxy/executor"
)
// claudeFastRequestError marks a Fast request failure as request-scoped. Fast
// errors must stop at the caller: they do not justify retrying another
// credential or changing the selected credential's availability.
// credential or changing the selected credential's availability, unless the failure
// is a genuine credential-level rate limit.
type claudeFastRequestError struct {
cause error
status int
cause error
status int
retryAfter *time.Duration
}
func (e *claudeFastRequestError) Error() string {
@@ -40,14 +45,43 @@ func (e *claudeFastRequestError) StatusCode() int {
}
func (e *claudeFastRequestError) IsRequestScoped() bool {
return e != nil
if e == nil {
return false
}
if e.IsCredentialScoped() {
return false
}
return true
}
func (e *claudeFastRequestError) IsCredentialScoped() bool {
if e == nil {
return false
}
type credentialScopedProvider interface {
IsCredentialScoped() bool
}
var csp credentialScopedProvider
if errors.As(e.cause, &csp) && csp != nil {
return csp.IsCredentialScoped()
}
return false
}
func (e *claudeFastRequestError) RetryAfter() *time.Duration {
if e == nil {
return nil
}
return e.retryAfter
}
// claudeFastDirectResponseError carries an upstream HTTP error response through
// the auth manager and protocol handlers without retrying or rebuilding its
// status and JSON body.
type claudeFastDirectResponseError struct {
response *cliproxyexecutor.RequestTerminatedError
response *cliproxyexecutor.RequestTerminatedError
retryAfter *time.Duration
credentialScoped bool
}
func (e *claudeFastDirectResponseError) Error() string {
@@ -65,14 +99,38 @@ func (e *claudeFastDirectResponseError) Unwrap() error {
}
func (e *claudeFastDirectResponseError) IsRequestScoped() bool {
return e != nil
if e == nil {
return false
}
if e.credentialScoped {
return false
}
return true
}
func (e *claudeFastDirectResponseError) IsCredentialScoped() bool {
if e == nil {
return false
}
return e.credentialScoped
}
func (e *claudeFastDirectResponseError) RetryAfter() *time.Duration {
if e == nil {
return nil
}
return e.retryAfter
}
func wrapClaudeFastRequestError(fastRequest bool, status int, err error) error {
if err == nil || !fastRequest {
return err
}
return &claudeFastRequestError{cause: err, status: status}
var retryAfter *time.Duration
if rap, ok := err.(interface{ RetryAfter() *time.Duration }); ok && rap != nil {
retryAfter = rap.RetryAfter()
}
return &claudeFastRequestError{cause: err, status: status, retryAfter: retryAfter}
}
func newClaudeFastDirectResponseError(resp *http.Response, body []byte) error {
@@ -84,11 +142,25 @@ func newClaudeFastDirectResponseError(resp *http.Response, body []byte) error {
// length headers that describe the compressed upstream bytes.
headers.Del("Content-Encoding")
headers.Del("Content-Length")
return &claudeFastDirectResponseError{response: &cliproxyexecutor.RequestTerminatedError{
HTTPStatus: resp.StatusCode,
Header: headers,
Body: bytes.Clone(body),
}}
var retryAfter *time.Duration
credentialScoped := false
if resp.StatusCode == http.StatusTooManyRequests {
retryAfter = helps.ParseClaudeRateLimitReset(resp.Header, time.Now())
if helps.ClaudeHeadersIndicateUnifiedRateLimitRejection(resp.Header) {
credentialScoped = true
}
}
return &claudeFastDirectResponseError{
response: &cliproxyexecutor.RequestTerminatedError{
HTTPStatus: resp.StatusCode,
Header: headers,
Body: bytes.Clone(body),
},
retryAfter: retryAfter,
credentialScoped: credentialScoped,
}
}
func claudeRequestIsFast(req *http.Request, body []byte) bool {

View File

@@ -0,0 +1,790 @@
package executor
import (
"context"
"errors"
"io"
"net/http"
"net/http/httptest"
"strconv"
"strings"
"sync/atomic"
"testing"
"time"
"github.com/google/uuid"
"github.com/router-for-me/CLIProxyAPI/v7/internal/config"
"github.com/router-for-me/CLIProxyAPI/v7/internal/registry"
cliproxyauth "github.com/router-for-me/CLIProxyAPI/v7/sdk/cliproxy/auth"
cliproxyexecutor "github.com/router-for-me/CLIProxyAPI/v7/sdk/cliproxy/executor"
sdktranslator "github.com/router-for-me/CLIProxyAPI/v7/sdk/translator"
)
type retryAfterProvider interface {
RetryAfter() *time.Duration
}
func TestClaudeExecutor_HonorsAnthropicRateLimitHeaders_Execute(t *testing.T) {
now := time.Now()
sevenDayReset := now.Add(7 * 24 * time.Hour).Unix()
fiveHourReset := now.Add(5 * time.Hour).Unix()
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "application/json")
w.Header().Set("Anthropic-Ratelimit-Unified-Status", "rejected")
w.Header().Set("Anthropic-Ratelimit-Unified-5h-Status", "allowed")
w.Header().Set("Anthropic-Ratelimit-Unified-5h-Reset", strconv.FormatInt(fiveHourReset, 10))
w.Header().Set("Anthropic-Ratelimit-Unified-7d-Status", "rejected")
w.Header().Set("Anthropic-Ratelimit-Unified-7d-Reset", strconv.FormatInt(sevenDayReset, 10))
w.Header().Set("Anthropic-Ratelimit-Unified-Representative-Claim", "seven_day")
w.Header().Set("Anthropic-Ratelimit-Unified-Reset", strconv.FormatInt(sevenDayReset, 10))
w.Header().Set("Retry-After", strconv.FormatInt(7*24*3600, 10))
w.WriteHeader(http.StatusTooManyRequests)
_, _ = w.Write([]byte(`{"type":"error","error":{"type":"rate_limit_error","message":"Number of requests has exceeded your 7-day rate limit."}}`))
}))
defer server.Close()
executor := NewClaudeExecutor(&config.Config{})
auth := &cliproxyauth.Auth{
ID: "claude-auth-1",
Provider: "claude",
Attributes: map[string]string{
"api_key": "test-key",
"base_url": server.URL,
},
}
payload := []byte(`{"messages":[{"role":"user","content":[{"type":"text","text":"hi"}]}]}`)
_, err := executor.Execute(context.Background(), auth, cliproxyexecutor.Request{
Model: "claude-3-5-sonnet-20241022",
Payload: payload,
}, cliproxyexecutor.Options{SourceFormat: sdktranslator.FormatClaude})
if err == nil {
t.Fatal("expected error from Execute, got nil")
}
var rap retryAfterProvider
if !errors.As(err, &rap) || rap == nil {
t.Fatalf("expected error %T to implement RetryAfter() *time.Duration", err)
}
retryAfter := rap.RetryAfter()
if retryAfter == nil {
t.Fatalf("expected non-nil RetryAfter, got nil")
}
// Should be at least 7 days (reported reset) and at most 7 days + 35s (fuzz upper bound).
minExpected := 7*24*time.Hour - 5*time.Second
maxExpected := 7*24*time.Hour + 35*time.Second
if *retryAfter < minExpected || *retryAfter > maxExpected {
t.Fatalf("RetryAfter = %v, want between %v and %v", *retryAfter, minExpected, maxExpected)
}
// Verify one-time fuzz stability: repeat calls return exact same value
if second := rap.RetryAfter(); second == nil || *second != *retryAfter {
t.Fatalf("RetryAfter changed across calls: %v vs %v", *second, *retryAfter)
}
}
func TestClaudeExecutor_HonorsAnthropicRateLimitHeaders_ExecuteStream(t *testing.T) {
now := time.Now()
fiveHourReset := now.Add(5 * time.Hour).Unix()
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "application/json")
w.Header().Set("Anthropic-Ratelimit-Unified-Status", "rejected")
w.Header().Set("Anthropic-Ratelimit-Unified-5h-Status", "rejected")
w.Header().Set("Anthropic-Ratelimit-Unified-5h-Reset", strconv.FormatInt(fiveHourReset, 10))
w.Header().Set("Anthropic-Ratelimit-Unified-7d-Status", "allowed")
w.Header().Set("Anthropic-Ratelimit-Unified-7d-Reset", strconv.FormatInt(now.Add(7*24*time.Hour).Unix(), 10))
w.WriteHeader(http.StatusTooManyRequests)
_, _ = w.Write([]byte(`{"type":"error","error":{"type":"rate_limit_error","message":"5-hour limit exceeded."}}`))
}))
defer server.Close()
executor := NewClaudeExecutor(&config.Config{})
auth := &cliproxyauth.Auth{
ID: "claude-auth-1",
Provider: "claude",
Attributes: map[string]string{
"api_key": "test-key",
"base_url": server.URL,
},
}
payload := []byte(`{"messages":[{"role":"user","content":[{"type":"text","text":"hi"}]}]}`)
_, err := executor.ExecuteStream(context.Background(), auth, cliproxyexecutor.Request{
Model: "claude-3-5-sonnet-20241022",
Payload: payload,
}, cliproxyexecutor.Options{SourceFormat: sdktranslator.FormatClaude})
if err == nil {
t.Fatal("expected error from ExecuteStream, got nil")
}
var rap retryAfterProvider
if !errors.As(err, &rap) || rap == nil {
t.Fatalf("expected error %T to implement RetryAfter() *time.Duration", err)
}
retryAfter := rap.RetryAfter()
if retryAfter == nil {
t.Fatalf("expected non-nil RetryAfter, got nil")
}
minExpected := 5*time.Hour - 5*time.Second
maxExpected := 5*time.Hour + 35*time.Second
if *retryAfter < minExpected || *retryAfter > maxExpected {
t.Fatalf("RetryAfter = %v, want between %v and %v (5h window)", *retryAfter, minExpected, maxExpected)
}
}
func TestClaudeExecutor_RateLimit_BothRejectedUsesLongest(t *testing.T) {
now := time.Now()
fiveHourReset := now.Add(5 * time.Hour).Unix()
sevenDayReset := now.Add(7 * 24 * time.Hour).Unix()
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "application/json")
w.Header().Set("Anthropic-Ratelimit-Unified-5h-Status", "rejected")
w.Header().Set("Anthropic-Ratelimit-Unified-5h-Reset", strconv.FormatInt(fiveHourReset, 10))
w.Header().Set("Anthropic-Ratelimit-Unified-7d-Status", "rejected")
w.Header().Set("Anthropic-Ratelimit-Unified-7d-Reset", strconv.FormatInt(sevenDayReset, 10))
w.WriteHeader(http.StatusTooManyRequests)
_, _ = w.Write([]byte(`{"type":"error","error":{"type":"rate_limit_error","message":"Both limits exceeded."}}`))
}))
defer server.Close()
executor := NewClaudeExecutor(&config.Config{})
auth := &cliproxyauth.Auth{
ID: "claude-auth-1",
Provider: "claude",
Attributes: map[string]string{
"api_key": "test-key",
"base_url": server.URL,
},
}
payload := []byte(`{"messages":[{"role":"user","content":[{"type":"text","text":"hi"}]}]}`)
_, err := executor.Execute(context.Background(), auth, cliproxyexecutor.Request{
Model: "claude-3-5-sonnet-20241022",
Payload: payload,
}, cliproxyexecutor.Options{SourceFormat: sdktranslator.FormatClaude})
if err == nil {
t.Fatal("expected error, got nil")
}
var rap retryAfterProvider
if !errors.As(err, &rap) || rap == nil {
t.Fatalf("expected error %T to implement RetryAfter() *time.Duration", err)
}
retryAfter := rap.RetryAfter()
if retryAfter == nil {
t.Fatalf("expected non-nil RetryAfter, got nil")
}
minExpected := 7*24*time.Hour - 5*time.Second
maxExpected := 7*24*time.Hour + 35*time.Second
if *retryAfter < minExpected || *retryAfter > maxExpected {
t.Fatalf("RetryAfter = %v, want between %v and %v", *retryAfter, minExpected, maxExpected)
}
}
func TestClaudeExecutor_RateLimit_CountTokensHonorsRateLimitReset(t *testing.T) {
now := time.Now()
sevenDayReset := now.Add(7 * 24 * time.Hour).Unix()
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "application/json")
w.Header().Set("Anthropic-Ratelimit-Unified-Status", "rejected")
w.Header().Set("Anthropic-Ratelimit-Unified-7d-Status", "rejected")
w.Header().Set("Anthropic-Ratelimit-Unified-7d-Reset", strconv.FormatInt(sevenDayReset, 10))
w.Header().Set("Retry-After", "604800")
w.WriteHeader(http.StatusTooManyRequests)
_, _ = w.Write([]byte(`{"type":"error","error":{"type":"rate_limit_error","message":"7-day rate limit exceeded."}}`))
}))
defer server.Close()
executor := NewClaudeExecutor(&config.Config{})
auth := &cliproxyauth.Auth{
ID: "claude-auth-1",
Provider: "claude",
Attributes: map[string]string{
"api_key": "test-key",
"base_url": server.URL,
},
}
payload := []byte(`{"messages":[{"role":"user","content":[{"type":"text","text":"hi"}]}]}`)
_, err := executor.countTokensUpstream(context.Background(), auth, cliproxyexecutor.Request{
Model: "claude-3-5-sonnet-20241022",
Payload: payload,
}, cliproxyexecutor.Options{SourceFormat: sdktranslator.FormatClaude})
if err == nil {
t.Fatal("expected error from CountTokens, got nil")
}
var rap retryAfterProvider
if !errors.As(err, &rap) || rap == nil || rap.RetryAfter() == nil {
t.Fatalf("expected CountTokens rate limit error to implement RetryAfter, got %v", err)
}
type credentialScopedProvider interface {
IsCredentialScoped() bool
}
var csp credentialScopedProvider
if !errors.As(err, &csp) || csp == nil || !csp.IsCredentialScoped() {
t.Fatalf("expected CountTokens rate limit error to be credential-scoped, got %v", err)
}
minExpected := 7*24*time.Hour - 5*time.Second
maxExpected := 7*24*time.Hour + 35*time.Second
if *rap.RetryAfter() < minExpected || *rap.RetryAfter() > maxExpected {
t.Fatalf("RetryAfter = %v, want between %v and %v", *rap.RetryAfter(), minExpected, maxExpected)
}
}
func TestClaudeExecutor_RateLimit_CaseInsensitiveRawHeaderMap(t *testing.T) {
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Header()["anthropic-ratelimit-unified-status"] = []string{"rejected"}
w.Header()["anthropic-ratelimit-unified-7d-status"] = []string{"rejected"}
w.Header()["anthropic-ratelimit-unified-7d-reset"] = []string{strconv.FormatInt(time.Now().Add(2*time.Hour).Unix(), 10)}
w.Header()["retry-after"] = []string{"7200"}
w.WriteHeader(http.StatusTooManyRequests)
_, _ = w.Write([]byte(`{"type":"error","error":{"type":"rate_limit_error","message":"Too many requests."}}`))
}))
defer server.Close()
executor := NewClaudeExecutor(&config.Config{})
auth := &cliproxyauth.Auth{
ID: "claude-auth-1",
Provider: "claude",
Attributes: map[string]string{
"api_key": "test-key",
"base_url": server.URL,
},
}
payload := []byte(`{"messages":[{"role":"user","content":[{"type":"text","text":"hi"}]}]}`)
_, err := executor.Execute(context.Background(), auth, cliproxyexecutor.Request{
Model: "claude-3-5-sonnet-20241022",
Payload: payload,
}, cliproxyexecutor.Options{SourceFormat: sdktranslator.FormatClaude})
if err == nil {
t.Fatal("expected error, got nil")
}
var rap retryAfterProvider
if !errors.As(err, &rap) || rap == nil || rap.RetryAfter() == nil {
t.Fatalf("expected RetryAfter for non-canonical header map, got %v", err)
}
minExpected := 2*time.Hour - 5*time.Second
maxExpected := 2*time.Hour + 35*time.Second
if *rap.RetryAfter() < minExpected || *rap.RetryAfter() > maxExpected {
t.Fatalf("RetryAfter = %v, want between %v and %v", *rap.RetryAfter(), minExpected, maxExpected)
}
}
func TestClaudeExecutor_RateLimit_FastModeAuthoritativeRejectionHeadersOverrideBody(t *testing.T) {
var attemptsCred1 atomic.Int32
var attemptsCred2 atomic.Int32
server1 := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
attemptsCred1.Add(1)
w.Header().Set("Content-Type", "application/json")
w.Header().Set("Anthropic-Ratelimit-Unified-Status", "rejected")
w.Header().Set("Anthropic-Ratelimit-Unified-7d-Status", "rejected")
w.Header().Set("Anthropic-Ratelimit-Unified-7d-Reset", strconv.FormatInt(time.Now().Add(7*24*time.Hour).Unix(), 10))
w.Header().Set("Retry-After", "604800")
w.WriteHeader(http.StatusTooManyRequests)
// Body text mentioning fast request rejected, but headers explicitly reject unified quota
_, _ = w.Write([]byte(`{"type":"error","error":{"type":"rate_limit_error","message":"Fast request rejected"}}`))
}))
defer server1.Close()
server2 := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
attemptsCred2.Add(1)
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(http.StatusOK)
_, _ = w.Write([]byte(`{"id":"msg-fast-ok","type":"message","role":"assistant","content":[{"type":"text","text":"hello from cred2"}]}`))
}))
defer server2.Close()
cfg := &config.Config{DisableCooling: false}
manager := cliproxyauth.NewManager(nil, nil, nil)
manager.SetRetryConfig(0, 0, 2)
executor := NewClaudeExecutor(cfg)
manager.RegisterExecutor(executor)
baseID := uuid.NewString()
auth1 := &cliproxyauth.Auth{ID: baseID + "-fast-override-1", Provider: "claude", Attributes: map[string]string{"api_key": "k1", "base_url": server1.URL}}
auth2 := &cliproxyauth.Auth{ID: baseID + "-fast-override-2", Provider: "claude", Attributes: map[string]string{"api_key": "k2", "base_url": server2.URL}}
reg := registry.GetGlobalRegistry()
reg.RegisterClient(auth1.ID, "claude", []*registry.ModelInfo{{ID: "claude-3-5-sonnet-20241022"}})
reg.RegisterClient(auth2.ID, "claude", []*registry.ModelInfo{{ID: "claude-3-5-sonnet-20241022"}})
t.Cleanup(func() {
reg.UnregisterClient(auth1.ID)
reg.UnregisterClient(auth2.ID)
})
if _, err := manager.Register(context.Background(), auth1); err != nil {
t.Fatalf("register auth1: %v", err)
}
if _, err := manager.Register(context.Background(), auth2); err != nil {
t.Fatalf("register auth2: %v", err)
}
payload := []byte(`{"speed":"fast","messages":[{"role":"user","content":[{"type":"text","text":"hi"}]}]}`)
resp, err := manager.Execute(context.Background(), []string{"claude"}, cliproxyexecutor.Request{
Model: "claude-3-5-sonnet-20241022",
Payload: payload,
}, cliproxyexecutor.Options{SourceFormat: sdktranslator.FormatClaude})
if err != nil {
t.Fatalf("expected failover to cred2 when authoritative rate limit headers present, got: %v", err)
}
if len(resp.Payload) == 0 {
t.Fatal("expected response from cred2")
}
if attemptsCred1.Load() != 1 {
t.Fatalf("attempts on cred1 = %d, want 1", attemptsCred1.Load())
}
if attemptsCred2.Load() != 1 {
t.Fatalf("attempts on cred2 = %d, want 1", attemptsCred2.Load())
}
// Verify cred1 was cooled down at credential level
registeredAuth, ok := manager.GetByID(auth1.ID)
if !ok || registeredAuth == nil {
t.Fatal("auth1 not found")
}
if !registeredAuth.Unavailable || !registeredAuth.Quota.Exceeded {
t.Fatalf("cred1 was not cooled down: unavailable=%v quota=%+v", registeredAuth.Unavailable, registeredAuth.Quota)
}
}
func TestClaudeExecutor_RateLimit_FastEntitlementWithRetryAfterRemainsRequestScoped(t *testing.T) {
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "application/json")
w.Header().Set("Retry-After", "120")
w.WriteHeader(http.StatusTooManyRequests)
_, _ = w.Write([]byte(`{"type":"error","error":{"type":"rate_limit_error","message":"Usage credits are required for fast mode."}}`))
}))
defer server.Close()
cfg := &config.Config{DisableCooling: false}
manager := cliproxyauth.NewManager(nil, nil, nil)
manager.SetRetryConfig(0, 0, 2)
executor := NewClaudeExecutor(cfg)
manager.RegisterExecutor(executor)
baseID := uuid.NewString()
auth := &cliproxyauth.Auth{ID: baseID + "-fast-entitlement", Provider: "claude", Attributes: map[string]string{"api_key": "k1", "base_url": server.URL}}
reg := registry.GetGlobalRegistry()
reg.RegisterClient(auth.ID, "claude", []*registry.ModelInfo{{ID: "claude-3-5-sonnet-20241022"}})
t.Cleanup(func() {
reg.UnregisterClient(auth.ID)
})
if _, err := manager.Register(context.Background(), auth); err != nil {
t.Fatalf("register auth: %v", err)
}
payload := []byte(`{"speed":"fast","messages":[{"role":"user","content":[{"type":"text","text":"hi"}]}]}`)
_, err := manager.Execute(context.Background(), []string{"claude"}, cliproxyexecutor.Request{
Model: "claude-3-5-sonnet-20241022",
Payload: payload,
}, cliproxyexecutor.Options{SourceFormat: sdktranslator.FormatClaude})
if err == nil {
t.Fatal("expected error, got nil")
}
registeredAuth, ok := manager.GetByID(auth.ID)
if !ok || registeredAuth == nil {
t.Fatal("auth not found")
}
if registeredAuth.Unavailable || registeredAuth.Quota.Exceeded {
t.Fatalf("fast entitlement refusal incorrectly cooled down the credential: unavailable=%v quota=%+v", registeredAuth.Unavailable, registeredAuth.Quota)
}
}
func TestClaudeExecutor_AuthManager_CredentialScopeBlocksAllModelsAndAliases(t *testing.T) {
var upstreamAttempts atomic.Int32
now := time.Now()
sevenDayReset := now.Add(7 * 24 * time.Hour).Unix()
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
upstreamAttempts.Add(1)
w.Header().Set("Content-Type", "application/json")
w.Header().Set("Anthropic-Ratelimit-Unified-Status", "rejected")
w.Header().Set("Anthropic-Ratelimit-Unified-7d-Status", "rejected")
w.Header().Set("Anthropic-Ratelimit-Unified-7d-Reset", strconv.FormatInt(sevenDayReset, 10))
w.Header().Set("Retry-After", strconv.FormatInt(7*24*3600, 10))
w.WriteHeader(http.StatusTooManyRequests)
_, _ = w.Write([]byte(`{"type":"error","error":{"type":"rate_limit_error","message":"7d limit rejected."}}`))
}))
defer server.Close()
cfg := &config.Config{
DisableCooling: false,
}
manager := cliproxyauth.NewManager(nil, nil, nil)
manager.SetRetryConfig(0, 0, 0)
executor := NewClaudeExecutor(cfg)
manager.RegisterExecutor(executor)
baseID := uuid.NewString()
auth := &cliproxyauth.Auth{
ID: baseID + "-claude-cred",
Provider: "claude",
Attributes: map[string]string{
"api_key": "test-key",
"base_url": server.URL,
},
}
reg := registry.GetGlobalRegistry()
reg.RegisterClient(auth.ID, "claude", []*registry.ModelInfo{
{ID: "claude-3-5-sonnet-20241022"},
{ID: "claude-3-opus-20240229"},
{ID: "claude-3-7-sonnet-20250219"},
})
t.Cleanup(func() {
reg.UnregisterClient(auth.ID)
})
if _, errRegister := manager.Register(context.Background(), auth); errRegister != nil {
t.Fatalf("failed to register auth: %v", errRegister)
}
// 1. Initial request on sonnet triggers 429 and records 7d cooldown
payload := []byte(`{"messages":[{"role":"user","content":[{"type":"text","text":"hi"}]}]}`)
_, err := manager.Execute(context.Background(), []string{"claude"}, cliproxyexecutor.Request{
Model: "claude-3-5-sonnet-20241022",
Payload: payload,
}, cliproxyexecutor.Options{SourceFormat: sdktranslator.FormatClaude})
if err == nil {
t.Fatal("expected error on first execute, got nil")
}
if attempts := upstreamAttempts.Load(); attempts != 1 {
t.Fatalf("upstream attempts = %d, want 1", attempts)
}
// 2. Try requesting a completely different model (opus) on the same credential -> must be blocked locally
_, errOpus := manager.Execute(context.Background(), []string{"claude"}, cliproxyexecutor.Request{
Model: "claude-3-opus-20240229",
Payload: payload,
}, cliproxyexecutor.Options{SourceFormat: sdktranslator.FormatClaude})
if errOpus == nil {
t.Fatal("expected error for opus, got nil")
}
if attempts := upstreamAttempts.Load(); attempts != 1 {
t.Fatalf("upstream attempts after opus = %d, want 1 (must be blocked locally)", attempts)
}
// 3. Try requesting a thinking suffix alias on the same credential -> must also be blocked locally
_, errThinking := manager.Execute(context.Background(), []string{"claude"}, cliproxyexecutor.Request{
Model: "claude-3-7-sonnet-20250219-thinking-16k",
Payload: payload,
}, cliproxyexecutor.Options{SourceFormat: sdktranslator.FormatClaude})
if errThinking == nil {
t.Fatal("expected error for thinking suffix, got nil")
}
if attempts := upstreamAttempts.Load(); attempts != 1 {
t.Fatalf("upstream attempts after thinking suffix = %d, want 1 (must be blocked locally)", attempts)
}
// 4. Try streaming execution for opus on the same cooling credential -> must also be blocked locally
_, errStream := manager.ExecuteStream(context.Background(), []string{"claude"}, cliproxyexecutor.Request{
Model: "claude-3-opus-20240229",
Payload: payload,
}, cliproxyexecutor.Options{SourceFormat: sdktranslator.FormatClaude})
if errStream == nil {
t.Fatal("expected error for streaming opus, got nil")
}
if attempts := upstreamAttempts.Load(); attempts != 1 {
t.Fatalf("upstream attempts after streaming opus = %d, want 1 (must be blocked locally)", attempts)
}
}
func TestClaudeExecutor_AuthManager_OrdinaryModel429DoesNotBlockSiblingModels(t *testing.T) {
var attemptsSonnet atomic.Int32
var attemptsOpus atomic.Int32
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
body, _ := io.ReadAll(r.Body)
if strings.Contains(string(body), "claude-3-5-sonnet") {
attemptsSonnet.Add(1)
w.Header().Set("Content-Type", "application/json")
w.Header().Set("Anthropic-Ratelimit-Unified-5h-Status", "allowed")
w.Header().Set("Anthropic-Ratelimit-Unified-7d-Status", "allowed")
w.WriteHeader(http.StatusTooManyRequests)
_, _ = w.Write([]byte(`{"type":"error","error":{"type":"rate_limit_error","message":"Model rate limit exceeded."}}`))
return
}
attemptsOpus.Add(1)
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(http.StatusOK)
_, _ = w.Write([]byte(`{"id":"msg-opus","type":"message","role":"assistant","content":[{"type":"text","text":"hello from opus"}]}`))
}))
defer server.Close()
cfg := &config.Config{DisableCooling: false}
manager := cliproxyauth.NewManager(nil, nil, nil)
manager.SetRetryConfig(0, 0, 0)
executor := NewClaudeExecutor(cfg)
manager.RegisterExecutor(executor)
baseID := uuid.NewString()
auth := &cliproxyauth.Auth{
ID: baseID + "-ordinary-429",
Provider: "claude",
Attributes: map[string]string{
"api_key": "test-key",
"base_url": server.URL,
},
}
reg := registry.GetGlobalRegistry()
reg.RegisterClient(auth.ID, "claude", []*registry.ModelInfo{
{ID: "claude-3-5-sonnet-20241022"},
{ID: "claude-3-opus-20240229"},
})
t.Cleanup(func() {
reg.UnregisterClient(auth.ID)
})
if _, err := manager.Register(context.Background(), auth); err != nil {
t.Fatalf("register auth: %v", err)
}
// 1. Initial request on sonnet triggers ordinary model 429
payloadSonnet := []byte(`{"messages":[{"role":"user","content":[{"type":"text","text":"hi"}]}]}`)
_, errSonnet := manager.Execute(context.Background(), []string{"claude"}, cliproxyexecutor.Request{
Model: "claude-3-5-sonnet-20241022",
Payload: payloadSonnet,
}, cliproxyexecutor.Options{SourceFormat: sdktranslator.FormatClaude})
if errSonnet == nil {
t.Fatal("expected error on sonnet execute, got nil")
}
if attemptsSonnet.Load() != 1 {
t.Fatalf("sonnet attempts = %d, want 1", attemptsSonnet.Load())
}
// 2. Request on opus MUST succeed on the same credential (not blocked by ordinary model-level 429)
payloadOpus := []byte(`{"messages":[{"role":"user","content":[{"type":"text","text":"hi opus"}]}]}`)
respOpus, errOpus := manager.Execute(context.Background(), []string{"claude"}, cliproxyexecutor.Request{
Model: "claude-3-opus-20240229",
Payload: payloadOpus,
}, cliproxyexecutor.Options{SourceFormat: sdktranslator.FormatClaude})
if errOpus != nil {
t.Fatalf("expected opus to succeed on same credential, got error: %v", errOpus)
}
if len(respOpus.Payload) == 0 {
t.Fatal("expected non-empty response for opus")
}
if attemptsOpus.Load() != 1 {
t.Fatalf("opus attempts = %d, want 1", attemptsOpus.Load())
}
}
func TestClaudeExecutor_AuthManager_AlternativeCredentialCanBeSelected(t *testing.T) {
var attemptsCred1 atomic.Int32
var attemptsCred2 atomic.Int32
now := time.Now()
sevenDayReset := now.Add(7 * 24 * time.Hour).Unix()
server1 := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
attemptsCred1.Add(1)
w.Header().Set("Content-Type", "application/json")
w.Header().Set("Anthropic-Ratelimit-Unified-Status", "rejected")
w.Header().Set("Anthropic-Ratelimit-Unified-7d-Status", "rejected")
w.Header().Set("Anthropic-Ratelimit-Unified-7d-Reset", strconv.FormatInt(sevenDayReset, 10))
w.WriteHeader(http.StatusTooManyRequests)
_, _ = w.Write([]byte(`{"type":"error","error":{"type":"rate_limit_error","message":"7d limit rejected."}}`))
}))
defer server1.Close()
server2 := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
attemptsCred2.Add(1)
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(http.StatusOK)
_, _ = w.Write([]byte(`{"id":"msg-123","type":"message","role":"assistant","content":[{"type":"text","text":"hello from cred2"}]}`))
}))
defer server2.Close()
cfg := &config.Config{
DisableCooling: false,
}
manager := cliproxyauth.NewManager(nil, nil, nil)
manager.SetRetryConfig(0, 0, 2)
executor := NewClaudeExecutor(cfg)
manager.RegisterExecutor(executor)
baseID := uuid.NewString()
auth1 := &cliproxyauth.Auth{
ID: baseID + "-claude-cred-1",
Provider: "claude",
Attributes: map[string]string{
"api_key": "test-key-1",
"base_url": server1.URL,
},
}
auth2 := &cliproxyauth.Auth{
ID: baseID + "-claude-cred-2",
Provider: "claude",
Attributes: map[string]string{
"api_key": "test-key-2",
"base_url": server2.URL,
},
}
reg := registry.GetGlobalRegistry()
reg.RegisterClient(auth1.ID, "claude", []*registry.ModelInfo{{ID: "claude-3-5-sonnet-20241022"}})
reg.RegisterClient(auth2.ID, "claude", []*registry.ModelInfo{{ID: "claude-3-5-sonnet-20241022"}})
t.Cleanup(func() {
reg.UnregisterClient(auth1.ID)
reg.UnregisterClient(auth2.ID)
})
if _, err := manager.Register(context.Background(), auth1); err != nil {
t.Fatalf("register auth1: %v", err)
}
if _, err := manager.Register(context.Background(), auth2); err != nil {
t.Fatalf("register auth2: %v", err)
}
payload := []byte(`{"messages":[{"role":"user","content":[{"type":"text","text":"hi"}]}]}`)
resp, err := manager.Execute(context.Background(), []string{"claude"}, cliproxyexecutor.Request{
Model: "claude-3-5-sonnet-20241022",
Payload: payload,
}, cliproxyexecutor.Options{SourceFormat: sdktranslator.FormatClaude})
if err != nil {
t.Fatalf("expected successful failover to cred2, got error: %v", err)
}
if attemptsCred1.Load() != 1 {
t.Fatalf("attempts on cred1 = %d, want 1", attemptsCred1.Load())
}
if attemptsCred2.Load() != 1 {
t.Fatalf("attempts on cred2 = %d, want 1", attemptsCred2.Load())
}
if len(resp.Payload) == 0 {
t.Fatal("expected non-empty response payload from cred2")
}
// Next request should directly use cred2 without attempting cred1 (which is cooling down)
resp2, err2 := manager.Execute(context.Background(), []string{"claude"}, cliproxyexecutor.Request{
Model: "claude-3-5-sonnet-20241022",
Payload: payload,
}, cliproxyexecutor.Options{SourceFormat: sdktranslator.FormatClaude})
if err2 != nil {
t.Fatalf("expected successful request on cred2, got error: %v", err2)
}
if attemptsCred1.Load() != 1 {
t.Fatalf("attempts on cred1 after 2nd request = %d, want 1 (must stay 1)", attemptsCred1.Load())
}
if attemptsCred2.Load() != 2 {
t.Fatalf("attempts on cred2 after 2nd request = %d, want 2", attemptsCred2.Load())
}
if len(resp2.Payload) == 0 {
t.Fatal("expected non-empty response payload from 2nd request")
}
}
func TestClaudeExecutor_AuthManager_MultiModelPoolStreamStopsProbingOn429(t *testing.T) {
var attemptsCred1 atomic.Int32
var attemptsCred2 atomic.Int32
server1 := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
attemptsCred1.Add(1)
w.Header().Set("Content-Type", "application/json")
w.Header().Set("Anthropic-Ratelimit-Unified-Status", "rejected")
w.Header().Set("Anthropic-Ratelimit-Unified-7d-Status", "rejected")
w.Header().Set("Anthropic-Ratelimit-Unified-7d-Reset", strconv.FormatInt(time.Now().Add(7*24*time.Hour).Unix(), 10))
w.Header().Set("Retry-After", "604800")
w.WriteHeader(http.StatusTooManyRequests)
_, _ = w.Write([]byte(`{"type":"error","error":{"type":"rate_limit_error","message":"rate limited"}}`))
}))
defer server1.Close()
server2 := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
attemptsCred2.Add(1)
w.Header().Set("Content-Type", "text/event-stream")
w.WriteHeader(http.StatusOK)
_, _ = w.Write([]byte("event: message_start\ndata: {\"type\":\"message_start\",\"message\":{\"id\":\"msg-1\",\"model\":\"claude-3-5-sonnet-20241022\"}}\n\nevent: message_stop\ndata: {\"type\":\"message_stop\"}\n\n"))
}))
defer server2.Close()
cfg := &config.Config{DisableCooling: false}
manager := cliproxyauth.NewManager(nil, nil, nil)
manager.SetRetryConfig(0, 0, 2)
manager.SetOAuthModelAlias(map[string][]config.OAuthModelAlias{
"claude": {
{Name: "claude-3-5-sonnet-20241022", Alias: "claude-pool-alias"},
{Name: "claude-3-opus-20240229", Alias: "claude-pool-alias"},
},
})
executor := NewClaudeExecutor(cfg)
manager.RegisterExecutor(executor)
baseID := uuid.NewString()
auth1 := &cliproxyauth.Auth{ID: baseID + "-pool-1", Provider: "claude", Attributes: map[string]string{"api_key": "k1", "base_url": server1.URL}}
auth2 := &cliproxyauth.Auth{ID: baseID + "-pool-2", Provider: "claude", Attributes: map[string]string{"api_key": "k2", "base_url": server2.URL}}
reg := registry.GetGlobalRegistry()
reg.RegisterClient(auth1.ID, "claude", []*registry.ModelInfo{
{ID: "claude-pool-alias"},
{ID: "claude-3-5-sonnet-20241022"},
{ID: "claude-3-opus-20240229"},
})
reg.RegisterClient(auth2.ID, "claude", []*registry.ModelInfo{
{ID: "claude-pool-alias"},
{ID: "claude-3-5-sonnet-20241022"},
{ID: "claude-3-opus-20240229"},
})
t.Cleanup(func() {
reg.UnregisterClient(auth1.ID)
reg.UnregisterClient(auth2.ID)
})
if _, err := manager.Register(context.Background(), auth1); err != nil {
t.Fatalf("register auth1: %v", err)
}
if _, err := manager.Register(context.Background(), auth2); err != nil {
t.Fatalf("register auth2: %v", err)
}
payload := []byte(`{"messages":[{"role":"user","content":[{"type":"text","text":"hi"}]}]}`)
res, err := manager.ExecuteStream(context.Background(), []string{"claude"}, cliproxyexecutor.Request{
Model: "claude-pool-alias",
Payload: payload,
}, cliproxyexecutor.Options{SourceFormat: sdktranslator.FormatClaude})
if err != nil {
t.Fatalf("ExecuteStream failed: %v", err)
}
for chunk := range res.Chunks {
if chunk.Err != nil {
t.Fatalf("unexpected chunk error: %v", chunk.Err)
}
}
// Must have tried cred1 exactly once (did NOT probe the 2nd model on cred1 after 429) and failed over to cred2
if got := attemptsCred1.Load(); got != 1 {
t.Fatalf("attempts on cred1 = %d, want 1 (must not probe other models on cooled cred)", got)
}
if got := attemptsCred2.Load(); got != 1 {
t.Fatalf("attempts on cred2 = %d, want 1", got)
}
}

View File

@@ -13,6 +13,7 @@ import (
"net/url"
"sort"
"strings"
"time"
"github.com/andybalholm/brotli"
"github.com/google/uuid"
@@ -262,6 +263,23 @@ func (claudeEntitlementError) IsRequestScoped() bool {
return true
}
func (claudeEntitlementError) IsCredentialScoped() bool {
return false
}
type claudeRateLimitError struct {
statusErr
credentialScoped bool
}
func (e claudeRateLimitError) IsCredentialScoped() bool {
return e.credentialScoped
}
func (e claudeRateLimitError) IsRequestScoped() bool {
return false
}
// classifyClaudeUpstreamError promotes upstream refusals that no other credential
// can satisfy into request-scoped errors.
//
@@ -272,10 +290,21 @@ func (claudeEntitlementError) IsRequestScoped() bool {
// next one, which returns the same 429. A single speed:"fast" request would walk
// the whole Claude pool and cool down every credential, all of which remain
// perfectly healthy for ordinary traffic. The refusal belongs to the request.
func classifyClaudeUpstreamError(statusCode int, body []byte) error {
err := statusErr{code: statusCode, msg: string(body)}
if statusCode == http.StatusTooManyRequests && claudeBodyIndicatesFastModeCredits(body) {
return claudeEntitlementError{err}
func classifyClaudeUpstreamError(statusCode int, headers http.Header, body []byte) error {
var retryAfter *time.Duration
if statusCode == http.StatusTooManyRequests || (statusCode >= 400 && statusCode < 600) {
retryAfter = helps.ParseClaudeRateLimitReset(headers, time.Now())
}
err := statusErr{code: statusCode, msg: string(body), retryAfter: retryAfter}
if statusCode == http.StatusTooManyRequests {
if helps.ClaudeHeadersIndicateUnifiedRateLimitRejection(headers) {
return claudeRateLimitError{statusErr: err, credentialScoped: true}
}
if claudeBodyIndicatesFastModeCredits(body) {
return claudeEntitlementError{err}
}
// Ordinary model-level Claude 429 (not a unified 5h/7d rejection)
return claudeRateLimitError{statusErr: err, credentialScoped: false}
}
return err
}
@@ -287,8 +316,9 @@ func claudeBodyIndicatesFastModeCredits(body []byte) bool {
if message == "" {
message = strings.ToLower(string(body))
}
return strings.Contains(message, "fast mode") &&
(strings.Contains(message, "usage credits") || strings.Contains(message, "credits are required"))
return strings.Contains(message, "fast request rejected") ||
(strings.Contains(message, "fast") &&
(strings.Contains(message, "usage credits") || strings.Contains(message, "credits are required")))
}
// claudeRequestedBetas collects every beta the caller asked for, from the

View File

@@ -232,7 +232,11 @@ func (e *ClaudeExecutor) ExecuteStream(ctx context.Context, auth *cliproxyauth.A
helps.RecordAPIResponseError(ctx, e.cfg, decErr)
msg := fmt.Sprintf("failed to decode error response body: %v", decErr)
helps.LogWithRequestID(ctx).Warn(msg)
return nil, wrapClaudeFastRequestError(fastRequest, httpResp.StatusCode, statusErr{code: httpResp.StatusCode, msg: msg})
errClassified := classifyClaudeUpstreamError(httpResp.StatusCode, httpResp.Header, []byte(msg))
if fastRequest {
return nil, wrapClaudeFastRequestError(fastRequest, httpResp.StatusCode, errClassified)
}
return nil, errClassified
}
b, readErr := io.ReadAll(errBody)
if readErr != nil {
@@ -249,7 +253,7 @@ func (e *ClaudeExecutor) ExecuteStream(ctx context.Context, auth *cliproxyauth.A
if fastRequest {
return nil, newClaudeFastDirectResponseError(httpResp, b)
}
return nil, classifyClaudeUpstreamError(httpResp.StatusCode, b)
return nil, classifyClaudeUpstreamError(httpResp.StatusCode, httpResp.Header, b)
}
decodedBody, err := decodeResponseBody(httpResp.Body, claudeResponseContentEncoding(httpResp.Header))
if err != nil {

View File

@@ -258,7 +258,7 @@ func (e *ClaudeExecutor) countTokensUpstream(ctx context.Context, auth *cliproxy
helps.RecordAPIResponseError(ctx, e.cfg, decErr)
msg := fmt.Sprintf("failed to decode error response body: %v", decErr)
helps.LogWithRequestID(ctx).Warn(msg)
return cliproxyexecutor.Response{}, statusErr{code: resp.StatusCode, msg: msg}
return cliproxyexecutor.Response{}, classifyClaudeUpstreamError(resp.StatusCode, resp.Header, []byte(msg))
}
b, readErr := io.ReadAll(errBody)
if readErr != nil {
@@ -271,7 +271,7 @@ func (e *ClaudeExecutor) countTokensUpstream(ctx context.Context, auth *cliproxy
if errClose := errBody.Close(); errClose != nil {
log.Errorf("response body close error: %v", errClose)
}
return cliproxyexecutor.Response{}, statusErr{code: resp.StatusCode, msg: string(b)}
return cliproxyexecutor.Response{}, classifyClaudeUpstreamError(resp.StatusCode, resp.Header, b)
}
decodedBody, err := decodeResponseBody(resp.Body, claudeResponseContentEncoding(resp.Header))
if err != nil {

View File

@@ -0,0 +1,229 @@
package helps
import (
cryptorand "crypto/rand"
"math/big"
"net/http"
"strconv"
"strings"
"time"
log "github.com/sirupsen/logrus"
)
const (
defaultClaudeRateLimitFuzzMinSeconds = 1
defaultClaudeRateLimitFuzzMaxSeconds = 30
)
// ClaudeHeadersIndicateUnifiedRateLimitRejection reports whether response headers explicitly
// declare an Anthropic unified 5h or 7d rate-limit rejection.
func ClaudeHeadersIndicateUnifiedRateLimitRejection(headers http.Header) bool {
if headers == nil {
return false
}
unifiedStatus := strings.ToLower(strings.TrimSpace(getHeaderCaseInsensitive(headers, "Anthropic-Ratelimit-Unified-Status")))
if unifiedStatus == "rejected" {
return true
}
status5h := strings.ToLower(strings.TrimSpace(getHeaderCaseInsensitive(headers, "Anthropic-Ratelimit-Unified-5h-Status")))
if status5h == "rejected" {
return true
}
status7d := strings.ToLower(strings.TrimSpace(getHeaderCaseInsensitive(headers, "Anthropic-Ratelimit-Unified-7d-Status")))
if status7d == "rejected" {
return true
}
return false
}
// ParseClaudeRateLimitReset inspects Anthropic response headers for unified rate-limit
// and standard Retry-After reset information, returning the conservative cooldown
// duration including a bounded non-negative random grace period.
// If no valid future reset information is present, it returns nil.
func ParseClaudeRateLimitReset(headers http.Header, now time.Time) *time.Duration {
return parseClaudeRateLimitResetWithFuzz(headers, now, defaultClaudeRateLimitFuzzMinSeconds, defaultClaudeRateLimitFuzzMaxSeconds)
}
func parseClaudeRateLimitResetWithFuzz(headers http.Header, now time.Time, minFuzzSec, maxFuzzSec int) *time.Duration {
if headers == nil {
return nil
}
unifiedStatus := strings.ToLower(strings.TrimSpace(getHeaderCaseInsensitive(headers, "Anthropic-Ratelimit-Unified-Status")))
status5h := strings.ToLower(strings.TrimSpace(getHeaderCaseInsensitive(headers, "Anthropic-Ratelimit-Unified-5h-Status")))
status7d := strings.ToLower(strings.TrimSpace(getHeaderCaseInsensitive(headers, "Anthropic-Ratelimit-Unified-7d-Status")))
var candidateDeadlines []time.Time
var rejectedWindows []string
if unifiedStatus == "rejected" {
rejectedWindows = append(rejectedWindows, "unified")
}
if status5h == "rejected" {
rejectedWindows = append(rejectedWindows, "5h")
}
if status7d == "rejected" {
rejectedWindows = append(rejectedWindows, "7d")
}
// 1. Retry-After header
if rawRetryAfter := getHeaderCaseInsensitive(headers, "Retry-After"); rawRetryAfter != "" {
if !containsString(rejectedWindows, "retry-after") {
rejectedWindows = append(rejectedWindows, "retry-after")
}
if t, ok := parseRetryAfterHeader(rawRetryAfter, now); ok && t.After(now) {
candidateDeadlines = append(candidateDeadlines, t)
}
}
// 2. 5-hour window reset (only when rejected)
if status5h == "rejected" {
if raw := getHeaderCaseInsensitive(headers, "Anthropic-Ratelimit-Unified-5h-Reset"); raw != "" {
if t, ok := parseUnixOrTimestamp(raw); ok && t.After(now) {
candidateDeadlines = append(candidateDeadlines, t)
}
}
}
// 3. 7-day window reset (only when rejected)
if status7d == "rejected" {
if raw := getHeaderCaseInsensitive(headers, "Anthropic-Ratelimit-Unified-7d-Reset"); raw != "" {
if t, ok := parseUnixOrTimestamp(raw); ok && t.After(now) {
candidateDeadlines = append(candidateDeadlines, t)
}
}
}
// 4. Unified reset header:
unifiedRejected := unifiedStatus == "rejected" || status5h == "rejected" || status7d == "rejected" ||
(unifiedStatus == "" && status5h != "allowed" && status7d != "allowed")
if unifiedRejected {
if raw := getHeaderCaseInsensitive(headers, "Anthropic-Ratelimit-Unified-Reset"); raw != "" {
if !containsString(rejectedWindows, "unified") {
rejectedWindows = append(rejectedWindows, "unified")
}
if t, ok := parseUnixOrTimestamp(raw); ok && t.After(now) {
candidateDeadlines = append(candidateDeadlines, t)
}
}
}
if len(candidateDeadlines) == 0 {
if len(rejectedWindows) > 0 {
log.WithFields(log.Fields{
"rejected_windows": strings.Join(rejectedWindows, ","),
"status": "fallback_exponential_backoff",
}).Info("Anthropic rate limit window rejected; falling back to generic exponential backoff")
}
return nil
}
// Pick the latest applicable deadline across rejected windows
var latestDeadline time.Time
for _, deadline := range candidateDeadlines {
if deadline.After(latestDeadline) {
latestDeadline = deadline
}
}
if latestDeadline.IsZero() || !latestDeadline.After(now) {
if len(rejectedWindows) > 0 {
log.WithFields(log.Fields{
"rejected_windows": strings.Join(rejectedWindows, ","),
"status": "fallback_exponential_backoff",
}).Info("Anthropic rate limit window rejected; falling back to generic exponential backoff")
}
return nil
}
baseDuration := latestDeadline.Sub(now)
fuzz := randomClaudeFuzzDuration(minFuzzSec, maxFuzzSec)
effectiveDuration := baseDuration + fuzz
log.WithFields(log.Fields{
"rejected_windows": strings.Join(rejectedWindows, ","),
"effective_cooldown": effectiveDuration.String(),
"base_cooldown": baseDuration.String(),
"fuzz": fuzz.String(),
"deadline": latestDeadline.Format(time.RFC3339),
}).Info("parsed Anthropic rate limit reset headers")
return &effectiveDuration
}
func containsString(list []string, target string) bool {
for _, item := range list {
if item == target {
return true
}
}
return false
}
func getHeaderCaseInsensitive(h http.Header, target string) string {
if h == nil {
return ""
}
if val := h.Get(target); val != "" {
return val
}
for k, v := range h {
if strings.EqualFold(k, target) && len(v) > 0 {
return v[0]
}
}
return ""
}
func parseUnixOrTimestamp(raw string) (time.Time, bool) {
raw = strings.TrimSpace(raw)
if raw == "" {
return time.Time{}, false
}
if sec, err := strconv.ParseFloat(raw, 64); err == nil && sec > 0 {
secInt := int64(sec)
nsec := int64((sec - float64(secInt)) * 1e9)
return time.Unix(secInt, nsec), true
}
if t, err := time.Parse(time.RFC3339, raw); err == nil {
return t, true
}
if t, err := http.ParseTime(raw); err == nil {
return t, true
}
return time.Time{}, false
}
func parseRetryAfterHeader(raw string, now time.Time) (time.Time, bool) {
raw = strings.TrimSpace(raw)
if raw == "" {
return time.Time{}, false
}
if sec, err := strconv.ParseFloat(raw, 64); err == nil && sec > 0 {
d := time.Duration(sec * float64(time.Second))
return now.Add(d), true
}
if t, err := http.ParseTime(raw); err == nil {
return t, true
}
if t, err := time.Parse(time.RFC3339, raw); err == nil {
return t, true
}
return time.Time{}, false
}
func randomClaudeFuzzDuration(minSec, maxSec int) time.Duration {
if maxSec <= minSec {
if minSec < 0 {
return 0
}
return time.Duration(minSec) * time.Second
}
nBig, err := cryptorand.Int(cryptorand.Reader, big.NewInt(int64(maxSec-minSec+1)))
if err != nil {
return time.Duration(minSec) * time.Second
}
return time.Duration(minSec+int(nBig.Int64())) * time.Second
}

View File

@@ -0,0 +1,141 @@
package helps
import (
"net/http"
"strconv"
"testing"
"time"
)
func TestParseClaudeRateLimitReset_AllCases(t *testing.T) {
now := time.Now()
t.Run("nil headers returns nil", func(t *testing.T) {
if got := ParseClaudeRateLimitReset(nil, now); got != nil {
t.Fatalf("expected nil, got %v", got)
}
})
t.Run("empty headers returns nil", func(t *testing.T) {
h := make(http.Header)
if got := ParseClaudeRateLimitReset(h, now); got != nil {
t.Fatalf("expected nil, got %v", got)
}
})
t.Run("retry-after only seconds", func(t *testing.T) {
h := make(http.Header)
h.Set("Retry-After", "60")
got := parseClaudeRateLimitResetWithFuzz(h, now, 0, 0)
if got == nil {
t.Fatal("expected non-nil RetryAfter")
}
if *got != 60*time.Second {
t.Fatalf("expected 60s, got %v", *got)
}
})
t.Run("retry-after HTTP date", func(t *testing.T) {
h := make(http.Header)
futureTime := now.Add(90 * time.Second).UTC().Truncate(time.Second)
h.Set("Retry-After", futureTime.Format(http.TimeFormat))
got := parseClaudeRateLimitResetWithFuzz(h, now, 0, 0)
if got == nil {
t.Fatal("expected non-nil RetryAfter")
}
if *got < 89*time.Second || *got > 91*time.Second {
t.Fatalf("expected ~90s, got %v", *got)
}
})
t.Run("5h rejected and 7d allowed with unified reset", func(t *testing.T) {
h := make(http.Header)
// Missing Anthropic-Ratelimit-Unified-Status, 5h is rejected, 7d is allowed
h.Set("Anthropic-Ratelimit-Unified-5h-Status", "rejected")
h.Set("Anthropic-Ratelimit-Unified-5h-Reset", strconv.FormatInt(now.Add(5*time.Hour).Unix(), 10))
h.Set("Anthropic-Ratelimit-Unified-7d-Status", "allowed")
h.Set("Anthropic-Ratelimit-Unified-7d-Reset", strconv.FormatInt(now.Add(7*24*time.Hour).Unix(), 10))
h.Set("Anthropic-Ratelimit-Unified-Reset", strconv.FormatInt(now.Add(5*time.Hour).Unix(), 10))
got := parseClaudeRateLimitResetWithFuzz(h, now, 0, 0)
if got == nil {
t.Fatal("expected non-nil RetryAfter")
}
if *got < 5*time.Hour-5*time.Second || *got > 5*time.Hour+5*time.Second {
t.Fatalf("expected ~5h, got %v", *got)
}
})
t.Run("7d rejected and 5h allowed", func(t *testing.T) {
h := make(http.Header)
h.Set("Anthropic-Ratelimit-Unified-Status", "rejected")
h.Set("Anthropic-Ratelimit-Unified-5h-Status", "allowed")
h.Set("Anthropic-Ratelimit-Unified-5h-Reset", strconv.FormatInt(now.Add(5*time.Hour).Unix(), 10))
h.Set("Anthropic-Ratelimit-Unified-7d-Status", "rejected")
h.Set("Anthropic-Ratelimit-Unified-7d-Reset", strconv.FormatInt(now.Add(7*24*time.Hour).Unix(), 10))
got := parseClaudeRateLimitResetWithFuzz(h, now, 0, 0)
if got == nil {
t.Fatal("expected non-nil RetryAfter")
}
if *got < 7*24*time.Hour-5*time.Second || *got > 7*24*time.Hour+5*time.Second {
t.Fatalf("expected ~7d, got %v", *got)
}
})
t.Run("both 5h and 7d rejected chooses longest", func(t *testing.T) {
h := make(http.Header)
h.Set("Anthropic-Ratelimit-Unified-5h-Status", "rejected")
h.Set("Anthropic-Ratelimit-Unified-5h-Reset", strconv.FormatInt(now.Add(5*time.Hour).Unix(), 10))
h.Set("Anthropic-Ratelimit-Unified-7d-Status", "rejected")
h.Set("Anthropic-Ratelimit-Unified-7d-Reset", strconv.FormatInt(now.Add(7*24*time.Hour).Unix(), 10))
got := parseClaudeRateLimitResetWithFuzz(h, now, 0, 0)
if got == nil {
t.Fatal("expected non-nil RetryAfter")
}
if *got < 7*24*time.Hour-5*time.Second || *got > 7*24*time.Hour+5*time.Second {
t.Fatalf("expected ~7d, got %v", *got)
}
})
t.Run("all allowed returns nil", func(t *testing.T) {
h := make(http.Header)
h.Set("Anthropic-Ratelimit-Unified-Status", "allowed")
h.Set("Anthropic-Ratelimit-Unified-5h-Status", "allowed")
h.Set("Anthropic-Ratelimit-Unified-5h-Reset", strconv.FormatInt(now.Add(5*time.Hour).Unix(), 10))
h.Set("Anthropic-Ratelimit-Unified-7d-Status", "allowed")
h.Set("Anthropic-Ratelimit-Unified-7d-Reset", strconv.FormatInt(now.Add(7*24*time.Hour).Unix(), 10))
got := ParseClaudeRateLimitReset(h, now)
if got != nil {
t.Fatalf("expected nil for allowed status, got %v", got)
}
})
t.Run("past timestamp returns nil", func(t *testing.T) {
h := make(http.Header)
h.Set("Anthropic-Ratelimit-Unified-5h-Status", "rejected")
h.Set("Anthropic-Ratelimit-Unified-5h-Reset", strconv.FormatInt(now.Add(-5*time.Hour).Unix(), 10))
got := ParseClaudeRateLimitReset(h, now)
if got != nil {
t.Fatalf("expected nil for past reset, got %v", got)
}
})
t.Run("fuzz is bounded and non-negative", func(t *testing.T) {
h := make(http.Header)
h.Set("Retry-After", "100")
for i := 0; i < 50; i++ {
got := ParseClaudeRateLimitReset(h, now)
if got == nil {
t.Fatal("expected non-nil")
}
diff := *got - 100*time.Second
if diff < 1*time.Second || diff > 30*time.Second {
t.Fatalf("fuzz %v out of bounds [1s, 30s]", diff)
}
}
})
}

View File

@@ -0,0 +1,315 @@
package auth
import (
"context"
"net/http"
"testing"
"time"
"github.com/google/uuid"
"github.com/router-for-me/CLIProxyAPI/v7/internal/registry"
)
func TestAuthManager_ConcurrentSuccessDoesNotClearActiveCredentialCooldown(t *testing.T) {
now := time.Now()
sevenDayReset := now.Add(7 * 24 * time.Hour)
manager := NewManager(nil, nil, nil)
baseID := uuid.NewString()
auth := &Auth{
ID: baseID + "-claude-concurrent",
Provider: "claude",
Attributes: map[string]string{
"api_key": "test-key",
},
}
reg := registry.GetGlobalRegistry()
reg.RegisterClient(auth.ID, "claude", []*registry.ModelInfo{
{ID: "claude-3-5-sonnet-20241022"},
{ID: "claude-3-opus-20240229"},
})
t.Cleanup(func() {
reg.UnregisterClient(auth.ID)
})
if _, err := manager.Register(context.Background(), auth); err != nil {
t.Fatalf("register auth: %v", err)
}
// 1. Request A fails with 7d credential-scoped cooldown
sevenDayDuration := 7 * 24 * time.Hour
manager.MarkResult(context.Background(), Result{
AuthID: auth.ID,
Provider: "claude",
Model: "claude-3-5-sonnet-20241022",
Success: false,
RetryAfter: &sevenDayDuration,
CredentialScope: true,
Error: &Error{HTTPStatus: http.StatusTooManyRequests, Message: "7d limit rejected"},
})
// 2. An earlier in-flight request on opus returns 200 OK after the 429
manager.MarkResult(context.Background(), Result{
AuthID: auth.ID,
Provider: "claude",
Model: "claude-3-opus-20240229",
Success: true,
})
// 3. The credential MUST still be blocked for all models
updatedAuth, ok := manager.GetByID(auth.ID)
if !ok || updatedAuth == nil {
t.Fatal("auth not found")
}
if !updatedAuth.Quota.Exceeded || !updatedAuth.Quota.NextRecoverAt.After(now.Add(6*24*time.Hour)) {
t.Fatalf("auth quota was cleared or shortened by concurrent success: quota=%+v", updatedAuth.Quota)
}
// Selecting any model on this credential must be blocked locally
for _, m := range []string{"claude-3-5-sonnet-20241022", "claude-3-opus-20240229", "claude-3-7-sonnet-20250219"} {
blocked, reason, next := isAuthBlockedForModel(updatedAuth, m, time.Now())
if !blocked {
t.Fatalf("model %q was unblocked despite active 7d credential cooldown", m)
}
if reason != blockReasonCooldown || next.Before(sevenDayReset.Add(-time.Minute)) {
t.Fatalf("model %q block reason=%v next=%v, want cooldown ~7d", m, reason, next)
}
}
}
func TestAuthManager_UpdatePreservesActiveCredentialCooldown(t *testing.T) {
now := time.Now()
manager := NewManager(nil, nil, nil)
baseID := uuid.NewString()
auth := &Auth{
ID: baseID + "-claude-update",
Provider: "claude",
Attributes: map[string]string{
"api_key": "test-key",
},
}
reg := registry.GetGlobalRegistry()
reg.RegisterClient(auth.ID, "claude", []*registry.ModelInfo{
{ID: "claude-3-5-sonnet-20241022"},
})
t.Cleanup(func() {
reg.UnregisterClient(auth.ID)
})
if _, err := manager.Register(context.Background(), auth); err != nil {
t.Fatalf("register auth: %v", err)
}
sevenDayDuration := 7 * 24 * time.Hour
manager.MarkResult(context.Background(), Result{
AuthID: auth.ID,
Provider: "claude",
Model: "claude-3-5-sonnet-20241022",
Success: false,
RetryAfter: &sevenDayDuration,
CredentialScope: true,
Error: &Error{HTTPStatus: http.StatusTooManyRequests, Message: "7d limit rejected"},
})
// Reload/update auth (e.g. config reload or token refresh)
updatedAuth := &Auth{
ID: auth.ID,
Provider: "claude",
Attributes: map[string]string{
"api_key": "test-key-updated",
},
}
if _, err := manager.Update(context.Background(), updatedAuth); err != nil {
t.Fatalf("update auth: %v", err)
}
persistedAuth, ok := manager.GetByID(auth.ID)
if !ok || persistedAuth == nil {
t.Fatal("auth not found after update")
}
if !persistedAuth.Quota.Exceeded || persistedAuth.Quota.Reason != "credential_quota" || !persistedAuth.Quota.NextRecoverAt.After(now.Add(6*24*time.Hour)) {
t.Fatalf("credential cooldown was lost after Update: quota=%+v", persistedAuth.Quota)
}
blocked, reason, _ := isAuthBlockedForModel(persistedAuth, "claude-3-5-sonnet-20241022", time.Now())
if !blocked || reason != blockReasonCooldown {
t.Fatalf("model unblocked after Update: blocked=%v reason=%v", blocked, reason)
}
}
func TestAuthManager_DisableCoolingDoesNotPermanentlyBlock(t *testing.T) {
SetQuotaCooldownDisabled(true)
t.Cleanup(func() { SetQuotaCooldownDisabled(false) })
manager := NewManager(nil, nil, nil)
baseID := uuid.NewString()
auth := &Auth{
ID: baseID + "-claude-disable-cooling",
Provider: "claude",
Attributes: map[string]string{
"api_key": "test-key",
},
}
reg := registry.GetGlobalRegistry()
reg.RegisterClient(auth.ID, "claude", []*registry.ModelInfo{
{ID: "claude-3-5-sonnet-20241022"},
{ID: "claude-3-opus-20240229"},
})
t.Cleanup(func() {
reg.UnregisterClient(auth.ID)
})
if _, err := manager.Register(context.Background(), auth); err != nil {
t.Fatalf("register auth: %v", err)
}
// 429 arrives while cooling is disabled
sevenDayDuration := 7 * 24 * time.Hour
manager.MarkResult(context.Background(), Result{
AuthID: auth.ID,
Provider: "claude",
Model: "claude-3-5-sonnet-20241022",
Success: false,
RetryAfter: &sevenDayDuration,
CredentialScope: true,
Error: &Error{HTTPStatus: http.StatusTooManyRequests, Message: "7d limit rejected"},
})
// Must NOT be blocked when cooling is disabled
for _, m := range []string{"claude-3-5-sonnet-20241022", "claude-3-opus-20240229"} {
updatedAuth, _ := manager.GetByID(auth.ID)
blocked, _, _ := isAuthBlockedForModel(updatedAuth, m, time.Now())
if blocked {
t.Fatalf("model %q was blocked even though cooling is disabled", m)
}
}
}
func TestAuthManager_NonClaudeProvider_Model429DoesNotBlockSiblingModels(t *testing.T) {
manager := NewManager(nil, nil, nil)
baseID := uuid.NewString()
auth := &Auth{
ID: baseID + "-openai-auth",
Provider: "openai",
Attributes: map[string]string{
"api_key": "test-key",
},
}
reg := registry.GetGlobalRegistry()
reg.RegisterClient(auth.ID, "openai", []*registry.ModelInfo{
{ID: "gpt-4o"},
{ID: "gpt-4o-mini"},
})
t.Cleanup(func() {
reg.UnregisterClient(auth.ID)
})
if _, err := manager.Register(context.Background(), auth); err != nil {
t.Fatalf("register auth: %v", err)
}
// Regular model 429 on gpt-4o (CredentialScope is false)
manager.MarkResult(context.Background(), Result{
AuthID: auth.ID,
Provider: "openai",
Model: "gpt-4o",
Success: false,
CredentialScope: false,
Error: &Error{HTTPStatus: http.StatusTooManyRequests, Message: "rate limit"},
})
// gpt-4o should be blocked
updatedAuth, _ := manager.GetByID(auth.ID)
blocked4o, _, _ := isAuthBlockedForModel(updatedAuth, "gpt-4o", time.Now())
if !blocked4o {
t.Fatal("gpt-4o should be blocked after 429")
}
// gpt-4o-mini MUST remain selectable (unaffected by sibling model 429)
blockedMini, _, _ := isAuthBlockedForModel(updatedAuth, "gpt-4o-mini", time.Now())
if blockedMini {
t.Fatal("gpt-4o-mini was incorrectly blocked by sibling model 429")
}
}
func TestAuthManager_CooldownPersistenceAcrossRestore(t *testing.T) {
manager := NewManager(nil, nil, nil)
baseID := uuid.NewString()
auth := &Auth{
ID: baseID + "-persistence-test",
Provider: "claude",
Attributes: map[string]string{
"api_key": "k",
},
}
reg := registry.GetGlobalRegistry()
reg.RegisterClient(auth.ID, "claude", []*registry.ModelInfo{{ID: "claude-3-5-sonnet-20241022"}})
t.Cleanup(func() {
reg.UnregisterClient(auth.ID)
})
if _, err := manager.Register(context.Background(), auth); err != nil {
t.Fatalf("register auth: %v", err)
}
futureCooldown := 7 * 24 * time.Hour
manager.MarkResult(context.Background(), Result{
AuthID: auth.ID,
Provider: "claude",
Model: "claude-3-5-sonnet-20241022",
Success: false,
RetryAfter: &futureCooldown,
CredentialScope: true,
Error: &Error{HTTPStatus: http.StatusTooManyRequests, Message: "7d rejected"},
})
records := manager.cooldownStateRecordsSnapshot()
if len(records) == 0 {
t.Fatal("expected cooldown state records to be captured")
}
// Create a new manager instance and restore state
newManager := NewManager(nil, nil, nil)
newAuth := &Auth{
ID: auth.ID,
Provider: "claude",
}
if _, err := newManager.Register(context.Background(), newAuth); err != nil {
t.Fatalf("register new auth: %v", err)
}
newManager.SetCooldownStateStore(&mockCooldownStateStore{records: records})
if err := newManager.RestoreCooldownStates(context.Background()); err != nil {
t.Fatalf("RestoreCooldownStates error: %v", err)
}
restoredAuth, ok := newManager.GetByID(auth.ID)
if !ok || restoredAuth == nil {
t.Fatal("restored auth not found")
}
if !restoredAuth.Quota.Exceeded || restoredAuth.Quota.NextRecoverAt.Before(time.Now().Add(6*24*time.Hour)) {
t.Fatalf("restored auth quota was not preserved: quota=%+v", restoredAuth.Quota)
}
}
type mockCooldownStateStore struct {
records []CooldownStateRecord
}
func (s *mockCooldownStateStore) Load(context.Context) ([]CooldownStateRecord, error) {
return s.records, nil
}
func (s *mockCooldownStateStore) Save(context.Context, []CooldownStateRecord) error {
return nil
}

View File

@@ -54,6 +54,8 @@ type Result struct {
Success bool
// RetryAfter carries a provider supplied retry hint (e.g. 429 retryDelay).
RetryAfter *time.Duration
// CredentialScope indicates that the failure affects the whole credential across models (e.g. Anthropic 5h/7d unified limits).
CredentialScope bool
// Error describes the failure when Success is false.
Error *Error
}

View File

@@ -724,7 +724,9 @@ func (m *Manager) MarkResult(ctx context.Context, result Result) {
}
if result.Success {
if modelKey != "" {
if auth.Quota.Reason == "credential_quota" && auth.Quota.NextRecoverAt.After(now) {
// Retain active credential-scoped cooldown
} else if modelKey != "" {
state := ensureModelState(auth, modelKey)
resetModelState(state, now)
updateAggregatedAvailability(auth, now)
@@ -819,6 +821,9 @@ func (m *Manager) MarkResult(ctx context.Context, result Result) {
} else {
next, backoffLevel = quotaCooldownAfterFailure(state.Quota, now)
}
if state.Quota.Exceeded && state.Quota.NextRecoverAt.After(next) {
next = state.Quota.NextRecoverAt
}
}
state.NextRetryAfter = next
state.Quota = QuotaState{
@@ -832,6 +837,34 @@ func (m *Manager) MarkResult(ctx context.Context, result Result) {
shouldSuspendModel = true
setModelQuota = true
}
if result.CredentialScope && !disableCooling {
for _, otherState := range auth.ModelStates {
if otherState != nil && otherState != state {
otherState.Unavailable = true
otherState.Status = StatusError
otherNext := next
if otherState.Quota.Exceeded && otherState.Quota.NextRecoverAt.After(otherNext) {
otherNext = otherState.Quota.NextRecoverAt
}
otherState.NextRetryAfter = otherNext
otherState.Quota = QuotaState{
Exceeded: true,
Reason: "credential_quota",
NextRecoverAt: otherNext,
BackoffLevel: backoffLevel,
}
}
}
auth.Unavailable = true
auth.Quota.Exceeded = true
auth.Quota.Reason = "credential_quota"
authNext := next
if auth.Quota.NextRecoverAt.After(authNext) {
authNext = auth.Quota.NextRecoverAt
}
auth.Quota.NextRecoverAt = authNext
auth.NextRetryAfter = authNext
}
case 408, 500, 502, 503, 504:
state.NextRetryAfter = recoverableFailureRetryAfter(now, disableCooling)
state.Unavailable = !state.NextRetryAfter.IsZero()
@@ -1066,6 +1099,10 @@ func updateAggregatedAvailability(auth *Auth, now time.Time) {
if auth == nil {
return
}
if auth.Quota.Exceeded && auth.Quota.Reason == "credential_quota" && auth.Quota.NextRecoverAt.After(now) {
auth.Unavailable = true
return
}
if len(auth.ModelStates) == 0 {
clearAggregatedAvailability(auth)
return
@@ -1123,8 +1160,13 @@ func updateAggregatedAvailability(auth *Auth, now time.Time) {
if quotaExceeded {
auth.Quota.Exceeded = true
auth.Quota.Reason = "quota"
if auth.Quota.NextRecoverAt.After(quotaRecover) {
quotaRecover = auth.Quota.NextRecoverAt
}
auth.Quota.NextRecoverAt = quotaRecover
auth.Quota.BackoffLevel = maxBackoffLevel
} else if auth.Quota.Exceeded && auth.Quota.NextRecoverAt.After(now) {
// Retain active auth-level quota cooldown
} else {
auth.Quota.Exceeded = false
auth.Quota.Reason = ""
@@ -1370,6 +1412,17 @@ func retryAfterFromError(err error) *time.Duration {
return &value
}
func isCredentialScopedError(err error) bool {
if err == nil {
return false
}
type credentialScopedProvider interface {
IsCredentialScoped() bool
}
var csp credentialScopedProvider
return errors.As(err, &csp) && csp != nil && csp.IsCredentialScoped()
}
func statusCodeFromResult(err *Error) int {
if err == nil {
return 0
@@ -1808,6 +1861,9 @@ func applyAuthFailureState(auth *Auth, resultErr *Error, retryAfter *time.Durati
} else {
next, auth.Quota.BackoffLevel = quotaCooldownAfterFailure(auth.Quota, now)
}
if auth.Quota.Exceeded && auth.Quota.NextRecoverAt.After(next) {
next = auth.Quota.NextRecoverAt
}
}
auth.Quota.NextRecoverAt = next
auth.NextRetryAfter = next

View File

@@ -374,11 +374,17 @@ func (m *Manager) executeMixedOnce(ctx context.Context, providers []string, req
if ra := retryAfterFromError(errExec); ra != nil {
result.RetryAfter = ra
}
if isCredentialScopedError(errExec) {
result.CredentialScope = true
}
m.MarkResult(execCtx, result)
if isRequestInvalidError(errExec) {
return cliproxyexecutor.Response{}, errExec
}
authErr = errExec
if result.CredentialScope {
break
}
continue
}
m.MarkResult(execCtx, result)
@@ -508,12 +514,18 @@ func (m *Manager) executeCountMixedOnce(ctx context.Context, providers []string,
if isCountTokensEndpointNotFoundError(errExec, execReq.Model) {
m.recordAvailabilityNeutralResult(execCtx, result)
} else {
if isCredentialScopedError(errExec) {
result.CredentialScope = true
}
m.MarkResult(execCtx, result)
}
if isRequestInvalidError(errExec) {
return cliproxyexecutor.Response{}, errExec
}
authErr = errExec
if result.CredentialScope {
break
}
continue
}
m.MarkResult(execCtx, result)

View File

@@ -1121,7 +1121,13 @@ func (m *Manager) tryAntigravityCreditsExecute(ctx context.Context, req cliproxy
if ra := retryAfterFromError(errExec); ra != nil {
result.RetryAfter = ra
}
if isCredentialScopedError(errExec) {
result.CredentialScope = true
}
m.MarkResult(creditsCtx, result)
if result.CredentialScope {
break
}
continue
}
m.MarkResult(creditsCtx, result)

View File

@@ -178,6 +178,9 @@ func (m *Manager) executeHome(ctx context.Context, providers []string, req clipr
}
result.Error = resultErrorFromError(errExecute)
result.RetryAfter = retryAfterFromError(errExecute)
if isCredentialScopedError(errExecute) {
result.CredentialScope = true
}
m.reportHomeResult(execCtx, result, preparedAuth)
lastErr = errExecute
if isRequestInvalidError(errExecute) {
@@ -185,6 +188,9 @@ func (m *Manager) executeHome(ctx context.Context, providers []string, req clipr
selection.End("request_invalid")
return cliproxyexecutor.Response{}, errExecute
}
if result.CredentialScope {
break
}
}
releaseAttempt()
if errEnd := m.endHomeSelectionBeforeRedispatch(ctx, selection, "execution_failed"); errEnd != nil {

View File

@@ -125,6 +125,14 @@ func (m *Manager) Update(ctx context.Context, auth *Auth) (*Auth, error) {
if len(auth.ModelStates) == 0 && len(existing.ModelStates) > 0 {
auth.ModelStates = existing.ModelStates
}
if existing.Quota.Exceeded && existing.Quota.Reason == "credential_quota" && existing.Quota.NextRecoverAt.After(time.Now()) {
auth.Unavailable = existing.Unavailable
auth.NextRetryAfter = existing.NextRetryAfter
auth.Quota = existing.Quota
if auth.Status == StatusActive {
auth.Status = existing.Status
}
}
}
now := time.Now()
cooldownStateChanged := normalizeModelStates(auth)

View File

@@ -268,11 +268,17 @@ func (m *Manager) executeStreamWithModelPool(ctx context.Context, executor Provi
rerr := resultErrorFromError(errStream)
result := Result{AuthID: auth.ID, Provider: provider, Model: resultModel, Success: false, Error: rerr}
result.RetryAfter = retryAfterFromError(errStream)
if isCredentialScopedError(errStream) {
result.CredentialScope = true
}
m.recordExecutionResult(ctx, result, auth, ephemeralResult)
if isRequestInvalidError(errStream) {
return nil, errStream
}
lastErr = errStream
if result.CredentialScope {
return nil, errStream
}
continue
}
@@ -328,6 +334,9 @@ func (m *Manager) executeStreamWithModelPool(ctx context.Context, executor Provi
rerr := resultErrorFromError(bootstrapErr)
result := Result{AuthID: auth.ID, Provider: provider, Model: resultModel, Success: false, Error: rerr}
result.RetryAfter = retryAfterFromError(bootstrapErr)
if isCredentialScopedError(bootstrapErr) {
result.CredentialScope = true
}
m.recordExecutionResult(ctx, result, auth, ephemeralResult)
discardStreamChunks(streamResult.Chunks)
return nil, bootstrapErr
@@ -336,14 +345,23 @@ func (m *Manager) executeStreamWithModelPool(ctx context.Context, executor Provi
rerr := resultErrorFromError(bootstrapErr)
result := Result{AuthID: auth.ID, Provider: provider, Model: resultModel, Success: false, Error: rerr}
result.RetryAfter = retryAfterFromError(bootstrapErr)
if isCredentialScopedError(bootstrapErr) {
result.CredentialScope = true
}
m.recordExecutionResult(ctx, result, auth, ephemeralResult)
discardStreamChunks(streamResult.Chunks)
lastErr = bootstrapErr
if result.CredentialScope {
return nil, newStreamBootstrapError(bootstrapErr, streamResult.Headers)
}
continue
}
rerr := resultErrorFromError(bootstrapErr)
result := Result{AuthID: auth.ID, Provider: provider, Model: resultModel, Success: false, Error: rerr}
result.RetryAfter = retryAfterFromError(bootstrapErr)
if isCredentialScopedError(bootstrapErr) {
result.CredentialScope = true
}
m.recordExecutionResult(ctx, result, auth, ephemeralResult)
discardStreamChunks(streamResult.Chunks)
return nil, newStreamBootstrapError(bootstrapErr, streamResult.Headers)

View File

@@ -541,6 +541,9 @@ func isAuthBlockedForModel(auth *Auth, model string, now time.Time) (bool, block
if auth.Disabled || auth.Status == StatusDisabled {
return true, blockReasonDisabled, time.Time{}
}
if auth.Quota.Exceeded && auth.Quota.Reason == "credential_quota" && auth.Quota.NextRecoverAt.After(now) {
return true, blockReasonCooldown, auth.Quota.NextRecoverAt
}
if model != "" {
if len(auth.ModelStates) > 0 {
modelKey := canonicalModelKey(model)
@@ -572,7 +575,6 @@ func isAuthBlockedForModel(auth *Auth, model string, now time.Time) (bool, block
if matched {
return blocked, blockedReason, nextRetry
}
// Auth-level availability can aggregate failures from other models.
return false, blockReasonNone, time.Time{}
}
return availabilityBlock(auth.Unavailable, auth.Quota.Exceeded, auth.NextRetryAfter, auth.Quota.NextRecoverAt, now)