Documentation
¶
Overview ¶
Package proxy 是 AI 请求热路径:分组 key 鉴权 → 调度器选号 → SDK 转发 → 用量采集。 规格 §6/§9。不变量:热路径零 DB、零 per-request 锁。
Index ¶
- func AIRouter(p *Proxy) http.Handler
- type Auth
- func (a *Auth) Acquire(meta domain.KeyMeta) (int, bool)
- func (a *Auth) Authenticate(r *http.Request) (domain.KeyMeta, bool)
- func (a *Auth) DeductQuota(keyID, tokens int64)
- func (a *Auth) Delete(raw string)
- func (a *Auth) InFlightUsers() map[int64]int64
- func (a *Auth) QuotaExhausted(meta domain.KeyMeta) bool
- func (a *Auth) Release(meta domain.KeyMeta, level int)
- func (a *Auth) Reload(ctx context.Context) error
- func (a *Auth) RemoveUser(userID int64)
- func (a *Auth) SetInstancesProvider(p InstancesProvider)
- func (a *Auth) Upsert(raw string, meta domain.KeyMeta)
- func (a *Auth) UpsertUser(userID int64, snap domain.UserSnapshot)
- func (a *Auth) UserSnapshot(userID int64) (domain.UserSnapshot, bool)
- type BillingHooks
- type ConcSyncWorker
- type Config
- type InstancesProvider
- type KeyLoader
- type PriceResolver
- type Proxy
- func (p *Proxy) CloseAllWS()
- func (p *Proxy) HandleAnthropic(w http.ResponseWriter, r *http.Request)
- func (p *Proxy) HandleChat(w http.ResponseWriter, r *http.Request)
- func (p *Proxy) HandleImagesEdits(w http.ResponseWriter, r *http.Request)
- func (p *Proxy) HandleImagesGenerations(w http.ResponseWriter, r *http.Request)
- func (p *Proxy) HandleModels(w http.ResponseWriter, r *http.Request)
- func (p *Proxy) HandleResponses(w http.ResponseWriter, r *http.Request)
- func (p *Proxy) HandleResponsesWS(w http.ResponseWriter, r *http.Request)
- func (p *Proxy) HandleSearch(w http.ResponseWriter, r *http.Request)
- func (p *Proxy) Inflight() int64
- func (p *Proxy) SetCodex(c *sdkbridge.Codex)
- func (p *Proxy) SetInstancesProvider(inst InstancesProvider)
- type QuotaUsedReader
- type ToolView
- type UpstreamCaller
- type UserStatusLoader
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
Types ¶
type Auth ¶
type Auth struct {
// contains filtered or unexported fields
}
Auth 鉴权快照:key_raw(明文)→ KeyMeta(含归属用户门禁字段)+ 用户快照表 (status+role)+ 两级并发/额度内存计数(gate)。热路径零 DB、零 per-request 锁(RWMutex 读多写少,规格 §10.3)。用户变更(禁用/降权/并发/额度调整) 走 invalidate 回调 → Reload 全量刷新(评审 I-2),JWT 24h 长时效仅作快照 失效后的最终兜底。
func NewAuth ¶
func NewAuth(loader KeyLoader, users UserStatusLoader, log *logx.Logger) *Auth
NewAuth 构造鉴权快照(空表——首载统一由快照注册表 ReloadAll 承担,单一启动 入口,消灭"构造即载 + 注册表再刷"双重加载冗余;构造到首刷之间无请求流量, 见 main 装配序)。
func (*Auth) Acquire ¶
Acquire 两级并发门禁:user → key 依次 CAS 抢占;key 失败回滚 user 计数 (评审 I-3:防泄漏)。返回已 acquire 层级位掩码(release 仅释放已 acquire 层级)。未设置上限(max=0)或计数器缺失(跨 reload 竞态窗口)→ 该层跳过。
func (*Auth) Authenticate ¶
Authenticate 解析网关 key 并返回 KeyMeta。兼容两种客户端口径: OpenAI 客户端发 Authorization: Bearer;Anthropic 官方 SDK / Claude Code 发 x-api-key 头。两者同时提供时以 Authorization 为准。 key 或归属用户被禁用 → 快照直接拒绝(401,即时失效)。 快照 map key = 明文,等值直查(零哈希);meta 无 key 字符串字段—— 鉴权失败日志天然不落明文。
func (*Auth) DeductQuota ¶
DeductQuota 请求结束扣减(后扣模型;usage 已知;无额度 key 无计数器 → no-op)。
func (*Auth) InFlightUsers ¶
InFlightUsers 门禁在途并发只读快照(/api/admin/users-top 端点用;spec 2026-08-14 P2-3:gateSnapshot.users 未导出,经本访问器只读暴露):gateSnapshot 整体 原子换入换出(reload/upsert 重建,不可变),store.Load() 零锁取当前引用后 遍历 + 原子读各计数器 → map[int64]int64 拷贝(含 0——过滤由调用方做)。 冷面调用(管理端聚合,不涉请求热路径);多实例部署下为本实例在途计数。
func (*Auth) QuotaExhausted ¶
QuotaExhausted 额度检查:本地预算快读(零锁零 DB);预算耗尽触发 DB 复核 认领(#14 §3.2——复核成功续预算继续放行,复核确认真尽才 429)。检查在并发 acquire 之前(评审提醒①:失败无并发槽副作用);未设置额度 key 短路零成本。
func (*Auth) Release ¶
Release 释放并发计数(仅释放 acquire 返回的层级;跨 reload 命中新快照的 继承计数,与 scheduler Release 同语义)。
func (*Auth) Reload ¶
Reload 全量刷新鉴权快照(注册表首刷/周期 auth-sync/用户变更 invalidate): keys 元数据 + 用户状态 + 门禁计数器(在途值跨 reload 继承)。 失败必打 Warn(含调用方是否忽略错误——invalidate 回调等吞错路径):加载 失败若被忽略,快照保持旧值/空表 → 鉴权全部 401 或用旧 key 放行,静默恶化 (IN 超限事故的"运行中静默失败"形态即此类)。
func (*Auth) RemoveUser ¶
RemoveUser 增量移除用户(暂未使用,防御性;与 UpsertUser 对称)。
func (*Auth) SetInstancesProvider ¶
func (a *Auth) SetInstancesProvider(p InstancesProvider)
SetInstancesProvider 注入集群实例数 N 提供者(#14 多实例预算分摊;discovery 装配——main 装配点,spec 2026-08-25-redis-instance-discovery-design §2.2)。 注入即触发预算重算(幂等 reload,在途值继承);此后 N 在每次预算分配现读, 心跳计数变化 ≤1 tick 天然生效。
func (*Auth) UpsertUser ¶
func (a *Auth) UpsertUser(userID int64, snap domain.UserSnapshot)
UpsertUser 增量刷新单个用户状态(本地立即可见,不等去抖窗口)。 供 admin 创建用户 / 注册 / 状态变更后本地实例立即对 RequireJWT 可见, 消除 200ms 窗口内新建用户 401;远端实例仍经 NOTIFY → 全量 Reload 收敛。
func (*Auth) UserSnapshot ¶
func (a *Auth) UserSnapshot(userID int64) (domain.UserSnapshot, bool)
UserSnapshot 用户快照(RequireJWT 状态校验 + adminAuth 快照 role 校验共用; 用户变更走 invalidate → Reload,不用 DB 直查)。单次查找同时取 status+role (热路径零分配)。
type BillingHooks ¶
type BillingHooks struct {
Resolver PriceResolver
Balances *billing.Balances // 余额只读快照(预检 + 扣费后定向刷新)
Flusher *billing.Flusher // 批量扣费落库(billed 路由终点)
TierPolicy func(tier billing.Tier) billing.TierPolicyMode
}
BillingHooks 计费钩子(proxy.New 参数;nil = 计费全关:不查价、不记 BillingTier、不处理 service_tier 转发策略、不做余额预检)。
type ConcSyncWorker ¶
type ConcSyncWorker struct {
// contains filtered or unexported fields
}
ConcSyncWorker 并发门双向同步 worker:实现 worker.Worker。客户端经 main 注入 (pkg/redisx 单构造点纪律,本包不自建连接);gate 引用取自 Auth(同包直访 私有字段,无需导出接口)。nil client = Start no-op(测试/降级形态零专门分支)。
func NewConcSyncWorker ¶
func NewConcSyncWorker(auth *Auth, client *redis.Client, self string, log *logx.Logger) *ConcSyncWorker
NewConcSyncWorker 构造(零副作用:Start 才起 goroutine,Close 未 Start 时安全)。
func (*ConcSyncWorker) Close ¶
func (w *ConcSyncWorker) Close(ctx context.Context) error
Close 幂等优雅停机:停 tick 即完事(spec §1.5 协调态可丢、无排空义务—— 在途字段由 ts 新鲜度 ≤4s 出局、整键 EXPIRE 自灭,无清理命令)。
func (*ConcSyncWorker) Name ¶
func (w *ConcSyncWorker) Name() string
Name worker 名(worker.Worker 契约)。
func (*ConcSyncWorker) Start ¶
func (w *ConcSyncWorker) Start(_ context.Context) error
Start 非阻塞启动同步循环(幂等;自持 ctx,同 discovery baseCtx 惯用法)。 nil client 直接短路:无视图装配能力即全额本地语义,循环无存在意义。
func (*ConcSyncWorker) Stats ¶
func (w *ConcSyncWorker) Stats() any
Stats worker 观测(协调面冻结可见性):fail-open 静默退化时这是运维面唯一 痕迹——视图停止换入则 last_tick_ok 翻 false、consecutive_errors 增长。
type Config ¶
type Config struct {
MaxBodySize int64
MaxInflight int64
UpstreamTimeout time.Duration // codex 非流式上游超时(resp/images 各自包 ctx——B-P2-7;HTTPClient.Timeout 不可用:流式/非流式四方法共享,覆盖整响应体读取会切断长流式 SSE)。同源同值 cfg.Proxy.UpstreamTimeout(aiclient.Config.UpstreamTimeout 管 typed 面)
UpstreamStreamTimeout time.Duration // 流式 backstop(非流式超时在 aiclient.Config/cfg.Proxy.UpstreamTimeout)
FailoverAttempts int
GroupKeyRPM int
UsageCapture bool
BillingCapture bool // 计费开关(config.Billing.Enabled 映射;余额预检门控 + billable 行 Billed 出生标记取反——F2 单写点后不再路由分流)
// BehindCDN 客户端 IP 识别开关(config.proxy.behind_cdn 映射;clientIP
// 提取门控——false 完全不读供应商头直取 RemoteAddr,true 按序采信三头)。
// 部署前提见 config.go 注释与 clientip.go:源站只对 CDN 暴露。
BehindCDN bool
}
type InstancesProvider ¶
type InstancesProvider interface {
ClusterInstances() int
}
InstancesProvider 集群实例数 N 提供者(discovery.Discovery 实现 ClusterInstances; 多实例预算分摊 #14 §3.1——N = Redis 心跳活体数,spec 2026-08-25-redis-instance-discovery-design;原 DB settings 手工设置已删)。 nil(未装配)按 N=1(单实例语义)。N 在每次预算分配现读(instancesN),心跳 计数变化 ≤1 tick 天然生效。
type PriceResolver ¶
type PriceResolver interface {
ResolvePrices(model string, promptTokens int64, tier string, at time.Time) (domain.ResolvedPrices, bool)
}
PriceResolver 统一价格解析(零 DB 快照读,首中即停变体解析)。
type Proxy ¶
type Proxy struct {
// contains filtered or unexported fields
}
func New ¶
func New(cfg Config, sched *scheduler.Scheduler, creds *credential.Registry, rec *usage.Recorder, clients *aiclient.Factory, auth *Auth, log *logx.Logger, bill *BillingHooks, errlog *usage.ErrLogWorker) *Proxy
New 构造代理。creds 为凭据注册表(评审 M2:直接参数注入,编译期强制; 不用 Config 字段——避免 nil 运行时才炸)。bill 为计费钩子(Phase 5; nil = 计费全关——现有调用点/测试兼容)。errlog 为错误明细落盘 worker (分表设计;nil = 未装配——拒绝/异常路径只聚统计不落 err_logs 明细)。
func (*Proxy) CloseAllWS ¶
func (p *Proxy) CloseAllWS()
CloseAllWS closes all hijacked WS client connections (F3). Uses CloseNow for immediate TCP close — shutdown path, no handshake. Closed sessions unwind through existing classify→finish→rec.Record and inflight drops naturally. Idempotent.
func (*Proxy) HandleAnthropic ¶
func (p *Proxy) HandleAnthropic(w http.ResponseWriter, r *http.Request)
HandleAnthropic 转发 /v1/messages(anthropic 格式)。全部逻辑在 通用骨架 handleFormat + anthropicCaller(Phase 2 转发骨架通用化),本方法 只是端点入口委托。测试按处理函数直接调用(forward_ext_test.go 等),故入口保留。
func (*Proxy) HandleChat ¶
func (p *Proxy) HandleChat(w http.ResponseWriter, r *http.Request)
HandleChat 转发 /v1/chat/completions(openai-chat 格式)。全部逻辑在 通用骨架 handleFormat + chatCaller(Phase 2 转发骨架通用化),本方法 只是端点入口委托。测试按处理函数直接调用(proxy_test.go 等),故入口保留。
func (*Proxy) HandleImagesEdits ¶
func (p *Proxy) HandleImagesEdits(w http.ResponseWriter, r *http.Request)
HandleImagesEdits 转发 POST /v1/images/edits(openai-images 格式;上游子路径 由 handleFormat 内 imagesCallerFor 按请求路径选择)。
func (*Proxy) HandleImagesGenerations ¶
func (p *Proxy) HandleImagesGenerations(w http.ResponseWriter, r *http.Request)
HandleImagesGenerations 转发 POST /v1/images/generations(openai-images 格式; 全部逻辑在 handleFormat + imagesCaller——端点入口委托,同 HandleChat 形态)。
func (*Proxy) HandleModels ¶
func (p *Proxy) HandleModels(w http.ResponseWriter, r *http.Request)
HandleModels GET /v1/models:OpenAI 兼容模型列表(OpenAI SDK / codex 等 客户端用网关 key 拉取可用模型)。数据源 = 调度器内存快照(零 DB): key → meta.GroupID → 组快照 routes → 模型去重排序(scheduler.GroupModels)。 端点冷面(非转发热路径):鉴权通过即放行——不走计费/限流/并发门禁(只读 列表端点,OpenAI 语义——声明于注释);不建新缓存/新快照。空组/无模型 → 空 data 数组(200 不 404);组不存在/快照未加载 → 404(对齐 Select 的 ErrGroupNotFound 语义——鉴权已过但组失效)。
func (*Proxy) HandleResponses ¶
func (p *Proxy) HandleResponses(w http.ResponseWriter, r *http.Request)
HandleResponses 转发 /v1/responses(openai-responses 格式)。全部逻辑在 通用骨架 handleFormat + responsesCaller(Phase 2 转发骨架通用化),本方法 只是端点入口委托。测试按处理函数直接调用(forward_ext_test.go 等),故入口保留。
func (*Proxy) HandleResponsesWS ¶
func (p *Proxy) HandleResponsesWS(w http.ResponseWriter, r *http.Request)
HandleResponsesWS 处理 resp-ws 升级请求(/v1/responses 带 upgrade 头—— 真实客户端无 /ws 后缀,WS 与 POST /v1/responses 同路径,按协议分流)。 与 handleFormat 同构:guardPipeline(鉴权 → 额度/余额预检 → 两级并发门禁 → 限流,见 pipeline.go)→ 升级 → 首帧(= 请求体)→ 模型提取 → 选号 → failoverLoop(wsAttempt + wsSink,precheck=true)→ 双向 relay(usage 嗅探) → 记录。差异段留本文件:
- 无 HTTP body:请求体 = 升级后首个 WS 帧(response.create)
- 选号在首帧之后(模型来自首帧;挂死不占账号槽)
- 本地拒绝在升级后无 HTTP 状态码 → 错误事件帧承载(wsWriteError)
- 门禁覆盖整个长会话:guard 释放 defer 到 relay 结束(同现状语义)
func (*Proxy) HandleSearch ¶
func (p *Proxy) HandleSearch(w http.ResponseWriter, r *http.Request)
HandleSearch 转发 codex /v1/alpha/search(spec 2026-08-13 v2):codex CLI 以 独立 unary POST 调 web search(模型发 web.run tool call 时触发,与主 /responses 流并发)。**透传语义:请求体/响应体原样**(opaque results/ encrypted_output 网关零解析——alpha 端点实验性,上游变更网关免疫)。
与主 handleFormat 的差异(search 专属语义,spec 边界声明):
- 账号选择:body.model → Scheduler.Select(groupID, openai-responses, model) (复用主流 resp 路由面——四类型全可达;**独立选号无会话绑定**——P2 裁 决:search 请求自包含,上游鉴权 = 有效 Bearer,无会话亲和机制)
- **不走计费预检**(余额/缺价 402 均不执行——search 无预检语义;按次价在 2xx 落账时结算,零余额透支扣费为产品语义,防实现期误当缺陷"修复")
- **四类型分派(用户裁决 2026-08-13)**:codex-oauth/codex-pat → codex-sdk Search(适配层 clientFor 缓存客户端直接复用——统一 client 形态; search URL 由 SDK 方法内派生,网关零拼装;Auth 注入/刷新/fatal 生命周期 复用既有 SDK 面);api_key/responses-special → 静态透传(Bearer upstream key 直连上游——aiclient 既有静态 key 通道零新增机制;URL 裸根派生 base/v1/alpha/search)。组内混合类型路由允许(任一类型均可用——不再本地 拒绝)
- **x-codex-turn-metadata 统一不转发**(两路径均不带上游——SDK 默认头面 无该头;静态 rawPostCT 构造全新 Header 只设 Content-Type + Authorization, 与主流静态路径现状一致)
- **不做 ModelMapping 改写(P3-3 显式取舍)**:请求体原样 = 映射对 search 不生效(上游收客户端模型名)——零解析是 spec 显式约束,自洽记录
- 计费:2xx → usage_logs 行(format=openai-search + call_count=1 + price_per_call_millis=PriceResolver call 档(codex-search 模型) + cost=按次价×整单 倍率,applyBilling search 分支);非 2xx/网络错误 → 不计费(cost=0,错误 行走既有 err_logs 面)
复用面(评审 P3-4 点名):guardPipeline(鉴权/配额/并发门禁/限流序列)、 Select + handleSelectError、信封分类(statusOf/upstreamBody)、failoverLoop (**每轮按当轮 sel.CredentialType 重新分派**——searchAttempt 对齐 P1-1 教训: 跨类型换账号复用旧调用器会把健康账号路由到错误凭据路径)、 recordRejected/finish/buildLog/MarkResult 全部既有机制零改动。
func (*Proxy) SetCodex ¶
SetCodex 注入 codex SDK 适配层(T2 §3 装配点——main 构造 sdkbridge.NewCodex(统一失效回调) 后注入;nil = 未装配 → codex 类型请求 501 显式拒绝)。
func (*Proxy) SetInstancesProvider ¶
func (p *Proxy) SetInstancesProvider(inst InstancesProvider)
SetInstancesProvider 注入集群实例数 N 提供者(#14 多实例预算分摊;discovery 构造后调用——main 装配点:px.SetInstancesProvider(disco),spec 2026-08-25-redis-instance-discovery-design §2.2)。转发给 auth(gate 预算 ceil(剩余/N))与 limit(RPM ceil(rpm/N));N 在每次预算分配现读,心跳计数 变化 ≤1 tick 天然生效。
type QuotaUsedReader ¶
QuotaUsedReader 预算复核的 DB 只读接口(repository.KeyRepo 实现):预算耗尽 时读 key 当前已用额度(DB 权威值——usage.Recorder 批量增量回写 quota_used, 复核时刻分配预算用)。热路径不调(仅复核慢路径)。
type ToolView ¶
type ToolView struct {
Raw []byte
Type []byte
Name []byte
Namespace []byte
Result []byte
ID []byte
}
ToolView 工具视图:Raw 为 body 子切片(真零拷贝——引用 tools 值区间字节); Type/Name/Namespace 为顶层键值(去引号、\uXXXX 已解码;未提取 = 空)。 Result/ID 为 resp 响应检测旁路(spec §6)复用提取:ID = id 字符串值; Result = result 值裸字节(字符串值去引号/\u 解码;非字符串值——数组/null/ 数字——取原始值区间,空语义在判定处 imageResultNonEmpty 处理)。
type UpstreamCaller ¶
type UpstreamCaller interface {
Call(ctx context.Context, w http.ResponseWriter, r *http.Request, reqID string, groupID int64,
start time.Time, sel *scheduler.Selection, cred string, body []byte, stream bool) (code int, respBody []byte, handled bool, err error)
}
UpstreamCaller 一格式一实现:完成单次上游调用(含流式写出、客户端断开判定 与 usage 记录)。记录职责全在 caller(finish/buildLog/recordStreamAbort/ MarkResult 直接可用——评审 I-1);骨架只做 code 分支(429/5xx 转移、4xx 透传记录)、handled 短路与耗尽 record。凭据值经 aiclient 格式方法传入 (头名 aiclient 内组装,Phase 1 正交延续——评审 M-2)。
语义:
- handled == true → 请求已处理完毕(成功/客户端断开/流中止已记录;本地拒绝 已写出无记录),骨架直接 return(不可转移)
- handled == false → 上游未接受,骨架接手: code 429 → MarkResult(Kind429) + Release + 转移 code >= 500 或 code == 0(连接级/凭据错)→ MarkResult(RuleKindOf(code)) + Release + 转移 code 4xx(err == nil)→ 骨架 finish(buildLog(Err4xx)) + 透传 respBody (空 → 网关文案 "upstream rejected request")
- err 非 nil 仅在错误路径返回(分类由 code 承载);骨架用它提取错误文本 (部署故障修复):code==0 → err.Error() 落 ErrorMessage/last_error + Warn(err 全文),4xx → respBody 原文落 ErrorMessage。成功路径 err 恒 nil(零新增分配)。
- 例外(首字节前客户端断连,分类正确性):code==0 且 r.Context().Err()!=nil (客户端已断开)→ 记 499+ErrAbort 立即返回——不 failover、不 MarkResult、 不冷却(否则连接级误分类把无辜账号冷却 + failover 空转)。
type UserStatusLoader ¶
type UserStatusLoader interface {
LoadUsers(ctx context.Context) (map[int64]domain.UserSnapshot, error)
}
UserStatusLoader 由 repository.UserRepo 实现(RequireJWT 用户状态 + adminAuth 快照 role 校验的数据源)。
Source Files
¶
- auth.go
- billing.go
- caller.go
- caller_anthropic.go
- caller_chat.go
- caller_converted.go
- caller_images.go
- caller_images_codex.go
- caller_images_stream.go
- caller_responses.go
- caller_responses_ws.go
- clientip.go
- codex_responses_http.go
- codex_responses_ws.go
- concsync.go
- forward.go
- forward_anthropic.go
- forward_chat.go
- forward_images.go
- forward_responses.go
- forward_search.go
- gate.go
- limit.go
- models.go
- pipeline.go
- router.go
- strip_image.go
- strip_scan.go
- usage_extract.go
- ws_registry.go
- ws_relay.go