Files
李建琦 6d655c9903
Release / update-version (push) Has been cancelled
Release / build-frontend (push) Has been cancelled
Release / release (push) Has been cancelled
Release / sync-version-file (push) Has been cancelled
CI / shell (push) Canceled after 0s
CI / test (push) Canceled after 0s
CI / frontend (push) Canceled after 0s
CI / golangci-lint (push) Canceled after 0s
Security Scan / backend-security (push) Canceled after 0s
Security Scan / frontend-security (push) Canceled after 0s
Sub2API v1.0 - AI API 网关(二开初始版本,基于上游 Wei-Shaw/sub2api)
2026-08-21 18:30:13 +08:00

309 lines
12 KiB
Go
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
package handler
import (
"context"
"net/http"
"time"
"github.com/gin-gonic/gin"
"github.com/Wei-Shaw/sub2api/internal/pkg/logger"
"github.com/Wei-Shaw/sub2api/internal/service"
"go.uber.org/zap"
)
// TempUnscheduler 用于 HandleFailoverError 中同账号重试耗尽后的临时封禁。
// GatewayService 隐式实现此接口。
type TempUnscheduler interface {
TempUnscheduleRetryableError(ctx context.Context, accountID int64, failoverErr *service.UpstreamFailoverError)
}
// FailoverAction 表示 failover 错误处理后的下一步动作
type FailoverAction int
const (
// FailoverContinue 继续循环(同账号重试或切换账号,调用方统一 continue)
FailoverContinue FailoverAction = iota
// FailoverExhausted 切换次数耗尽(调用方应返回错误响应)
FailoverExhausted
// FailoverCanceled context 已取消(调用方应直接 return)
FailoverCanceled
)
const (
// maxSameAccountRetries 同账号重试次数默认上限(针对 RetryableOnSameAccount 错误)。
// 生产调用方通常传入账号级配置 account.GetPoolModeRetryCount(),该常量仅作兜底/测试默认值。
maxSameAccountRetries = 3
// sameAccountRetryDelay 同账号重试间隔
sameAccountRetryDelay = 500 * time.Millisecond
// maxRequestScopedRetryDelay 限制请求级瞬时错误的指数退避上限,避免高重试配置
// 将单次请求拖入分钟级等待。
maxRequestScopedRetryDelay = 8 * time.Second
// singleAccountBackoffDelay 单账号分组 503 退避重试固定延时。
// Service 层在 SingleAccountRetry 模式下已做充分原地重试(最多 3 次、总等待 30s),
// Handler 层只需短暂间隔后重新进入 Service 层即可。
singleAccountBackoffDelay = 2 * time.Second
// maxProfitVetoAttempts 单次请求内允许的分组利润门终检否决次数上限。
// 利润否决不产生上游请求,因此不会推进 SwitchCount;没有独立上限的话,
// 「选号 → 终检否决 → 重选」在候选池与账号快照短暂不一致时可以空转很久。
// 取值与 maxAccountSwitches 默认值一致:混合定价的大分组仍有充分重选机会,
// 同时把整池越线时的无谓选号开销限制在常数级。
maxProfitVetoAttempts = 10
)
// profitVetoExhaustedMessage 是利润否决次数耗尽时返回给客户端的文案。
// 语义上等同于「无可用账号」:候选账号都不满足分组的利润约束。
const profitVetoExhaustedMessage = "No available accounts: all candidates rejected by group profit control"
func sameAccountRetryDelayFor(failoverErr *service.UpstreamFailoverError, retryCount int) time.Duration {
if failoverErr == nil || !failoverErr.RequestScopedTransient || retryCount <= 1 {
return sameAccountRetryDelay
}
delay := sameAccountRetryDelay
for i := 1; i < retryCount; i++ {
if delay >= maxRequestScopedRetryDelay/2 {
return maxRequestScopedRetryDelay
}
delay *= 2
}
return delay
}
// FailoverState 跨循环迭代共享的 failover 状态
type FailoverState struct {
SwitchCount int
MaxSwitches int
FailedAccountIDs map[int64]struct{}
SameAccountRetryCount map[int64]int
LastFailoverErr *service.UpstreamFailoverError
ForceCacheBilling bool
hasBoundSession bool
// profitVetoedAccountIDs 记录被分组利润门终检否决的账号,是 FailedAccountIDs
// 的子集。之所以单独维护:HandleSelectionExhausted 的 503 退避分支会清空
// FailedAccountIDs,而利润否决在同一请求内的判定不会改变(下游倍率 D 已在
// 请求开始冻结),被清空的账号会被立即重选并再次否决,形成没有任何上游请求、
// SwitchCount 也不前进的活锁。清空后必须把它们放回排除集。
profitVetoedAccountIDs map[int64]struct{}
// profitVetoCount 本次请求累计的利润否决次数,用于 maxProfitVetoAttempts 上限。
profitVetoCount int
}
// NewFailoverState 创建 failover 状态
func NewFailoverState(maxSwitches int, hasBoundSession bool) *FailoverState {
return &FailoverState{
MaxSwitches: maxSwitches,
FailedAccountIDs: make(map[int64]struct{}),
SameAccountRetryCount: make(map[int64]int),
hasBoundSession: hasBoundSession,
profitVetoedAccountIDs: make(map[int64]struct{}),
}
}
// RecordProfitVeto 记录一次分组利润门终检否决:把账号加入排除列表(同时登记到
// 利润否决集,使其不被 503 退避分支清掉)并递增否决计数。
//
// 返回 FailoverContinue 表示调用方可以继续重选下一个账号;返回 FailoverExhausted
// 表示本次请求的利润否决次数已达上限,调用方应按「无可用账号」终止,
// 不得继续 continue。
func (s *FailoverState) RecordProfitVeto(accountID int64) FailoverAction {
s.FailedAccountIDs[accountID] = struct{}{}
if s.profitVetoedAccountIDs == nil {
s.profitVetoedAccountIDs = make(map[int64]struct{})
}
s.profitVetoedAccountIDs[accountID] = struct{}{}
s.profitVetoCount++
if s.profitVetoCount >= maxProfitVetoAttempts {
return FailoverExhausted
}
return FailoverContinue
}
// ProfitVetoCount 返回本次请求累计的利润否决次数(供日志使用)。
func (s *FailoverState) ProfitVetoCount() int { return s.profitVetoCount }
// allExclusionsAreProfitVetoed 判断排除列表是否已全部由利润门否决贡献。
// 此时清空 FailedAccountIDs 会被原样恢复,退避重试不会带来任何新候选。
func (s *FailoverState) allExclusionsAreProfitVetoed() bool {
if len(s.profitVetoedAccountIDs) == 0 || len(s.FailedAccountIDs) == 0 {
return false
}
for id := range s.FailedAccountIDs {
if _, ok := s.profitVetoedAccountIDs[id]; !ok {
return false
}
}
return true
}
// HandleFailoverError 处理 UpstreamFailoverError,返回下一步动作。
// 包含:缓存计费判断、同账号重试、临时封禁、切换计数、Antigravity 延时。
func (s *FailoverState) HandleFailoverError(
ctx context.Context,
gatewayService TempUnscheduler,
accountID int64,
platform string,
retryLimit int,
failoverErr *service.UpstreamFailoverError,
) FailoverAction {
// 客户端已断开:failover 只会用已取消的 context 重新选号并必然失败,
// 不应再被当成账号耗尽处理(误报 502)。
if ctx != nil && ctx.Err() != nil {
return FailoverCanceled
}
s.LastFailoverErr = failoverErr
if failoverErr == nil || !failoverErr.ShouldRetryNextAccount() {
return FailoverExhausted
}
// 同账号重试不算切换账号,粘性会话仅在实际切换时强制缓存计费。
sameAccountRetry := failoverErr.RetryableOnSameAccount && s.SameAccountRetryCount[accountID] < retryLimit
if needForceCacheBilling(s.hasBoundSession, failoverErr, sameAccountRetry) {
s.ForceCacheBilling = true
}
// 同账号重试:对 RetryableOnSameAccount 的临时性错误,先在同一账号上重试。
// 重试次数上限 retryLimit 由调用方传入(账号级 pool_mode_retry_count 配置)。
if failoverErr.RetryableOnSameAccount && s.SameAccountRetryCount[accountID] < retryLimit {
s.SameAccountRetryCount[accountID]++
retryDelay := sameAccountRetryDelayFor(failoverErr, s.SameAccountRetryCount[accountID])
logger.FromContext(ctx).Warn("gateway.failover_same_account_retry",
zap.Int64("account_id", accountID),
zap.Int("upstream_status", failoverErr.StatusCode),
zap.Int("same_account_retry_count", s.SameAccountRetryCount[accountID]),
zap.Int("same_account_retry_max", retryLimit),
zap.Duration("retry_delay", retryDelay),
)
if !sleepWithContext(ctx, retryDelay) {
return FailoverCanceled
}
return FailoverContinue
}
// 同账号重试用尽,执行临时封禁
if failoverErr.RetryableOnSameAccount {
gatewayService.TempUnscheduleRetryableError(ctx, accountID, failoverErr)
}
// 加入失败列表
s.FailedAccountIDs[accountID] = struct{}{}
// 检查是否耗尽
if s.SwitchCount >= s.MaxSwitches {
return FailoverExhausted
}
// 递增切换计数
s.SwitchCount++
logger.FromContext(ctx).Warn("gateway.failover_switch_account",
zap.Int64("account_id", accountID),
zap.Int("upstream_status", failoverErr.StatusCode),
zap.Int("switch_count", s.SwitchCount),
zap.Int("max_switches", s.MaxSwitches),
)
// Antigravity 平台换号线性递增延时
if platform == service.PlatformAntigravity {
delay := time.Duration(s.SwitchCount-1) * time.Second
if !sleepWithContext(ctx, delay) {
return FailoverCanceled
}
}
return FailoverContinue
}
// HandleSelectionExhausted 处理选号失败(所有候选账号都在排除列表中)时的退避重试决策。
// 针对 Antigravity 单账号分组的 503 (MODEL_CAPACITY_EXHAUSTED) 场景:
// 清除排除列表、等待退避后重新选号。
//
// 返回 FailoverContinue 时,调用方应设置 SingleAccountRetry context 并 continue。
// 返回 FailoverExhausted 时,调用方应返回错误响应。
// 返回 FailoverCanceled 时,调用方应直接 return。
func (s *FailoverState) HandleSelectionExhausted(ctx context.Context) FailoverAction {
// 客户端已断开时选号失败是 context canceled 的必然结果,
// 不代表账号耗尽,直接按取消终止。
if ctx.Err() != nil {
return FailoverCanceled
}
if s.LastFailoverErr != nil &&
s.LastFailoverErr.StatusCode == http.StatusServiceUnavailable &&
s.SwitchCount <= s.MaxSwitches {
// 排除列表全由利润门否决贡献时,清空后会被原样恢复:退避重试拿不到
// 任何新候选,而利润否决不推进 SwitchCount,退避条件将永远成立。
// 这里直接判定耗尽,避免每 2s 空转一轮的活锁。
if s.allExclusionsAreProfitVetoed() {
logger.FromContext(ctx).Warn("gateway.failover_selection_exhausted_by_profit_veto",
zap.Int("profit_veto_count", s.profitVetoCount),
zap.Int("excluded_accounts", len(s.FailedAccountIDs)),
)
return FailoverExhausted
}
logger.FromContext(ctx).Warn("gateway.failover_single_account_backoff",
zap.Duration("backoff_delay", singleAccountBackoffDelay),
zap.Int("switch_count", s.SwitchCount),
zap.Int("max_switches", s.MaxSwitches),
)
if !sleepWithContext(ctx, singleAccountBackoffDelay) {
return FailoverCanceled
}
logger.FromContext(ctx).Warn("gateway.failover_single_account_retry",
zap.Int("switch_count", s.SwitchCount),
zap.Int("max_switches", s.MaxSwitches),
)
s.FailedAccountIDs = make(map[int64]struct{})
// 利润门否决的账号不参与退避重试的解除:判定依据(冻结的下游倍率)在
// 同一请求内不变,放它们回池只会被再次否决。
for id := range s.profitVetoedAccountIDs {
s.FailedAccountIDs[id] = struct{}{}
}
return FailoverContinue
}
return FailoverExhausted
}
// needForceCacheBilling 判断 failover 时是否需要强制缓存计费。
// 粘性会话实际切换账号、或上游明确标记时,将 input_tokens 转为 cache_read 计费。
func needForceCacheBilling(hasBoundSession bool, failoverErr *service.UpstreamFailoverError, sameAccountRetry bool) bool {
return (hasBoundSession && !sameAccountRetry) || (failoverErr != nil && failoverErr.ForceCacheBilling)
}
// failoverClientGone 判断下游客户端是否已断开(请求 context 已取消)。
// 客户端断开后 failover 必须静默终止:用已取消的 context 重新选号只会得到
// context.Canceled,并被误报成账号耗尽(通用 502);上游 detach 的在途请求
// 照常完成计费,但不再为无人接收的响应启动新的上游尝试。
// 响应尚未提交时把状态码标记为 499client closed request),供访问日志归类。
func failoverClientGone(c *gin.Context) bool {
if c == nil || c.Request == nil || c.Request.Context().Err() == nil {
return false
}
// 先停 compact 心跳(接管 ResponseWriter,建立 happens-before),与
// handleStreamingAwareError/errorResponse 等终结路径对齐,避免心跳
// goroutine 与下面的状态标记并发触碰同一 writer。心跳已提交 200 时
// 状态码已固化,不再标 499。
if service.StopOpenAICompactSSEKeepaliveCommitted(c) {
return true
}
if !c.Writer.Written() {
c.Status(statusClientClosedRequest)
}
return true
}
// sleepWithContext 等待指定时长,返回 false 表示 context 已取消。
func sleepWithContext(ctx context.Context, d time.Duration) bool {
if d <= 0 {
return true
}
select {
case <-ctx.Done():
return false
case <-time.After(d):
return true
}
}