mirror of
https://github.com/Wei-Shaw/sub2api.git
synced 2026-09-24 16:05:44 +08:00
fix: preserve Anthropic window cooldowns
This commit is contained in:
@@ -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.
|
||||
//
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
Reference in New Issue
Block a user