fix(kimi): keep concurrency 403 recoverable

This commit is contained in:
wucm667
2026-08-25 19:40:39 +08:00
parent aa2c4e8d13
commit 3802268e2d
4 changed files with 146 additions and 1 deletions
@@ -11,6 +11,7 @@ package service
import (
"context"
"errors"
"net/http"
"net/http/httptest"
"testing"
@@ -142,3 +143,112 @@ func TestHandle403_CNProviderStructured403TempUnschedulableFirstHit(t *testing.T
require.Equal(t, 1, repo.tempCalls)
require.Contains(t, repo.lastTempReason, "(1/3)")
}
func TestIsCNProviderConcurrencyLimit403_ExactClassification(t *testing.T) {
kimi := &Account{Platform: PlatformKimi}
require.True(t, isCNProviderConcurrencyLimit403(kimi, kimiConcurrentRequestLimitMessage))
require.True(t, isCNProviderConcurrencyLimit403(kimi, " "+kimiConcurrentRequestLimitMessage+"\n"))
for name, tc := range map[string]struct {
account *Account
message string
}{
"permission denied": {kimi, "You do not have permission to access this resource."},
"generic concurrency wording": {kimi, "concurrent request limit reached"},
"near match missing punctuation": {kimi, "You've reached your concurrent request limit. Please wait for your ongoing requests to finish and try again"},
"other CN provider": {&Account{Platform: PlatformZhipu}, kimiConcurrentRequestLimitMessage},
"non CN provider": {&Account{Platform: PlatformOpenAI}, kimiConcurrentRequestLimitMessage},
"nil account": {nil, kimiConcurrentRequestLimitMessage},
} {
t.Run(name, func(t *testing.T) {
require.False(t, isCNProviderConcurrencyLimit403(tc.account, tc.message))
})
}
}
func TestHandle403_OtherCNProviderWithKimiConcurrencyMessageUsesNormalPolicy(t *testing.T) {
repo := &rateLimitAccountRepoStub{}
counter := &openAI403CounterCacheStub{counts: []int64{openAI403DisableThreshold}}
blocker := &runtimeBlockRecorder{}
service := NewRateLimitService(repo, nil, &config.Config{}, nil, nil)
service.SetOpenAI403CounterCache(counter)
service.SetAccountRuntimeBlocker(blocker)
account := &Account{ID: 405, Platform: PlatformZhipu, Type: AccountTypeAPIKey}
shouldDisable := service.HandleUpstreamError(
context.Background(), account, http.StatusForbidden, http.Header{},
[]byte(`{"error":{"message":"You've reached your concurrent request limit. Please wait for your ongoing requests to finish and try again."}}`),
)
require.True(t, shouldDisable)
require.Equal(t, 1, repo.setErrorCalls, "non-Kimi CN provider must retain the normal permanent-error policy")
require.Equal(t, 0, repo.tempCalls)
require.Empty(t, counter.counts, "normal CN 403 policy must consume the counter result")
require.Equal(t, []string{"auth_error"}, blocker.reasons, "the Kimi-specific runtime block must not apply")
}
func TestHandle403_CNProviderConcurrencyLimitAlwaysUsesTemporaryCooldown(t *testing.T) {
repo := &rateLimitAccountRepoStub{}
counter := &openAI403CounterCacheStub{counts: []int64{openAI403DisableThreshold}}
blocker := &runtimeBlockRecorder{}
service := NewRateLimitService(repo, nil, &config.Config{}, nil, nil)
service.SetOpenAI403CounterCache(counter)
service.SetAccountRuntimeBlocker(blocker)
account := &Account{ID: 403, Platform: PlatformKimi, Type: AccountTypeAPIKey}
shouldDisable := service.HandleUpstreamError(
context.Background(), account, http.StatusForbidden, http.Header{},
[]byte(`{"error":{"message":"You've reached your concurrent request limit. Please wait for your ongoing requests to finish and try again."}}`),
)
require.True(t, shouldDisable, "the request must still fail over to another account")
require.Equal(t, 0, repo.setErrorCalls)
require.Equal(t, 1, repo.tempCalls)
require.Contains(t, repo.lastTempReason, cnConcurrencyLimitReasonPrefix)
require.Equal(t, []int64{openAI403DisableThreshold}, counter.counts, "transient concurrency 403 must bypass the permanent-error counter")
require.Len(t, blocker.accounts, 1)
require.Equal(t, cnConcurrencyLimitReasonPrefix, blocker.reasons[0])
require.True(t, blocker.until[0].After(time.Now()))
}
func TestHandle403_KimiConcurrencyLimitRepositoryFailureKeepsRuntimeBlock(t *testing.T) {
repo := &rateLimitAccountRepoStub{tempErr: errors.New("repository unavailable")}
counter := &openAI403CounterCacheStub{counts: []int64{openAI403DisableThreshold}}
blocker := &runtimeBlockRecorder{}
service := NewRateLimitService(repo, nil, &config.Config{}, nil, nil)
service.SetOpenAI403CounterCache(counter)
service.SetAccountRuntimeBlocker(blocker)
account := &Account{ID: 406, Platform: PlatformKimi, Type: AccountTypeAPIKey}
shouldDisable := service.HandleUpstreamError(
context.Background(), account, http.StatusForbidden, http.Header{},
[]byte(`{"error":{"message":"You've reached your concurrent request limit. Please wait for your ongoing requests to finish and try again."}}`),
)
require.True(t, shouldDisable, "the current request must fail over even when persistence fails")
require.Equal(t, 1, repo.tempCalls, "the temporary cooldown should still be persisted when possible")
require.Equal(t, 0, repo.setErrorCalls, "persistence failure must not fall back to permanent account error")
require.Equal(t, []int64{openAI403DisableThreshold}, counter.counts, "persistence failure must not enter the permanent-error counter path")
require.Len(t, blocker.accounts, 1, "the in-memory runtime block must survive repository failure")
require.Same(t, account, blocker.accounts[0])
require.Equal(t, cnConcurrencyLimitReasonPrefix, blocker.reasons[0])
require.True(t, blocker.until[0].After(time.Now()))
}
func TestHandle403_CNProviderNearMatchRetainsNormalPermanentErrorPolicy(t *testing.T) {
repo := &rateLimitAccountRepoStub{}
counter := &openAI403CounterCacheStub{counts: []int64{openAI403DisableThreshold}}
service := NewRateLimitService(repo, nil, &config.Config{}, nil, nil)
service.SetOpenAI403CounterCache(counter)
account := &Account{ID: 404, Platform: PlatformKimi, Type: AccountTypeAPIKey}
shouldDisable := service.HandleUpstreamError(
context.Background(), account, http.StatusForbidden, http.Header{},
[]byte(`{"error":{"message":"You've reached your concurrent request limit. Please contact support."}}`),
)
require.True(t, shouldDisable)
require.Equal(t, 1, repo.setErrorCalls, "non-exact 403 must retain existing permission/auth protection")
require.Equal(t, 0, repo.tempCalls)
}
@@ -26,6 +26,33 @@ const cnBalanceExtraSuffixLow = "balance_low"
// 其他子系统(阈值/限流/401)写入的临时停调。
const cnBalanceLowReasonPrefix = "cn_balance_low"
const kimiConcurrentRequestLimitMessage = "You've reached your concurrent request limit. Please wait for your ongoing requests to finish and try again."
const cnConcurrencyLimitReasonPrefix = "cn_concurrency_limit"
func isCNProviderConcurrencyLimit403(account *Account, upstreamMsg string) bool {
return account != nil && account.Platform == PlatformKimi &&
strings.TrimSpace(upstreamMsg) == kimiConcurrentRequestLimitMessage
}
func (s *RateLimitService) handleCNProviderConcurrencyLimit403(
ctx context.Context,
account *Account,
) {
until := time.Now().Add(time.Duration(openAI403CooldownMinutesDefault) * time.Minute)
reason := cnConcurrencyLimitReasonPrefix + ": " + kimiConcurrentRequestLimitMessage
s.notifyAccountSchedulingBlocked(account, until, cnConcurrencyLimitReasonPrefix)
if err := s.accountRepo.SetTempUnschedulable(ctx, account.ID, until, reason); err != nil {
slog.Warn("cn_concurrency_limit_set_temp_unschedulable_failed", "account_id", account.ID, "error", err)
return
}
slog.Info("cn_provider_concurrency_limited",
"account_id", account.ID,
"platform", account.Platform,
"until", until.UTC(),
)
}
// cnBalanceLowReason 构造余额不足临时停调的 reason(带稳定前缀)。
func cnBalanceLowReason(upstreamMsg string) string {
if upstreamMsg = strings.TrimSpace(upstreamMsg); upstreamMsg != "" {
@@ -936,6 +936,13 @@ func (s *RateLimitService) handle403(ctx context.Context, account *Account, upst
if account.Platform == PlatformAntigravity {
return s.handleAntigravity403(ctx, account, upstreamMsg, responseBody)
}
// Kimi reports its transient per-account concurrency/business limit as a 403.
// Keep the normal 403 failover signal (true), but never feed this exact message
// into the escalating 403 counter that can permanently mark the account error.
if isCNProviderConcurrencyLimit403(account, upstreamMsg) {
s.handleCNProviderConcurrencyLimit403(ctx, account)
return true
}
// 国产供应商与 openai 同口径:HTML 403(CDN/代理拦截页)不构成账号失效证据,
// 且 403 在 failover 状态集里会被逐账号重放——直接 SetError 会让一个坏请求/
// 一层坏代理连环永久禁用整组账号。走 HTML 豁免 + N 次累计 + 临时冷却。
@@ -25,6 +25,7 @@ type rateLimitAccountRepoStub struct {
lastTempReason string
lastErrorID int64
lastTempID int64
tempErr error
}
func (r *rateLimitAccountRepoStub) SetError(ctx context.Context, id int64, errorMsg string) error {
@@ -38,7 +39,7 @@ func (r *rateLimitAccountRepoStub) SetTempUnschedulable(ctx context.Context, id
r.tempCalls++
r.lastTempID = id
r.lastTempReason = reason
return nil
return r.tempErr
}
func (r *rateLimitAccountRepoStub) UpdateCredentials(ctx context.Context, id int64, credentials map[string]any) error {