proxy

package
v0.0.1-beta.5 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Aug 26, 2026 License: AGPL-3.0 Imports: 44 Imported by: 0

Documentation

Overview

Package proxy 是 AI 请求热路径:分组 key 鉴权 → 调度器选号 → SDK 转发 → 用量采集。 规格 §6/§9。不变量:热路径零 DB、零 per-request 锁。

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

func AIRouter

func AIRouter(p *Proxy) http.Handler

AIRouter 挂载 AI 端点(规格 §6.1/§9):路径决定请求格式,全部走通用转发 骨架 handleFormat(Phase 2:UpstreamCaller 注册表分发);resp-ws 与 HTTP responses 同路径(真实客户端无 /ws 后缀)——/v1/responses 带 upgrade 头 → 按 resp-ws 处理,走专用编排 HandleResponsesWS(caller_responses_ws.go)。

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

func (a *Auth) Acquire(meta domain.KeyMeta) (int, bool)

Acquire 两级并发门禁:user → key 依次 CAS 抢占;key 失败回滚 user 计数 (评审 I-3:防泄漏)。返回已 acquire 层级位掩码(release 仅释放已 acquire 层级)。未设置上限(max=0)或计数器缺失(跨 reload 竞态窗口)→ 该层跳过。

func (*Auth) Authenticate

func (a *Auth) Authenticate(r *http.Request) (domain.KeyMeta, bool)

Authenticate 解析网关 key 并返回 KeyMeta。兼容两种客户端口径: OpenAI 客户端发 Authorization: Bearer;Anthropic 官方 SDK / Claude Code 发 x-api-key 头。两者同时提供时以 Authorization 为准。 key 或归属用户被禁用 → 快照直接拒绝(401,即时失效)。 快照 map key = 明文,等值直查(零哈希);meta 无 key 字符串字段—— 鉴权失败日志天然不落明文。

func (*Auth) DeductQuota

func (a *Auth) DeductQuota(keyID, tokens int64)

DeductQuota 请求结束扣减(后扣模型;usage 已知;无额度 key 无计数器 → no-op)。

func (*Auth) Delete

func (a *Auth) Delete(raw string)

func (*Auth) InFlightUsers

func (a *Auth) InFlightUsers() map[int64]int64

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

func (a *Auth) QuotaExhausted(meta domain.KeyMeta) bool

QuotaExhausted 额度检查:本地预算快读(零锁零 DB);预算耗尽触发 DB 复核 认领(#14 §3.2——复核成功续预算继续放行,复核确认真尽才 429)。检查在并发 acquire 之前(评审提醒①:失败无并发槽副作用);未设置额度 key 短路零成本。

func (*Auth) Release

func (a *Auth) Release(meta domain.KeyMeta, level int)

Release 释放并发计数(仅释放 acquire 返回的层级;跨 reload 命中新快照的 继承计数,与 scheduler Release 同语义)。

func (*Auth) Reload

func (a *Auth) Reload(ctx context.Context) error

Reload 全量刷新鉴权快照(注册表首刷/周期 auth-sync/用户变更 invalidate): keys 元数据 + 用户状态 + 门禁计数器(在途值跨 reload 继承)。 失败必打 Warn(含调用方是否忽略错误——invalidate 回调等吞错路径):加载 失败若被忽略,快照保持旧值/空表 → 鉴权全部 401 或用旧 key 放行,静默恶化 (IN 超限事故的"运行中静默失败"形态即此类)。

func (*Auth) RemoveUser

func (a *Auth) RemoveUser(userID int64)

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) Upsert

func (a *Auth) Upsert(raw string, meta domain.KeyMeta)

Upsert 增量刷新单个 key(key 创建/轮换/更新后调用;门禁计数器同步)。

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 KeyLoader

type KeyLoader interface {
	LoadKeys(ctx context.Context) (map[string]domain.KeyMeta, error)
}

KeyLoader 由 repository.KeyRepo 实现(keys 独立表鉴权快照)。

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) Inflight

func (p *Proxy) Inflight() int64

func (*Proxy) SetCodex

func (p *Proxy) SetCodex(c *sdkbridge.Codex)

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

type QuotaUsedReader interface {
	QuotaUsed(ctx context.Context, keyID int64) (int64, error)
}

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 校验的数据源)。

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL