Revert "feat(429): add configurable cooldown and retry strategies"

This reverts commit 6c3edc0956.
This commit is contained in:
IanShaw
2026-08-20 05:44:21 -07:00
parent 6c3edc0956
commit e62ec2c42f
12 changed files with 555 additions and 693 deletions
@@ -112,23 +112,15 @@ func (h *SettingHandler) GetRateLimit429CooldownSettings(c *gin.Context) {
}
response.Success(c, dto.RateLimit429CooldownSettings{
Enabled: settings.Enabled,
CooldownSeconds: settings.CooldownSeconds,
Strategy: settings.Strategy,
RetryIntervalMs: settings.RetryIntervalMs,
RetryMaxDurationSeconds: settings.RetryMaxDurationSeconds,
MaxAccountSwitches: settings.MaxAccountSwitches,
Enabled: settings.Enabled,
CooldownSeconds: settings.CooldownSeconds,
})
}
// UpdateRateLimit429CooldownSettingsRequest 更新429默认回避配置请求
type UpdateRateLimit429CooldownSettingsRequest struct {
Strategy string `json:"strategy"`
RetryIntervalMs int `json:"retry_interval_ms"`
RetryMaxDurationSeconds int `json:"retry_max_duration_seconds"`
MaxAccountSwitches int `json:"max_account_switches"`
Enabled bool `json:"enabled"`
CooldownSeconds int `json:"cooldown_seconds"`
Enabled bool `json:"enabled"`
CooldownSeconds int `json:"cooldown_seconds"`
}
// UpdateRateLimit429CooldownSettings 更新429默认回避配置
@@ -141,12 +133,8 @@ func (h *SettingHandler) UpdateRateLimit429CooldownSettings(c *gin.Context) {
}
settings := &service.RateLimit429CooldownSettings{
Strategy: req.Strategy,
RetryIntervalMs: req.RetryIntervalMs,
RetryMaxDurationSeconds: req.RetryMaxDurationSeconds,
MaxAccountSwitches: req.MaxAccountSwitches,
Enabled: req.Enabled,
CooldownSeconds: req.CooldownSeconds,
Enabled: req.Enabled,
CooldownSeconds: req.CooldownSeconds,
}
if err := h.settingService.SetRateLimit429CooldownSettings(c.Request.Context(), settings); err != nil {
@@ -161,12 +149,8 @@ func (h *SettingHandler) UpdateRateLimit429CooldownSettings(c *gin.Context) {
}
response.Success(c, dto.RateLimit429CooldownSettings{
Enabled: updatedSettings.Enabled,
CooldownSeconds: updatedSettings.CooldownSeconds,
Strategy: updatedSettings.Strategy,
RetryIntervalMs: updatedSettings.RetryIntervalMs,
RetryMaxDurationSeconds: updatedSettings.RetryMaxDurationSeconds,
MaxAccountSwitches: updatedSettings.MaxAccountSwitches,
Enabled: updatedSettings.Enabled,
CooldownSeconds: updatedSettings.CooldownSeconds,
})
}
+71 -101
View File
@@ -148,23 +148,21 @@ type SystemSettings struct {
GoogleOAuthRedirectURL string `json:"google_oauth_redirect_url"`
GoogleOAuthFrontendRedirectURL string `json:"google_oauth_frontend_redirect_url"`
SiteName string `json:"site_name"`
SiteLogo string `json:"site_logo"`
SiteSubtitle string `json:"site_subtitle"`
APIBaseURL string `json:"api_base_url"`
ContactInfo string `json:"contact_info"`
SupportQRCodes []service.SupportQRCodeEntry `json:"support_qr_codes"`
DownloadToolsURL string `json:"download_tools_url"`
DocURL string `json:"doc_url"`
HomeContent string `json:"home_content"`
CompactHomeEnabled bool `json:"compact_home_enabled"`
HideCcsImportButton bool `json:"hide_ccs_import_button"`
PurchaseSubscriptionEnabled bool `json:"purchase_subscription_enabled"`
PurchaseSubscriptionURL string `json:"purchase_subscription_url"`
TableDefaultPageSize int `json:"table_default_page_size"`
TablePageSizeOptions []int `json:"table_page_size_options"`
CustomMenuItems []CustomMenuItem `json:"custom_menu_items"`
CustomEndpoints []CustomEndpoint `json:"custom_endpoints"`
SiteName string `json:"site_name"`
SiteLogo string `json:"site_logo"`
SiteSubtitle string `json:"site_subtitle"`
APIBaseURL string `json:"api_base_url"`
ContactInfo string `json:"contact_info"`
DocURL string `json:"doc_url"`
HomeContent string `json:"home_content"`
CompactHomeEnabled bool `json:"compact_home_enabled"`
HideCcsImportButton bool `json:"hide_ccs_import_button"`
PurchaseSubscriptionEnabled bool `json:"purchase_subscription_enabled"`
PurchaseSubscriptionURL string `json:"purchase_subscription_url"`
TableDefaultPageSize int `json:"table_default_page_size"`
TablePageSizeOptions []int `json:"table_page_size_options"`
CustomMenuItems []CustomMenuItem `json:"custom_menu_items"`
CustomEndpoints []CustomEndpoint `json:"custom_endpoints"`
DefaultConcurrency int `json:"default_concurrency"`
DefaultBalance float64 `json:"default_balance"`
@@ -172,9 +170,6 @@ type SystemSettings struct {
AffiliateRebateFreezeHours int `json:"affiliate_rebate_freeze_hours"`
AffiliateRebateDurationDays int `json:"affiliate_rebate_duration_days"`
AffiliateRebatePerInviteeCap float64 `json:"affiliate_rebate_per_invitee_cap"`
AffiliateRebateCap float64 `json:"affiliate_rebate_cap"`
AffiliateRebateInviteeLimit int `json:"affiliate_rebate_invitee_limit"`
AffiliateSignupBonus float64 `json:"affiliate_signup_bonus"`
AdminRechargeRebateEnabled bool `json:"affiliate_admin_recharge_enabled"`
DefaultUserRPMLimit int `json:"default_user_rpm_limit"`
DefaultSubscriptions []DefaultSubscriptionSetting `json:"default_subscriptions"`
@@ -310,23 +305,13 @@ type SystemSettings struct {
ChannelMonitorMode string `json:"channel_monitor_mode"`
ChannelMonitorDefaultIntervalSeconds int `json:"channel_monitor_default_interval_seconds"`
ChannelMonitorHideThroughput bool `json:"channel_monitor_hide_throughput"`
ChannelMonitorShowQuota bool `json:"channel_monitor_show_quota"`
// Grok model mapping policy (admin settings; empty account mapping falls back to these).
GrokDefaultTextModel string `json:"grok_default_text_model"`
GrokCrossClientModelMapEnabled bool `json:"grok_cross_client_model_map_enabled"`
GrokDefaultBaseURLMode string `json:"grok_default_base_url_mode"`
// Kiro runtime defaults (admin-only; not exposed on PublicSettings).
KiroDefaultVersion string `json:"kiro_version"`
KiroDefaultCommit string `json:"kiro_commit"`
KiroDefaultSystemVersion string `json:"system_version"`
KiroDefaultNodeVersion string `json:"node_version"`
KiroCacheHitRateScale int `json:"cache_hit_rate_scale"`
KiroCacheMinBlockTokens int `json:"cache_min_block_tokens"`
KiroCacheIndependentTTLSeconds int `json:"cache_independent_ttl_seconds"`
KiroCachePrefixTTLSeconds int `json:"cache_prefix_ttl_seconds"`
KiroCodeExecutionSandboxCommand string `json:"kiro_code_execution_sandbox_command"`
// Available Channels feature switch (user-facing aggregate view)
AvailableChannelsEnabled bool `json:"available_channels_enabled"`
@@ -345,14 +330,6 @@ type SystemSettings struct {
// Affiliate (邀请返利) feature switch
AffiliateEnabled bool `json:"affiliate_enabled"`
// Ticket feature switch (default enabled)
TicketEnabled bool `json:"ticket_enabled"`
IPMultiAccountBanEnabled bool `json:"ip_multi_account_ban_enabled"`
IPMultiAccountBanWindowMinutes int `json:"ip_multi_account_ban_window_minutes"`
IPMultiAccountBanThreshold int `json:"ip_multi_account_ban_threshold"`
IPMultiAccountBanLearningUntil string `json:"ip_multi_account_ban_learning_until"`
// OpenAI fast/flex policy
OpenAIFastPolicySettings *OpenAIFastPolicySettings `json:"openai_fast_policy_settings,omitempty"`
@@ -372,61 +349,58 @@ type DefaultSubscriptionSetting struct {
}
type PublicSettings struct {
RegistrationEnabled bool `json:"registration_enabled"`
EmailVerifyEnabled bool `json:"email_verify_enabled"`
ForceEmailOnThirdPartySignup bool `json:"force_email_on_third_party_signup"`
RegistrationEmailSuffixWhitelist []string `json:"registration_email_suffix_whitelist"`
RegistrationEmailDomainQuotaEnabled bool `json:"registration_email_domain_quota_enabled"`
PromoCodeEnabled bool `json:"promo_code_enabled"`
PasswordResetEnabled bool `json:"password_reset_enabled"`
InvitationCodeEnabled bool `json:"invitation_code_enabled"`
TotpEnabled bool `json:"totp_enabled"` // TOTP 双因素认证
PasskeyEnabled bool `json:"passkey_enabled"`
LoginAgreementEnabled bool `json:"login_agreement_enabled"`
LoginAgreementMode string `json:"login_agreement_mode"`
LoginAgreementUpdatedAt string `json:"login_agreement_updated_at"`
LoginAgreementRevision string `json:"login_agreement_revision"`
LoginAgreementDocuments []LoginAgreementDocument `json:"login_agreement_documents"`
TurnstileEnabled bool `json:"turnstile_enabled"`
TurnstileSiteKey string `json:"turnstile_site_key"`
TencentCaptchaEnabled bool `json:"tencent_captcha_enabled"`
TencentCaptchaAppID string `json:"tencent_captcha_app_id"`
TencentCaptchaRegion string `json:"tencent_captcha_region"`
AliyunCaptchaEnabled bool `json:"aliyun_captcha_enabled"`
AliyunCaptchaSceneID string `json:"aliyun_captcha_scene_id"`
AliyunCaptchaPrefix string `json:"aliyun_captcha_prefix"`
AliyunCaptchaRegion string `json:"aliyun_captcha_region"`
SiteName string `json:"site_name"`
SiteLogo string `json:"site_logo"`
SiteSubtitle string `json:"site_subtitle"`
APIBaseURL string `json:"api_base_url"`
ContactInfo string `json:"contact_info"`
SupportQRCodes []service.SupportQRCodeEntry `json:"support_qr_codes"`
DownloadToolsURL string `json:"download_tools_url"`
DocURL string `json:"doc_url"`
HomeContent string `json:"home_content"`
CompactHomeEnabled bool `json:"compact_home_enabled"`
HideCcsImportButton bool `json:"hide_ccs_import_button"`
PurchaseSubscriptionEnabled bool `json:"purchase_subscription_enabled"`
PurchaseSubscriptionURL string `json:"purchase_subscription_url"`
TableDefaultPageSize int `json:"table_default_page_size"`
TablePageSizeOptions []int `json:"table_page_size_options"`
CustomMenuItems []CustomMenuItem `json:"custom_menu_items"`
CustomEndpoints []CustomEndpoint `json:"custom_endpoints"`
DingTalkOAuthEnabled bool `json:"dingtalk_oauth_enabled"`
LinuxDoOAuthEnabled bool `json:"linuxdo_oauth_enabled"`
WeChatOAuthEnabled bool `json:"wechat_oauth_enabled"`
WeChatOAuthOpenEnabled bool `json:"wechat_oauth_open_enabled"`
WeChatOAuthMPEnabled bool `json:"wechat_oauth_mp_enabled"`
WeChatOAuthMobileEnabled bool `json:"wechat_oauth_mobile_enabled"`
OIDCOAuthEnabled bool `json:"oidc_oauth_enabled"`
OIDCOAuthProviderName string `json:"oidc_oauth_provider_name"`
GitHubOAuthEnabled bool `json:"github_oauth_enabled"`
GoogleOAuthEnabled bool `json:"google_oauth_enabled"`
SoraClientEnabled bool `json:"sora_client_enabled"`
BackendModeEnabled bool `json:"backend_mode_enabled"`
PaymentEnabled bool `json:"payment_enabled"`
Version string `json:"version"`
RegistrationEnabled bool `json:"registration_enabled"`
EmailVerifyEnabled bool `json:"email_verify_enabled"`
ForceEmailOnThirdPartySignup bool `json:"force_email_on_third_party_signup"`
RegistrationEmailSuffixWhitelist []string `json:"registration_email_suffix_whitelist"`
RegistrationEmailDomainQuotaEnabled bool `json:"registration_email_domain_quota_enabled"`
PromoCodeEnabled bool `json:"promo_code_enabled"`
PasswordResetEnabled bool `json:"password_reset_enabled"`
InvitationCodeEnabled bool `json:"invitation_code_enabled"`
TotpEnabled bool `json:"totp_enabled"` // TOTP 双因素认证
PasskeyEnabled bool `json:"passkey_enabled"`
LoginAgreementEnabled bool `json:"login_agreement_enabled"`
LoginAgreementMode string `json:"login_agreement_mode"`
LoginAgreementUpdatedAt string `json:"login_agreement_updated_at"`
LoginAgreementRevision string `json:"login_agreement_revision"`
LoginAgreementDocuments []LoginAgreementDocument `json:"login_agreement_documents"`
TurnstileEnabled bool `json:"turnstile_enabled"`
TurnstileSiteKey string `json:"turnstile_site_key"`
TencentCaptchaEnabled bool `json:"tencent_captcha_enabled"`
TencentCaptchaAppID string `json:"tencent_captcha_app_id"`
TencentCaptchaRegion string `json:"tencent_captcha_region"`
AliyunCaptchaEnabled bool `json:"aliyun_captcha_enabled"`
AliyunCaptchaSceneID string `json:"aliyun_captcha_scene_id"`
AliyunCaptchaPrefix string `json:"aliyun_captcha_prefix"`
AliyunCaptchaRegion string `json:"aliyun_captcha_region"`
SiteName string `json:"site_name"`
SiteLogo string `json:"site_logo"`
SiteSubtitle string `json:"site_subtitle"`
APIBaseURL string `json:"api_base_url"`
ContactInfo string `json:"contact_info"`
DocURL string `json:"doc_url"`
HomeContent string `json:"home_content"`
CompactHomeEnabled bool `json:"compact_home_enabled"`
HideCcsImportButton bool `json:"hide_ccs_import_button"`
PurchaseSubscriptionEnabled bool `json:"purchase_subscription_enabled"`
PurchaseSubscriptionURL string `json:"purchase_subscription_url"`
TableDefaultPageSize int `json:"table_default_page_size"`
TablePageSizeOptions []int `json:"table_page_size_options"`
CustomMenuItems []CustomMenuItem `json:"custom_menu_items"`
CustomEndpoints []CustomEndpoint `json:"custom_endpoints"`
DingTalkOAuthEnabled bool `json:"dingtalk_oauth_enabled"`
LinuxDoOAuthEnabled bool `json:"linuxdo_oauth_enabled"`
WeChatOAuthEnabled bool `json:"wechat_oauth_enabled"`
WeChatOAuthOpenEnabled bool `json:"wechat_oauth_open_enabled"`
WeChatOAuthMPEnabled bool `json:"wechat_oauth_mp_enabled"`
WeChatOAuthMobileEnabled bool `json:"wechat_oauth_mobile_enabled"`
OIDCOAuthEnabled bool `json:"oidc_oauth_enabled"`
OIDCOAuthProviderName string `json:"oidc_oauth_provider_name"`
GitHubOAuthEnabled bool `json:"github_oauth_enabled"`
GoogleOAuthEnabled bool `json:"google_oauth_enabled"`
BackendModeEnabled bool `json:"backend_mode_enabled"`
PaymentEnabled bool `json:"payment_enabled"`
Version string `json:"version"`
// 服务器全局时区(IANA 名称与当前 UTC 偏移,如 "Asia/Shanghai" / "+08:00")。
// 高峰时段等按服务器本地时间判定的窗口,前端展示时据此标注,避免用户按浏览器本地时间误读。
ServerTimezone string `json:"server_timezone"`
@@ -440,6 +414,7 @@ type PublicSettings struct {
ChannelMonitorMode string `json:"channel_monitor_mode"`
ChannelMonitorDefaultIntervalSeconds int `json:"channel_monitor_default_interval_seconds"`
ChannelMonitorHideThroughput bool `json:"channel_monitor_hide_throughput"`
ChannelMonitorShowQuota bool `json:"channel_monitor_show_quota"`
AvailableChannelsEnabled bool `json:"available_channels_enabled"`
@@ -447,7 +422,6 @@ type PublicSettings struct {
ModelPlazaRequireAuth bool `json:"model_plaza_require_auth"`
AffiliateEnabled bool `json:"affiliate_enabled"`
TicketEnabled bool `json:"ticket_enabled"`
RiskControlEnabled bool `json:"risk_control_enabled"`
@@ -468,12 +442,8 @@ type OverloadCooldownSettings struct {
// RateLimit429CooldownSettings 429默认回避配置 DTO
type RateLimit429CooldownSettings struct {
Strategy string `json:"strategy"`
RetryIntervalMs int `json:"retry_interval_ms"`
RetryMaxDurationSeconds int `json:"retry_max_duration_seconds"`
MaxAccountSwitches int `json:"max_account_switches"`
Enabled bool `json:"enabled"`
CooldownSeconds int `json:"cooldown_seconds"`
Enabled bool `json:"enabled"`
CooldownSeconds int `json:"cooldown_seconds"`
}
// PanelRateLimitSettings 面板 API 限流配置 DTO
+16 -75
View File
@@ -56,16 +56,13 @@ const (
const profitVetoExhaustedMessage = "No available accounts: all candidates rejected by group profit control"
func sameAccountRetryDelayFor(failoverErr *service.UpstreamFailoverError, retryCount int) time.Duration {
// An OAuth 429 with a request-scoped deadline intentionally retries
// immediately; zero is meaningful here and must not fall back to 500ms.
if failoverErr != nil && failoverErr.StatusCode == http.StatusTooManyRequests &&
!failoverErr.SameAccountRetryDeadline.IsZero() && failoverErr.SameAccountRetryDelay <= 0 {
return 0
if failoverErr == nil {
return sameAccountRetryDelay
}
if failoverErr != nil && failoverErr.SameAccountRetryDelay > 0 {
if failoverErr.SameAccountRetryDelay > 0 {
return failoverErr.SameAccountRetryDelay
}
if failoverErr == nil || !failoverErr.RequestScopedTransient || retryCount <= 1 {
if !failoverErr.RequestScopedTransient || retryCount <= 1 {
return sameAccountRetryDelay
}
@@ -79,43 +76,10 @@ func sameAccountRetryDelayFor(failoverErr *service.UpstreamFailoverError, retryC
return delay
}
func sameAccountRetryAllowed(failoverErr *service.UpstreamFailoverError, retryCount, retryLimit int) bool {
if failoverErr == nil || !failoverErr.RetryableOnSameAccount {
return false
}
if retryLimit > 0 && retryCount >= retryLimit {
return false
}
if !failoverErr.SameAccountRetryDeadline.IsZero() {
return time.Now().Before(failoverErr.SameAccountRetryDeadline)
}
return retryCount < retryLimit
}
func pinSameAccountRetryContext(
ctx context.Context,
fs *FailoverState,
accountID int64,
groupID *int64,
failoverErr *service.UpstreamFailoverError,
prevRetryCount int,
prevSwitchCount int,
bridgeOldKeys bool,
) context.Context {
if ctx == nil || fs == nil || failoverErr == nil || !failoverErr.RetryableOnSameAccount {
return ctx
}
if fs.SwitchCount != prevSwitchCount {
return ctx
}
if fs.SameAccountRetryCount[accountID] != prevRetryCount+1 {
return ctx
}
prefetchedGroupID := int64(0)
if groupID != nil {
prefetchedGroupID = *groupID
}
return service.WithPrefetchedStickySession(ctx, accountID, prefetchedGroupID, bridgeOldKeys)
// sameAccountRetryDeadlineAllows prevents a retry from starting after the
// service-provided same-account retry window has elapsed.
func sameAccountRetryDeadlineAllows(failoverErr *service.UpstreamFailoverError) bool {
return failoverErr == nil || failoverErr.SameAccountRetryDeadline.IsZero() || time.Now().Before(failoverErr.SameAccountRetryDeadline)
}
// FailoverState 跨循环迭代共享的 failover 状态
@@ -168,24 +132,6 @@ func (s *FailoverState) RecordProfitVeto(accountID int64) FailoverAction {
return FailoverContinue
}
// RecordConcurrencyTimeout excludes a busy account after the slot ladder
// (deadline two-shot or wait-queue full) and continues onto another account
// while the original sticky binding is preserved by the caller.
func (s *FailoverState) RecordConcurrencyTimeout(accountID int64) FailoverAction {
if s == nil {
return FailoverExhausted
}
if s.FailedAccountIDs == nil {
s.FailedAccountIDs = make(map[int64]struct{})
}
s.FailedAccountIDs[accountID] = struct{}{}
if s.SwitchCount >= s.MaxSwitches {
return FailoverExhausted
}
s.SwitchCount++
return FailoverContinue
}
// ProfitVetoCount 返回本次请求累计的利润否决次数(供日志使用)。
func (s *FailoverState) ProfitVetoCount() int { return s.profitVetoCount }
@@ -223,27 +169,27 @@ func (s *FailoverState) HandleFailoverError(
return FailoverExhausted
}
retryMax := retryLimit
if failoverErr.SameAccountRetryMax > 0 {
retryMax = failoverErr.SameAccountRetryMax
}
// 同账号重试不算切换账号,粘性会话仅在实际切换时强制缓存计费。
sameAccountRetry := sameAccountRetryAllowed(failoverErr, s.SameAccountRetryCount[accountID], retryMax)
retryCount := s.SameAccountRetryCount[accountID]
sameAccountRetryAllowed := failoverErr.RetryableOnSameAccount && retryLimit > 0 && retryCount < retryLimit
if sameAccountRetryAllowed && !failoverErr.SameAccountRetryDeadline.IsZero() {
sameAccountRetryAllowed = time.Now().Before(failoverErr.SameAccountRetryDeadline)
}
sameAccountRetry := sameAccountRetryAllowed
if needForceCacheBilling(s.hasBoundSession, failoverErr, sameAccountRetry) {
s.ForceCacheBilling = true
}
// 同账号重试:对 RetryableOnSameAccount 的临时性错误,先在同一账号上重试。
// 重试次数上限 retryLimit 由调用方传入(账号级 pool_mode_retry_count 配置)。
if sameAccountRetryAllowed(failoverErr, s.SameAccountRetryCount[accountID], retryMax) {
if sameAccountRetryAllowed {
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", retryMax),
zap.Int("same_account_retry_max", retryLimit),
zap.Duration("retry_delay", retryDelay),
)
if !sleepWithContext(ctx, retryDelay) {
@@ -259,11 +205,6 @@ func (s *FailoverState) HandleFailoverError(
// 加入失败列表
s.FailedAccountIDs[accountID] = struct{}{}
for _, excludedAccountID := range failoverErr.ExcludedAccountIDs {
if excludedAccountID > 0 {
s.FailedAccountIDs[excludedAccountID] = struct{}{}
}
}
// 检查是否耗尽
if s.SwitchCount >= s.MaxSwitches {
@@ -11,34 +11,13 @@ import (
const (
openAIAccountStateUpdateTimeout = 5 * time.Second
openAIOAuth429RetryWindow = 2 * time.Minute
openAIOAuth429RetryDelay = 0
openAIOAuth429FallbackCooldown = 5 * time.Second
openAIStopSchedulingBridgeCooldown = 2 * time.Minute
openAIOAuth429MaxAccountAttempts = 3
openAIOAuth429StormWindow = 10 * time.Second
openAIOAuth429StormThreshold = 20
openAIOAuth429StormMaxAccountSwitches = 1
)
func (s *OpenAIGatewayService) rateLimit429StrategySettings() RateLimit429CooldownSettings {
defaults := DefaultRateLimit429CooldownSettings()
if s == nil || s.settingService == nil {
return *defaults
}
s.openai429StrategyMu.Lock()
defer s.openai429StrategyMu.Unlock()
if time.Since(s.openai429StrategyCachedAt) < 5*time.Second {
return s.openai429StrategyCached
}
settings := *defaults
if loaded, err := s.settingService.GetRateLimit429CooldownSettings(context.Background()); err == nil && loaded != nil {
settings = *loaded
}
s.openai429StrategyCached = settings
s.openai429StrategyCachedAt = time.Now()
return settings
}
// OpenAIOAuth429FailoverState tracks the request-local follow-up budget after
// the first Grok OAuth 429. Once that 429 occurs, exactly one different account
// may be attempted; any failure from that follow-up account ends failover.
@@ -76,6 +55,11 @@ func (s *OpenAIGatewayService) handleOpenAIAccountUpstreamError(ctx context.Cont
if s != nil {
scheduleOllamaCloudUsageActivity(s.deferredService, account)
}
// Capacity shedding describes this request, not account health. Keep the
// account schedulable while the request-local retry budget handles recovery.
if account != nil && account.Platform == PlatformOpenAI && isOpenAIRequestScopedCapacityShed("", responseBody) {
return false
}
stateCtx, cancel := openAIAccountStateContext(ctx)
defer cancel()
@@ -93,6 +77,10 @@ func (s *OpenAIGatewayService) handleOpenAIAccountUpstreamError(ctx context.Cont
if s == nil || account == nil {
return false
}
// Team 联动熔断必须先于 model-not-found 与账户级临时不可调度规则的早退。
if s.rateLimitService != nil {
s.rateLimitService.maybeHandleOpenAITeamLinkedError(stateCtx, account, statusCode, responseBody)
}
stateCtx = withTempUnschedulableModel(stateCtx, canonicalModel)
if s.rateLimitService != nil && len(canonicalModel) > 0 && s.rateLimitService.HandleUpstreamModelNotFound(stateCtx, account, canonicalModel[0], statusCode, responseBody) {
return true
@@ -161,164 +149,20 @@ func (s *OpenAIGatewayService) markOpenAIOAuth429RateLimited(ctx context.Context
return
}
s.recordOpenAIOAuth429()
if s.ShouldRetryOpenAIOAuth429(account, headers, responseBody) {
return
}
cooldownUntil := time.Time{}
hasCooldown := false
cooldownUntil := time.Now().Add(openAIOAuth429FallbackCooldown)
if s.rateLimitService != nil {
if resetAt := s.rateLimitService.calculateOpenAI429ResetTime(headers); resetAt != nil && resetAt.After(time.Now()) {
cooldownUntil = *resetAt
hasCooldown = true
} else if resetUnix := parseOpenAIRateLimitResetTime(responseBody); resetUnix != nil {
if resetAt := time.Unix(*resetUnix, 0); resetAt.After(time.Now()) {
cooldownUntil = resetAt
hasCooldown = true
}
} else if cooldown, ok := s.rateLimitService.get429FallbackCooldown(ctx, account); ok && cooldown > 0 {
cooldownUntil = time.Now().Add(cooldown)
hasCooldown = true
}
}
if !hasCooldown {
// The request-local retry window has expired without an upstream reset
// signal. Keep the account out of new selections while this request
// switches to another candidate, rather than immediately selecting it
// again on a concurrent request.
cooldownUntil = time.Now().Add(openAIStopSchedulingBridgeCooldown)
}
s.BlockAccountScheduling(account, cooldownUntil, "429")
s.openaiOAuth429RetryStartedAt.Delete(account.ID)
}
// shouldRetryOpenAIOAuth429OnSameAccount keeps an OAuth account pinned while
// a transient 429 is still inside its retry window. API-key accounts keep the
// existing pool-mode behavior.
func (s *OpenAIGatewayService) shouldRetryOpenAIOAuth429OnSameAccount(account *Account, statusCode int, shouldDisable bool) bool {
if shouldDisable || account == nil {
return false
}
if statusCode == http.StatusTooManyRequests && isOpenAIOAuthAccount(account) && !account.IsShadow() {
if s.settingService != nil && s.rateLimit429StrategySettings().Strategy != "same_account_retry" {
return false
}
// A prior retry window may already have expired and parked this account.
// Do not create a fresh window while that runtime block is active.
if s.isOpenAIAccountRuntimeBlocked(account) {
return false
}
return s.openAIOAuth429RetryWindowActive(account)
}
return account.IsPoolMode() && account.IsPoolModeRetryableStatus(statusCode)
}
// ShouldRetryOpenAIOAuth429 is used before persisting a scheduler block. An
// upstream-provided reset takes precedence; only temporary 429s without one
// stay on the same OAuth account during the retry window.
func (s *OpenAIGatewayService) ShouldRetryOpenAIOAuth429(account *Account, headers http.Header, responseBody []byte) bool {
if s == nil || !isOpenAIOAuthAccount(account) || account.IsShadow() {
return false
}
if s.isOpenAIAccountRuntimeBlocked(account) {
return false
}
if s.settingService != nil && s.rateLimit429StrategySettings().Strategy != "same_account_retry" {
return false
}
if s.rateLimitService != nil && s.rateLimitService.calculateOpenAI429ResetTime(headers) != nil {
return false
}
if parseOpenAIRateLimitResetTime(responseBody) != nil {
return false
}
return s.openAIOAuth429RetryWindowActive(account)
}
func (s *OpenAIGatewayService) openAIOAuth429RetryWindowActive(account *Account) bool {
if s == nil || !isOpenAIOAuthAccount(account) || account.IsShadow() {
return false
}
now := time.Now()
value, _ := s.openaiOAuth429RetryStartedAt.LoadOrStore(account.ID, now)
startedAt, ok := value.(time.Time)
if !ok {
s.openaiOAuth429RetryStartedAt.Store(account.ID, now)
startedAt = now
}
window := openAIOAuth429RetryWindow
if s.settingService != nil {
window = time.Duration(s.rateLimit429StrategySettings().RetryMaxDurationSeconds) * time.Second
}
return now.Sub(startedAt) < window
}
func openAIOAuth429SameAccountRetryDelay(statusCode int, account *Account) time.Duration {
if statusCode == http.StatusTooManyRequests && isOpenAIOAuthAccount(account) && !account.IsShadow() {
return openAIOAuth429RetryDelay
}
return 0
}
func (s *OpenAIGatewayService) openAIOAuth429SameAccountRetryDelay(statusCode int, account *Account) time.Duration {
if statusCode == http.StatusTooManyRequests && isOpenAIOAuthAccount(account) && !account.IsShadow() && s != nil && s.settingService != nil {
return time.Duration(s.rateLimit429StrategySettings().RetryIntervalMs) * time.Millisecond
}
return openAIOAuth429SameAccountRetryDelay(statusCode, account)
}
// openAIOAuth429RetryDeadline returns the request-local retry window end that
// was established when the account first saw a temporary OAuth 429.
func (s *OpenAIGatewayService) openAIOAuth429RetryDeadline(account *Account) time.Time {
if s == nil || !isOpenAIOAuthAccount(account) || account.IsShadow() {
return time.Time{}
}
value, ok := s.openaiOAuth429RetryStartedAt.Load(account.ID)
if !ok {
return time.Time{}
}
startedAt, ok := value.(time.Time)
if !ok {
return time.Time{}
}
window := openAIOAuth429RetryWindow
if s.settingService != nil {
window = time.Duration(s.rateLimit429StrategySettings().RetryMaxDurationSeconds) * time.Second
}
return startedAt.Add(window)
}
// SameAccountRetryLimit returns the request-local retry budget. OAuth 429s
// deliberately use a time-derived budget rather than an account pool setting.
func SameAccountRetryLimit(account *Account, failoverErr *UpstreamFailoverError) int {
if failoverErr != nil && failoverErr.StatusCode == http.StatusTooManyRequests &&
isOpenAIOAuthAccount(account) && !account.IsShadow() {
if failoverErr.SameAccountRetryMax > 0 {
return failoverErr.SameAccountRetryMax
}
return 24
}
if account == nil {
return 0
}
return account.GetPoolModeRetryCount()
}
func (s *OpenAIGatewayService) openAIOAuth429SameAccountRetryMax() int {
if s == nil || s.settingService == nil {
return 24
}
settings := s.rateLimit429StrategySettings()
interval := time.Duration(settings.RetryIntervalMs) * time.Millisecond
window := time.Duration(settings.RetryMaxDurationSeconds) * time.Second
max := int(window / interval)
if max < 1 {
max = 1
}
if max > 240 {
max = 240
}
return max
}
func (s *OpenAIGatewayService) BlockAccountScheduling(account *Account, until time.Time, reason string) {
@@ -501,11 +345,7 @@ func (s *OpenAIGatewayService) isOpenAIOAuth429Storm() bool {
}
func (s *OpenAIGatewayService) ShouldStopOpenAIOAuth429Failover(account *Account, statusCode int, failedSwitches int, state *OpenAIOAuth429FailoverState) bool {
maxSwitches := openAIOAuth429StormMaxAccountSwitches
if s != nil && s.settingService != nil {
maxSwitches = s.rateLimit429StrategySettings().MaxAccountSwitches
}
if failedSwitches < maxSwitches {
if failedSwitches < openAIOAuth429StormMaxAccountSwitches {
return false
}
if state != nil && state.grokOAuth429FollowupPending {
@@ -528,9 +368,5 @@ func (s *OpenAIGatewayService) ShouldStopOpenAIOAuth429Failover(account *Account
if statusCode != http.StatusTooManyRequests || !isOpenAIOAuthAccount(account) {
return false
}
// failedSwitches is incremented after each exhausted candidate. Therefore,
// a value of three means this request has already given three distinct OAuth
// accounts their full same-account retry window. A 429 storm is diagnostic
// only; it must not skip those candidates and return a client 429 early.
return failedSwitches >= maxSwitches+1
return s.isOpenAIOAuth429Storm()
}
@@ -125,19 +125,13 @@ func (s *OpenAIGatewayService) failoverOpenAIUpstreamHTTPError(
if account.Platform != PlatformGrok && !tempUnscheduled {
shouldDisable = s.handleOpenAIAccountUpstreamError(ctx, account, resp.StatusCode, resp.Header, respBody, upstreamModel)
}
failoverErr := newOpenAIUpstreamFailoverError(
return newOpenAIUpstreamFailoverError(
resp.StatusCode,
resp.Header,
respBody,
upstreamMsg,
s.shouldRetryOpenAIOAuth429OnSameAccount(account, resp.StatusCode, shouldDisable) || (!shouldDisable && account.IsPoolMode() && isOpenAITransientProcessingError(resp.StatusCode, upstreamMsg, respBody)),
!shouldDisable && account.IsPoolMode() && (account.IsPoolModeRetryableStatus(resp.StatusCode) || isOpenAITransientProcessingError(resp.StatusCode, upstreamMsg, respBody)),
)
if failoverErr.RetryableOnSameAccount {
failoverErr.SameAccountRetryDelay = s.openAIOAuth429SameAccountRetryDelay(resp.StatusCode, account)
failoverErr.SameAccountRetryDeadline = s.openAIOAuth429RetryDeadline(account)
failoverErr.SameAccountRetryMax = s.openAIOAuth429SameAccountRetryMax()
}
return failoverErr
}
// openAIChatCompletionsTargetURL 解析账号的(非 Grok)Chat Completions 上游端点。
@@ -156,7 +150,7 @@ func (s *OpenAIGatewayService) openAIChatCompletionsTargetURL(account *Account)
// resolveCCFallbackTarget 解析两条 CC 回退路径共用的账号凭证与上游端点
// (回退路径仅面向 APIKey 账号,凭证恒为 openai api_key)。
func (s *OpenAIGatewayService) resolveCCFallbackTarget(account *Account) (apiKey string, targetURL string, err error) {
apiKey = account.GetOpenAIApiKey()
apiKey = strings.TrimSpace(account.GetOpenAIProtocolAPIKey())
if apiKey == "" {
return "", "", fmt.Errorf("account %d missing api_key", account.ID)
}
@@ -214,9 +208,7 @@ func (s *OpenAIGatewayService) sendCCUpstreamRequest(
if account.Platform == PlatformGrok {
if account.IsGrokOAuth() {
if err := applyGrokInteractiveUpstreamHeadersFromAccount(ctx, upstreamReq, account); err != nil {
return nil, err
}
applyGrokCLIHeaders(upstreamReq.Header)
}
applyGrokCacheHeaders(upstreamReq.Header, grokCacheIdentity)
}
@@ -228,7 +220,7 @@ func (s *OpenAIGatewayService) sendCCUpstreamRequest(
if account.Proxy != nil {
proxyURL = account.Proxy.URL()
}
resp, err := s.httpUpstream.Do(upstreamReq, proxyURL, account.ID, account.EffectiveConcurrency())
resp, err := s.httpUpstream.Do(upstreamReq, proxyURL, account.ID, account.Concurrency)
if err != nil {
return nil, s.handleOpenAIUpstreamTransportError(ctx, c, account, err, false)
}
@@ -21,6 +21,7 @@ import (
func (s *OpenAIGatewayService) Forward(ctx context.Context, c *gin.Context, account *Account, body []byte) (*OpenAIForwardResult, error) {
beginUpstreamResponseModelObservation(c)
clearGrokResponsesClientToolMapping(c)
clearOpenAIResponsesClientToolMapping(c)
clearOpenAIResponsesNamespaceNames(c)
startTime := time.Now()
// 固定渠道映射后的请求级 canonical body;账号 normalize/strip 不得改写跨 failover hint。
@@ -39,9 +40,6 @@ func (s *OpenAIGatewayService) Forward(ctx context.Context, c *gin.Context, acco
})
return nil, errors.New("codex_cli_only restriction: only codex official clients are allowed")
}
if c != nil && c.Request != nil {
maybeLearnOfficialDeviceProfile(ctx, account, c.Request.Header)
}
normalizedBody, normalized, err := normalizeOpenAICodexCompactReasoningEffortForAccount(c, account, body)
if err != nil {
@@ -106,12 +104,21 @@ func (s *OpenAIGatewayService) Forward(ctx context.Context, c *gin.Context, acco
requestView := newOpenAIRequestView(body)
reqModel, reqStream, promptCacheKey := requestView.Model, requestView.Stream, requestView.PromptCacheKey
originalModel := reqModel
nativeDeepSeekResponses := account.Platform == PlatformDeepseek &&
(account.GetAPIProtocol() == APIProtocolResponses || account.IsAdaptiveAPIProtocol())
if account.Platform == PlatformGrok {
return s.forwardGrokResponses(ctx, c, account, body, originalModel, reqStream, startTime)
}
if account.Type == AccountTypeAPIKey && !openai_compat.ShouldUseResponsesAPI(account.Extra) {
// CN 供应商 anthropic 协议账号:/v1/responses 入站是交叉协议组合
// (Responses 客户端 × Anthropic 上游),转成 Anthropic 请求走原生端点。
// 不能落到下面的 raw-CC 分支——其 URL 构造会把 anthropic base 当 CC base 用。
if account.IsAnthropicProtocol() {
return s.forwardResponsesViaNativeAnthropic(ctx, c, account, body, reqModel)
}
if shouldForwardOpenAIResponsesViaRawChatCompletions(account) {
return s.forwardResponsesViaRawChatCompletions(ctx, c, account, body)
}
if account.Platform == PlatformOpenAI && account.Type == AccountTypeAPIKey {
@@ -281,7 +288,7 @@ func (s *OpenAIGatewayService) Forward(ctx context.Context, c *gin.Context, acco
instructions := gjson.GetBytes(body, "instructions")
instructionsEmpty := !instructions.Exists() || instructions.Type != gjson.String || strings.TrimSpace(instructions.String()) == ""
if instructionsEmpty && !compatMessagesBridge {
if instructionsEmpty && !compatMessagesBridge && !nativeDeepSeekResponses {
markPatchSet("instructions", defaultCodexSynthInstructions(reqModel))
}
@@ -413,22 +420,35 @@ func (s *OpenAIGatewayService) Forward(ctx context.Context, c *gin.Context, acco
if codexResult.Modified {
markDecodedModified()
}
// 一次加载档案:请求体 client_metadata 与出站头共享同一份 fingerprint IDs。
// compact 形态不同,跳过。fpIDs == nil 时不单独盖 installation id。
// 带真实 device_id 时补齐 client_metadata 安装标识,与真实 Codex 对齐(compact 形态不同,跳过)。
if !isCompactRequest && applyCodexClientMetadata(decoded, account) {
markDecodedModified()
}
stageCodexFingerprintIDs(c, nil)
// 指纹收敛:一次性解析收敛 ID,请求体和出站头共享同一份 IDs(保证 turn_id 等随机字段一致)。
// fingerprintIDs 在此处解析,后续 buildUpstreamRequest 中使用同一份。
if !isCompactRequest {
var clientHeaders http.Header
if c != nil && c.Request != nil {
clientHeaders = c.Request.Header
}
fpIDs := applyCodexForwardRequestIdentity(ctx, c, decoded, account, clientHeaders)
fpIDs := resolveCodexFingerprintIDsFromRequest(account, clientHeaders)
if fpIDs != nil {
markDecodedModified()
if applyCodexFingerprintClientMetadata(decoded, fpIDs) {
markDecodedModified()
}
}
// 将 fpIDs 存入 gin context,供 buildUpstreamRequest 中头改写使用。
// 无条件覆写(含 nil):failover 从收敛账号切到 off 账号时,上一
// 账号的 IDs 不得残留(stageCodexFingerprintIDs 注释)。
stageCodexFingerprintIDs(c, fpIDs)
}
if codexResult.NormalizedModel != "" {
upstreamModel = codexResult.NormalizedModel
}
if codexResult.PromptCacheKey != "" {
if currentPromptCacheKey, ok := decoded["prompt_cache_key"].(string); ok && currentPromptCacheKey != "" {
promptCacheKey = currentPromptCacheKey
} else if codexResult.PromptCacheKey != "" {
promptCacheKey = codexResult.PromptCacheKey
}
}
@@ -441,7 +461,7 @@ func (s *OpenAIGatewayService) Forward(ctx context.Context, c *gin.Context, acco
maxOutputTokens := gjson.GetBytes(body, "max_output_tokens")
if maxOutputTokens.Exists() {
switch account.Platform {
case PlatformOpenAI:
case PlatformOpenAI, PlatformDeepseek:
// Preserve Responses-native output limits unless the selected upstream
// explicitly rejects the field in the bounded HTTP retry loop below.
case PlatformAnthropic:
@@ -835,7 +855,7 @@ func (s *OpenAIGatewayService) Forward(ctx context.Context, c *gin.Context, acco
// Send request
upstreamStart := time.Now()
resp, err := s.doAccountHTTP(ctx, c, account, upstreamReq, proxyURL, "responses")
resp, err := s.httpUpstream.Do(upstreamReq, proxyURL, account.ID, account.Concurrency)
SetOpsLatencyMs(c, OpsUpstreamLatencyMsKey, time.Since(upstreamStart).Milliseconds())
if headerGuard != nil && headerGuard.stopHeaderWait() {
if resp != nil && resp.Body != nil {
@@ -929,19 +949,13 @@ func (s *OpenAIGatewayService) Forward(ctx context.Context, c *gin.Context, acco
})
shouldDisable := s.handleFailoverSideEffects(ctx, resp, account, respBody, upstreamModel)
failoverErr := newOpenAIUpstreamFailoverError(
return nil, newOpenAIUpstreamFailoverError(
resp.StatusCode,
resp.Header,
respBody,
upstreamMsg,
s.shouldRetryOpenAIOAuth429OnSameAccount(account, resp.StatusCode, shouldDisable) || (!shouldDisable && account.IsPoolMode() && isOpenAITransientProcessingError(resp.StatusCode, upstreamMsg, respBody)),
!shouldDisable && account.IsPoolMode() && (account.IsPoolModeRetryableStatus(resp.StatusCode) || isOpenAITransientProcessingError(resp.StatusCode, upstreamMsg, respBody)),
)
if failoverErr.RetryableOnSameAccount {
failoverErr.SameAccountRetryDelay = s.openAIOAuth429SameAccountRetryDelay(resp.StatusCode, account)
failoverErr.SameAccountRetryDeadline = s.openAIOAuth429RetryDeadline(account)
failoverErr.SameAccountRetryMax = s.openAIOAuth429SameAccountRetryMax()
}
return nil, failoverErr
}
return s.handleErrorResponse(ctx, resp, c, account, body, billingModel)
}
@@ -1027,6 +1041,25 @@ func (s *OpenAIGatewayService) Forward(ctx context.Context, c *gin.Context, acco
}
}
func shouldForwardOpenAIResponsesViaRawChatCompletions(account *Account) bool {
if account == nil || account.Type != AccountTypeAPIKey {
return false
}
if account.IsCNProvider() {
// CN 的显式协议配置优先于异步探针 Extra;adaptive 仅 DeepSeek 有原生
// Responses,Kimi/GLM 回退 Chat Completions。
switch account.GetAPIProtocol() {
case APIProtocolChatCompletions:
return true
case APIProtocolAdaptive:
return account.Platform != PlatformDeepseek
default:
return false
}
}
return !openai_compat.ShouldUseResponsesAPI(account.Extra)
}
func (s *OpenAIGatewayService) buildUpstreamRequest(ctx context.Context, c *gin.Context, account *Account, body []byte, token string, isStream bool, promptCacheKey string, isCodexCLI bool) (*http.Request, error) {
// Determine target URL based on account type
var targetURL string
@@ -1037,6 +1070,9 @@ func (s *OpenAIGatewayService) buildUpstreamRequest(ctx context.Context, c *gin.
case AccountTypeAPIKey:
// API Key accounts use Platform API or custom base URL
baseURL := account.GetOpenAIBaseURL()
if account.Platform == PlatformDeepseek && account.IsAdaptiveAPIProtocol() {
baseURL = account.GetCNProtocolBaseURL(APIProtocolResponses)
}
if baseURL == "" {
targetURL = openaiPlatformAPIURL
} else {
@@ -1044,13 +1080,17 @@ func (s *OpenAIGatewayService) buildUpstreamRequest(ctx context.Context, c *gin.
if err != nil {
return nil, err
}
targetURL = buildOpenAIResponsesURL(validatedURL)
targetURL = buildOpenAIResponsesURLForPlatform(account.Platform, validatedURL)
}
default:
targetURL = openaiPlatformAPIURL
}
targetURL = appendOpenAIResponsesRequestPathSuffix(targetURL, openAIResponsesRequestPathSuffix(c))
// DeepSeek 原生 Responses 端点为无状态实现:强制 store=false、清除
// previous_response_id,避免携带状态字段被上游拒绝。
body = normalizeDeepSeekResponsesRequestBody(account, body)
req, err := http.NewRequestWithContext(ctx, "POST", targetURL, bytes.NewReader(body))
if err != nil {
return nil, err
@@ -1087,6 +1127,9 @@ func (s *OpenAIGatewayService) buildUpstreamRequest(ctx context.Context, c *gin.
}
}
}
// 客户端回带的 x-codex-turn-state 若已知由其他账号铸造(failover 换号),
// 剥离后再出站——异账号 blob 与本账号的(指纹收敛后)出站身份自相矛盾。
s.guardOpenAICodexTurnStateEcho(c, account, req.Header)
if account.Type == AccountTypeOAuth {
compatMessagesBridge := isOpenAICompatMessagesBridgeContext(c) || isOpenAICompatMessagesBridgeBody(body)
// 清除客户端透传的 session 头,后续用隔离后的值重新设置,防止跨用户会话碰撞。
@@ -1101,19 +1144,18 @@ func (s *OpenAIGatewayService) buildUpstreamRequest(ctx context.Context, c *gin.
req.Header.Set("originator", resolveOpenAIUpstreamOriginator(c, isCodexCLI))
}
apiKeyID := getAPIKeyIDFromContext(c)
profile := resolveOpenAIOutboundDeviceProfile(ctx, c, account)
if isOpenAIResponsesCompactPath(c) {
req.Header.Set("accept", "application/json")
if req.Header.Get("version") == "" {
req.Header.Set("version", codexCLIVersion)
req.Header.Set("version", CodexCanonicalClientVersion())
}
compactSession := resolveOpenAICompactSessionID(c)
req.Header.Set("session_id", openaiOutboundSessionIDFromProfile(profile, apiKeyID, compactSession))
req.Header.Set("session_id", isolateOpenAISessionID(apiKeyID, compactSession))
} else {
req.Header.Set("accept", "text/event-stream")
}
if promptCacheKey != "" {
isolated := openaiOutboundSessionIDFromProfile(profile, apiKeyID, promptCacheKey)
isolated := isolateOpenAISessionID(apiKeyID, promptCacheKey)
req.Header.Set("session_id", isolated)
if !compatMessagesBridge || clientConversationID != "" {
req.Header.Set("conversation_id", isolated)
@@ -1131,28 +1173,19 @@ func (s *OpenAIGatewayService) buildUpstreamRequest(ctx context.Context, c *gin.
req.Header.Set("user-agent", customUA)
}
// 若开启 ForceCodexCLI,则强制将上游 User-Agent 伪装为 Codex CLI。
// 若开启 ForceCodexCLI,则强制将上游 User-Agent 伪装为规范 Codex 身份。
// 用于网关未透传/改写 User-Agent 时,仍能命中 Codex 侧识别逻辑。
if s.cfg != nil && s.cfg.Gateway.ForceCodexCLI {
req.Header.Set("user-agent", codexCLIUserAgent)
req.Header.Set("user-agent", CodexCanonicalUserAgent())
}
// 指纹收敛:使用 Forward() 中预计算的收敛 ID 改写出站头,与请求体使用同一份 IDs。
// leftover 5 session/full 模式会用账号级恒定 session_id 覆盖上面的
// isolate+namespace 值;那是「一号一安装」收敛,不是 leftover 11 的缺口。
// leftover 11 在 off/device、以及不走指纹的 passthrough/WS/compat 路径生效。
if account.Type == AccountTypeOAuth && c != nil {
if fpIDs, ok := c.Get("codex_fingerprint_ids"); ok {
if ids, ok := fpIDs.(*codexFingerprintIDs); ok && fingerprintIDsBelongToAccount(ids, account) {
applyCodexFingerprintHeaders(req.Header, ids)
}
}
}
applyStagedCodexFingerprintHeaders(c, account, req.Header)
// 终态收口:强制统一 OAuth 出站身份(User-Agent / originator / version 同源自洽)。
// 客户端自报身份不参与构造,浏览器型 UA 也因此不会再到达上游(原浏览器 UA 兜底已被吸收)。
if account.Type == AccountTypeOAuth {
s.enforceCodexIdentityFromLoadedProfile(req.Header, account, outboundDeviceProfileFromGin(c, account))
enforceCodexIdentityHeadersWithUA(req.Header, s.codexIdentityOverrideUA(account))
}
// Ensure required headers exist
@@ -1162,6 +1195,9 @@ func (s *OpenAIGatewayService) buildUpstreamRequest(ctx context.Context, c *gin.
// 账号级请求头覆写(仅 openai api_key 账号启用时生效;OAuth 路径 no-op)
account.ApplyHeaderOverrides(req.Header)
// x-codex-beta-features:按真实 Codex 的会话级行为补注(在账号级覆写之后,
// 保证不被覆盖丢失)。
applyOpenAICodexBetaFeatures(c, account, req.Header)
setOpenAICodexRoutingHintFromBody(req.Header, account, body)
logOpenAIRoutingDiagnosticsFromBody(ctx, account, "http", req.Header, body, "not_applicable")
@@ -1177,25 +1213,3 @@ func (s *OpenAIGatewayService) codexIdentityOverrideUA(account *Account) string
}
return account.GetOpenAIUserAgent()
}
func (s *OpenAIGatewayService) enforceCodexIdentityFromAccount(ctx context.Context, h http.Header, account *Account) {
if account == nil {
s.enforceCodexIdentityFromLoadedProfile(h, account, nil)
return
}
profile, err := LoadOutboundDeviceProfile(ctx, account)
if err != nil {
profile = nil
}
s.enforceCodexIdentityFromLoadedProfile(h, account, profile)
}
func (s *OpenAIGatewayService) enforceCodexIdentityFromLoadedProfile(h http.Header, account *Account, profile *AccountDeviceProfile) {
fallback := s.codexIdentityOverrideUA(account)
if profile == nil {
enforceCodexIdentityHeadersWithUA(h, fallback)
return
}
identity := resolveCodexOutboundIdentityFromProfile(profile, fallback)
enforceCodexIdentityHeadersWithUA(h, identity.userAgent)
}
@@ -16,8 +16,8 @@ import (
"strings"
"time"
"github.com/Wei-Shaw/sub2api/internal/pkg/apicompat"
"github.com/Wei-Shaw/sub2api/internal/pkg/logger"
"github.com/Wei-Shaw/sub2api/internal/pkg/openai"
"github.com/Wei-Shaw/sub2api/internal/util/responseheaders"
"github.com/gin-gonic/gin"
"github.com/tidwall/gjson"
@@ -25,6 +25,84 @@ import (
"go.uber.org/zap"
)
const openAIResponsesClientToolMappingContextKey = "openai_responses_client_tool_mapping"
func hasOpenAIResponsesClientToolMapping(mapping apicompat.ResponsesClientToolMapping) bool {
return len(mapping.CustomTools) > 0 || mapping.ToolSearch || len(mapping.NamespaceTools) > 0
}
func adaptOpenAIResponsesClientTools(body []byte) ([]byte, apicompat.ResponsesClientToolMapping, error) {
if !needsOpenAIResponsesClientToolAdaptation(body) {
return body, apicompat.ResponsesClientToolMapping{}, nil
}
decoder := json.NewDecoder(bytes.NewReader(body))
decoder.UseNumber()
var requestBody map[string]any
if err := decoder.Decode(&requestBody); err != nil {
return body, apicompat.ResponsesClientToolMapping{}, fmt.Errorf("decode OpenAI Responses client tools: %w", err)
}
var trailingValue any
if err := decoder.Decode(&trailingValue); !errors.Is(err, io.EOF) {
if err == nil {
err = errors.New("multiple JSON values")
}
return body, apicompat.ResponsesClientToolMapping{}, fmt.Errorf("decode OpenAI Responses client tools trailing data: %w", err)
}
mapping, changed, err := apicompat.AdaptResponsesClientTools(requestBody)
if err != nil || !changed {
return body, mapping, err
}
rebuilt, err := marshalOpenAIUpstreamJSON(requestBody)
if err != nil {
return body, apicompat.ResponsesClientToolMapping{}, fmt.Errorf("encode OpenAI Responses client tools: %w", err)
}
return rebuilt, mapping, nil
}
func needsOpenAIResponsesClientToolAdaptation(body []byte) bool {
needsAdaptation := false
var visit func(gjson.Result) bool
visit = func(value gjson.Result) bool {
if value.IsObject() {
switch strings.TrimSpace(value.Get("type").String()) {
case "custom", "custom_tool_call", "custom_tool_call_output",
"tool_search", "tool_search_call", "tool_search_output":
needsAdaptation = true
return false
}
}
if value.IsObject() || value.IsArray() {
value.ForEach(func(_, child gjson.Result) bool {
return visit(child)
})
}
return !needsAdaptation
}
visit(gjson.ParseBytes(body))
return needsAdaptation
}
func openAIResponsesClientToolMapping(c *gin.Context) (apicompat.ResponsesClientToolMapping, bool) {
if c == nil {
return apicompat.ResponsesClientToolMapping{}, false
}
value, ok := c.Get(openAIResponsesClientToolMappingContextKey)
mapping, typed := value.(apicompat.ResponsesClientToolMapping)
return mapping, ok && typed && hasOpenAIResponsesClientToolMapping(mapping)
}
// clearOpenAIResponsesClientToolMapping removes mapping state from the prior
// forwarding attempt. Forward retries accounts on the same Gin context.
func clearOpenAIResponsesClientToolMapping(c *gin.Context) {
if c == nil {
return
}
if _, exists := c.Get(openAIResponsesClientToolMappingContextKey); exists {
c.Set(openAIResponsesClientToolMappingContextKey, apicompat.ResponsesClientToolMapping{})
}
}
func (s *OpenAIGatewayService) forwardOpenAIPassthrough(
ctx context.Context,
c *gin.Context,
@@ -80,6 +158,39 @@ func (s *OpenAIGatewayService) forwardOpenAIPassthrough(
body = normalizedBody
}
reqStream = gjson.GetBytes(body, "stream").Bool()
stageCodexFingerprintIDs(c, nil)
// 指纹收敛:与非透传路径同门控(仅 OAuth、legacy compact 形态跳过)。
// 一次性解析收敛 ID:请求体 client_metadata 在此改写(raw 字节外科
// 手术,透传热路径禁全量 Unmarshal),出站头改写由请求构造器读取
// context 中的同一份 IDs 完成(turn_id 等随机字段两侧必须一致)。
if !isOpenAIResponsesCompactPath(c) {
var clientHeaders http.Header
if c != nil && c.Request != nil {
clientHeaders = c.Request.Header
}
fpIDs := resolveCodexFingerprintIDsFromRequest(account, clientHeaders)
if fpIDs != nil {
fpBody, fpChanged, fpErr := applyCodexFingerprintClientMetadataRaw(body, fpIDs)
if fpErr != nil {
return nil, fpErr
}
if fpChanged {
body = fpBody
}
}
stageCodexFingerprintIDs(c, fpIDs)
}
}
if account != nil && account.Platform == PlatformOpenAI && account.Type == AccountTypeAPIKey &&
!isOpenAIResponsesCompactPath(c) && needsOpenAIResponsesClientToolAdaptation(body) {
adaptedBody, mapping, adaptErr := adaptOpenAIResponsesClientTools(body)
if adaptErr != nil {
return nil, adaptErr
}
body = adaptedBody
c.Set(openAIResponsesClientToolMappingContextKey, mapping)
}
sanitizedBody, sanitized, err := sanitizeEmptyBase64InputImagesInOpenAIBody(body)
@@ -202,7 +313,7 @@ func (s *OpenAIGatewayService) forwardOpenAIPassthrough(
}
upstreamStart := time.Now()
resp, err = s.doAccountHTTP(ctx, c, account, upstreamReq, proxyURL, "responses")
resp, err = s.httpUpstream.Do(upstreamReq, proxyURL, account.ID, account.Concurrency)
SetOpsLatencyMs(c, OpsUpstreamLatencyMsKey, time.Since(upstreamStart).Milliseconds())
if err != nil {
// Transport-level failure (proxy/DNS/TCP/TLS — no HTTP response). Convert to
@@ -236,9 +347,23 @@ func (s *OpenAIGatewayService) forwardOpenAIPassthrough(
return nil, s.handleErrorResponsePassthrough(ctx, resp, c, account, body, probeBody)
}
defer func() { _ = resp.Body.Close() }()
if mapping, ok := openAIResponsesClientToolMapping(c); ok && isEventStreamResponse(resp.Header) {
maxLineSize := defaultMaxLineSize
if s.cfg != nil && s.cfg.Gateway.MaxLineSize > 0 {
maxLineSize = s.cfg.Gateway.MaxLineSize
}
resp.Body = newGrokResponsesClientToolStreamBody(resp.Body, mapping, maxLineSize)
}
serviceTier := extractOpenAIServiceTierFromBody(body)
// x-codex-turn-state 溯源:下游回传由 writeOpenAIPassthroughResponseHeaders
// 在各 handler 的写头点强制放行,铸造账号在此统一记录,供出站守卫剥离
// failover 换号后的跨账号回带(openai_codex_turn_state.go)。
if extractOpenAICodexTurnState(resp.Header) != "" {
s.noteOpenAICodexTurnStateProvenance(c, account)
}
var usage *OpenAIUsage
var firstTokenMs *int
responseID := ""
@@ -351,11 +476,14 @@ func (s *OpenAIGatewayService) buildUpstreamRequestOpenAIPassthrough(
if err != nil {
return nil, err
}
targetURL = buildOpenAIResponsesURL(validatedURL)
targetURL = buildOpenAIResponsesURLForPlatform(account.Platform, validatedURL)
}
}
targetURL = appendOpenAIResponsesRequestPathSuffix(targetURL, openAIResponsesRequestPathSuffix(c))
// DeepSeek 原生 Responses 端点为无状态实现(见 normalizeDeepSeekResponsesRequestBody)。
body = normalizeDeepSeekResponsesRequestBody(account, body)
req, err := http.NewRequestWithContext(ctx, http.MethodPost, targetURL, bytes.NewReader(body))
if err != nil {
return nil, err
@@ -376,6 +504,10 @@ func (s *OpenAIGatewayService) buildUpstreamRequestOpenAIPassthrough(
}
}
// 客户端回带的 x-codex-turn-state 若已知由其他账号铸造(failover 换号),
// 剥离后再出站(openai_codex_turn_state.go)。
s.guardOpenAICodexTurnStateEcho(c, account, req.Header)
// 覆盖入站鉴权残留,并注入上游认证
req.Header.Del("authorization")
req.Header.Del("x-api-key")
@@ -408,7 +540,7 @@ func (s *OpenAIGatewayService) buildUpstreamRequestOpenAIPassthrough(
if isOpenAIResponsesCompactPath(c) {
req.Header.Set("accept", "application/json")
if req.Header.Get("version") == "" {
req.Header.Set("version", codexCLIVersion)
req.Header.Set("version", CodexCanonicalClientVersion())
}
if clientSessionID == "" {
clientSessionID = resolveOpenAICompactSessionID(c)
@@ -417,7 +549,7 @@ func (s *OpenAIGatewayService) buildUpstreamRequestOpenAIPassthrough(
req.Header.Set("accept", "text/event-stream")
}
if req.Header.Get("originator") == "" {
req.Header.Set("originator", openai.CodexDefaultOriginator)
req.Header.Set("originator", resolveCodexOutboundIdentity("").originator)
}
// 用隔离后的 session 标识符覆盖客户端透传值,防止跨用户会话碰撞。
if clientSessionID == "" {
@@ -426,14 +558,11 @@ func (s *OpenAIGatewayService) buildUpstreamRequestOpenAIPassthrough(
if clientConversationID == "" {
clientConversationID = promptCacheKey
}
if clientSessionID != "" || clientConversationID != "" {
sessionID, conversationID := openaiOutboundSessionPair(ctx, account, apiKeyID, clientSessionID, clientConversationID)
if clientSessionID != "" {
req.Header.Set("session_id", sessionID)
}
if clientConversationID != "" {
req.Header.Set("conversation_id", conversationID)
}
if clientSessionID != "" {
req.Header.Set("session_id", isolateOpenAISessionID(apiKeyID, clientSessionID))
}
if clientConversationID != "" {
req.Header.Set("conversation_id", isolateOpenAISessionID(apiKeyID, clientConversationID))
}
} else if isOpenAIResponsesCompactPath(c) {
// 透传白名单会放行客户端的 Accept: text/event-stream;compact 上游是
@@ -448,12 +577,16 @@ func (s *OpenAIGatewayService) buildUpstreamRequestOpenAIPassthrough(
req.Header.Set("user-agent", customUA)
}
if s.cfg != nil && s.cfg.Gateway.ForceCodexCLI {
req.Header.Set("user-agent", codexCLIUserAgent)
req.Header.Set("user-agent", CodexCanonicalUserAgent())
}
// 指纹收敛:使用 forwardOpenAIPassthrough 中预计算的收敛 ID 改写出站头,
// 与请求体 client_metadata 共享同一份 IDs(与非透传路径相同的相对位置:
// 会话隔离之后、终态身份收口之前)。
applyStagedCodexFingerprintHeaders(c, account, req.Header)
// 终态收口:透传路径的 OAuth 与非透传完全一致,同样强制统一出站身份
// (User-Agent / originator / version 同源自洽),客户端自报身份不会到达上游。
if account.Type == AccountTypeOAuth {
s.enforceCodexIdentityFromAccount(ctx, req.Header, account)
enforceCodexIdentityHeadersWithUA(req.Header, s.codexIdentityOverrideUA(account))
}
if req.Header.Get("content-type") == "" {
@@ -462,6 +595,9 @@ func (s *OpenAIGatewayService) buildUpstreamRequestOpenAIPassthrough(
// 账号级请求头覆写(仅 openai api_key 账号启用时生效;OAuth 路径 no-op)
account.ApplyHeaderOverrides(req.Header)
// x-codex-beta-features:按真实 Codex 的会话级行为补注(在账号级覆写之后,
// 保证不被覆盖丢失)。
applyOpenAICodexBetaFeatures(c, account, req.Header)
setOpenAICodexRoutingHintFromBody(req.Header, account, body)
logOpenAIRoutingDiagnosticsFromBody(ctx, account, "http_passthrough", req.Header, body, "not_applicable")
@@ -638,19 +774,13 @@ func (s *OpenAIGatewayService) handleFailoverErrorResponsePassthrough(
Detail: upstreamDetail,
UpstreamResponseBody: upstreamDetail,
})
failoverErr := newOpenAIUpstreamFailoverError(
return newOpenAIUpstreamFailoverError(
resp.StatusCode,
resp.Header,
body,
upstreamMsg,
s.shouldRetryOpenAIOAuth429OnSameAccount(account, resp.StatusCode, shouldDisable),
!shouldDisable && account.IsPoolMode() && account.IsPoolModeRetryableStatus(resp.StatusCode),
)
if failoverErr.RetryableOnSameAccount {
failoverErr.SameAccountRetryDelay = s.openAIOAuth429SameAccountRetryDelay(resp.StatusCode, account)
failoverErr.SameAccountRetryDeadline = s.openAIOAuth429RetryDeadline(account)
failoverErr.SameAccountRetryMax = s.openAIOAuth429SameAccountRetryMax()
}
return failoverErr
}
func (s *OpenAIGatewayService) handleErrorResponsePassthrough(
@@ -778,6 +908,19 @@ type openaiNonStreamingResultPassthrough struct {
imageOutputSizes []string
}
const openAIStreamKeepaliveBytesKey = "openai_stream_keepalive_bytes"
func recordOpenAIStreamKeepaliveBytes(c *gin.Context, written int) {
if c == nil || written <= 0 {
return
}
current := 0
if value, ok := c.Get(openAIStreamKeepaliveBytesKey); ok {
current, _ = value.(int)
}
c.Set(openAIStreamKeepaliveBytesKey, current+written)
}
func openAIStreamClientOutputStarted(c *gin.Context, localStarted bool) bool {
if localStarted {
return true
@@ -800,6 +943,85 @@ func openAIStreamEventIsPreamble(eventType string) bool {
}
}
func openAIStreamAddedEventStartsClientOutput(payload []byte, eventType string) bool {
if len(payload) == 0 || !gjson.ValidBytes(payload) {
return true
}
switch strings.TrimSpace(eventType) {
case "response.output_item.added":
item := gjson.GetBytes(payload, "item")
if !item.Exists() || !item.IsObject() {
return true
}
switch strings.TrimSpace(item.Get("type").String()) {
case "reasoning":
if item.Get("encrypted_content").String() != "" {
return true
}
summary := item.Get("summary")
if !summary.IsArray() {
return false
}
for _, part := range summary.Array() {
if strings.TrimSpace(part.Get("type").String()) != "summary_text" || part.Get("text").String() != "" {
return true
}
}
return false
case "message":
content := item.Get("content")
if !content.IsArray() {
return false
}
for _, part := range content.Array() {
switch strings.TrimSpace(part.Get("type").String()) {
case "output_text":
if part.Get("text").String() != "" {
return true
}
case "refusal":
if part.Get("refusal").String() != "" {
return true
}
default:
return true
}
}
return false
case "function_call":
return item.Get("arguments").String() != ""
case "custom_tool_call":
return item.Get("input").String() != ""
case "compaction":
return item.Get("encrypted_content").String() != ""
default:
return true
}
case "response.content_part.added":
part := gjson.GetBytes(payload, "part")
if !part.Exists() || !part.IsObject() {
return true
}
switch strings.TrimSpace(part.Get("type").String()) {
case "output_text":
return part.Get("text").String() != ""
case "refusal":
return part.Get("refusal").String() != ""
default:
return true
}
case "response.reasoning_summary_part.added":
part := gjson.GetBytes(payload, "part")
if !part.Exists() || !part.IsObject() || strings.TrimSpace(part.Get("type").String()) != "summary_text" {
return true
}
return part.Get("text").String() != ""
default:
return true
}
}
func openAIStreamDataStartsClientOutput(data, eventType string) bool {
trimmed := strings.TrimSpace(data)
if trimmed == "" {
@@ -816,6 +1038,8 @@ func openAIStreamDataStartsClientOutput(data, eventType string) bool {
// (content_policy / invalid_request 等)维持原样转发,保留上游错误细节。
payload := []byte(trimmed)
return !openAIStreamFailedEventShouldFailover(payload, extractOpenAISSEErrorMessage(payload))
case "response.output_item.added", "response.content_part.added", "response.reasoning_summary_part.added":
return openAIStreamAddedEventStartsClientOutput([]byte(trimmed), eventType)
}
return !openAIStreamEventIsPreamble(eventType)
}
@@ -894,9 +1118,34 @@ func isOpenAIUpstreamCapacityShedEvent(payload []byte) bool {
switch openAIStreamFailedEventErrorCode(payload) {
case "server_is_overloaded", "slow_down":
return true
default:
return false
}
for _, path := range []string{"response.error.message", "error.message", "message"} {
if isOpenAICapacityShedMessage(gjson.GetBytes(payload, path).String()) {
return true
}
}
return false
}
func logOpenAICapacityFailoverSuppressed(
ctx context.Context,
account *Account,
path string,
upstreamRequestID string,
eventType string,
) {
fields := []zap.Field{
zap.String("path", path),
zap.String("event_type", strings.TrimSpace(eventType)),
zap.String("upstream_request_id", strings.TrimSpace(upstreamRequestID)),
}
if account != nil {
fields = append(fields,
zap.Int64("account_id", account.ID),
zap.String("platform", account.Platform),
)
}
logger.FromContext(ctx).Warn("gateway.failover_suppressed_after_semantic_output", fields...)
}
// openAICapacityShedRetryableClientCode 是把上游容量降载错误转发给客户端时改写
@@ -919,9 +1168,12 @@ func sanitizeOpenAICapacityShedErrorCodeForClient(payload []byte) ([]byte, bool)
updated := payload
changed := false
for _, path := range []string{"response.error.code", "error.code"} {
switch strings.ToLower(strings.TrimSpace(gjson.GetBytes(updated, path).String())) {
case "server_is_overloaded", "slow_down":
default:
parent := strings.TrimSuffix(path, ".code")
if !gjson.GetBytes(updated, parent).Exists() {
continue
}
code := strings.ToLower(strings.TrimSpace(gjson.GetBytes(updated, path).String()))
if code != "" && code != "server_is_overloaded" && code != "slow_down" {
continue
}
next, err := sjson.SetBytes(updated, path, openAICapacityShedRetryableClientCode)
@@ -954,7 +1206,7 @@ func openAIStreamFailedEventSemanticStatus(payload []byte, message string) int {
return http.StatusUnauthorized
case strings.Contains(combined, "permission") || strings.Contains(combined, "forbidden") || strings.Contains(combined, "access denied"):
return http.StatusForbidden
case code == "server_is_overloaded" || code == "slow_down":
case isOpenAIUpstreamCapacityShedEvent(payload):
return http.StatusServiceUnavailable
default:
return http.StatusBadGateway
@@ -1088,6 +1340,16 @@ func openAIStreamFailedEventShouldFailover(payload []byte, message string) bool
return true
}
func openAIStreamErrorEventShouldFailover(payload []byte, message string) bool {
if strings.TrimSpace(gjson.GetBytes(payload, "type").String()) != "error" {
return false
}
if isOpenAIContextWindowError(message, payload) {
return false
}
return isOpenAITransientProcessingError(http.StatusBadRequest, message, payload)
}
func openAIStreamFailedEventRetryableOnSameAccount(account *Account, payload []byte, message string) bool {
if account == nil {
return false
@@ -1230,6 +1492,7 @@ func (s *OpenAIGatewayService) handleStreamingResponsePassthrough(
sawTerminalEvent := false
sawFailedEvent := false
semanticOutputSeen := false
capacityFailoverSuppressedLogged := false
failedMessage := ""
clientOutputStarted := false
upstreamRequestID := strings.TrimSpace(resp.Header.Get("x-request-id"))
@@ -1316,6 +1579,32 @@ func (s *OpenAIGatewayService) handleStreamingResponsePassthrough(
}
}
eventType := strings.TrimSpace(gjson.Get(trimmedData, "type").String())
if !capacityFailoverSuppressedLogged && account != nil && account.Platform == PlatformOpenAI &&
(eventType == "error" || eventType == "response.failed") &&
openAIStreamClientOutputStarted(c, clientOutputStarted) &&
isOpenAIUpstreamCapacityShedEvent(dataBytes) {
logOpenAICapacityFailoverSuppressed(ctx, account, "passthrough_sse", upstreamRequestID, eventType)
capacityFailoverSuppressedLogged = true
}
if eventType == "error" && !openAIStreamClientOutputStarted(c, clientOutputStarted) {
errorMessage := extractOpenAISSEErrorMessage(dataBytes)
if status, errType, errMsg, matched := applyOpenAIStreamFailedErrorPassthroughRule(c, account.Platform, dataBytes, errorMessage); matched {
s.recordOpenAIStreamUpstreamError(c, account, true, upstreamRequestID, "http_error", dataBytes, errorMessage)
MarkResponseCommitted(c)
c.Writer.Header().Set("Content-Type", "application/json; charset=utf-8")
c.JSON(status, gin.H{
"error": gin.H{
"type": errType,
"message": errMsg,
},
})
return resultWithUsage(), fmt.Errorf("upstream error event: passthrough rule matched message=%s", errMsg)
}
if openAIStreamErrorEventShouldFailover(dataBytes, errorMessage) {
return resultWithUsage(),
s.newOpenAIStreamFailoverError(c, account, true, upstreamRequestID, dataBytes, errorMessage, resp.Header)
}
}
if eventType == "response.failed" {
failedMessage = extractOpenAISSEErrorMessage(dataBytes)
// response.failed 自带上游已消耗的 usage(input token 通常已扣);必须先解析
@@ -1529,6 +1818,12 @@ func (s *OpenAIGatewayService) handleNonStreamingResponsePassthrough(
if err != nil {
return nil, fmt.Errorf("restore OpenAI passthrough namespace response: %w", err)
}
if mapping, ok := openAIResponsesClientToolMapping(c); ok && json.Valid(body) {
body, _, err = apicompat.RestoreResponsesClientToolPayload(body, mapping)
if err != nil {
return nil, fmt.Errorf("restore OpenAI Responses client tools: %w", err)
}
}
if !writeOpenAICompactSSEBridge(c, resp.StatusCode, body) {
c.Data(resp.StatusCode, contentType, body)
}
@@ -1659,4 +1954,13 @@ func writeOpenAIPassthroughResponseHeaders(dst http.Header, src http.Header, fil
dst.Add(key, v)
}
}
// x-codex-turn-state:Codex 回合状态头,客户端会在同回合后续请求回带。
// 与上面的用量头不同,这里在上游缺失时也主动清除——failover 换号后残留
// 上一账号的 blob 会构成跨账号矛盾(openai_codex_turn_state.go)。
turnStateKey := http.CanonicalHeaderKey(openAICodexTurnStateHeader)
dst.Del(turnStateKey)
for _, v := range getCaseInsensitiveValues(src, openAICodexTurnStateHeader) {
dst.Add(turnStateKey, v)
}
}
@@ -61,7 +61,7 @@ const (
// 陈旧版本会被优先丢弃(HTTP 200 + 流内 server_is_overloaded);非官方客户端配不出
// 官方身份时整体回退到本常量,因此它必须跟随官方 CLI 的当前发布版本,
// 落后多个版本会让这些请求稳定落在被优先丢弃的一侧。
codexCLIVersion = "0.147.0"
codexCLIVersion = "0.146.0"
// Codex 限额快照仅用于后台展示/诊断,不需要每个成功请求都立即落库。
openAICodexSnapshotPersistMinInterval = 30 * time.Second
// 配额自动暂停时,超过该时长仍未刷新的 used% 快照视为陈旧,不再据此暂停账号。
@@ -284,8 +284,9 @@ type OpenAIForwardResult struct {
// AudioUsage carries Voice billing units when present.
AudioUsage *AudioUsage
wsReplayInput []json.RawMessage
wsReplayInputExists bool
wsReplayInput []json.RawMessage
wsReplayInputExists bool
wsAccountFailoverReplayInput []json.RawMessage
}
// SucceededForScheduling reports whether this result is an upstream success
@@ -398,8 +399,8 @@ func (t *accountWriteThrottle) Allow(id int64, now time.Time) bool {
var defaultOpenAICodexSnapshotPersistThrottle = newAccountWriteThrottle(openAICodexSnapshotPersistMinInterval)
// ErrNoAvailableCompactAccounts indicates the request needs /responses/compact
// support but no compatible account is available.
// ErrNoAvailableCompactAccounts indicates a legacy /responses/compact request
// needs compact support but no compatible account is available.
var ErrNoAvailableCompactAccounts = errors.New("no available accounts support /responses/compact")
// OpenAIGatewayService handles OpenAI API gateway operations
@@ -428,8 +429,6 @@ type OpenAIGatewayService struct {
channelService *ChannelService
balanceNotifyService *BalanceNotifyService
settingService *SettingService
tlsFPProfileService *TLSFingerprintProfileService
tlsFPRouterService *TLSFingerprintRouterService
userPlatformQuotaRepo UserPlatformQuotaRepository
liveAttestation liveattestation.Provider
liveAttestationCipher SecretEncryptor
@@ -455,10 +454,6 @@ type OpenAIGatewayService struct {
openaiAccountRuntimeBlockLocks sync.Map // key: int64(accountID), value: *sync.Mutex
openaiAccountRuntimeBlockGeneration sync.Map // key: int64(accountID), value: uint64
openaiAccountRuntimeBlockSequence atomic.Uint64
openaiOAuth429RetryStartedAt sync.Map // key: int64(accountID), value: time.Time
openai429StrategyMu sync.Mutex
openai429StrategyCachedAt time.Time
openai429StrategyCached RateLimit429CooldownSettings
grokCredentialMutationLocks sync.Map // key: int64(accountID), value: *sync.Mutex
openaiOAuth429WindowStartUnixNano atomic.Int64
openaiOAuth429WindowCount atomic.Int64
@@ -468,6 +463,11 @@ type OpenAIGatewayService struct {
codexModelsManifestCache codexModelsManifestCache
openaiCompatSessionResponses sync.Map
openaiCompatAnthropicDigestSessions sync.Map
// openaiCodexTurnStateOrigins: 下游会话 seed → openAICodexTurnStateOrigin,
// 记录最近一次向该会话下发 x-codex-turn-state 的铸造账号,供出站守卫
// 剥离跨账号回带(openai_codex_turn_state.go)。
openaiCodexTurnStateOrigins sync.Map
openaiCodexTurnStateWrites atomic.Uint64
}
// NewOpenAIGatewayService creates a new OpenAIGatewayService
@@ -548,27 +548,6 @@ func NewOpenAIGatewayService(
return svc
}
func (s *OpenAIGatewayService) SetTLSFingerprintServices(profile *TLSFingerprintProfileService, router *TLSFingerprintRouterService) {
if s == nil {
return
}
s.tlsFPProfileService = profile
s.tlsFPRouterService = router
}
func (s *OpenAIGatewayService) doAccountHTTP(ctx context.Context, c *gin.Context, account *Account, req *http.Request, proxyURL, protocol string) (*http.Response, error) {
if account != nil && account.Platform == PlatformGrok {
return doLeftoverAccountHTTP(ctx, s.httpUpstream, req, proxyURL, account, s.tlsFPProfileService, s.tlsFPRouterService, inboundUserAgentFromGin(c), "http", protocol)
}
// Last mutation before send: callers may Header.Set/Get after buildUpstreamRequest
// (images Content-Type, messages identity + turn-state). leftover 5 session_id
// values are unchanged; originator stays lowercase.
if account != nil && account.Type == AccountTypeOAuth && req != nil {
applyCodexHeaderWireCasing(req.Header)
}
return doAccountHTTPUpstreamFromGin(ctx, c, s.httpUpstream, req, proxyURL, account, s.tlsFPProfileService, s.tlsFPRouterService, "http", protocol)
}
// ResolveChannelMapping 解析渠道级模型映射(代理到 ChannelService)
func (s *OpenAIGatewayService) ResolveChannelMapping(ctx context.Context, groupID int64, model string) ChannelMappingResult {
if s.channelService == nil {
@@ -625,6 +604,10 @@ func (s *OpenAIGatewayService) isUpstreamModelRestrictedByChannel(ctx context.Co
if s.channelService == nil {
return false
}
if compactForwardModel, ok := openAIForwardModelFromContext(ctx); ok {
requestedModel = compactForwardModel.model
requireCompact = compactForwardModel.useCompactModelMapping
}
upstreamModel := resolveOpenAIAccountUpstreamModelForRequest(account, requestedModel, requireCompact)
if upstreamModel == "" {
return false
@@ -1073,9 +1056,6 @@ func getAPIKeyIDFromContext(c *gin.Context) int64 {
// isolateOpenAISessionID 将 apiKeyID 混入 session 标识符,
// 确保不同 API Key 的用户即使使用相同的原始 session_id/conversation_id,
// 到达上游的标识符也不同,防止跨用户会话碰撞。
//
// Outbound session/conversation headers should use openaiOutboundSessionID
// or openaiOutboundSessionUUID so a valid device-profile namespace is folded in.
func isolateOpenAISessionID(apiKeyID int64, raw string) string {
raw = strings.TrimSpace(raw)
if raw == "" {
@@ -1087,60 +1067,6 @@ func isolateOpenAISessionID(apiKeyID int64, raw string) string {
return fmt.Sprintf("%016x", h.Sum64())
}
func loadOpenAIOutboundSessionProfile(ctx context.Context, account *Account) *AccountDeviceProfile {
if account == nil || !account.IsOpenAIOAuth() {
return nil
}
return loadOutboundCodexProfile(ctx, account)
}
func deriveOpenAIOutboundSessionIDFromProfile(profile *AccountDeviceProfile, isolated string) string {
if profile == nil || isolated == "" {
return ""
}
sessionID, _, _, err := DeriveSessionIDs(profile.SessionNamespace, isolated)
if err != nil || sessionID == "" {
return ""
}
return sessionID
}
func deriveOpenAIOutboundSessionID(ctx context.Context, account *Account, apiKeyID int64, raw string) string {
isolated := isolateOpenAISessionID(apiKeyID, raw)
if isolated == "" {
return ""
}
return deriveOpenAIOutboundSessionIDFromProfile(loadOpenAIOutboundSessionProfile(ctx, account), isolated)
}
func openaiOutboundSessionIDFromProfile(profile *AccountDeviceProfile, apiKeyID int64, raw string) string {
isolated := isolateOpenAISessionID(apiKeyID, raw)
if derived := deriveOpenAIOutboundSessionIDFromProfile(profile, isolated); derived != "" {
return derived
}
return isolated
}
func openaiOutboundSessionID(ctx context.Context, account *Account, apiKeyID int64, raw string) string {
if derived := deriveOpenAIOutboundSessionID(ctx, account, apiKeyID, raw); derived != "" {
return derived
}
return isolateOpenAISessionID(apiKeyID, raw)
}
func openaiOutboundSessionUUID(ctx context.Context, account *Account, apiKeyID int64, raw string) string {
if derived := deriveOpenAIOutboundSessionID(ctx, account, apiKeyID, raw); derived != "" {
return derived
}
return generateSessionUUID(isolateOpenAISessionID(apiKeyID, raw))
}
func openaiOutboundSessionPair(ctx context.Context, account *Account, apiKeyID int64, sessionRaw, conversationRaw string) (sessionID, conversationID string) {
profile := loadOpenAIOutboundSessionProfile(ctx, account)
return openaiOutboundSessionIDFromProfile(profile, apiKeyID, sessionRaw),
openaiOutboundSessionIDFromProfile(profile, apiKeyID, conversationRaw)
}
func logCodexCLIOnlyDetection(ctx context.Context, c *gin.Context, account *Account, apiKeyID int64, result CodexClientRestrictionDetectionResult, body []byte) {
if !result.Enabled {
return
@@ -1280,7 +1206,7 @@ func (s *OpenAIGatewayService) GetAccessToken(ctx context.Context, account *Acco
}
return apiKey, "apikey", nil
}
apiKey := account.GetOpenAIApiKey()
apiKey := strings.TrimSpace(account.GetOpenAIProtocolAPIKey())
if apiKey == "" {
return "", "", errors.New("api_key not found in credentials")
}
@@ -739,27 +739,6 @@ func (s *SettingService) GetRateLimit429CooldownSettings(ctx context.Context) (*
if settings.CooldownSeconds > 7200 {
settings.CooldownSeconds = 7200
}
if settings.Strategy != "same_account_retry" {
settings.Strategy = "cooldown"
}
if settings.RetryIntervalMs < 100 {
settings.RetryIntervalMs = 500
}
if settings.RetryIntervalMs > 60000 {
settings.RetryIntervalMs = 60000
}
if settings.RetryMaxDurationSeconds < 1 {
settings.RetryMaxDurationSeconds = 120
}
if settings.RetryMaxDurationSeconds > 600 {
settings.RetryMaxDurationSeconds = 600
}
if settings.MaxAccountSwitches < 0 {
settings.MaxAccountSwitches = 0
}
if settings.MaxAccountSwitches > 10 {
settings.MaxAccountSwitches = 10
}
return &settings, nil
}
@@ -769,10 +748,6 @@ func (s *SettingService) SetRateLimit429CooldownSettings(ctx context.Context, se
if settings == nil {
return fmt.Errorf("settings cannot be nil")
}
if settings.Strategy == "" { settings.Strategy = "cooldown" }
if settings.RetryIntervalMs == 0 { settings.RetryIntervalMs = 500 }
if settings.RetryMaxDurationSeconds == 0 { settings.RetryMaxDurationSeconds = 120 }
if settings.MaxAccountSwitches < 0 { settings.MaxAccountSwitches = 0 }
if settings.CooldownSeconds < 1 || settings.CooldownSeconds > 7200 {
if settings.Enabled {
@@ -780,18 +755,6 @@ func (s *SettingService) SetRateLimit429CooldownSettings(ctx context.Context, se
}
settings.CooldownSeconds = 5
}
if settings.Strategy != "cooldown" && settings.Strategy != "same_account_retry" {
return fmt.Errorf("strategy must be cooldown or same_account_retry")
}
if settings.RetryIntervalMs < 100 || settings.RetryIntervalMs > 60000 {
return fmt.Errorf("retry_interval_ms must be between 100-60000")
}
if settings.RetryMaxDurationSeconds < 1 || settings.RetryMaxDurationSeconds > 600 {
return fmt.Errorf("retry_max_duration_seconds must be between 1-600")
}
if settings.MaxAccountSwitches < 0 || settings.MaxAccountSwitches > 10 {
return fmt.Errorf("max_account_switches must be between 0-10")
}
data, err := json.Marshal(settings)
if err != nil {
+18 -60
View File
@@ -154,8 +154,6 @@ type SystemSettings struct {
SiteSubtitle string
APIBaseURL string
ContactInfo string
SupportQRCodes string
DownloadToolsURL string
DocURL string
HomeContent string
CompactHomeEnabled bool
@@ -167,36 +165,19 @@ type SystemSettings struct {
CustomMenuItems string // JSON array of custom menu items
CustomEndpoints string // JSON array of custom endpoints
DefaultConcurrency int
DefaultBalance float64
RiskControlEnabled bool
CyberSessionBlockEnabled bool
CyberSessionBlockTTLSeconds int
AffiliateEnabled bool
AffiliateRebateRate float64
AffiliateRebateFreezeHours int
AffiliateRebateDurationDays int
AffiliateRebatePerInviteeCap float64
AffiliateRebateCap float64
AffiliateRebateInviteeLimit int
AffiliateSignupBonus float64
AdminRechargeRebateEnabled bool
TicketEnabled bool
KiroDefaultVersion string
KiroDefaultCommit string
KiroDefaultSystemVersion string
KiroDefaultNodeVersion string
KiroCacheHitRateScale int
KiroCacheMinBlockTokens int
KiroCacheIndependentTTLSeconds int
KiroCachePrefixTTLSeconds int
KiroCodeExecutionSandboxCommand string
IPMultiAccountBanEnabled bool
IPMultiAccountBanWindowMinutes int
IPMultiAccountBanThreshold int
IPMultiAccountBanLearningUntil string
DefaultUserRPMLimit int
DefaultSubscriptions []DefaultSubscriptionSetting
DefaultConcurrency int
DefaultBalance float64
RiskControlEnabled bool
CyberSessionBlockEnabled bool
CyberSessionBlockTTLSeconds int
AffiliateEnabled bool
AffiliateRebateRate float64
AffiliateRebateFreezeHours int
AffiliateRebateDurationDays int
AffiliateRebatePerInviteeCap float64
AdminRechargeRebateEnabled bool
DefaultUserRPMLimit int
DefaultSubscriptions []DefaultSubscriptionSetting
// Model fallback configuration
EnableModelFallback bool `json:"enable_model_fallback"`
@@ -220,6 +201,7 @@ type SystemSettings struct {
ChannelMonitorMode string `json:"channel_monitor_mode"`
ChannelMonitorDefaultIntervalSeconds int `json:"channel_monitor_default_interval_seconds"`
ChannelMonitorHideThroughput bool `json:"channel_monitor_hide_throughput"`
ChannelMonitorShowQuota bool `json:"channel_monitor_show_quota"`
// Grok model mapping policy (admin settings; empty mapping falls back to these).
GrokDefaultTextModel string `json:"grok_default_text_model"`
@@ -331,18 +313,6 @@ type DefaultSubscriptionSetting struct {
ValidityDays int `json:"validity_days"`
}
type DefaultAccountModelConfig struct {
ModelWhitelist []string `json:"model_whitelist,omitempty"`
ModelMapping map[string]string `json:"model_mapping,omitempty"`
CompactModelMapping map[string]string `json:"compact_model_mapping,omitempty"`
KiroSubscriptionTypeModelMap map[string]DefaultAccountModelConfig `json:"kiro_subscription_type_model_config,omitempty"`
TempUnschedulableEnabled bool `json:"temp_unschedulable_enabled,omitempty"`
TempUnschedulableRules []TempUnschedulableRule `json:"temp_unschedulable_rules,omitempty"`
CustomErrorCodesEnabled bool `json:"custom_error_codes_enabled,omitempty"`
CustomErrorCodes []int `json:"custom_error_codes,omitempty"`
}
type PublicSettings struct {
RegistrationEnabled bool
EmailVerifyEnabled bool
@@ -373,8 +343,6 @@ type PublicSettings struct {
SiteSubtitle string
APIBaseURL string
ContactInfo string
SupportQRCodes string
DownloadToolsURL string
DocURL string
HomeContent string
CompactHomeEnabled bool
@@ -411,6 +379,7 @@ type PublicSettings struct {
ChannelMonitorMode string `json:"channel_monitor_mode"`
ChannelMonitorDefaultIntervalSeconds int `json:"channel_monitor_default_interval_seconds"`
ChannelMonitorHideThroughput bool `json:"channel_monitor_hide_throughput"`
ChannelMonitorShowQuota bool `json:"channel_monitor_show_quota"`
// Grok model mapping policy (admin settings).
GrokDefaultTextModel string `json:"grok_default_text_model"`
@@ -427,9 +396,6 @@ type PublicSettings struct {
// Affiliate (邀请返利) feature toggle
AffiliateEnabled bool `json:"affiliate_enabled"`
// Ticket feature toggle (default enabled)
TicketEnabled bool `json:"ticket_enabled"`
// 风控中心功能开关
RiskControlEnabled bool `json:"risk_control_enabled"`
@@ -594,11 +560,7 @@ type RateLimit429CooldownSettings struct {
// Enabled 是否在无法解析上游重置时间时应用默认429回避
Enabled bool `json:"enabled"`
// CooldownSeconds 默认回避时长(秒)
CooldownSeconds int `json:"cooldown_seconds"`
Strategy string `json:"strategy"`
RetryIntervalMs int `json:"retry_interval_ms"`
RetryMaxDurationSeconds int `json:"retry_max_duration_seconds"`
MaxAccountSwitches int `json:"max_account_switches"`
CooldownSeconds int `json:"cooldown_seconds"`
}
// DefaultOverloadCooldownSettings 返回默认的过载冷却配置(启用,10分钟)
@@ -612,12 +574,8 @@ func DefaultOverloadCooldownSettings() *OverloadCooldownSettings {
// DefaultRateLimit429CooldownSettings 返回默认的429回避配置(启用,5秒)
func DefaultRateLimit429CooldownSettings() *RateLimit429CooldownSettings {
return &RateLimit429CooldownSettings{
Enabled: true,
CooldownSeconds: 5,
Strategy: "cooldown",
RetryIntervalMs: 500,
RetryMaxDurationSeconds: 120,
MaxAccountSwitches: 2,
Enabled: true,
CooldownSeconds: 5,
}
}
-4
View File
@@ -1289,10 +1289,6 @@ export async function updateOverloadCooldownSettings(
export interface RateLimit429CooldownSettings {
enabled: boolean;
cooldown_seconds: number;
strategy: "cooldown" | "same_account_retry";
retry_interval_ms: number;
retry_max_duration_seconds: number;
max_account_switches: number;
}
export async function getRateLimit429CooldownSettings(): Promise<RateLimit429CooldownSettings> {
+1 -23
View File
@@ -340,16 +340,8 @@
<Toggle v-model="rateLimit429CooldownForm.enabled" />
</div>
<div class="border-t border-gray-100 pt-4 dark:border-dark-700">
<label class="mb-2 block text-sm font-medium text-gray-700 dark:text-gray-300">处理策略</label>
<select v-model="rateLimit429CooldownForm.strategy" class="input w-64">
<option value="cooldown">账号冷却后恢复</option>
<option value="same_account_retry">原号重试后切号</option>
</select>
</div>
<div
v-if="rateLimit429CooldownForm.enabled && rateLimit429CooldownForm.strategy === 'cooldown'"
v-if="rateLimit429CooldownForm.enabled"
class="space-y-4 border-t border-gray-100 pt-4 dark:border-dark-700"
>
<div>
@@ -379,12 +371,6 @@
</div>
</div>
<div v-if="rateLimit429CooldownForm.strategy === 'same_account_retry'" class="grid grid-cols-1 gap-4 border-t border-gray-100 pt-4 dark:border-dark-700 md:grid-cols-3">
<label class="text-sm">重试间隔(ms)<input v-model.number="rateLimit429CooldownForm.retry_interval_ms" type="number" min="100" max="60000" class="input mt-2 w-full" /></label>
<label class="text-sm">单号最大重试时长(秒)<input v-model.number="rateLimit429CooldownForm.retry_max_duration_seconds" type="number" min="1" max="600" class="input mt-2 w-full" /></label>
<label class="text-sm">最多切号次数<input v-model.number="rateLimit429CooldownForm.max_account_switches" type="number" min="0" max="10" class="input mt-2 w-full" /></label>
</div>
<div
class="flex justify-end border-t border-gray-100 pt-4 dark:border-dark-700"
>
@@ -8951,10 +8937,6 @@ const rateLimit429CooldownSaving = ref(false);
const rateLimit429CooldownForm = reactive({
enabled: true,
cooldown_seconds: 5,
strategy: "cooldown" as "cooldown" | "same_account_retry",
retry_interval_ms: 500,
retry_max_duration_seconds: 120,
max_account_switches: 2,
});
// Panel API Rate Limit 状态
@@ -11855,10 +11837,6 @@ async function saveRateLimit429CooldownSettings() {
const updated = await adminAPI.settings.updateRateLimit429CooldownSettings({
enabled: rateLimit429CooldownForm.enabled,
cooldown_seconds: rateLimit429CooldownForm.cooldown_seconds,
strategy: rateLimit429CooldownForm.strategy,
retry_interval_ms: rateLimit429CooldownForm.retry_interval_ms,
retry_max_duration_seconds: rateLimit429CooldownForm.retry_max_duration_seconds,
max_account_switches: rateLimit429CooldownForm.max_account_switches,
});
Object.assign(rateLimit429CooldownForm, updated);
appStore.showSuccess(t("admin.settings.rateLimit429Cooldown.saved"));