diff --git a/backend/internal/service/ratelimit_service.go b/backend/internal/service/ratelimit_service.go index e2704f3c20..abdcec5cbb 100644 --- a/backend/internal/service/ratelimit_service.go +++ b/backend/internal/service/ratelimit_service.go @@ -181,6 +181,16 @@ func (s *RateLimitService) HandleUpstreamError(ctx context.Context, account *Acc return true } + // Anthropic official 5h / 7d window exhaustion is a hard account limit. + // It must take precedence over user-configured 429 temp-unsched rules, + // otherwise a broad "rate limit" keyword rule can shorten a multi-hour + // cooldown to a local temporary pause. + if statusCode == http.StatusTooManyRequests && account.Platform == PlatformAnthropic { + if s.persistAnthropicExhaustedWindowLimit(ctx, account, headers) { + return false + } + } + // 先尝试临时不可调度规则(401除外) // 如果匹配成功,直接返回,不执行后续禁用逻辑 if statusCode != 401 { @@ -1083,6 +1093,115 @@ type anthropic429Result struct { fiveHourReset *time.Time // 5h window reset timestamp (for session window calculation), nil if not available } +type anthropicWindowLimit struct { + window string + resetAt time.Time + reason string +} + +func selectAnthropicExhaustedWindow(headers http.Header, now time.Time) *anthropicWindowLimit { + reset5h, ok5hReset := parseAnthropicWindowReset(headers, "5h", now) + reset7d, ok7dReset := parseAnthropicWindowReset(headers, "7d", now) + + exceeded5h := isAnthropic5hRejected(headers) || isAnthropicWindowExceeded(headers, "5h") + exceeded7d := isAnthropicWindowExceeded(headers, "7d") + + if exceeded7d && ok7dReset { + return &anthropicWindowLimit{ + window: "7d", + resetAt: reset7d, + reason: "anthropic_7d_window_exhausted", + } + } + if exceeded5h && ok5hReset { + return &anthropicWindowLimit{ + window: "5h", + resetAt: reset5h, + reason: "anthropic_5h_window_exhausted", + } + } + return nil +} + +func isAnthropic5hRejected(headers http.Header) bool { + return strings.EqualFold(strings.TrimSpace(headers.Get("anthropic-ratelimit-unified-5h-status")), "rejected") +} + +func parseAnthropicWindowReset(headers http.Header, window string, now time.Time) (time.Time, bool) { + raw := strings.TrimSpace(headers.Get("anthropic-ratelimit-unified-" + window + "-reset")) + if raw == "" { + return time.Time{}, false + } + ts, err := strconv.ParseInt(raw, 10, 64) + if err != nil { + return time.Time{}, false + } + if ts > 1e11 { + ts = ts / 1000 + } + resetAt := time.Unix(ts, 0) + if !resetAt.After(now) { + return time.Time{}, false + } + + maxAge := 8 * 24 * time.Hour + if window == "5h" { + maxAge = 6 * time.Hour + } + if resetAt.After(now.Add(maxAge)) { + return time.Time{}, false + } + return resetAt, true +} + +func shouldPersistAnthropicWindowLimit(account *Account, limit *anthropicWindowLimit, now time.Time) bool { + if account == nil || limit == nil || !limit.resetAt.After(now) { + return false + } + if account.RateLimitResetAt == nil { + return true + } + if !account.RateLimitResetAt.After(now) { + return true + } + return limit.resetAt.After(*account.RateLimitResetAt) +} + +func (s *RateLimitService) persistAnthropicExhaustedWindowLimit(ctx context.Context, account *Account, headers http.Header) bool { + if s == nil || s.accountRepo == nil || account == nil { + return false + } + now := time.Now() + limit := selectAnthropicExhaustedWindow(headers, now) + if limit == nil { + return false + } + if !shouldPersistAnthropicWindowLimit(account, limit, now) { + slog.Info("anthropic_window_rate_limit_kept", + "account_id", account.ID, + "window", limit.window, + "reset_at", limit.resetAt, + "existing_reset_at", account.RateLimitResetAt) + return true + } + + s.notifyAccountSchedulingBlocked(account, limit.resetAt, limit.reason) + if err := s.accountRepo.SetRateLimited(ctx, account.ID, limit.resetAt); err != nil { + slog.Warn("anthropic_window_rate_limit_set_failed", + "account_id", account.ID, + "window", limit.window, + "reset_at", limit.resetAt, + "error", err) + return true + } + slog.Info("anthropic_window_rate_limited", + "account_id", account.ID, + "window", limit.window, + "reset_at", limit.resetAt, + "reset_in", time.Until(limit.resetAt).Truncate(time.Second)) + return true +} + // calculateAnthropic429ResetTime parses Anthropic's per-window rate-limit headers // to determine which window (5h or 7d) actually triggered the 429. // diff --git a/backend/internal/service/ratelimit_service_anthropic_window_limit_test.go b/backend/internal/service/ratelimit_service_anthropic_window_limit_test.go new file mode 100644 index 0000000000..6168b250a6 --- /dev/null +++ b/backend/internal/service/ratelimit_service_anthropic_window_limit_test.go @@ -0,0 +1,68 @@ +//go:build unit + +package service + +import ( + "context" + "net/http" + "strconv" + "testing" + "time" + + "github.com/stretchr/testify/require" +) + +type anthropicWindowLimitRepo struct { + mockAccountRepoForGemini + rateLimitCalls int + tempUnschedCalls int + lastRateLimitReset time.Time +} + +func (r *anthropicWindowLimitRepo) SetRateLimited(_ context.Context, _ int64, resetAt time.Time) error { + r.rateLimitCalls++ + r.lastRateLimitReset = resetAt + return nil +} + +func (r *anthropicWindowLimitRepo) SetTempUnschedulable(_ context.Context, _ int64, _ time.Time, _ string) error { + r.tempUnschedCalls++ + return nil +} + +func TestHandleUpstreamError_AnthropicWindowLimitPreemptsTempUnschedRule(t *testing.T) { + resetAt := time.Now().Add(3 * time.Hour).Truncate(time.Second) + headers := http.Header{} + headers.Set("anthropic-ratelimit-unified-5h-utilization", "1.02") + headers.Set("anthropic-ratelimit-unified-5h-reset", strconv.FormatInt(resetAt.Unix(), 10)) + + repo := &anthropicWindowLimitRepo{} + svc := NewRateLimitService(repo, nil, nil, nil, nil) + account := &Account{ + ID: 42, + Type: AccountTypeOAuth, + Platform: PlatformAnthropic, + Credentials: map[string]any{ + "temp_unschedulable_enabled": true, + "temp_unschedulable_rules": []any{ + map[string]any{ + "error_code": float64(http.StatusTooManyRequests), + "keywords": []any{"rate limit"}, + "duration_minutes": float64(10), + }, + }, + }, + } + + svc.HandleUpstreamError( + context.Background(), + account, + http.StatusTooManyRequests, + headers, + []byte(`{"type":"error","error":{"type":"rate_limit_error","message":"This request would exceed your account's rate limit. Please try again later."}}`), + ) + + require.Zero(t, repo.tempUnschedCalls, "official Anthropic window limits should not be shortened by local temp-unsched rules") + require.Equal(t, 1, repo.rateLimitCalls) + require.Equal(t, resetAt, repo.lastRateLimitReset) +}