gcli2api / internal /server /handler.go
a3216's picture
chore: 同步到上游 1.12.0-panel + 凭证同步/独立启动器/保活
6d60378 verified
Raw History Blame Contribute Delete
73.4 kB
// Package server 暴露 OpenAI 兼容 HTTP 接口,内部驱动 pool 挑号 + upstream 转发。
package server
import (
"context"
"encoding/json"
"errors"
"fmt"
"io"
"log"
"net"
"net/http"
"os"
"strings"
"sync"
"time"
"github.com/linguo2625469/workbuddy2api-panel/internal/auth"
"github.com/linguo2625469/workbuddy2api-panel/internal/httpauth"
"github.com/linguo2625469/workbuddy2api-panel/internal/livecfg"
"github.com/linguo2625469/workbuddy2api-panel/internal/logfmt"
"github.com/linguo2625469/workbuddy2api-panel/internal/pool"
"github.com/linguo2625469/workbuddy2api-panel/internal/prompt"
"github.com/linguo2625469/workbuddy2api-panel/internal/reqlog"
"github.com/linguo2625469/workbuddy2api-panel/internal/session"
"github.com/linguo2625469/workbuddy2api-panel/internal/upstream"
"github.com/linguo2625469/workbuddy2api-panel/internal/usage"
)
// Config handler 依赖。
type Config struct {
Pool *pool.Pool
Upstream *upstream.Client
APIKey string // 空 = 不鉴权(静态值;与 Live 同时给出时 Live 优先)
MaxRotate int // 单请求最多换号次数,默认 3
// Session 会话粘性路由器(可选;nil = 关闭粘性,纯 Pick 轮换)。
Session *session.Router
// StickyCount 返回当前粘性会话绑定数(供 /status);nil 时报告 0。
StickyCount func() int
// RedisMode 观测字段("upstash" / "noop"),供 /status 透出。
RedisMode string
SoftCooldown time.Duration // 429/限流文案软冷却基数,默认 600s(连续触发指数退避,封顶 soft_rate_max)
RefreshSkew time.Duration // token 提前刷新窗口,默认 10m
// Panel 管理面板 handler(可选;nil = 不挂载)。挂载在 /panel/ 前缀下,
// 面板自带 Bearer 鉴权(同一 api_key)与内嵌静态资源,主路由只做转发。
Panel http.Handler
// Live 运行期可变配置(面板在线改 api_key / soft_rate / 脱敏开关时立即生效)。
// nil 时回退静态字段(测试与裸用场景)。
Live *livecfg.Holder
// PromptMode "custom"(网关用自有提示词替换 system)/ "passthrough"(透传)。
PromptMode string
// PromptText custom 模式下注入的系统提示词文本(来自 config.PromptText)。
PromptText string
// GlobalEnabled global realm 路由开关(config global.enabled,缺省 true)。
// handler 侧第三道闸(与 main 注入 auth 开关、upstream.GlobalEnabled 呼应):
// false(显式逃生门)时即便 auth realm=global 也不提供 global: 模型名
// (modelList 不列 global 名单)。
GlobalEnabled bool
// Usage 逐请求用量记录器(可选;nil = 不记录)。
// 在 recordAttempt 这一唯一汇聚点调用,因此流式/非流式、成功/失败都会计入,
// 且与 pool 的每账号累计器同源,两条口径不会漂移。
Usage *usage.Recorder
// RequestLog 请求指标与脱敏 JSONL 归档(可选;nil = 不记录)。
RequestLog *reqlog.Recorder
// RecordClientInfo 是否在请求日志里记录调用来源(客户端 IP / User-Agent)。
// 来自 logging.request_client_info(缺省 true);关闭时 reqlog 事件的来源字段
// 保持为空,归档与面板都不出现来源信息。
RecordClientInfo bool
}
// loadLive 返回当前运行期快照;Live 为 nil 时用静态字段合成。
func (h *Handler) loadLive() livecfg.Snapshot {
if h.cfg.Live != nil {
return h.cfg.Live.Load()
}
return livecfg.Snapshot{
APIKey: h.cfg.APIKey,
SoftCooldown: h.cfg.SoftCooldown,
RecordClientInfo: h.cfg.RecordClientInfo,
}
}
// softCooldown 返回当前生效的软冷却基数(热改优先,<=0 回退默认)。
func (h *Handler) softCooldown() time.Duration {
if d := h.loadLive().SoftCooldown; d > 0 {
return d
}
if h.cfg.SoftCooldown > 0 {
return h.cfg.SoftCooldown
}
return 600 * time.Second
}
// notFoundCooldown 上游 404 的固定短冷却时长。
// 与 SoftCooldown 分流的原因:404 是上游**偶发**路径缺失,不是"本账号在限流",
// 若共用 soft_rate(600s 起 + 指数升级),一次偶发 404 会把好账号罚 10 分钟并逐次加倍。
// 故固定 60s 防雪崩即可,不随 soft_rate 配置、也不参与软退避指数。
const notFoundCooldown = 60 * time.Second
// ServiceName 网关身份标识。经 /healthz 响应体 service 字段与 X-Service 头同时透出:
// 宿主(如 workbuddy-switch 托管网关子进程)探测同端口的旧服务/其他服务时,对方即使
// 返回 2xx 也不带本标识,宿主据此可识别"假成功"。
const ServiceName = "workbuddy2api"
// Handler 主路由。
type Handler struct {
cfg Config
mux *http.ServeMux
degrade degradeGate
// wafIP WAF IP 级拦截状态机(fail-fast,wafip.go):短窗多号 WAF 403 →
// 激活期轮转遇 WAF 403 直接终止(不放大请求量)。进程内状态、重启清零。
wafIP wafIPGate
}
// NewHandler 构建 handler。
func NewHandler(cfg Config) *Handler {
if cfg.MaxRotate <= 0 {
cfg.MaxRotate = 3
}
if cfg.SoftCooldown <= 0 {
cfg.SoftCooldown = 600 * time.Second // 软限流基数(连续触发按指数退避放大)
}
if cfg.RefreshSkew <= 0 {
cfg.RefreshSkew = 10 * time.Minute
}
if cfg.PromptMode == "" {
cfg.PromptMode = "custom" // 缺省 custom:网关自有提示词
}
h := &Handler{cfg: cfg, mux: http.NewServeMux()}
h.mux.HandleFunc("POST /v1/chat/completions", h.withAuth(h.chatCompletions))
h.mux.HandleFunc("GET /v1/models", h.withAuth(h.models))
h.mux.HandleFunc("GET /status", h.withAuth(h.status))
h.mux.HandleFunc("GET /healthz", h.healthz)
if cfg.Panel != nil {
h.mux.Handle("/panel/", cfg.Panel) // /panel → /panel/ 由 ServeMux 自动重定向
}
return h
}
func (h *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
if h.cfg.RequestLog != nil && r.Method == http.MethodPost && r.URL.Path == "/v1/chat/completions" {
trace := &requestTrace{id: reqlog.NewRequestID(), start: time.Now()}
if h.loadLive().RecordClientInfo {
trace.captureClientInfo(r)
}
r = r.WithContext(context.WithValue(r.Context(), requestTraceKey{}, trace))
obs := &responseObserver{ResponseWriter: w}
w.Header().Set("X-Request-Id", trace.id)
h.cfg.RequestLog.Begin()
defer func() {
status := obs.status
if status == 0 {
status = http.StatusOK
}
h.cfg.RequestLog.Record(trace.event(status))
}()
h.mux.ServeHTTP(obs, r)
return
}
h.mux.ServeHTTP(w, r)
}
func (h *Handler) withAuth(next http.HandlerFunc) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
if !httpauth.VerifyBearer(r, h.loadLive().APIKey) {
writeOpenAIError(w, http.StatusUnauthorized, "invalid_api_key", "missing or invalid API key")
return
}
next(w, r)
}
}
func (h *Handler) healthz(w http.ResponseWriter, r *http.Request) {
total, healthy, _, _, _ := h.cfg.Pool.CountsDetailed()
// 用 ServableNow 判定:healthy>0 但全占满在途时 chat 会 503,探活必须同口径,
// 否则负载均衡器会把流量持续打进无法受理的实例。
status := http.StatusOK
if !h.cfg.Pool.ServableNow() {
status = http.StatusServiceUnavailable
}
// realm_servable 域可服务维度:不改判活语义(存在性探活保持不变),
// 只新增 CN/global 各自可达性供双域部署运维观察(任一域不可用单独告警)。
realmServable := map[string]bool{
"cn": h.cfg.Pool.ServableForRealm("cn"),
"global": h.cfg.Pool.ServableForRealm("global"),
}
// 恒无鉴权(负载均衡/编排探活只需 2xx/503 语义),身份靠 service 字段 + X-Service 头双保险。
w.Header().Set("X-Service", ServiceName)
writeJSON(w, status, map[string]any{
"healthy": healthy,
"total": total,
"service": ServiceName,
"realm_servable": realmServable,
})
}
func (h *Handler) status(w http.ResponseWriter, r *http.Request) {
total, healthy, cooling, disabled, inFlightFull := h.cfg.Pool.CountsDetailed()
sticky := 0
if h.cfg.StickyCount != nil {
sticky = h.cfg.StickyCount()
}
redisMode := h.cfg.RedisMode
if redisMode == "" {
redisMode = "noop"
}
// cost_explore 探索台账(issue #136 §5 可观测性):累计探索事件数 + 各
// (域, 模型) 的最近探索时刻(键 "realm|model")。与 accounts[].model_costs
// 行对照即可读出「探索→毕业」全链路(单一事实来源,不做双表示)。零回归只增键。
exploreEvents, exploreLast := h.cfg.Pool.CostExploreStatus()
writeJSON(w, http.StatusOK, map[string]any{
"accounts": h.cfg.Pool.List(),
"total": total,
"healthy": healthy,
"cooling": cooling,
"disabled": disabled,
"in_flight_full": inFlightFull,
// realm_totals 按域分组的计数汇总(双 realm 并存时运维一眼看到各域可用性):
// 只新增字段,既有 total/healthy/cooling/disabled/in_flight_full 汇总键不变(零回归)。
"realm_totals": map[string]map[string]int{
"cn": countsMapFrom(h.cfg.Pool.CountsDetailedForRealm("cn")),
"global": countsMapFrom(h.cfg.Pool.CountsDetailedForRealm("global")),
},
"sticky_sessions": sticky,
"redis_mode": redisMode,
// model_locks 当前有未过期模型级限流的 (域, 模型) 全清单:账号池视图回答
// 「哪些号不能用」,本键回答「哪些模型不能用、锁了几个号、还要锁多久」。
// 与 ModelBlocked(请求失败时的单模型判定)互补;无锁时为 null。零回归只增键。
"model_locks": h.cfg.Pool.ModelLockView(),
// credit_floor 生效的积分保底值(0 = 关闭)。与 accounts[].credits +
// model_costs 对照即可判定「某号为何对某模型不出票」。零值也显式写出
// (运维口径:缺失会让人误以为没记录)。
"credit_floor": h.cfg.Pool.CreditFloor(),
// cost_explore 事件与 per-model 时间戳(时间值由 encoding/json 写 RFC3339)。
"cost_explore": map[string]any{
"events_total": exploreEvents,
"per_model": exploreLast,
},
})
}
// countsMapFrom 把 CountsDetailed 五元组打包成 /status 的域分组建模。
func countsMapFrom(total, healthy, cooling, disabled, inFlightFull int) map[string]int {
return map[string]int{
"total": total,
"healthy": healthy,
"cooling": cooling,
"disabled": disabled,
"in_flight_full": inFlightFull,
}
}
// dynamicModelsCache 动态模型缓存。
var dynamicModelsCache struct {
sync.RWMutex
ids []upstream.ModelInfo
fetched time.Time // 最近一次成功拉取时间
lastFail time.Time // 最近一次拉取失败时间(负缓存)
}
const (
// dynamicModelsTTL 模型目录缓存时长。曾是 1h;缩到 10min 对齐「面板实时、
// API 缓存」的漂移痛点(PR #38 报告):目录新增模型时面板立即可见,公开
// /v1/models 最多滞后一个 TTL。再短就不值得——每次失效都是 2 次上游探测。
dynamicModelsTTL = 10 * time.Minute
modelsFetchFailCooldown = 5 * time.Minute
)
// models 返回模型列表:纯动态(缓存 10min),失败/无号返回空列表(无静态兜底——
// 拉不出目录即意味着上游不可用,假名单只会让客户端选到 11102 的模型)。
func (h *Handler) models(w http.ResponseWriter, r *http.Request) {
writeJSON(w, http.StatusOK, map[string]any{
"object": "list",
"data": h.modelList(),
})
}
// fmtCreditsPrefix 从上游 credits 原文提取倍率并格式化为 "[x0.05 credit]"。
// 上游格式不统一:"x0.05 credits" / "x0.29" / "x0.00 credits" 等,
// 统一提取 x数字 部分,去 "credits" 后缀。
func fmtCreditsPrefix(raw string) string {
s := strings.TrimSpace(raw)
s = strings.TrimSuffix(s, "credits")
s = strings.TrimSpace(s)
if s == "" {
return ""
}
return "[" + s + " credit]"
}
// applyModelInfoFields 把上游模型对象全字段(ModelInfo)按「空值省略」写出规则
// 合入 /v1/models 条目:name/description/credits/tags/vendor/能力旗标/
// max_allowed_size/reasoning_effort/reasoning_summary。CN 动态分支与 global
// 探测命中分支共用(两域模型对象同构),保证输出字段集一致。
// 不覆盖 id/object/created/owned_by 及调用方先前写好的基础字段;上游未下发的
// 字段(零值)整体省略——不编造。
func applyModelInfoFields(entry map[string]any, mi upstream.ModelInfo) map[string]any {
if mi.Name != "" {
entry["name"] = mi.Name
}
if mi.Description != "" {
// 积分倍率前缀:从 "x0.05 credits" / "x0.29" 等格式提取纯数字,
// 统一为 "[x0.05 credit]" 前缀拼入 description,方便下游面板直接展示。
if mi.Credits != "" {
entry["description"] = fmtCreditsPrefix(mi.Credits) + " " + mi.Description
} else {
entry["description"] = mi.Description // descriptionZh 中文描述
}
}
if mi.Credits != "" {
entry["credits"] = mi.Credits // 积分倍率原文(如 "x0.05"),仅展示
}
if len(mi.Tags) > 0 {
entry["tags"] = mi.Tags
}
if mi.Vendor != "" {
entry["vendor"] = mi.Vendor
}
if mi.IsDefault {
entry["is_default"] = true
}
if mi.SupportsImages {
entry["supports_images"] = true // 多模态能力透出
}
if mi.SupportsReasoning {
entry["supports_reasoning"] = true
if mi.CanDisableThinking {
entry["can_disable_thinking"] = true
}
}
if mi.SupportsToolCall {
entry["supports_tool_call"] = true
}
if mi.OnlyReasoning {
entry["only_reasoning"] = true
}
if mi.MaxAllowedSize > 0 {
entry["max_allowed_size"] = mi.MaxAllowedSize
}
if mi.ReasoningEffort != "" {
entry["reasoning_effort"] = mi.ReasoningEffort
}
if mi.ReasoningSummary != "" {
entry["reasoning_summary"] = mi.ReasoningSummary
}
return entry
}
// modelList 模型列表:CN 模型输出统一加 "cn:" 前缀(gateway 路由协议,与 resolveModel
// 对称);global.enabled=true 时追加 global: 前缀的国际版名单。
// 纯动态:动态拉取失败/无号 → 该域空列表,无静态兜底。
func (h *Handler) modelList() []map[string]any {
out := make([]map[string]any, 0)
for _, mi := range h.fetchDynamicModels() {
entry := map[string]any{
"id": "cn:" + mi.ID,
"object": "model",
"created": 1753600000,
"owned_by": "workbuddy",
}
// context_length / max_output_tokens 四级查找(upstream.model_catalog):
// 上游动态值(maxInputTokens/maxOutputTokens)权威 → 静态种子表 →
// model.json 本地缓存 → models.dev 按需拉取(异步不阻塞本次响应,拉到后
// 写 model.json 供下次命中)→ 1M 兜底 / max_output_tokens 省略。
// 上游零值不再透出假 131072(误导 Codex/ZCode 等按 context_length 提前
// 截断、白白丢上下文)。
entry["context_length"] = upstream.ContextWindowListingV4(mi.ID, mi.ContextWindow, h.cfg.Upstream.HTTP)
if mo, ok := upstream.MaxOutputTokensListingV4(mi.ID, mi.MaxTokens, h.cfg.Upstream.HTTP); ok {
entry["max_output_tokens"] = mo
}
// 上游模型对象全字段透出(name/描述/标签/倍率/能力旗标等,空值省略)。
entry = applyModelInfoFields(entry, mi)
// effort 能力透出——远端 supportedEfforts 权威,缺失落到 CN 静态兜底表
// (客户端可发现档位,不再盲传)。无档位 → 省略字段。
if efforts, def := upstream.EffortListing("cn", mi.ID, mi.Efforts, mi.DefaultEffort); efforts != nil {
entry["reasoning_supported_efforts"] = efforts
if def != "" {
entry["reasoning_default_effort"] = def
}
}
out = append(out, entry)
}
// global 模型名单:仅 GlobalEnabled=true 时列出(逃生门)。
// 名单 = 纯动态探测结果(fetchGlobalModels,失败/无号 → 空)。
if h.cfg.GlobalEnabled {
// global 域 effort 能力三级查找:探测下发桶(权威)→ 静态兜底表 → 省略。
// 先 fetchGlobalModels(内部探测并落 effort 桶),再按 id 取快照。
globalIDs, globalAccount := h.fetchGlobalModels()
// 探测对象形态的全字段条目(与 fetchGlobalModels 共享同一次探测缓存):
// 命中 id 才透出富字段;窄表/失败 → nil,按裸 ID 条目输出(不编造字段)。
// globalAccount 为 nil(无 global 号)时返回 nil,跳过富字段映射。
globalInfos := map[string]upstream.ModelInfo{}
for _, mi := range h.cfg.Upstream.FetchGlobalModelInfos(globalAccount) {
globalInfos[mi.ID] = mi
}
globalEfforts, globalDefaults := h.cfg.Upstream.GlobalEffortSnapshot()
for _, id := range globalIDs {
entry := map[string]any{
"id": "global:" + id,
"object": "model",
"created": 1753600000,
"owned_by": "workbuddy",
}
// context_length / max_output_tokens 四级查找(与 CN 动态分支同口径)。
var remoteCtx, remoteOut int64
if mi, ok := globalInfos[id]; ok {
entry = applyModelInfoFields(entry, mi)
remoteCtx, remoteOut = mi.ContextWindow, mi.MaxTokens
}
entry["context_length"] = upstream.ContextWindowListingV4(id, remoteCtx, h.cfg.Upstream.HTTP)
if mo, ok := upstream.MaxOutputTokensListingV4(id, remoteOut, h.cfg.Upstream.HTTP); ok {
entry["max_output_tokens"] = mo
}
if efforts, def := upstream.EffortListing("global", id, globalEfforts[id], globalDefaults[id]); efforts != nil {
entry["reasoning_supported_efforts"] = efforts
if def != "" {
entry["reasoning_default_effort"] = def
}
}
out = append(out, entry)
}
}
return out
}
// fetchGlobalModels 返回 global 模型名单(纯动态探测结果)及被探测账号。
// 缓存/失败回落封在 upstream.FetchGlobalModels(内部 1h + 5min 负缓存)。
// 本方法只负责"何时探测":池中无 global 账号 → 空名单 + nil 账号(零上游调用)。
// 返回的 acct 供调用方在同一账号上取富 ModelInfo(FetchGlobalModelInfos 与
// FetchGlobalModels 共享缓存,不会触发第二次上游探测)。
// GlobalEnabled=false 时 modelList 已不进入本分支(逃生门在调用方 gate)。
func (h *Handler) fetchGlobalModels() ([]string, *auth.Auth) {
acct := h.cfg.Pool.PickExcludingForRealm(nil, "", "global")
if acct == nil {
return nil, nil
}
return h.cfg.Upstream.FetchGlobalModels(acct), acct
}
// fetchDynamicModels 从第一个可用 CN 账号拉模型列表(含 contextWindow/maxTokens),
// 缓存 10min。
// 选号与 /panel/api/models 完全同口径(AvailableUIDsForRealm("cn") 首个 + AuthByUID),
// 而非 Pool.Pick():Pick 无 realm 过滤,混合池里可能选中 global 号去打 CN 端点,
// 表现为偶发失败/面板与 /v1/models 两套目录(PR #38 报告并给出的选号修复)。
// 缓存 + 5min 负缓存按既有语义**保留**(#38 原案整体删除缓存被拒):公开端点逐请求
// 实时拉取 = 每次 2 个上游探测,客户端周期性刷新模型列表会持续打上游;上游故障时
// 无冷却窗口,客户端重试即放大请求量——负缓存正是为此设计(见 handler_test 吸收
// 上游 9832283 的注释);且 cachedModelsSnapshot(gateway_hint 判定)依赖缓存写入。
func (h *Handler) fetchDynamicModels() []upstream.ModelInfo {
dynamicModelsCache.RLock()
if len(dynamicModelsCache.ids) > 0 && time.Since(dynamicModelsCache.fetched) < dynamicModelsTTL {
out := dynamicModelsCache.ids
dynamicModelsCache.RUnlock()
return out
}
// 失败负缓存:冷却期内不再请求上游。
if !dynamicModelsCache.lastFail.IsZero() && time.Since(dynamicModelsCache.lastFail) < modelsFetchFailCooldown {
dynamicModelsCache.RUnlock()
return nil
}
dynamicModelsCache.RUnlock()
uids := h.cfg.Pool.AvailableUIDsForRealm("cn")
if len(uids) == 0 {
return nil
}
acct := h.cfg.Pool.AuthByUID(uids[0])
if acct == nil {
return nil
}
infos, err := h.cfg.Upstream.FetchModels(acct)
if err != nil || len(infos) == 0 {
// 拉取失败只进负缓存(5min lastFail),不 NoteError:NoteError 喂的是 chat
// 熔断器,models 端点偶发 5xx 跨界惩罚 chat 通道健康的账号;
// models 拉取失败 ≠ 账号 chat 不可用。
dynamicModelsCache.Lock()
dynamicModelsCache.lastFail = time.Now()
dynamicModelsCache.Unlock()
return nil
}
dynamicModelsCache.Lock()
dynamicModelsCache.ids = infos
dynamicModelsCache.fetched = time.Now()
dynamicModelsCache.lastFail = time.Time{} // 成功则清空负缓存
dynamicModelsCache.Unlock()
return infos
}
// cachedModelsSnapshot 只读模型目录缓存(TTL 内快照);缓存冷/空 → nil。
// 不发起任何上游调用(hint 判定用:错误路径加一次 FetchModels 网络调用既拖慢
// 错误响应、又污染上游调用语义)。
func cachedModelsSnapshot() []upstream.ModelInfo {
dynamicModelsCache.RLock()
defer dynamicModelsCache.RUnlock()
if len(dynamicModelsCache.ids) == 0 || time.Since(dynamicModelsCache.fetched) >= dynamicModelsTTL {
return nil
}
return dynamicModelsCache.ids
}
func (h *Handler) chatCompletions(w http.ResponseWriter, r *http.Request) {
// 客户端 IP 提取(按请求传递到 ChatStream,不透传时 upstream 侧忽略);
// 消除早年共享字段方案的并发交叉污染(issue:ClientIP 竞态)。
clientIP := upstream.ExtractClientIP(r)
// 请求体无大小上限(max_body_mb 已移除,对齐上游):完整读入,超限类问题交由
// 上游自然返回错误(其响应经既有错误分类链路透出,信息量更大)。#41 的截断
// 防御语义保留在读错误路径——移除预拦截后,截断只可能来自客户端自己断流,
// 读 body 出错就地 400,不把半截 JSON 喂上游 unmarshal 报 unexpected EOF 冤枉罚号。
body, err := io.ReadAll(r.Body)
if err != nil {
writeOpenAIError(w, http.StatusBadRequest, "invalid_request", "read body: "+err.Error())
return
}
var peek struct {
Stream bool `json:"stream"`
Model string `json:"model"`
}
_ = json.Unmarshal(body, &peek)
// realm 前缀解析(D6):model 名可能带 "[realm:]" 前缀。剥出 realm + bareModel,
// bareModel 用于选号/粘性/出站 body 重写(前缀是网关侧路由协议,上游只认裸名)。
// 裸名 → ("cn", 原串),CN 现状零回归。
realm, bareModel := resolveModel(peek.Model)
modelRate := ""
if h.cfg.Upstream != nil {
modelRate = h.cfg.Upstream.ModelRate(realm, bareModel)
}
// 请求级统计:出口即打一行表格日志(任何路径都会走到)。
st := newChatStat(time.Now(), body, peek.Stream)
if tr := requestTraceFrom(r); tr != nil {
tr.stat = st
// 来源在 ServeHTTP 入口采集(此时才知道开关与请求头),此处转交给统计对象,
// 让 stdout 流水行与归档事件共用同一份来源值,两处不会漂移。
st.clientIP, st.userAgent = tr.clientIP, tr.userAgent
}
defer st.done()
tried := map[string]bool{}
var lastErr error
// modelBlock 记录本次是否因「模型级冷却」而选不到号(下方 acct==nil 分支填充)。
// 默认零值 Blocked=false = 按既有的"没有可用账号"口径报错。
var modelBlock pool.ModelBlockStatus
// 会话粘性:从请求体提取会话键并解析绑定号(找不到/无效则 stickyUID 为空,走普通轮换)。
// ExtractKey 与粘性开关解耦(issue #35 侧):关闭粘性时会话头族的聚合主键仍按
// 会话级(RequestIDForKey(sessKey)),不悄悄退化成轮级——提取本身与粘性无关。
sessKey := session.ExtractKey(body)
stickyUID := ""
if h.cfg.Session != nil && sessKey != "" {
// 按模型解析:绑定号在**当前模型**被 6004 限额时视为不可用 → 重新分配,
// 而不是钉在限额号上反复失败("限额后换不动号"的正解)。
if uid, ok := h.cfg.Session.ResolveForModel(sessKey, peek.Model); ok {
stickyUID = uid
}
}
// 轮级聚合键:按 body 里最后一条 user 消息派生(同轮内所有上游调用同键,
// 换 user 消息换键)。#170 起带会话键的客户端也统一走轮级(对齐官方桌面 CLI
// 的 X-Conversation-Request-ID 轮级语义——TraceStartHook 每次 USER_PROMPT_SUBMIT
// 清空重生成),故不再限 sessKey=="" 才计算;sessKey 由下方派生处以复合键方式
// 入键(防不同会话同轮文本互撞)。
// 必须在下方 prompt.Rewrite 之前取——改写会动 messages 内容,之后取会让键漂移。
turnKey := session.TurnKey(body)
// gateway_hint 判定所需的请求形态(image_url part):在改写前取(与 turnKey
// 同理)。11133「模型不支持图片」指向的前提。
reqHasImage := hasImagePart(body)
// 在途租约:成功选中即占名额;函数出口(含成功 return 与 panic)统一释放。
var heldUID string
defer func() {
if heldUID != "" {
h.cfg.Pool.Release(heldUID)
}
}()
releaseHeld := func() {
if heldUID != "" {
h.cfg.Pool.Release(heldUID)
heldUID = ""
}
}
// unbindSticky 解绑当前会话粘性号(stickyUID 非空时)。供「粘性号不可用/被抢」与 fail 共用。
// 幂等:stickyUID 已空则空操作;不会误解绑其他轮的绑定。仅当 Session != nil 时 stickyUID 才会非空。
unbindSticky := func() {
if stickyUID != "" {
h.cfg.Session.Unbind(sessKey)
stickyUID = ""
}
}
// fail 在轮转失败分支统一:释放租约 + 若失败号正是粘性号则解绑(下次请求重新分配)。
fail := func(uid string) {
releaseHeld()
if stickyUID != "" && uid == stickyUID {
unbindSticky()
}
}
// ttfb 首 token 等待(仅流式有观测;非流式传 0 = 无观测,速率不扣减)。
// 显式入参而不是读 st.ttfb:后者在流式分支里是**调用之后**才赋值的,
// 靠顺序传递会让将来重排代码时静默把速率算回旧的错口径。
recordAttempt := func(uid string, delta pool.TokenUsageDelta, credit float64, hasCredit bool, started time.Time, ttfb time.Duration) {
st.attempts++
if delta.HasPromptTokens {
st.promptTokens = delta.PromptTokens
}
if delta.HasCompletionTokens {
st.completionTokens = delta.CompletionTokens
}
if delta.HasTotalTokens {
st.totalTokens = delta.TotalTokens
} else if delta.HasPromptTokens || delta.HasCompletionTokens {
st.totalTokens = st.promptTokens + st.completionTokens
}
delta.Model = bareModel
if delta.Model == "" {
delta.Model = peek.Model
}
latency := time.Since(started)
latencyMs := latency.Milliseconds()
if latencyMs < 1 {
latencyMs = 1
}
delta.HasLatencyMs = true
delta.LatencyMs = latencyMs
// HasCompletionTokens 是「本次有没有 token 观测」的独立判断,与速率怎么算
// 无关,所以留在调用点;速率本身的边界(TTFB 缺失/超界)收在 tokensPerSecond。
if delta.HasCompletionTokens {
if tps, ok := tokensPerSecond(delta.CompletionTokens,
time.Duration(latencyMs)*time.Millisecond, ttfb); ok {
delta.HasTokensPerSecond = true
delta.TokensPerSecond = tps
}
}
h.cfg.Pool.RecordTokenUsage(uid, delta)
// 用量时序记录。ok 以「上游是否给了 usage」判定:空 delta 意味着这次尝试
// 没拿到任何 token 统计(传输错误 / >=400 / 解析失败),计为失败尝试。
// 失败也计入请求数——否则重试放大在「用量」视图里看不见。
if h.cfg.Usage != nil {
realm := "cn"
if a, ok := h.cfg.Pool.Status(uid); ok && a.Realm != "" {
realm = a.Realm
}
h.cfg.Usage.Add(time.Now(), realm, uid, delta.Model, usage.Delta{
PromptTokens: delta.PromptTokens,
HasPromptTokens: delta.HasPromptTokens,
CompletionTokens: delta.CompletionTokens,
HasCompletion: delta.HasCompletionTokens,
TotalTokens: delta.TotalTokens,
HasTotal: delta.HasTotalTokens,
Credit: credit,
HasCredit: hasCredit,
HasCacheTokens: st.hasCache,
CacheHitTokens: st.cacheHit,
CacheMissTokens: st.cacheMiss,
ModelRate: modelRate,
LatencyMs: delta.LatencyMs,
HasLatency: delta.HasLatencyMs,
TokensPerSecond: delta.TokensPerSecond,
HasTPS: delta.HasTokensPerSecond,
}, delta.HasTotalTokens || delta.HasCompletionTokens || delta.HasPromptTokens)
}
}
// 系统提示词改写(出站前、轮转前;每个请求一次)。
// - custom:用自有提示词替换客户端 system/developer(从源头消灭 system 指纹误报)。
// - append:开头连续 system/developer 块后插自有提示词,既有消息逐字不动
// (客户端项目规范/工具约定与网关提示词并用,issue #129)。
// - passthrough + 降级期:换 Degraded 中性提示词直达,不再先撞 400。
// - passthrough / append 非降级期:透传客户端原始 system(append 则再插一条网关 system)。
// 降级裁决:append 在降级期退化为 replace(Rewrite(Degraded))——append 带
// 指纹原文重试是确定性再撞墙,replace 是一次性最小抢救(issue #129 设计 §4)。
degradedApplied := false
if h.cfg.PromptMode == "custom" && h.cfg.PromptText != "" {
body = prompt.Rewrite(body, h.cfg.PromptText)
} else if h.cfg.PromptMode == "append" && h.cfg.PromptText != "" && !h.degrade.Active() {
body = prompt.Append(body, h.cfg.PromptText)
} else if (h.cfg.PromptMode == "passthrough" || h.cfg.PromptMode == "append") && h.degrade.Active() {
body = prompt.Rewrite(body, prompt.Degraded)
degradedApplied = true
}
// outbound model 名重写为 bareModel(D6):realm 前缀是网关侧路由协议,
// 上游不认前缀(global 账号也请求裸模型名)。裸名时 bareModel==peek.Model 恒等。
if bareModel != peek.Model {
body = rewriteModel(body, bareModel)
}
// 会话头族(issue #35):后台按 X-Conversation-Request-ID(对话轮级)聚合请求,
// 官方客户端一次 user send 内所有 tool call/重试/换号复用同一个 ID。此处**轮转
// 循环外**生成一次,循环内每次出站原样复用 → 换号/重试/降级全部同 ID,后台不再
// 碎片化(此前网关一个都不发,上游按 HTTP 请求逐条记账,同一对话几十上百个
// RequestID)。
// - conversationID:body 提取(透传客户端原值,缺省空串——不伪造);
// - conversationRequestID:入站 X-Conversation-Request-ID 透传优先,否则按
// 粘性 key 进程内稳定生成;粘性 key 也空时走轮级兜底(TurnKey/TurnRequestID),
// 无 user 消息时退化成本请求级随机——轮转内捕获一次即共享;
// - messageID 在 ChatHeaders 内每条消息生成(消息级独立,无需外部可见)。
chatMeta := upstream.ChatMeta{ConversationID: session.ResolveConversationID(body)}
if v := r.Header.Get("X-Conversation-Request-ID"); v != "" {
chatMeta.ConversationRequestID = v
} else if turnKey != "" && sessKey != "" {
// 轮级复合键:sessKey 入键防跨会话同轮文本互撞(#170 统一轮级)。
chatMeta.ConversationRequestID = session.TurnRequestID(sessKey + ":" + turnKey)
} else if turnKey != "" {
// 无会话键客户端:纯轮级键(既有兜底语义不变,存量会话键值零漂移)。
chatMeta.ConversationRequestID = session.TurnRequestID(turnKey)
} else if sessKey != "" {
// 残留空态兜底(无 user 消息/无可签名内容):会话级聚合,好于请求级随机。
chatMeta.ConversationRequestID = session.RequestIDForKey(sessKey)
} else {
// 无会话键也无轮级键:请求级随机(轮转内捕获一次即共享)。
chatMeta.ConversationRequestID = session.TurnRequestID("")
}
chatMeta.TraceID = r.Header.Get("X-Trace-ID")
for i := 0; i < h.cfg.MaxRotate; i++ {
// 选号:粘性号优先(PickByUIDForModel 已校验该模型可用性 + 在途未满),否则普通轮换。
var acct *auth.Auth
if stickyUID != "" {
acct = h.cfg.Pool.PickByUIDForModel(stickyUID, bareModel)
if acct == nil || (realm != "" && acct.Realm() != realm) {
// 粘性号在当前模型不可用(冷却/占满/该模型被 6004 限额)或 realm 不符 → 解绑,
// 本次回落普通轮换。
unbindSticky()
acct = nil
}
}
if acct == nil {
// 模型感知 + realm 感知选号:模型非空时启用 6004 模型级冷却豁免
// (healthyForModel),realm 谓词过滤跨域账号。
acct = h.cfg.Pool.PickExcludingForRealm(tried, bareModel, realm)
}
if acct == nil {
st.status = http.StatusServiceUnavailable
// 记下"是不是模型级阻塞"。选号返回 nil 有两种完全不同的成因:
// (a) 池子真的没有可用号(账号级冷却/在途占满/积分保底);
// (b) 号都在,但每个号都对这个模型处于模型级冷却(11102/6004)。
// 两者此前都报 no_healthy_account,导致「模型不可用」被读成「号全挂了」。
// 在这里取一次快照,供下方错误构造区分(issue #102 附带发现 1)。
modelBlock = h.cfg.Pool.ModelBlocked(bareModel)
break
}
st.uid = acct.UID
// 同步昵称:请求流水行只写 uid8 时无法直观看是哪一号,昵称随本次选号带入日志行。
st.nick = acct.Nickname
tried[acct.UID] = true
// 占用在途名额:Pick 已跳过满额账号,此处 CAS 兜底并发抢名额的竞态。
if !h.cfg.Pool.Acquire(acct.UID) {
// 若被抢的正是粘性号,立即解绑并回落普通轮换,避免下一轮仍撞同一个
// 满载粘性号再浪费一次 PickByUID 往返(语义与 fail()/PickByUID-nil 的解绑一致)。
if stickyUID != "" && acct.UID == stickyUID {
unbindSticky()
}
if !rotateBackoff(i, r.Context()) {
// 客户端已断连:换号重试无意义,终止轮转走末端错误透传。
break
}
continue // 最后一个名额被并发抢走 → 换号
}
heldUID = acct.UID
// token 临近过期 → 先 refresh(失败冷却换号)
if acct.NeedsRefresh(h.cfg.RefreshSkew) {
if err := h.cfg.Upstream.RefreshToken(acct); err != nil {
lastErr = err
var ue *upstream.Error
if errors.As(err, &ue) && ue.Kind == upstream.ErrSessionDead {
h.cfg.Pool.Disable(acct.UID, "refresh session dead")
} else {
h.cfg.Pool.NoteError(acct.UID)
}
fail(acct.UID)
if !rotateBackoff(i, r.Context()) {
break // ctx 取消:终止轮转(refresh 失败换号退避)
}
continue
}
if err := acct.SaveAtomic(); err != nil {
// 刷新成功但落盘失败:下次启动会用旧 token,必须暴露
log.Printf("ERR: [server] chat refresh acct=%s: save auth failed: %v", logfmt.Label(acct.UID, acct.Nickname), err)
}
}
// 客户端 IP 按请求传递(PassthroughIP 开启时注入;消除共享字段竞态)。
attemptStarted := time.Now()
rc, status, respBody, terr := h.cfg.Upstream.ChatStreamContext(r.Context(), acct, body, clientIP, chatMeta)
// 分类信封一次成型:upstream 已在错误路径返回 *upstream.Error(Kind +
// Retry-After 头解析)。传输层错误(非 *Error)走抖动换号分支;防御分支
// (terr 为 nil 但 status>=400,如 ErrNone 兜底)回落本地 Classify,双保险。
var uerr *upstream.Error
if errors.As(terr, &uerr) {
status = uerr.Status
}
if uerr == nil && terr != nil {
// 上游超时 / 停滞:**不换号、不罚号**。
//
// 超时不是账号的问题:同一份请求换到别的号,撞上的是同一个慢上游,
// 只会把客户端拖到 MaxRotate × header_timeout(部署值 600s 时最坏
// 约半小时),期间还给一串健康号喂连败计数。此前全仓没有任何超时
// 识别,超时和"网络抖动"共用同一条换号路径。
//
// 判定三态:net.Error.Timeout()(ResponseHeaderTimeout / Client.Timeout)、
// 显式 deadline(DeadlineExceeded / os.ErrDeadlineExceeded)、以及
// **客户端仍在但 ctx 被取消**——那只能是我们自己的空闲看门狗掐的流,
// 也就是上游停滞。客户端主动断连时 r.Context() 已取消,走下面的抖动分支。
if isUpstreamTimeout(terr, r.Context().Err() != nil) {
recordAttempt(acct.UID, pool.TokenUsageDelta{}, 0, false, attemptStarted, 0)
st.status = http.StatusServiceUnavailable
lastErr = fmt.Errorf("%w: %v", errUpstreamTimeout, terr)
log.Printf("WARN: [server] upstream timeout acct=%s: %v (rotation stopped, account not penalized)",
logfmt.Label(acct.UID, acct.Nickname), terr)
break
}
// 网络层抖动:只换号,不喂熔断计数(传输层错误对连续失败连坐熔断过于严苛)。
// 连败兜底(issue #114):喂连败计数——连不上上游是「不知道原因的失败」,
// 连败 N 次临时出池,单次/偶发不罚(NoteFailures 内部达阈才动作)。
// 上游 client 已打 transport error 日志。
recordAttempt(acct.UID, pool.TokenUsageDelta{}, 0, false, attemptStarted, 0)
st.status = http.StatusServiceUnavailable
lastErr = terr
h.cfg.Pool.NoteFailures(acct.UID)
fail(acct.UID)
if !rotateBackoff(i, r.Context()) {
break // ctx 取消:终止轮转(传输层错误换号退避)
}
continue
}
if status >= 400 {
recordAttempt(acct.UID, pool.TokenUsageDelta{}, 0, false, attemptStarted, 0)
st.status = status
var kind upstream.ErrKind
if uerr != nil {
kind = uerr.Kind
} else {
kind = upstream.Classify(status, string(respBody))
uerr = &upstream.Error{Kind: kind, Status: status, Msg: string(respBody)}
}
// 内容拦截误报(passthrough/append 模式首遇):判定为 system 指纹误报,
// 触发降级到次日 00:00 CST,换 Degraded 中性提示词同请求内重试(append
// 降级重试同样退化为 replace——原文在场只会确定性再撞 400)。
// 第二次仍被拦(用户内容本身触发审核)→ 回内容防火墙错误(见下分支)。
// 内容问题非账号问题:applyErrorPolicy 不罚账号(见 ErrContentBlocked 分支)。
if kind == upstream.ErrContentBlocked && (h.cfg.PromptMode == "passthrough" || h.cfg.PromptMode == "append") && !degradedApplied {
h.degrade.Trigger()
body = prompt.Rewrite(body, prompt.Degraded)
degradedApplied = true
delete(tried, acct.UID) // 单账号池也能拿到重试机会(降级重试占一次名额)
releaseHeld()
log.Printf("content-blocked (likely fingerprint false positive) -> degraded prompt retry")
continue
}
if kind == upstream.ErrContentBlocked {
// 内容命中网关内容防火墙:立即回客户端,**不轮转**——换任何账号都会撞同一
// 审核,轮转纯属浪费时间。不罚账号(ErrContentBlocked 分支无冷却/熔断/NoteError)。
// error-passthrough:message 装上游 body 原文(code/msg/requestId 原样),
// 不再改写成网关固定文案——客户端必须看到真实错误才能排查。
h.applyErrorPolicy(acct.UID, kind, string(respBody), bareModel, uerr)
fail(acct.UID)
msg := string(respBody)
if strings.TrimSpace(msg) == "" {
// 空 body 兜底:无上游原文可透传,保留可读分类文案(不编造原文)。
msg = "content blocked by upstream content firewall"
}
writeOpenAIErrorHint(w, http.StatusBadRequest, "content_blocked", msg,
h.hintOf(upstream.ErrContentBlocked, string(respBody), bareModel, reqHasImage, uerr))
st.status = http.StatusBadRequest
st.outcome = reqlog.OutcomeHTTPError
return
}
// 11115「prompt is too long」:立即透传上游原文回客户端,**不罚号不轮转**
// ——上下文超限是请求的问题(同一 body 换任何号都超限,白扔健康号配额;
// 与 WAF IP fail-fast 同哲学:确定与账号无关的错误直接终止轮转)。
// applyErrorPolicy ErrPromptTooLong 分支零动作,fail 只释放租约。
// message 装上游 body 原文(含真实 token 数与上限值——上游原文是最有价值
// 的错误信息,客户端必须看到,禁止固定词覆盖)。
if kind == upstream.ErrPromptTooLong {
h.applyErrorPolicy(acct.UID, kind, string(respBody), bareModel, uerr)
fail(acct.UID)
writeOpenAIErrorHint(w, http.StatusBadRequest, "prompt_too_long", promptTooLongMessage(string(respBody)),
h.hintOf(upstream.ErrPromptTooLong, string(respBody), bareModel, reqHasImage, uerr))
st.status = http.StatusBadRequest
st.outcome = reqlog.OutcomeHTTPError
return
}
// 图片格式/数据无效:立即透传上游原文回客户端,不罚号不轮转。
// 同一 body 换账号仍是同样的解析结果,轮转只会放大无效请求。
if kind == upstream.ErrImageInvalid {
h.applyErrorPolicy(acct.UID, kind, string(respBody), bareModel, uerr)
fail(acct.UID)
msg := string(respBody)
if strings.TrimSpace(msg) == "" {
msg = "image request was rejected by upstream"
}
writeOpenAIErrorHint(w, http.StatusBadRequest, "image_invalid", msg,
h.hintOf(upstream.ErrImageInvalid, string(respBody), bareModel, reqHasImage, uerr))
st.status = http.StatusBadRequest
st.outcome = reqlog.OutcomeHTTPError
return
}
// 请求体解析失败(11101):与 11115 / 图片无效同一哲学——同一 body 换任何
// 账号都是同样的解析结果,轮转只会放大无效请求(每号一次上游调用 +
// rotateBackoff 占用在途名额)。更要紧的是:继续轮转后末端会落到
// 「其余保持 503」,把确定失败的请求伪装成"账号不可用、稍后再试",
// 客户端于是对必然失败的请求无限重试。立即透传上游原文回 400。
if kind == upstream.ErrBadParams {
h.applyErrorPolicy(acct.UID, kind, string(respBody), bareModel, uerr)
fail(acct.UID)
msg := string(respBody)
if strings.TrimSpace(msg) == "" {
msg = "chat request body was rejected by upstream"
}
writeOpenAIErrorHint(w, http.StatusBadRequest, "bad_params", msg,
h.hintOf(upstream.ErrBadParams, string(respBody), bareModel, reqHasImage, uerr))
st.status = http.StatusBadRequest
st.outcome = reqlog.OutcomeHTTPError
return
}
// lastErr 携带完整 body(uerr.Msg 在 upstream 侧截断 200 字符,透传语义
// 要求原文全量)+ Kind/RetryAfter(末端映射与冷却时长共用)。
lastErr = &upstream.Error{Kind: kind, Status: status, Msg: string(respBody), RetryAfter: uerr.RetryAfter}
h.applyErrorPolicy(acct.UID, kind, string(respBody), bareModel, uerr)
fail(acct.UID)
// WAF IP 级 fail-fast(优先于 rotateBackoff 退避——IP 级拦截时退避无意义):
// 该次 WAF 403 喂入 IP 级状态机,若激活(短窗多号命中,IP 被拦而非账号)
// 则立即终止轮转——继续换号只会把请求放大 MaxRotate 倍打同一出口 IP,
// 加重风控。账号级软冷却已在上方 applyErrorPolicy 照常记账。
if kind == upstream.ErrWafBlock && h.wafIP.noteWaf(acct.UID) {
break
}
if !rotateBackoff(i, r.Context()) {
break // ctx 取消:终止轮转(分类错误换号退避)
}
continue
}
// 成功判定与粘性绑定一律**延后到这一跳真正成功之后**(见下方流式/非流式分支):
// 上游「200 已开流 + 一帧 error」是真实形态(6004 限流、内容拦截、审核),
// 此前在读第一帧之前就 NoteSuccess + 清 11102 负缓存 + 绑粘性 → 被限流的号
// 记成健康、粘性把会话钉死在它身上,后续每一轮都打同一个限流号。
if peek.Stream {
// 流式:透传结束后立即关闭上游 body,避免 defer 在轮转场景下堆积 fd。
st.status = http.StatusOK
stats := newChatStatsReaderSince(rc, st.start)
// gateway_hint(SSE):成功状态 200 已开流,中途 error 帧透传时附加
// hint 字段(hintFn 惰性求值——正常流零开销,只有真撞到 error 帧才
// 组装请求上下文做判定)。
// errFrame:上游 error 帧原文(观察者旁路采集),用于流尾的账号处置。
var errFrame string
sErr := upstream.StreamHint(w, stats, upstream.FrameHintFunc(func() upstream.HintContext {
return h.hintContext(bareModel, reqHasImage)
}), upstream.WithErrorFrameObserver(func(payload string) { errFrame = payload }))
switch {
case upstream.IsEmptyStreamError(sErr):
// 上游 200 但空流(0 有效帧):StreamHint 已写 error 帧 + [DONE]
// 兜底(HTTP 头已发出只能 200),但这是上游缺陷不是成功——日志/
// 状态收敛到 502 观测,与非流式 Aggregate 空流→502 upstream_parse
// 同语义(此前 `_ =` 吞错把失败流记成 200,运维看到假成功)。
// 只认 IsEmptyStreamError:客户端断连的写失败不误标(人已走,
// 502 观测没有意义)。
st.status = http.StatusBadGateway
st.outcome = reqlog.OutcomeStreamError
log.Printf("WARN: [server] stream acct=%s model=%s: empty upstream stream (200+0 frames)", logfmt.Label(acct.UID, acct.Nickname), bareModel)
case errFrame != "":
// 上游以 error 帧报错(6004 限流 / 内容拦截 / 审核):按帧内容分类并
// 处置账号——**不记成功、不清 11102 负缓存、不绑粘性**。此前这些动作
// 在流开始前就做了,于是一个正在限流的号被当成健康号,粘性还会把
// 整个会话钉在它身上,后续每轮都失败。
kind := upstream.FrameKind(errFrame)
h.applyErrorPolicy(acct.UID, kind, errFrame, bareModel, nil)
st.status = http.StatusServiceUnavailable
st.outcome = reqlog.OutcomeStreamError
log.Printf("WARN: [server] stream acct=%s model=%s: upstream error frame kind=%s payload=%s",
logfmt.Label(acct.UID, acct.Nickname), bareModel, kind, logfmt.Truncate(errFrame, 200))
case sErr != nil:
// 客户端写失败(断连):上游帧无恙,账号健康——账号侧照常记成功
// (与 default 同语义),请求日志归为 Interrupted(人已走,未完成)。
st.outcome = reqlog.OutcomeInterrupted
h.cfg.Pool.NoteSuccess(acct.UID)
h.cfg.Pool.BlockModelClear(acct.UID, bareModel)
if sessKey != "" && h.cfg.Session != nil {
h.cfg.Session.Bind(sessKey, acct.UID)
}
default:
// 真成功:这一跳读完且上游没有报错,才记成功并让粘性跟上。
st.outcome = reqlog.OutcomeSuccess
h.cfg.Pool.NoteSuccess(acct.UID)
// 11102 负缓存清命:该账号该模型实测成功,立即解除避让(不必等 TTL 到期)。
// BlockModelClear 按 "11102" reason 前缀识别,只清 11102 条目、不碰 6004 独立冷却。
h.cfg.Pool.BlockModelClear(acct.UID, bareModel)
// 粘性跟随最终成功号:本轮成功的账号成为该会话的粘性绑定(覆盖旧绑定)。
// 若 sticky 号失败、轮换到别的号成功,这里把会话重绑到新号,多轮对话下一跳不再随机抽。
if sessKey != "" && h.cfg.Session != nil {
h.cfg.Session.Bind(sessKey, acct.UID)
}
}
credit, hasCredit := stats.Credit()
if hit, miss, ok := stats.CacheTokens(); ok {
st.cacheHit, st.cacheMiss, st.hasCache = hit, miss, true
}
recordAttempt(acct.UID, stats.Usage(), credit, hasCredit, attemptStarted, stats.TTFB())
// WARN 信号放 recordAttempt 之后:st.promptTokens 此时才是本次的观测值。
if st.hasCache {
cacheMissWarn.noteCacheTokens(bareModel, st.promptTokens, st.cacheHit, st.cacheMiss)
}
st.ttfb = stats.TTFB()
// usage 缺失时保留 chatStat.toks 的 -1 哨兵(观测缺失 → 显示 "-"),
// 不写入零值——否则「没观测到 usage」被伪造成「测得 0 token」,
// 与非流式走 completionTokens 返回 -1 的口径不一致。
if toks, hasUsage := stats.Tokens(); hasUsage {
st.toks = toks
}
// 成本账本:末帧 usage 带 credit 与 token 总数时记录实测单价,
// 供下次选号把免费/便宜的号排在前面。
if hasCredit {
st.credit = credit
st.hasCredit = true
if total, tok := stats.TotalTokens(); tok && total > 0 {
h.cfg.Pool.NoteModelCost(acct.UID, bareModel, credit, total)
}
}
rc.Close()
return
}
resp, err := upstream.Aggregate(rc)
rc.Close()
if err != nil {
recordAttempt(acct.UID, pool.TokenUsageDelta{}, 0, false, attemptStarted, 0)
// 上游流解析失败:客户端还没看到任何输出,回 502 并告知原因。
writeOpenAIError(w, http.StatusBadGateway, "upstream_parse", err.Error())
st.status = http.StatusBadGateway
st.outcome = reqlog.OutcomeHTTPError
return
}
credit, total, hasCredit := usageCreditTotal(resp)
if usage, ok := resp["usage"].(map[string]any); ok {
if hit, okH := upstream.UsageCacheHitTokens(usage); okH {
miss, _ := upstream.UsageCacheMissTokens(usage)
st.cacheHit, st.cacheMiss, st.hasCache = int64(hit), int64(miss), true
}
}
recordAttempt(acct.UID, usageDeltaFromResponse(resp), credit, hasCredit, attemptStarted, 0)
if st.hasCache {
cacheMissWarn.noteCacheTokens(bareModel, st.promptTokens, st.cacheHit, st.cacheMiss)
}
writeJSON(w, http.StatusOK, resp)
st.status = http.StatusOK
st.outcome = reqlog.OutcomeSuccess
st.toks = completionTokens(resp)
// 非流式同理:聚合成功(无 error 帧、非空流)才算这一跳成功,事后才记成功/绑粘性。
h.cfg.Pool.NoteSuccess(acct.UID)
h.cfg.Pool.BlockModelClear(acct.UID, bareModel)
if sessKey != "" && h.cfg.Session != nil {
h.cfg.Session.Bind(sessKey, acct.UID)
}
// 成本账本(非流式):从聚合响应的 usage 取 credit 与 token 总数。
if hasCredit {
st.credit = credit
st.hasCredit = true
h.cfg.Pool.NoteModelCost(acct.UID, bareModel, credit, total)
}
return
}
// 末端错误透传(error-passthrough):上游返回的错误原样透传,不再规范化成固定文案。
// 上游返回(*upstream.Error)→ error.message 装**上游 body 原文**(code/msg/
// requestId 原样保留)。HTTP 状态码按 OpenAI 兼容口径映射类别:ErrSoftRate → 429
// (限流语义、客户端应等待重试),其余保持 503。本地调度类错误(无可用账号/
// 传输层抖动/非上游返回的 lastErr)→ 保留自有文案 no_healthy_account(本地错误
// 没有上游原文可透传,不编造)。
status := http.StatusServiceUnavailable
code := "no_healthy_account"
msg := "all accounts are temporarily unavailable, please retry later"
// 上游超时:轮转已在传输层分支止损(见 isUpstreamTimeout),这里给一条**能区分**
// 的文案,别混进"没有可用账号"——两者的排查方向完全不同。
if errors.Is(lastErr, errUpstreamTimeout) {
code = "upstream_timeout"
msg = "upstream timed out: rotation stopped (another account would hit the same slow upstream), please retry later"
}
// gateway_hint(末端透传):上游错误按 Kind + 原文 + 请求形态判定;本地调度类
// 错误(无上游原文)固定 no_healthy_account hint。
hint := upstream.NoHealthyAccountHint()
// upstreamMsgPassed 记录 error.message 是否已被上游原文占据:模型级阻塞分支
// 据此决定要不要覆盖 msg(上游原文优先,含 requestId)。
upstreamMsgPassed := false
var ue *upstream.Error
if errors.As(lastErr, &ue) {
hint = h.hintOf(ue.Kind, ue.Msg, bareModel, reqHasImage, ue)
switch ue.Kind {
case upstream.ErrSoftRate:
status = http.StatusTooManyRequests
code = "rate_limit_exceeded"
msg = "rate limited: all accounts are cooling down, please wait a moment and try again"
case upstream.ErrWafBlock:
if h.wafIP.active() {
// IP 级拦截措辞(fail-fast 终止路径):网关出口 IP 被 WAF 拦截、
// 轮转已止损、窗口过后自动解除。客户端提前重试无意义(换号不换 IP);
// 有上游原文时原文优先(下方统一)。
code = "waf_ip_blocked"
msg = "waf ip-level block: upstream firewall is blocking the gateway IP, rotation stopped; retry after the block window expires"
}
case upstream.ErrModelBlocked:
// 11102「该后端无此模型」:上游原文(下方统一透传)已经写清了原因,但
// code 此前停在默认的 no_healthy_account —— 那是"服务端过载"的语义,
// 客户端据此会不断重试(#81 里就是 503 → 客户端自动重试 5 次),而真正
// 该做的是换个模型。用 400 + 明确的 code 把不可重试的性质讲清楚。
code = "model_unavailable"
status = http.StatusBadRequest
}
if s := strings.TrimSpace(ue.Msg); s != "" {
// 上游原文优先:透传 code/msg/requestId,不拼接本地前缀。
msg = s
upstreamMsgPassed = true
}
}
// 模型级阻塞覆盖上面的通用文案(放在最后 = 优先级最高)。
//
// 为什么能覆盖 code 而不丢信息:modelBlock.Reason 就是 BlockModelBackoff 存下的
// 上游原因,已在 hint 里原样带出,所以覆盖 message 不会丢掉上游信息,反而补上
// 了上游不会告诉客户端的两件事——「有几个号被挡」和「最早什么时候解封」。
//
// 用 400 而非 503:这是"你选的这个模型当前不可用",不是"服务端暂时过载"。
// 换号重试必然同样失败,客户端不该按可重试错误处理。
if modelBlock.Blocked {
code = "model_unavailable"
// 有上游原文时保留它(含 requestId,用户要拿去向上游反馈),只把"有几个号
// 被挡、最早何时解封"这类本地调度信息放进 gateway_hint。
if !upstreamMsgPassed {
msg = "model is unavailable on every account (per-model cooldown), try another model"
}
hint = fmt.Sprintf("model_blocked: %d account(s) cooling down this model", modelBlock.Count)
if !modelBlock.Until.IsZero() {
hint += "; earliest unblock at " + modelBlock.Until.Format(time.RFC3339)
}
if s := strings.TrimSpace(modelBlock.Reason); s != "" {
hint += "; upstream: " + s
}
status = http.StatusBadRequest
}
writeOpenAIErrorHint(w, status, code, msg, hint)
st.status = status
st.outcome = reqlog.OutcomeHTTPError
}
// tokensPerSecond 计算吐字速率(token/s),返回 (速率, 是否有意义)。
//
// 分母用「生成耗时」= 端到端耗时 - 首 token 等待(TTFB),不是端到端耗时。
//
// 为什么要减:不减的话,首 token 等待越长、报告速率被压得越低。同一模型换个
// 上游或网络,TTFB 从 0.3s 涨到 2s,速率能凭空掉一半——读起来像"模型变慢了",
// 其实只是排队久了。issue #34 报的正是这个,维护者也确认「速率计算时并没有减去
// 首token到达时间」。
//
// ttfb<=0 表示没有观测:非流式回复天然没有「首个 data 帧」(日志里那一列记的是
// "-")。此时不猜、不扣——凭空假定一个 TTFB 会把分母推向零、把速率抬成虚高,
// 比不扣更糟。只有真测到才扣。
//
// 同理,ttfb 不小于总耗时时(时钟粒度、或 TTFB 落在计时终点之后)退回端到端耗时,
// 避免零/负分母。token 数为负哨兵值(-1 = 观测缺失)时返回 false。
//
// 用量账本(handler)与控制台流水行(logging.go)都走这一个函数:两处各算一遍时
// 口径漂移过一次(流水行漏扣 TTFB、与面板数字对不上),共用是防再次分叉的唯一办法。
func tokensPerSecond(completionTokens int64, total, ttfb time.Duration) (float64, bool) {
if completionTokens < 0 || total <= 0 {
return 0, false
}
gen := total
if ttfb > 0 {
if g := total - ttfb; g > 0 {
gen = g
}
}
return float64(completionTokens) / gen.Seconds(), true
}
// promptTooLongMessage 11115 透传 message:上游 body 原文(含真实 token 数/
// 上限值/requestId,客户端自行排查);空 body 兜底为可读分类短文案(不编造原文)。
func promptTooLongMessage(body string) string {
if strings.TrimSpace(body) == "" {
return "prompt is too long"
}
return body
}
// usageCreditTotal 从聚合响应取 usage.credit 与 total_tokens(成本台账非流式入口)。
// 任一字段缺失/非法 → ok=false(不记录)。
func usageCreditTotal(resp map[string]any) (credit float64, total int, ok bool) {
usage, _ := resp["usage"].(map[string]any)
if usage == nil {
return 0, 0, false
}
c, _ := usage["credit"].(float64)
t, _ := usage["total_tokens"].(float64)
if t <= 0 {
return 0, 0, false
}
return c, int(t), true
}
// rotateBackoff 轮转间指数退避 + 抖动(WAF 403 修复 P0-2):第 i 次轮转失败
// (continue 换号前)等待 backoffAfter(i)(500ms·2^i 封顶 8s,±25% 抖动),
// ctx 取消(客户端断连/优雅停机)返回 false——调用方立即终止轮转(客户端已走,
// 换号重试无意义)。退避是「换号前歇一下」让上游频控窗口滑过;正常单号请求
// (首次成功)不经过本函数,零开销。
func rotateBackoff(i int, ctx context.Context) bool {
d := backoffAfter(i)
if d <= 0 {
return ctx.Err() == nil
}
if !sleepCtx(ctx, d) {
log.Printf("WARN: [server] rotate backoff aborted: ctx cancelled")
return false
}
return true
}
// errUpstreamTimeout 上游超时的哨兵:末端出口据此给出与「没号可用」可区分的文案。
var errUpstreamTimeout = errors.New("upstream timeout")
// isUpstreamTimeout 判断这一跳的失败是否属于「上游超时 / 停滞」。超时不是账号的
// 问题,换号注定白换(同一份请求撞同一个慢上游),必须止损:不轮转、不罚号。
// 三态判定见传输层错误分支的注释。
func isUpstreamTimeout(err error, clientGone bool) bool {
if err == nil {
return false
}
var ne net.Error
if errors.As(err, &ne) && ne.Timeout() {
return true
}
if errors.Is(err, context.DeadlineExceeded) || errors.Is(err, os.ErrDeadlineExceeded) {
return true
}
return !clientGone && errors.Is(err, context.Canceled)
}
// applyErrorPolicy 按错误分类对账号施加冷却/禁用/熔断策略(最终版状态机)。
// kind 是唯一权威分类(来自 upstream.Classify / ChatStreamContext 的 *Error 信封),
// 此处不再按原始 status 二次判断。仅在 chatCompletions 轮转循环内调用:内容拦截
// 会立即 400 返回,其余种类 continue 换号(continue 前由 rotateBackoff 退避)。
//
// 十条路径,各司其职:
// - ErrHardCredit → CooldownUntilTomorrow4AM:即时硬冷却到次日 04:00(等签到恢复)。
// - ErrSoftRate → 优先对齐上游重置墙钟(带「将在 … 重置」时 6004 走模型级豁免、
// 非 6004 走账号级,均不指数堆加);无重置时间才走有界退避。冷却时长优先采信
// Retry-After 头(uerr.RetryAfter,body 文案墙钟之外的头形态来源)。
// - ErrWafBlock → 账号级软冷却:**不 Disable**——WAF 403 是 IP/指纹维频控信号,
// 罚过即走、到期自愈。时长优先 Retry-After 头;缺失按 wafCooldownBase(60s)
// 起 · softStreak 指数、封顶 soft_rate_max 的既有 CooldownSoftRate 有界退避。
// 基数经 jitterDur 抖动(防多账号同相位冷却到期再聚团)。
// - ErrNotFound → Cooldown(CoolSoft, notFoundCooldown 固定 60s):短冷却防雪崩。
// - ErrSessionDead → Disable:session 死亡,永久禁用(需人工重登)。
// - ErrContentBlocked → 不罚账号;passthrough 首遇触发降级重试,最终仍拦则回 400。
// - ErrBadParams → 不罚账号,且与 ErrPromptTooLong/ErrImageInvalid 同待遇:
// 调用方在轮转循环内即刻 400 透传原文终止(换号必然同样失败)。
// - ErrPromptTooLong → 11115:请求的问题不是账号的问题。零动作(不冷却/不熔断/
// 不 NoteError、不喂连败),chatCompletions 已直接透传原文返回不轮转。
// - ErrImageInvalid → 图片格式/数据无效:请求的问题不是账号的问题(同一 body
// 换任何号都会得到相同的解析错误)。零动作(不冷却/不熔断/不 NoteError、
// 不喂连败),chatCompletions 已直接透传原文返回不轮转。
// - ErrModelBlocked → BlockModelBackoff:(账号, 模型) 11102 负缓存避让。
// - ErrServer → NoteError:喂单一连续失败计数器 fails + 累计错误 errTotal,
// 达到 breakerThreshold 触发熔断(指数退避)。
// - 其他(default:ErrClient/ErrNone)→ 只换号不罚(防雪崩),不喂熔断;ErrClient
// 额外喂连败计数(NoteFailures,issue #114):未知 4xx 连败 N 次临时出池。
//
// body 仅在 ErrSoftRate/ErrAccountFault 分支用于解析重置时间/分野;model 为请求
// 携带的模型名。uerr 是 ChatStreamContext 返回的分类信封(可携带 RetryAfter);
// 零值/防御路径下为 nil,冷却时长回落既有计算。
func (h *Handler) applyErrorPolicy(uid string, kind upstream.ErrKind, body, model string, uerr *upstream.Error) {
switch kind {
case upstream.ErrHardCredit:
// 402 + 余额关键词即积分耗尽:同步冷却到次日 04:00(签到任务 09/21 点恢复),
// 不需要异步核查(冗余)。立即换号。
h.cfg.Pool.CooldownUntilTomorrow4AM(uid, "余额不足")
case upstream.ErrSoftRate:
// 统一对齐上游重置时间:只要 body 带「将在 … 重置」,无论业务 code 是
// 6004 还是 11140 rate-limiting 等形态,都精确冷却到该墙钟、绝不指数堆加。
// - 模型级(6004)→ CooldownSoftForModel:写 modelCooldowns[model],切模型豁免。
// - 账号级(非 6004)→ CooldownSoftRate:写账号级 until,不产生模型豁免。
modelRateLimited := upstream.IsModelRateLimit(body)
if resetAt, ok := upstream.ParseRateReset(body); ok {
if modelRateLimited {
h.cfg.Pool.CooldownSoftForModel(uid, h.softCooldown(), resetAt, model, "6004 model rate limit")
return
}
h.cfg.Pool.CooldownSoftRate(uid, h.softCooldown(), resetAt, "429 rate limit")
return
}
// body 无重置文案但带 Retry-After 头 → 冷却到该时刻(不做指数堆加)。
// 头优先于「有界退避」,但低于 body 重置文案(文案是上游更权威的口径)。
if uerr != nil && uerr.RetryAfter > 0 {
h.cfg.Pool.CooldownSoftRate(uid, h.softCooldown(), time.Now().Add(uerr.RetryAfter), "429 rate limit (retry-after)")
if modelRateLimited {
h.cfg.Pool.RecordModelRateLimitAudit(uid, model, "6004 model rate limit (reset unknown)")
}
return
}
// 无重置时间 → 账号级有界退避(soft_rate 基数起、softStreak 翻倍、封顶
// soft_rate_max;已在冷却中的兜底探测不翻倍)。基数取 h.softCooldown()
// (热改优先),管理面板改 soft_rate 后立即生效。
h.cfg.Pool.CooldownSoftRate(uid, h.softCooldown(), time.Time{}, "429 rate limit")
if modelRateLimited {
h.cfg.Pool.RecordModelRateLimitAudit(uid, model, "6004 model rate limit (reset unknown)")
}
case upstream.ErrWafBlock:
// WAF 403(无业务信封拦截形态)。软冷却复用 CooldownSoftRate 家族:基数
// wafCooldownBase(60s,抖动后落 [45s,75s])、softStreak 指数升级、封顶
// soft_rate_max、冷却中兜底探测不翻倍——全部继承既有语义。
// Retry-After 头优先(WAF 拦截页可能带该头)。不 Disable。
if uerr != nil && uerr.RetryAfter > 0 {
h.cfg.Pool.CooldownSoftRate(uid, jitterDur(wafCooldownBase), time.Now().Add(uerr.RetryAfter), "waf 403 block (retry-after)")
return
}
h.cfg.Pool.CooldownSoftRate(uid, jitterDur(wafCooldownBase), time.Time{}, "waf 403 block")
case upstream.ErrSessionDead:
h.cfg.Pool.Disable(uid, "12153 session dead")
case upstream.ErrNotFound:
// 404 短冷却(软冷却),防雪崩。固定 notFoundCooldown,不随 soft_rate 退避:
// 偶发路径缺失不是限流信号,不该按限流惩罚升级。
h.cfg.Pool.Cooldown(uid, pool.CoolSoft, notFoundCooldown, "upstream 404")
case upstream.ErrAccountFault:
// 账号级授权/配额故障按 msg 分野(口径与 Classify 的 accountFaultMarkers 一致):
// - "request illegal"(code 11140)→ 账号级**授权封禁**:硬禁用(Disable)。
// - 14017(trial not activated)→ register 未完成,补完 register 后可能自愈,
// **保持软冷却**(禁用会让用户补完 register 后仍无法用)。
// 大小写不敏感(与 Classify 的 marker 匹配同口径)。
if strings.Contains(strings.ToLower(body), "request illegal") {
h.cfg.Pool.Disable(uid, "account banned by upstream (11140 request illegal), re-login required")
return
}
h.cfg.Pool.Cooldown(uid, pool.CoolSoft, h.softCooldown(), "account fault (14017)")
case upstream.ErrServer:
// 5xx 上游故障:Classify 已把 ≥500 判为 ErrServer,在此喂熔断计数(不再手写 status>=500)。
h.cfg.Pool.NoteError(uid)
case upstream.ErrContentBlocked:
// 内容策略拦截(误报):内容问题非账号问题,不罚账号(无冷却/熔断/NoteError)。
// passthrough 模式由 chatCompletions 内降级重试处理;custom 模式本不会到此分支。
case upstream.ErrPromptTooLong:
// 11115「prompt is too long」:请求的问题不是账号的问题(同一 body 换任何
// 号都超限)。零动作(不冷却/不熔断/不 NoteError,同 ErrContentBlocked 待遇),
// chatCompletions 已直接透传原文返回不轮转——该分支只为文档完备。
case upstream.ErrImageInvalid:
// 图片格式/数据无效:请求的问题不是账号的问题(同一 body 换任何号都会
// 得到相同解析错误)。零动作,chatCompletions 已 fail-fast 透传。
case upstream.ErrBadParams:
// 请求体解析失败(400 + Unmarshal chat params failed / 11101):发给上游的 body
// 有问题(网关侧不再截断,均为客户端畸形 JSON)。换了账号照样 400,
// 不罚账号(无冷却/熔断/NoteError,同 ErrContentBlocked 待遇);chatCompletions
// 已 fail-fast 400 透传原文、终止轮转——「换号可能有不同模型权限」属 11102
// (ErrModelBlocked)的分类域,与本类无关。
case upstream.ErrModelBlocked:
// 11102「该后端无此模型」:(账号, 模型) 负缓存避让。复用 modelCooldowns 机制
// (与 6004 同域),选号侧 healthyForModel 对该账号自动避开该模型。
// 立即换号(本轮 continue),该账号该模型冷却,下次选号避开。
h.cfg.Pool.BlockModelBackoff(uid, model, upstream.ModelBlockReason)
default:
// 其余(ErrClient/ErrNone):只换号不罚(防雪崩),不喂熔断。
// ErrClient(未知 4xx)喂连败计数(issue #114):连续 N 次该形态失败 →
// 账号临时出池(NoteFailures 达阈降权),单次/偶发不罚(不误伤)。ErrNone
// 到这里属防御路径(status>=400 但分类成功),语义不明不喂。
if kind == upstream.ErrClient {
h.cfg.Pool.NoteFailures(uid)
}
}
}
// ---------------------------------------------------------------------------
// helpers
// ---------------------------------------------------------------------------
func writeJSON(w http.ResponseWriter, status int, v any) {
raw, _ := json.Marshal(v)
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(status)
_, _ = w.Write(raw)
}
func writeOpenAIError(w http.ResponseWriter, status int, code, msg string) {
writeJSON(w, status, map[string]any{
"error": map[string]any{
"message": msg,
"type": "api_error",
"code": code,
},
})
}
// wafCooldownBase WAF 403 软冷却基数(建议 60s 起;抖动 ±25% 后落 [45s,75s],
// 实际进入 CooldownSoftRate 后再按 softStreak 指数、封顶 soft_rate_max)。
// 与 SoftCooldown 分流的原因:WAF 403 是 IP/指纹维频控,信号比 429「账号级限流」轻
// (账号本身健康),但比 404 重(带粘性会连环);60s 级的快速避让已足够让频控窗口
// 滑过。抖动复用 backoff.go jitterDur(单一来源)。
const wafCooldownBase = 60 * time.Second
// writeOpenAIErrorHint 同 writeOpenAIError,另在 error 对象上附加
// error.gateway_hint(hint 为空串时不带字段——未覆盖形态不编造)。
// message 仍是上游原文透传(hint 只做并列补充,绝不替换/包装 message)。
func writeOpenAIErrorHint(w http.ResponseWriter, status int, code, msg, hint string) {
if hint == "" {
writeOpenAIError(w, status, code, msg)
return
}
writeJSON(w, status, map[string]any{
"error": map[string]any{
"message": msg,
"type": "api_error",
"code": code,
"gateway_hint": hint,
},
})
}
// hasImagePart 报告聊天请求体是否携带多模态 image_url part(OpenAI 兼容形态
// messages[].content[] {type:"image_url"})。畸形/其他形态一律 false(hint 侧
// 宁缺勿滥:判不出带图就不给「模型不支持图片」指向)。
func hasImagePart(body []byte) bool {
var peek struct {
Messages []struct {
Content []struct {
Type string `json:"type"`
} `json:"content"`
} `json:"messages"`
}
if json.Unmarshal(body, &peek) != nil {
return false
}
for _, m := range peek.Messages {
for _, p := range m.Content {
if p.Type == "image_url" {
return true
}
}
}
return false
}
// hintContext 组装 chatCompletions 的 gateway_hint 判定上下文:请求裸模型名 +
// 是否带图 + 模型目录 supports_images 声明(目录未收录 → ModelInCatalog=false,
// 不做「不支持」判定,防查不到误判)。仅错误路径调用(成功请求零开销)。
//
// 目录查询只读既有缓存快照(cachedModelsSnapshot),**不触发上游拉取**:错误路径
// 加一次 FetchModels 网络调用既拖慢错误响应、又污染上游调用语义(错误风暴时放大
// 请求量——与 WAF IP fail-fast 的「不放大请求量」哲学相悖)。缓存冷(最近 10min 未
// 拉过)→ ModelInCatalog=false,11133 退中性 hint(宁缺勿滥,不编造能力事实)。
func (h *Handler) hintContext(bareModel string, hasImage bool) upstream.HintContext {
ctx := upstream.HintContext{Model: bareModel, HasImage: hasImage}
if bareModel == "" {
return ctx
}
for _, mi := range cachedModelsSnapshot() {
if mi.ID == bareModel {
ctx.ModelInCatalog = true
ctx.ModelSupportsImages = mi.SupportsImages
return ctx
}
}
return ctx
}
// hintOf 末端错误透传的统一 hint 入口:kind + 上游原文 + 请求上下文 →
// gateway_hint 文案(upstream.GatewayHint 单一事实来源)。uerr 为 nil 时回落
// body 原文判定(防御路径)。transport 层错误(lastErr 非 *upstream.Error 且
// 上游没回 body)→ 无 hint(不编造)。
func (h *Handler) hintOf(kind upstream.ErrKind, body, bareModel string, hasImage bool, uerr *upstream.Error) string {
msg := body
if uerr != nil && uerr.Msg != "" {
msg = uerr.Msg
}
return upstream.GatewayHint(kind, msg, h.hintContext(bareModel, hasImage))
}