From ac82bedfaf5b74e40f4530e3fcb61099cf6d582e Mon Sep 17 00:00:00 2001 From: Luis Pater Date: Sat, 15 Aug 2026 14:08:00 +0800 Subject: [PATCH] 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 --- .../claude_executor_beta_policy_test.go | 6 +- .../executor/claude_executor_execute.go | 8 +- .../executor/claude_executor_fast_error.go | 96 ++- .../claude_executor_ratelimit_test.go | 790 ++++++++++++++++++ .../executor/claude_executor_request.go | 42 +- .../executor/claude_executor_stream.go | 8 +- .../executor/claude_executor_tokens.go | 4 +- .../executor/helps/claude_ratelimit.go | 229 +++++ .../executor/helps/claude_ratelimit_test.go | 141 ++++ .../auth/claude_ratelimit_cooldown_test.go | 315 +++++++ sdk/cliproxy/auth/conductor.go | 2 + sdk/cliproxy/auth/conductor_cooldown.go | 58 +- sdk/cliproxy/auth/conductor_execution.go | 12 + sdk/cliproxy/auth/conductor_home.go | 6 + sdk/cliproxy/auth/conductor_home_execution.go | 6 + sdk/cliproxy/auth/conductor_lifecycle.go | 8 + sdk/cliproxy/auth/conductor_stream.go | 18 + sdk/cliproxy/auth/selector.go | 4 +- 18 files changed, 1724 insertions(+), 29 deletions(-) create mode 100644 internal/runtime/executor/claude_executor_ratelimit_test.go create mode 100644 internal/runtime/executor/helps/claude_ratelimit.go create mode 100644 internal/runtime/executor/helps/claude_ratelimit_test.go create mode 100644 sdk/cliproxy/auth/claude_ratelimit_cooldown_test.go diff --git a/internal/runtime/executor/claude_executor_beta_policy_test.go b/internal/runtime/executor/claude_executor_beta_policy_test.go index 957186e4b..454d2e621 100644 --- a/internal/runtime/executor/claude_executor_beta_policy_test.go +++ b/internal/runtime/executor/claude_executor_beta_policy_test.go @@ -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") } diff --git a/internal/runtime/executor/claude_executor_execute.go b/internal/runtime/executor/claude_executor_execute.go index 53a54dc0a..5f8dacec6 100644 --- a/internal/runtime/executor/claude_executor_execute.go +++ b/internal/runtime/executor/claude_executor_execute.go @@ -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 { diff --git a/internal/runtime/executor/claude_executor_fast_error.go b/internal/runtime/executor/claude_executor_fast_error.go index 5b895c8f0..6ce411bd5 100644 --- a/internal/runtime/executor/claude_executor_fast_error.go +++ b/internal/runtime/executor/claude_executor_fast_error.go @@ -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 { diff --git a/internal/runtime/executor/claude_executor_ratelimit_test.go b/internal/runtime/executor/claude_executor_ratelimit_test.go new file mode 100644 index 000000000..0c1a852fb --- /dev/null +++ b/internal/runtime/executor/claude_executor_ratelimit_test.go @@ -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) + } +} diff --git a/internal/runtime/executor/claude_executor_request.go b/internal/runtime/executor/claude_executor_request.go index 6f143b319..f74320413 100644 --- a/internal/runtime/executor/claude_executor_request.go +++ b/internal/runtime/executor/claude_executor_request.go @@ -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 diff --git a/internal/runtime/executor/claude_executor_stream.go b/internal/runtime/executor/claude_executor_stream.go index 6807d50f6..1e8136b7c 100644 --- a/internal/runtime/executor/claude_executor_stream.go +++ b/internal/runtime/executor/claude_executor_stream.go @@ -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 { diff --git a/internal/runtime/executor/claude_executor_tokens.go b/internal/runtime/executor/claude_executor_tokens.go index 428650676..3765cd56e 100644 --- a/internal/runtime/executor/claude_executor_tokens.go +++ b/internal/runtime/executor/claude_executor_tokens.go @@ -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 { diff --git a/internal/runtime/executor/helps/claude_ratelimit.go b/internal/runtime/executor/helps/claude_ratelimit.go new file mode 100644 index 000000000..bd091cc51 --- /dev/null +++ b/internal/runtime/executor/helps/claude_ratelimit.go @@ -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 +} diff --git a/internal/runtime/executor/helps/claude_ratelimit_test.go b/internal/runtime/executor/helps/claude_ratelimit_test.go new file mode 100644 index 000000000..5a9e928f5 --- /dev/null +++ b/internal/runtime/executor/helps/claude_ratelimit_test.go @@ -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) + } + } + }) +} diff --git a/sdk/cliproxy/auth/claude_ratelimit_cooldown_test.go b/sdk/cliproxy/auth/claude_ratelimit_cooldown_test.go new file mode 100644 index 000000000..e994d707e --- /dev/null +++ b/sdk/cliproxy/auth/claude_ratelimit_cooldown_test.go @@ -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 +} diff --git a/sdk/cliproxy/auth/conductor.go b/sdk/cliproxy/auth/conductor.go index 2c08f1f71..c8a498da6 100644 --- a/sdk/cliproxy/auth/conductor.go +++ b/sdk/cliproxy/auth/conductor.go @@ -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 } diff --git a/sdk/cliproxy/auth/conductor_cooldown.go b/sdk/cliproxy/auth/conductor_cooldown.go index 2c06182bf..f7d0040c0 100644 --- a/sdk/cliproxy/auth/conductor_cooldown.go +++ b/sdk/cliproxy/auth/conductor_cooldown.go @@ -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 diff --git a/sdk/cliproxy/auth/conductor_execution.go b/sdk/cliproxy/auth/conductor_execution.go index cb68fae1d..511171204 100644 --- a/sdk/cliproxy/auth/conductor_execution.go +++ b/sdk/cliproxy/auth/conductor_execution.go @@ -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) diff --git a/sdk/cliproxy/auth/conductor_home.go b/sdk/cliproxy/auth/conductor_home.go index 9da186475..4df0413c5 100644 --- a/sdk/cliproxy/auth/conductor_home.go +++ b/sdk/cliproxy/auth/conductor_home.go @@ -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) diff --git a/sdk/cliproxy/auth/conductor_home_execution.go b/sdk/cliproxy/auth/conductor_home_execution.go index 89ff72847..bd89f5e35 100644 --- a/sdk/cliproxy/auth/conductor_home_execution.go +++ b/sdk/cliproxy/auth/conductor_home_execution.go @@ -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 { diff --git a/sdk/cliproxy/auth/conductor_lifecycle.go b/sdk/cliproxy/auth/conductor_lifecycle.go index a8688529a..bc9c81d66 100644 --- a/sdk/cliproxy/auth/conductor_lifecycle.go +++ b/sdk/cliproxy/auth/conductor_lifecycle.go @@ -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) diff --git a/sdk/cliproxy/auth/conductor_stream.go b/sdk/cliproxy/auth/conductor_stream.go index 09850f305..cbe9a0d79 100644 --- a/sdk/cliproxy/auth/conductor_stream.go +++ b/sdk/cliproxy/auth/conductor_stream.go @@ -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) diff --git a/sdk/cliproxy/auth/selector.go b/sdk/cliproxy/auth/selector.go index c77513727..5a940d1c5 100644 --- a/sdk/cliproxy/auth/selector.go +++ b/sdk/cliproxy/auth/selector.go @@ -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)