Three of them were the same shape as the bug this branch exists to fix: something goes wrong, the library swallows it, and the caller is left with a number that means the opposite of what happened. The throttle key had no source in it. Five sources on one model is the normal case here, so the first one to break would warn once and silence the other four for the life of the process, and the message never said which gateway to look at. An unknown verdict in a cached entry threw away the whole response. The rehydrator tolerates unknown fields but not unknown values of a known field, so two library versions sharing a Redis would each invalidate the other's entries: halved hit rate, and the only log line says the cache rebuild failed. A purely observational field should not be able to void a response whose content is intact. Normalising for telemetry now degrades instead of raising, both for a bare string and for a value outside the domain. Either one used to reach the same except and cost the whole row, which is exactly how 1.3.0 lost nineteen calls without anyone noticing.
PolyGateway
实验室统一的大语言模型调度与中转库:LLM / VLM / OCR / Embedding 四类调用共用同一套生产级治理栈——多源多账号、限流、错误分类重试、熔断、响应缓存、流式看门狗、遥测与成本。治理单位是一次模型调用;任务编排、业务解析、图像预处理都留在业务侧。
由三个真实项目(GovDoc-SaaS / CHSAnalyzer / Video-Tree-TRM5)各自手写的治理栈提炼而来,并以"能否全量迁移回这三个项目"作为验收标准。v1.0.0 已通过 GovDoc 与 CHSAnalyzer 两项目的全量迁移验收(约 −6800 行项目侧治理代码由本库继任)。
为什么需要它
每个接入大模型的项目都会重写同一批东西:重试循环、429 处理、熔断器、SSE 解析、遥测埋点——写三遍就有三份 bug。本库把这些收敛为一份经过压测验证的实现:
| 能力 | 说明 |
|---|---|
| 多源多账号 | {SCOPE}__{PROVIDER}__{N}__* 配置任意多源;健康感知选源(EWMA×在途 P2C)自动避开坏源 |
| 限流 | 并发/RPM/TPM × 全局/单源六道闸;TPM 预扣入场、按实际用量结算退款;Redis 后端跨进程原子(Lua) |
| 错误分类重试 | 一切失败落入四分类(见下),由分类决定重试/换源/熔断;429 属 pushback 不消耗重试预算;退避含 jitter 且尊重 Retry-After |
| 熔断 | 双通道(连续失败 + 失败率窗口,健康证据抑制误熔);半开单探针带租约(持有者死亡自动回收);epoch fencing 拒绝迟到写回;开路时长指数递增;开路时当场失败还是等冷却可配(CIRCUIT_OPEN,单源 scope 应配 wait) |
| 自适应并发 | AIMD:429 削减、成功缓升,防止打爆上游 |
| 背压与判死 | 配额满与熔断开路各自可选等待或快速失败(QUOTA_FULL / CIRCUIT_OPEN,两键不可互相替代);等待期按双条件判死(本地非生产性等待与全局无进展同时超窗)。stall 窗口只计非生产性等待(429 退避/配额轮询/熔断冷却),与 TIMEOUT_S 无耦合 |
| 响应缓存 | Redis/内存;key 含 model + messages 摘要 + namespace(缓存隔离单位)+ salt + 采样参数,多模态 content 先摘要再 hash(防毒化);可 per-call 绕过(科研重采样) |
| 流式看门狗 | TTFT / inter-token / 总超时三层活性;thinking token 刷活性不计结果;截断流(缺 [DONE])判瞬时不入缓存 |
| 推理可观测性 | "这次到底推理没推理"由多信号裁定(推理正文压倒 usage 明细),三态落在 LLMResponse.thinking_observation:observed / absent / unknown——unknown 是"本次判不出",不是"没推理";请求方向与实测观测矛盾时按 (模型, 方向) 各告警一次(能力表过期、开启未生效、注入了却观测不到);裁定结果随遥测落库 |
| 遥测与成本 | 每次调用(含缓存命中与失败)必录 25 字段;SQLite / Postgres 后端(表已存在时不需要 schema 建表权限,最小权限账号可直接用);按价格表折算成本落库(注意 LLMResponse.cost 本身恒为 None,成本只进遥测);多模态内容摘要落库不存原图 |
| 遥测的资源与降级 | Postgres 池闲时占 0 条连接、忙时上限可配(PGW_TELEMETRY_PG_POOL_MAX,缺省 4),每次写入有硬预算(PGW_TELEMETRY_PG_WRITE_TIMEOUT_S,缺省 5s);后端不可用是可恢复的降级(冷却 60s 后自动重试,DBA 建完表/放开权限即自愈),永久失能只留给 DSN 本身写错;降级状态可编程查询——client.telemetry_status 给出 degraded/fatal/reason/dropped_rows 等只读快照,不必再靠人工对账。对账要同时看 degraded 与 dropped_rows: 池饱和超预算丢的行走行级丢弃,degraded 保持 False(后端没挂,是本进程并发超了),只按 degraded 告警会看不见这一类丢行——而它恰是 pool_max 配小了的唯一信号 |
| 调用方维度 | 每次调用可带 tenant_id(遥测表的真实列,可挂 RLS、可建复合索引)与 meta(≤16 个自定义 KV);四个公共方法全覆盖,校验超限即报错;库只交付列,不启用 RLS、不建索引 |
| 遥测表治理 | llm_calls 是下游的表:PG 侧缺省不再自动 ALTER 补列(PGW_TELEMETRY_SCHEMA_MODE 三态,不设则 sqlite→auto、postgres→manual),manual 档点名缺列并按现有列裁剪写入;telemetry_schema_sql(backend) 自取可粘进迁移文件的建表/补列 SQL;PGW_TELEMETRY_TEXT_CAP 限正文长度(不设 = 存全文);保留期与访问控制走生产部署 DDL 模板加 tools/telemetry_retention.py |
| 结构化输出 | json_repair 修复 / 原生 schema 双策略 + 校验失败有界带反馈重问 |
| OCR | MonkeyOCR 双端点(文本转录 + 版面解析),bbox 数值防御下沉,逐源健康预检 check_health() |
| Embedding | 分批、维度校验、与 chat 同一治理栈 |
降级方向是铁律:缓存/遥测后端掉线 → 降级而不冒泡(业务调用照常返回);限流/熔断后端掉线 → 报错而非放行(防击穿上游)。遥测的降级不是静默的——进入/恢复各一条日志、期间按行数与时间节流复述,并随时可经 client.telemetry_status 读到。asyncio.CancelledError 全链路穿透,in-flight 资源在 finally 释放;资源所有权的纪律是「谁建的谁关」——aclose() 只关自己 from_env()/from_settings() 建出来的组件,注入进来的 transport / recorder / limiter / breaker / cache 一律不碰(由注入方自己关)。
安装
发布在实验室 Gitea PyPI(公开包,匿名可装):
pip install --extra-index-url https://gitea.iomgaa.online/api/packages/iomgaa/pypi/simple/ \
"polygateway[redis,postgres,structured]>=1.3.0,<2"
核心仅依赖 httpx + pydantic;按需选 extras:
| extra | 内容 | 何时需要 |
|---|---|---|
redis |
redis-py | Redis 限流/熔断/缓存后端 |
postgres |
asyncpg | Postgres 遥测后端 |
structured |
json-repair | 结构化输出的修复策略 |
sdk |
openai | 可选的 SDK transport(默认手写 httpx,不需要) |
要求 Python ≥ 3.12。
快速开始
1. 配置 .env
LLM__MINIMAX__1__BASE_URL=https://your-gateway/v1
LLM__MINIMAX__1__API_KEY=sk-xxx
LLM__MINIMAX__1__MODEL=MiniMax-M3
LLM__MINIMAX__1__TIMEOUT_S=120
LLM_MAX_RETRIES=3
LLM_RETRY_BASE_DELAY=2.0
LLM_RETRY_MAX_DELAY=30.0
LLM_CIRCUIT_BREAKER_THRESHOLD=5
LLM_CIRCUIT_BREAKER_COOLDOWN=60
PGW_LIMITER_BACKEND=memory
PGW_BREAKER_BACKEND=memory
PGW_CACHE_BACKEND=none
PGW_TELEMETRY_BACKEND=none
缺任何关键键都会在装配时报错——本库禁止默认值兜底掩盖配置缺失。
2. 发起治理调用
from polygateway import GatewayClient
async def main() -> None:
client = GatewayClient.from_env("LLM") # 读 .env 装配整套治理栈
try:
resp = await client.chat([{"role": "user", "content": "你好"}])
print(resp.content, resp.source_name, resp.latency_ms)
finally:
await client.aclose() # 归还连接与治理后端资源
chat() 原生接受 OpenAI 多模态 content 数组(image_url data URL),VLM 调用无需专门客户端;session_id / parent_call_id / cache_salt 关键字参数用于链路追踪与缓存控制,tenant_id / meta 用于遥测归属(见下文「多租户与自定义维度」);overlay 传采样参数(temperature / seed / max_tokens 等,恒定值宜配在源的 EXTRA_BODY 上)——它会进缓存 key,故逐次变化的 seed 天然不命中缓存。
3. OCR 与 Embedding
from polygateway import EmbeddingClient
from polygateway.ocr import OcrClient
ocr = OcrClient.from_env("OCR") # OCR__MONKEY__1__* 多源
text = await ocr.recognize_text(image_bytes) # 文本转录
layout = await ocr.parse_layout(image_bytes) # 版面解析(带 bbox 的元素列表)
health = await ocr.check_health() # 逐源预检 {"monkey_1": True, ...}
embed = EmbeddingClient.from_env("EMBED") # EMBED__*__* + EMBED__BATCH_SIZE
vectors = (await embed.embed(["文本 a", "文本 b"])).vectors
4. 业务侧异常处理
from polygateway import GatewayUnavailableError, RequestRejectedError
try:
resp = await client.chat(messages)
except GatewayUnavailableError as exc:
# 整个 scope 暂时无源可用: 延期重投,不消耗业务失败预算
schedule_retry(after_s=exc.retry_after_s) # exc.reason / exc.per_source_reasons 供诊断
except RequestRejectedError:
... # 请求本身有问题(400/格式拒绝): 不重试,直接失败
5. 多租户与自定义维度
resp = await client.chat(
messages,
tenant_id="acme-corp", # 遥测表的真实列,可挂 RLS
meta={"batch_id": "b-42", "stage": "extract"}, # 任意自定义 KV,进 meta 列
)
1.2.1 起,四个公共方法(chat / embed / recognize_text / parse_layout)都接受这两个关键字参数,都可省略,既有调用点无需改动。校验在入口收口、超限报 ValueError 而非静默丢弃:tenant_id ≤128 字符、非空串、不含首尾空白(空白拒绝而非 strip——" t1" 与 "t1" 在 RLS 等值比较下是两个租户);meta 最多 16 个键,键须匹配 [a-z0-9_.]{1,64}(pg_ 前缀保留给库),值仅限 str / int / float / bool,字符串值 ≤256 字符、float 须有限(nan / inf 不是合法 JSON,JSONB 会拒收)。两者都不进缓存 key——缓存隔离由 cache_namespace 负责。
存储上 tenant_id 两端都是 TEXT NOT NULL DEFAULT '',meta 在 Postgres 是 JSONB、在 SQLite 是 TEXT;老表要补上这两列(补列是否由库自动执行取决于 PGW_TELEMETRY_SCHEMA_MODE,见遥测表 schema 与升级纪律),补列后老行读出是空串而非 NULL(NULL 在任何 RLS policy 下都对所有人不可见,空串则可用一条 SQL 审出还有多少行待归属)。
库只提供列,不启用 RLS、不建索引。 数据库层的强制隔离是下游 DBA 的职责,库不会代劳;不执行则 tenant_id 只是一个可查可过滤的普通列,没有任何数据库层强制。库不代劳的原因是 default-deny:启用 RLS 而没有匹配的 policy = 零行可写且静默不报错,会让非多租户部署的遥测全量写失败。三角色、RLS policy、分区与保留期的完整可执行模板见生产部署 DDL 模板。
遥测表 schema 与升级纪律
llm_calls 是下游的表,不是库的私有存储。库对它发出的语句只有三类,别的一概不发:
| 库会发 | 库不发 |
|---|---|
列/表探测:PG 走 to_regclass + pg_attribute,SQLite 走 PRAGMA table_info(都只读 catalog) |
SELECT 表数据——库只写不读,故你加多少列、建多少索引、怎么分区都不影响它 |
INSERT,永远显式列名,冲突处理不绑定具体约束(PG ON CONFLICT DO NOTHING / SQLite INSERT OR IGNORE) |
UPDATE / DELETE / TRUNCATE / DROP——保留期与清理全归下游 |
表不存在时 CREATE TABLE IF NOT EXISTS(PG 侧先探测,表在就不发) |
ALTER TABLE,除非该后端处于 auto 档(见下);manual 档一条 DDL 都不发 |
补列档位 PGW_TELEMETRY_SCHEMA_MODE
| 取值 | 含义 |
|---|---|
| 不设(缺省) | 按后端派生:sqlite → auto、postgres → manual |
auto |
旧表缺列时库逐列 ALTER TABLE ADD COLUMN 补齐 |
manual |
库一条 ALTER 都不发;缺列只发一条 warning(点名缺的维度 + 附上可直接执行的 SQL),并按现有列裁剪 INSERT 继续写 |
缺省为什么两端不对称:PG 侧是共享的生产表,ALTER TABLE ADD COLUMN 取 ACCESS EXCLUSIVE 锁,会排在长事务后阻塞该表其后的所有查询,而遥测是业务路径上的内联 await;这类部署有 DBA、有迁移工具、讲最小权限,DDL 的执行时机该由他们挑。SQLite 侧是下游自己的本地文件(现有下游典型是 runs/*.db):没有 DBA、没有迁移工具、没有第二个系统碰它,ALTER 是毫秒级元数据操作,要求"升级后手工跑一条 SQL"是给零运维场景强加运维步骤。调研过的 11 个同类系统(Celery / APScheduler / Alembic / Django contrib / Hangfire / Quartz.NET / dbt / Airbyte / Fivetran / Prefect / Airflow)里,没有一个把"库在下游库里自动 ALTER 出列"作为默认行为。同一个键两侧都可显式覆盖。
| 表状态 | auto |
manual |
|---|---|---|
| 不存在 | 建表 | 仍然建表(新表无既有数据、无并发访问者,不存在锁队列风险;停掉它会让"零配置起步"断掉) |
| 存在、列齐 | 不发任何 DDL | 不发任何 DDL |
| 存在、缺列 | 逐列 ALTER;失败不裁剪,缺列以逐行 warning 暴露(承诺的是"把列补上",补不上就让问题可见;要降级写入请显式选 manual) |
不发 DDL,裁剪写入,缺的维度不落库 |
无论哪档,遥测的失败方向都是静默降级:缺列、补列失败、写入失败都只 warning,绝不冒泡打断业务调用。
自取建表脚本
telemetry_schema_sql 输出与库运行时执行的 DDL 同源(同一份常量),照它建完表,库探测到的列就是齐的:
import polygateway
print(polygateway.telemetry_schema_sql("postgres")) # 或 "sqlite";非法值抛 ValueError
# 直接落成迁移文件:注释头 + CREATE TABLE IF NOT EXISTS(全量列)+ 各补列语句
python -c "import polygateway; print(polygateway.telemetry_schema_sql('postgres'))" \
> migrations/001_llm_calls.sql
PG 变体的补列语句带 ADD COLUMN IF NOT EXISTS,整段可重复执行(它即便列已存在也会先取 ACCESS EXCLUSIVE 锁,故请挑低峰);SQLite 没有该语法,脚本以注释标明"仅当该列不存在时执行"。注意这与库内部执行的 ALTER 是两份文本:库侧一律先探测后 ALTER,不用 IF NOT EXISTS,正是为了在稳态下一条排他锁都不取。
Expand/Contract 承诺
这张表的演进只走 expand,不走 contract。以下五条既是当前实现,也是库对下游的承诺——库此后的演进受它们约束:
| 承诺 | 你可以据此做什么 |
|---|---|
| 新列只增不删不改名,一律追加在既有列之后 | 已有的视图、报表、ETL 不会因升级而失效 |
| 新列必可空,或带非易失常量默认值 | PG 11+ 补列不重写全表,SQLite 补列是元数据操作——大表升级也是秒级 |
INSERT 永远显式写出列名 |
你可以自行加列(业务维度、生成列),库的写入不受影响 |
库从不 SELECT *,也从不读回这张表的数据 |
库侧根本没有读路径,你加索引、加自己的列、挂 RLS 都影响不到它 |
| 写入的冲突处理不绑定具体约束 | 你可以把 llm_calls 建成 PARTITION BY RANGE (created_at) 的分区表(此时主键必须是 (call_id, created_at),PG 要求分区表唯一约束含分区键),库的探测、补列与写入照常工作 |
生产部署 DDL 模板(PostgreSQL)
上一节讲的是库怎么对待这张表(只探测、只 INSERT、可选建表);本节讲的是你该把这张表部署成什么样:谁能读、谁能写、写进去的行能不能被改、存多久。这些库一件都不代劳——它没有、也不该有这些权限。
模板按下表顺序执行,标识符(角色名、schema、分区月份、密码)按你的环境改;llm_calls 一律不写 schema 限定,靠 search_path 解析,与库的写入口径一致。
| # | 锚点 | 做什么 |
|---|---|---|
| 1 | roles |
建三角色并授 schema 级权限 |
| 2 | table |
把 llm_calls 改造成按 created_at 的 RANGE 分区表,属主归 polygateway_owner |
| 3 | partition |
建一个月分区(生产用 pg_partman 自动滚动) |
| 4 | grants |
授表级权限并 REVOKE UPDATE, DELETE |
| 5 | immutable |
触发器兜底(只防误操作) |
| 6 | rls |
启用并 FORCE RLS + 两条 policy |
| 7 | index |
(tenant_id, created_at) 复合索引 |
1. 三角色
| 角色 | 拿到什么 | 谁在用 |
|---|---|---|
polygateway_owner |
表属主:DDL、加分区、删分区 | DBA / 定时任务;不用它连库跑业务 |
polygateway_app |
INSERT + 受 RLS 约束的 SELECT |
库的连接串用这个 |
polygateway_report |
受 RLS 约束的 SELECT |
BI、对账、成本报表 |
CREATE ROLE polygateway_owner NOLOGIN;
CREATE ROLE polygateway_app LOGIN PASSWORD 'CHANGE_ME_APP';
CREATE ROLE polygateway_report LOGIN PASSWORD 'CHANGE_ME_REPORT';
GRANT polygateway_owner TO CURRENT_USER; -- 下一块要把表属主改过去,须先成为它的成员
GRANT USAGE ON SCHEMA public TO polygateway_owner, polygateway_app, polygateway_report;
GRANT CREATE ON SCHEMA public TO polygateway_owner; -- 滚动分区要在该 schema 里建表
2. 分区表
分区表必须下游先手工建:库的 CREATE TABLE 只会建普通表。列不在这里重抄一份——抄了就会漂移,故先用库自带脚本建出普通表,再原地改造:
python -c "import polygateway; print(polygateway.telemetry_schema_sql('postgres'))" \
| psql "$PGW_TELEMETRY_PG_DSN"
ALTER TABLE llm_calls RENAME TO llm_calls_seed; -- 上一步建出的普通表当模子
CREATE TABLE llm_calls (
LIKE llm_calls_seed INCLUDING DEFAULTS, -- 列/类型/NOT NULL/DEFAULT 全照搬
PRIMARY KEY (call_id, created_at) -- 分区表的唯一约束必须含分区键
) PARTITION BY RANGE (created_at);
DROP TABLE llm_calls_seed;
ALTER TABLE llm_calls OWNER TO polygateway_owner;
CREATE TABLE llm_calls_2026_01 PARTITION OF llm_calls
FOR VALUES FROM ('2026-01-01 00:00:00+00') TO ('2026-02-01 00:00:00+00');
ALTER TABLE llm_calls_2026_01 OWNER TO polygateway_owner;
生产不要手工滚月份,交给 pg_partman:5.x 用 create_parent(p_parent_table := 'public.llm_calls', p_control := 'created_at', p_interval := '1 month')(4.x 的参数序不同,以你装的版本文档为准),再把 part_config.retention 设成 '6 months'、retention_keep_table 设成 false,run_maintenance_proc() 就会到期 DROP 整个分区。清理必须走 DETACH/DROP PARTITION 而不是 DELETE——这不是性能偏好,是权限张力的唯一解:下一块要对应用角色 REVOKE DELETE,而 DROP PARTITION 是属主的 DDL,两者不冲突,DELETE 则必然冲突。
分区部署改变了幂等键,按 cache_hit 出报表的下游必须知道:普通表上主键是 call_id,分区表上是 (call_id, created_at)。库的写入是无冲突目标的 ON CONFLICT DO NOTHING,两种表形态都合法;但 emit_cache_hit 复用的是响应里的历史 call_id,于是同一次缓存命中的重复回放,在普通表上第二次起被 DO NOTHING 吞掉、在分区表上每次都落一行(created_at 由 DEFAULT now() 生成,主键不再重复)。逐次尝试行不受影响(每次尝试都是新 call_id)。
3. 权限与不可变性
llm_calls 按不可变审计表对待:写进去的行谁都不许改、不许删,过期数据靠 DROP PARTITION 整块消失。
GRANT INSERT, SELECT ON llm_calls TO polygateway_app;
GRANT SELECT ON llm_calls TO polygateway_report;
REVOKE UPDATE, DELETE, TRUNCATE ON llm_calls FROM polygateway_app, polygateway_report;
CREATE FUNCTION llm_calls_reject_mutation() RETURNS trigger LANGUAGE plpgsql AS $$
BEGIN
RAISE EXCEPTION 'llm_calls 是不可变审计表,% 被拒绝', TG_OP;
END;
$$;
CREATE TRIGGER llm_calls_immutable BEFORE UPDATE OR DELETE ON llm_calls
FOR EACH ROW EXECUTE FUNCTION llm_calls_reject_mutation();
触发器只防误操作,不防恶意:表属主可以 ALTER TABLE llm_calls DISABLE TRIGGER llm_calls_immutable 把它关掉。真正的强制是上一块的 REVOKE——权限检查发生在触发器之前,应用角色连触发器都碰不到。要防属主本人,需要的是数据库之外的手段(WAL 归档、只追加的外部存证),不是本表能解决的。
DROP PARTITION 与 DETACH PARTITION 是 DDL,不会触发行级触发器,故保留期清理不受这一块影响。
4. 行级安全与多租户隔离
照抄过 1.2.1 那份 RLS 模板的部署请先查一遍:那份模板把写侧也绑在
app.tenant_id这个 GUC 上,而库从不设这个 GUC,于是它的每一条INSERT都被 policy 拒绝——遥测的失败方向是静默降级,表现不是报错而是整张表零行。用能绕过 RLS 的角色(superuser 或带BYPASSRLS)执行SELECT count(*) FROM llm_calls;,并在应用日志里搜Postgres 遥测写入失败(丢弃该行):。下面这份是修正后的模板。
ALTER TABLE llm_calls ENABLE ROW LEVEL SECURITY;
ALTER TABLE llm_calls FORCE ROW LEVEL SECURITY; -- 属主不豁免
CREATE POLICY llm_calls_app_write ON llm_calls FOR INSERT TO polygateway_app
WITH CHECK (true);
CREATE POLICY llm_calls_app_read ON llm_calls FOR SELECT TO polygateway_app
USING (tenant_id = NULLIF(current_setting('app.tenant_id', true), ''));
CREATE POLICY llm_calls_report_read ON llm_calls FOR SELECT TO polygateway_report
USING (tenant_id = NULLIF(current_setting('app.tenant_id', true), ''));
CREATE INDEX idx_llm_calls_tenant_created ON llm_calls (tenant_id, created_at);
current_setting(..., true) 的第二参数令 GUC 未设时返回 NULL 而非抛错,外层 NULLIF 把空串归一为 NULL——合起来使未设租户 = 零行(fail-closed)而不是全部行。索引列序不可颠倒:启用 RLS 后 policy 给每条查询隐式追加 tenant_id 等值谓词,它出现在 100% 的谓词里,必然是前导列。分区表上不能用 CREATE INDEX CONCURRENTLY(PG 不支持在分区父表上并发建索引);父表此时还没有数据,直接建即可,给已有数据的普通表补索引才需要逐个分区 CONCURRENTLY。
写侧 policy 为什么是 WITH CHECK (true) 而不是等值比较:库用一个连接池给所有租户写遥测,且从不发 set_config('app.tenant_id', ...)(源码里没有这条语句)。把写侧也绑到 GUC 上,库的每一条 INSERT 都会被 policy 拒绝——而遥测的失败方向是静默降级,表现是逐行 warning + 整表零行。隔离在这个模型里由读侧承担:写入方是库自己(可信),读取方才是要隔离的人。若你的调用点保证每次调用都带 tenant_id,可把写侧收紧成 WITH CHECK (tenant_id <> ''),代价是漏传 tenant_id 的调用点会丢遥测行(只留一条 warning)。
四个陷阱,每一个的失败形态都是静默的:
| 陷阱 | 后果 |
|---|---|
| 表属主默认豁免 RLS | 只写 ENABLE 而漏 FORCE,用属主角色连库时隔离形同虚设,且查询一切正常看不出来 |
FORCE 之后属主自己也被 policy 管 |
模板没给 polygateway_owner 任何 policy,故它读不到、也写不进任何行——这是有意的(它只用来做 DDL),但别拿它跑报表 |
租户上下文必须在显式事务内用 set_config('app.tenant_id', ..., true) |
asyncpg 默认 autocommit,单发 SET LOCAL 会当场失效,而 PG 只发 warning 不报错;表现是 policy 永远拿不到租户 → fail-closed 到零行 |
读侧 policy 漏写 USING |
FOR SELECT 的 policy 只认 USING;写成 WITH CHECK 不报错也不生效,隔离直接落空 |
5. 库本身需要的最小权限
按上面的模板部署后,库的连接串用 polygateway_app,它需要的权限恰好是下表这些——多一分都不必给:
| 库会发的语句 | 需要什么 |
|---|---|
| 连库 | 数据库 CONNECT + schema USAGE |
SELECT to_regclass('llm_calls')、查 pg_attribute(列探测) |
无需额外授权(系统 catalog 默认对 PUBLIC 可读) |
INSERT INTO llm_calls (...) |
表 INSERT;RLS 打开后还须有一条允许写的 policy |
CREATE TABLE IF NOT EXISTS(仅当表不存在) |
schema CREATE。生产建议不给:表由 owner 先建好,库探测到表在就不发这条 |
ALTER TABLE ADD COLUMN(仅 PGW_TELEMETRY_SCHEMA_MODE=auto) |
表属主——PG 的 ALTER TABLE 只认属主,这一项无法单独 GRANT。PG 侧缺省就是 manual,补列交给 DBA |
6. 合规下游的推荐配置
三件事(截断、保留期、访问控制)要一起上才有意义,故给一份可直接照抄的组合,而不是让你自己拼:
PGW_TELEMETRY_BACKEND=postgres
PGW_TELEMETRY_PG_DSN=postgresql://polygateway_app:...@db:5432/telemetry
PGW_TELEMETRY_SCHEMA_MODE=manual # PG 侧本就是缺省;写出来是为了不依赖缺省
PGW_TELEMETRY_TEXT_CAP=2000 # 落库正文的字符上限;不设 = 存全文
| 层 | 配置 |
|---|---|
| 正文体量 | PGW_TELEMETRY_TEXT_CAP=2000(按需调);超出部分头部硬切并附 …(略 N 字) |
| 保留期 | 上面的分区模板 + pg_partman 的 retention,过期分区整块 DROP |
| 访问控制 | 上面的三角色 + REVOKE UPDATE, DELETE + FORCE RLS |
| 存量兜底 | 已经攒成一张大普通表、来不及改造分区时,用 tools/telemetry_retention.py(默认 dry-run,--apply 才动手;探测到分区表会直接退出让路给 DROP PARTITION) |
PGW_TELEMETRY_TEXT_CAP 的覆盖面必须说清,否则合规判断会出错。 cap 落在四处:messages 里每条消息的字符串 content、多模态 content 数组中 type == "text" 的 part 的 text,以及 response 与 thinking 两列。消息侧的这个面与缓存摘要函数 digest_messages 一致——只碰 content,消息里别的字段一概不碰。所以调用方自己塞进 tool_calls.function.arguments、name 等字段的内容不在覆盖范围内:开了 cap 不等于表里没有全文残留。另需知道:缺省是不截断(存全文),而截断之后遥测不再是可复现重放的证据。
7. SQLite 侧的保留期
SQLite 侧不建议对着一个大库文件跑 DELETE + VACUUM,而应按天/按实验轮转库文件——runs/<date>.db、runs/<experiment>.db 这样,到期直接删文件。这是三个现有下游(Video-Tree-TRM5 / CHSAnalyzer / dissect)天然就有的形态,比删行省事也安全得多:删文件是 O(1) 且不可能删错行,而 VACUUM 会重写整库、期间需要一倍磁盘空间,还会把并发写入方挡在外面。
tools/telemetry_retention.py 的 SQLite 分支是给存量场景兜底的——已经攒成一个大库、来不及改轮转时用它,不是推荐路径。
该脚本随仓库分发,不在 pip 包内(它是运维工具而非库能力,库本体不 import 它,也不该拿到 DELETE 权限),请从仓库的 tools/telemetry_retention.py 取,用维护角色跑。
错误模型(四分类)
一切失败在 transport 层翻译为四类之一,治理行为由分类决定,业务侧不需要判断状态码:
| 分类 | 含义 | 库内行为 |
|---|---|---|
TransientError |
超时/5xx/网络抖动/截断流 | 换源重试 + 退避 |
SourceDeadError |
401/403/欠费(429+insufficient_quota) | 立即熔断该源 + 换源 |
RequestRejectedError |
400/内容拒绝/本地格式拒绝 | 不重试不换源,快速失败 |
ResultInvalidError |
调用成功但结果不合格(坏 JSON/维度不符/坏 bbox) | 不熔断("坏结果 ≠ 坏服务"),按策略有界重问或上抛 |
预算耗尽/全源熔断时抛 GatewayUnavailableError 族(CircuitOpenError / AllSourcesExhausted),携带 scope / reason / retry_after_s / per_source_reasons,供任务队列做延期重投。
网关拒绝的理由不会丢失(1.2.0 起):非 2xx 的响应体经折叠与截断后同时进入异常 message 与 exc.body_text,故遥测表的 error 列里就能看到网关的原话——不必再为查一次 400 单独埋点。截断保头保尾(总长 2048 字符),JSON 错误体尾部的 code / request_id 不会被切掉。经中转部署时请注意:第三方中转服务自身抖动也会回 400,从状态码上与"你的输入有问题"无法区分;库仍按确定性失败处理(直连供应商时重试只会白烧配额),批处理下游宜据 body_text 自备兜底分类。
哪些异常会到达调用方
上表的"库内行为"一列描述的是治理动作,不是调用方要处理的东西。四类里有两类根本到不了调用方——它们被重试循环接住,预算耗尽时统一包成 AllSourcesExhausted。这个区分只看类型树和 docstring 是读不出来的,曾让下游据此写错整段设计文档,故在此列明:
| 会到达调用方 | 库内吸收(不必 catch) |
|---|---|
GatewayUnavailableError 族——CircuitOpenError / AllSourcesExhausted / GovernanceBackendError |
TransientError(退避后换源重试,耗尽即转为 AllSourcesExhausted) |
RequestRejectedError |
SourceDeadError(立即熔断该源并换源,同上) |
ResultInvalidError |
|
SourceNotConfiguredError |
GovernanceBackendError 属于第一列: 限流/熔断的状态后端(如 Redis)自身故障时库 fail-closed——一个请求都发不出去,这就是"整个 scope 暂时不可用"。它继承 GatewayUnavailableError,所以 §4 那段 except GatewayUnavailableError 一条即覆盖完整,无需为它单列分支。retry_after_s 默认 5 秒(后端恢复时间不可知,取 0 会让积压任务零延迟冲击已挂掉的后端)。
SourceNotConfiguredError 有意不在第一列的族内: 源名不在限流后端的配置字典中是装配缺陷而非暂时故障,它应当消耗失败预算、进死信、让人看见——归入可重投家族只会让配置写错的任务永远重投且无人告警。
配置参考
配置只有两条装配路径:from_env()(读 .env/环境变量)或构造函数全量注入(测试/高级);库内部任何组件不自读环境变量。键名全集见 .env.example,约定速览:
| 键形态 | 作用 |
|---|---|
{SCOPE}__{PROVIDER}__{N}__{FIELD} |
第 N 个源;FIELD 全集 = BASE_URL/API_KEY/MODEL/TIMEOUT_S/MAX_CONCURRENCY/RPM/TPM/EST_TOKENS/TTFT_TIMEOUT_S/INTER_TOKEN_TIMEOUT_S/ENABLE_THINKING/MISSING_DONE/TRUST_ENV/EXTRA_BODY(表外的 FIELD 直接报错) |
{SCOPE}__GLOBAL__* |
scope 级全局限额(跨源并发/RPM/TPM) |
{SCOPE}__RETRY__* / BREAKER__* / BACKPRESSURE__* / SELECTOR / QUOTA_FULL / CIRCUIT_OPEN |
per-scope 韧性参数;缺省回落平铺键(LLM_MAX_RETRIES 等,兼容旧项目习惯) |
{SCOPE}__BATCH_SIZE / NORMALIZE / EXPECTED_DIM |
仅 EmbeddingClient 消费;BATCH_SIZE 必填(分批是行为关键,不设默认) |
PGW_LIMITER_BACKEND / PGW_BREAKER_BACKEND |
memory(单进程)或 redis(跨进程共享,需 REDIS_URL) |
PGW_CACHE_BACKEND |
none / memory / redis;非 none 时需 PGW_CACHE_NAMESPACE + PGW_CACHE_TTL_S(须 > 0) |
PGW_TELEMETRY_BACKEND |
none / sqlite(需 PGW_TELEMETRY_SQLITE_PATH)/ postgres(需 PGW_TELEMETRY_PG_DSN) |
PGW_TELEMETRY_SCHEMA_MODE |
可选:auto / manual;不设则按后端派生(sqlite→auto、postgres→manual),显式设置则两侧都可覆盖。决定库是否给已存在的旧表自动 ALTER 补列,详见遥测表 schema 与升级纪律 |
PGW_TELEMETRY_TEXT_CAP |
可选正整数:遥测落库正文的字符上限(作用于每条消息的文本 content、多模态 part 的 text、response、thinking);不设 = 不截断,详见合规下游的推荐配置 |
PGW_TELEMETRY_PG_POOL_MAX |
可选正整数(缺省 4):Postgres 遥测池的连接上限。池按需建连,闲时占 0 条,这一格是忙时天花板而非常驻量。调参按实测折算而非按 pool_max / RTT 估算——跨内网 RTT ≈ 123ms 上 pool_max=4 实测约 15.6 行/秒(一次 INSERT 的往返比一次 SELECT 1 重一倍);多个 client 共享同一 recorder 时并发在此汇聚,应相应放大 |
PGW_TELEMETRY_PG_WRITE_TIMEOUT_S |
可选正数(缺省 5.0):一次遥测写入的硬预算,同时用作建连、acquire 与「准备 + 取连接 + 执行」整段的上界;超时即丢该行,绝不让遥测无界地挂在业务路径上 |
PGW_PRICING_PATH / PGW_STRUCTURED_MAX_RETRIES / PGW_LEASE_TTL_S |
可选:价格表(缺省则成本恒 None)/ 结构化重问上限(缺省 2)/ permit 租约秒数(缺省 1500,须 ≥ 最大源 TIMEOUT_S) |
{SCOPE}__CIRCUIT_OPEN=fail_fast|wait(缺省 fail_fast)——单源 scope 请配 wait
熔断的设计前提是"这个源坏了,把流量导到别的源"。只配了一个源时这个前提不成立,同一段代码做的事变成"这个源坏了,所以整个 scope 停止服务":开路期间每一次调用都在几毫秒内失败,MAX_ATTEMPTS 一格用不上,一个网络包都没发出去。中转抖动几十秒就足以打断一条跑了几小时的长任务。
wait 档改变的只是"调用方当场失败还是排队等":等待期间照样一个请求都不发,熔断对配额和钱包的保护完整保留。代价是单次调用的最坏墙钟被拉长,上限为 {SCOPE}__BACKPRESSURE__STALL_WINDOW_S(缺省 300 秒)。wait 不豁免重试预算——冷却结束后放行的探针是一次真实尝试,失败照样烧一格 MAX_ATTEMPTS;因此密钥失效(401/403)这类一击即熔的源通常更早以 reason=retry_exhausted 失败,而非等满窗口的 stalled。库无法区分"密钥坏了"和"中转抖了",选 wait 就是声明"宁可等也不要当场死"。多源部署保持 fail_fast:有源可换时,换源比等待快。
该键与 {SCOPE}__QUOTA_FULL 同形但不可互相替代:配额满是"排队等自己的份额"(必然轮到),熔断开路是"等这个源恢复"(未必恢复),所以两者分开配置。
两个易被忽略的源级键:MISSING_DONE 决定 SSE 缺 [DONE] 时的处置(retry 默认判瞬时重试 / salvage 收下已收内容并把用量可信度降为 estimated;零内容恒 retry,不受该键影响);EXTRA_BODY 是该源恒定的采样参数(JSON 对象串,并入请求体,优先级低于 chat(overlay=...)),禁用键 model / messages / stream / stream_options 配了直接报错,OCR 与 EMBED scope 不消费该键(配了忽略并 warning)。
SCOPE 是逻辑角色(LLM/VLM/OCR/EMBED/JUDGE/SEARCH…任意大写名),同一进程可按角色装配多个 client,各自独立配置与治理状态。
架构
端口适配器 + 中间件洋葱:决策逻辑一份,状态存储可插拔。
graph LR
A[业务代码] --> B[GatewayClient]
B --> C[缓存 MW] --> D[遥测 MW] --> E[重试/选源/限流/熔断 MW]
E --> F[Transport httpx]
F --> G[(上游网关)]
E -.端口.-> H[(内存 / Redis 后端)]
D -.端口.-> I[(SQLite / Postgres)]
| 模块 | 职责 |
|---|---|
types.py / errors.py / ports.py |
内核:冻结类型、四分类异常、全部 Protocol(最内层,不依赖任何实现) |
middleware/ |
治理算法(重试/限流/熔断/缓存/遥测),只面向端口 |
transports/ |
协议细节:OpenAI 兼容 SSE、MonkeyOCR 双端点;错误翻译在此层 |
backends/ |
限流/熔断/缓存的内存与 Redis 实现(同一契约测试套件双后端共用) |
telemetry/ |
SQLite / Postgres 遥测后端 |
structured/ |
结构化输出策略 |
依赖纪律由 import-linter 机械化执法(make lint)。完整架构决策(D1-D15 含论证过程)见 research-wiki/ARCHITECTURE.md。
可靠性证据
行为不是宣称出来的,是压测出来的(数字见 research-wiki/findings/):
| 场景 | 结果 |
|---|---|
| 故障混编 soak(坏 key/黑洞/慢源/限流源混合,8000 调用) | 成功率 98.96%,坏源吸流被压制,真实源零误熔 |
| OCR 故障池 soak(1500 调用,redis 双后端跨进程) | 成功率 99.73%,13 项不变量全过(租约归零/探针不悬挂/零取消泄漏等) |
| 两项目全量迁移回归 | 原测试全绿 + 真实链路冒烟 + 50 样本批跑 100% 解析 |
时间语义测试(租约过期、窗口滚动、半开探针)全部真实等待不缩放;Redis/Postgres 测试打真实实验室后端,不 mock Lua。
开发
conda create -n PolyGateway python=3.12 && conda activate PolyGateway
make install # editable 安装(dev + 全部 extras)
make test # pytest + 覆盖率(目标 ≥80%)
make lint # ruff + import-linter
make ci # 只读全量验证
测试组织:tests/{unit,integration,e2e} + 双后端契约测试;并发/取消/降级方向是一等测试对象。压测 harness 在 tools/soak/。贡献流程与项目纪律见 CLAUDE.md。
文档导航
| 想了解 | 看 |
|---|---|
| 全部架构决策及理由(单一事实源) | research-wiki/ARCHITECTURE.md |
| 里程碑与状态 | research-wiki/ROADMAP.md |
| 项目迁移指南(删除清单/组件映射/行为审计) | research-wiki/migrations/ |
| 每个功能的设计与验收记录 | research-wiki/designs/、research-wiki/findings/ |
| 版本变更 | CHANGELOG.md |
兼容性承诺
LLMResponse 等被下游消费的公共类型,字段只增不删不改名且新增字段必带默认值;{SCOPE}__{PROVIDER}__{N}__{FIELD} 与平铺韧性键名(LLM_TIMEOUT 等)沿用三项目既有习惯,不做破坏性改名。实验室内部库,随实验室项目需求演进。