2af445cfc2
Dating the changelog and bumping both version strings is the cheap half. The README is the half that gets frozen into the sdist, so it is fixed first: the install pin now names 1.2.1 (1.2.0 has no tenant dimension), and the per-source FIELD table finally lists MISSING_DONE and EXTRA_BODY - the latter was already referenced elsewhere in the same file. The same table was missing QUOTA_FULL, the embedding-only keys, the memory cache backend and three optional PGW_* keys; all are reconciled against _SOURCE_FIELDS and _load_pgw rather than from memory.
266 lines
17 KiB
Markdown
266 lines
17 KiB
Markdown
# 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 拒绝迟到写回;开路时长指数递增 |
|
|
| 自适应并发 | AIMD:429 削减、成功缓升,防止打爆上游 |
|
|
| 背压与判死 | 配额满可选等待或快速失败;等待期按双条件判死(本地非生产性等待与全局无进展**同时**超窗)。stall 窗口只计**非生产性**等待(429 退避/配额轮询/熔断冷却),与 `TIMEOUT_S` 无耦合 |
|
|
| 响应缓存 | Redis/内存;key 含 model + messages 摘要 + namespace(缓存隔离单位)+ salt + 采样参数,多模态 content 先摘要再 hash(防毒化);可 per-call 绕过(科研重采样) |
|
|
| 流式看门狗 | TTFT / inter-token / 总超时三层活性;thinking token 刷活性不计结果;截断流(缺 `[DONE]`)判瞬时不入缓存 |
|
|
| 遥测与成本 | 每次调用(含缓存命中与失败)必录 24 字段;SQLite / Postgres 后端(表已存在时**不需要** schema 建表权限,最小权限账号可直接用);按价格表折算成本落库(注意 `LLMResponse.cost` 本身恒为 `None`,成本只进遥测);多模态内容摘要落库不存原图 |
|
|
| 调用方维度 | 每次调用可带 `tenant_id`(遥测表的真实列,可挂 RLS、可建复合索引)与 `meta`(≤16 个自定义 KV);四个公共方法全覆盖,校验超限即报错;**库只交付列,不启用 RLS、不建索引** |
|
|
| 结构化输出 | json_repair 修复 / 原生 schema 双策略 + 校验失败有界带反馈重问 |
|
|
| OCR | MonkeyOCR 双端点(文本转录 + 版面解析),bbox 数值防御下沉,逐源健康预检 `check_health()` |
|
|
| Embedding | 分批、维度校验、与 chat 同一治理栈 |
|
|
|
|
**降级方向是铁律**:缓存/遥测后端掉线 → 静默降级(warning);限流/熔断后端掉线 → 报错而非放行(防击穿上游)。`asyncio.CancelledError` 全链路穿透,in-flight 资源在 finally 释放。
|
|
|
|
## 安装
|
|
|
|
发布在实验室 Gitea PyPI(公开包,匿名可装):
|
|
|
|
```bash
|
|
pip install --extra-index-url https://gitea.iomgaa.online/api/packages/iomgaa/pypi/simple/ \
|
|
"polygateway[redis,postgres,structured]>=1.2.1,<2"
|
|
```
|
|
|
|
核心仅依赖 `httpx` + `pydantic`;按需选 extras:
|
|
|
|
| extra | 内容 | 何时需要 |
|
|
|---|---|---|
|
|
| `redis` | redis-py | Redis 限流/熔断/缓存后端 |
|
|
| `postgres` | asyncpg | Postgres 遥测后端 |
|
|
| `structured` | json-repair | 结构化输出的修复策略 |
|
|
| `sdk` | openai | 可选的 SDK transport(默认手写 httpx,不需要) |
|
|
|
|
要求 Python ≥ 3.11。
|
|
|
|
## 快速开始
|
|
|
|
### 1. 配置 `.env`
|
|
|
|
```bash
|
|
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. 发起治理调用
|
|
|
|
```python
|
|
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
|
|
|
|
```python
|
|
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. 业务侧异常处理
|
|
|
|
```python
|
|
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. 多租户与自定义维度
|
|
|
|
```python
|
|
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`;老表自动补列,**老行读出是空串而非 NULL**(NULL 在任何 RLS policy 下都对所有人不可见,空串则可用一条 SQL 审出还有多少行待归属)。
|
|
|
|
**库只提供列,不启用 RLS、不建索引。** 要数据库层的强制隔离,以下 DDL 是**下游 DBA 的职责,库不会代劳**;不执行则 `tenant_id` 只是一个可查可过滤的普通列,没有任何数据库层强制。库不代劳的原因是 default-deny:启用 RLS 而没有匹配的 policy = 零行可写且静默不报错,会让非多租户部署的遥测全量写失败。
|
|
|
|
```sql
|
|
ALTER TABLE llm_calls ENABLE ROW LEVEL SECURITY;
|
|
ALTER TABLE llm_calls FORCE ROW LEVEL SECURITY; -- 属主不豁免
|
|
CREATE POLICY llm_calls_tenant_isolation ON llm_calls TO polygateway_app
|
|
USING (tenant_id = NULLIF(current_setting('app.tenant_id', true), ''))
|
|
WITH CHECK (tenant_id = NULLIF(current_setting('app.tenant_id', true), ''));
|
|
```
|
|
|
|
```sql
|
|
CREATE INDEX CONCURRENTLY 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% 的谓词里,必然是前导列。
|
|
|
|
三个陷阱,每一个的失败形态都是**静默的**:
|
|
|
|
| 陷阱 | 后果 |
|
|
|---|---|
|
|
| 表属主默认**豁免** RLS | 只写 `ENABLE` 而漏 `FORCE`,用属主角色连库时隔离形同虚设,且查询一切正常看不出来 |
|
|
| 租户上下文必须在**显式事务内**用 `set_config('app.tenant_id', ..., true)` | asyncpg 默认 autocommit,单发 `SET LOCAL` 会当场失效,而 PG **只发 warning 不报错**;表现是 policy 永远拿不到租户 → fail-closed 到零行 |
|
|
| policy 必须同时写 `USING` 与 `WITH CHECK` | 只写前者则租户 A 读不到 B 的行,却**能插入标着 B 的行**——污染发生在写入侧,读侧查不出来 |
|
|
|
|
## 错误模型(四分类)
|
|
|
|
一切失败在 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](.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` | 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_PRICING_PATH` / `PGW_STRUCTURED_MAX_RETRIES` / `PGW_LEASE_TTL_S` | 可选:价格表(缺省则成本恒 `None`)/ 结构化重问上限(缺省 2)/ permit 租约秒数(缺省 1500,须 ≥ 最大源 `TIMEOUT_S`) |
|
|
|
|
两个易被忽略的源级键:`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,各自独立配置与治理状态。
|
|
|
|
## 架构
|
|
|
|
端口适配器 + 中间件洋葱:决策逻辑一份,状态存储可插拔。
|
|
|
|
```mermaid
|
|
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-D14 含论证过程)见 [research-wiki/ARCHITECTURE.md](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。
|
|
|
|
## 开发
|
|
|
|
```bash
|
|
conda create -n PolyGateway python=3.11 && 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](CLAUDE.md)。
|
|
|
|
## 文档导航
|
|
|
|
| 想了解 | 看 |
|
|
|---|---|
|
|
| 全部架构决策及理由(单一事实源) | `research-wiki/ARCHITECTURE.md` |
|
|
| 里程碑与状态 | `research-wiki/ROADMAP.md` |
|
|
| 项目迁移指南(删除清单/组件映射/行为审计) | `research-wiki/migrations/` |
|
|
| 每个功能的设计与验收记录 | `research-wiki/designs/`、`research-wiki/findings/` |
|
|
| 版本变更 | [CHANGELOG.md](CHANGELOG.md) |
|
|
|
|
## 兼容性承诺
|
|
|
|
`LLMResponse` 等被下游消费的公共类型,字段**只增不删不改名**且新增字段必带默认值;`{SCOPE}__{PROVIDER}__{N}__{FIELD}` 与平铺韧性键名(`LLM_TIMEOUT` 等)沿用三项目既有习惯,不做破坏性改名。实验室内部库,随实验室项目需求演进。
|