Merge branch 'feat/issue-13-schema-mode'

issue #13: the library no longer alters a downstream Postgres table on
its own. PGW_TELEMETRY_SCHEMA_MODE is tri-state and defaults by backend
-- SQLite keeps auto-migrating a local file, Postgres switches to manual,
where a stale table gets a named warning with runnable SQL and the INSERT
is trimmed to the columns that exist rather than dropping every row.

Schema constants now live in telemetry/schema.py so the SQL the library
prints cannot drift from the DDL it runs, and telemetry_schema_sql is
exported for downstreams writing their own migrations. The PG write drops
its conflict target, which partitioned tables require and which issue #12
depends on.
This commit is contained in:
2026-08-19 13:24:05 -04:00
19 changed files with 1462 additions and 271 deletions
+9
View File
@@ -56,6 +56,15 @@ PGW_BREAKER_BACKEND=memory # memory | redis
PGW_CACHE_BACKEND=none # redis | memory | none(必填,显式优于隐式)
PGW_TELEMETRY_BACKEND=none # sqlite | postgres | none(必填)
# PGW_TELEMETRY_SQLITE_PATH=logs/telemetry.db # sqlite 时必填
# PGW_TELEMETRY_SCHEMA_MODE=manual # auto | manual;三态: 不设 = 按后端派生(sqlite→auto、postgres→manual),
# # 显式设置则两侧都可覆盖。auto = 库给已存在的旧表自动 ALTER 补列;
# # manual = 库不发 ALTER,只 warning 点名缺列并打印可执行 SQL,
# # 按现有列裁剪 INSERT 继续写(遥测不会因缺列而全线丢失)。
# # 缺省为何不对称: postgres 是共享生产表,ALTER 取 ACCESS EXCLUSIVE 锁,
# # 会排在长事务后阻塞该表其后的所有查询,而遥测是业务路径上的内联 await;
# # 且这类部署有 DBA、有迁移工具、讲最小权限,DDL 该由他们择时执行。
# # sqlite 则是下游自己的本地文件(runs/*.db):没有 DBA、没有迁移工具、
# # 没有第二个系统碰它,ALTER 是毫秒级元数据操作,强加手工 SQL 步骤是净损失。
# PGW_TELEMETRY_PG_DSN=postgresql://user:pass@host:5432/polygateway # postgres 时必填;严禁指向在用业务库(实验室约定: 专用库 polygateway)
# PGW_PRICING_PATH=config/prices.json # 可选: {"<model>": {"input_per_1m": x, "output_per_1m": y}};缺省 cost 恒 None
# # 可选第三档 "cached_input_per_1m": z —— 供应商 prompt cache 命中部分的单价;
+41
View File
@@ -1,5 +1,46 @@
# Changelog
## 未发布
遥测表 `llm_calls` 的结构变更从此**由下游掌控**(issue #13)。此前两个后端都会在初始化期对下游数据库发 DDL:表不存在则建表,表存在但缺列则逐列 `ALTER TABLE ADD COLUMN`,而补列**没有任何开关**——库一升级、下次调用即自动执行。在共享的生产 Postgres 上这有三重问题:`ALTER` 取 ACCESS EXCLUSIVE 锁会排在长事务后阻塞该表其后的所有查询(而遥测是业务路径上的内联 `await`),多进程多版本共存时谁先补列是竞态,且这些 DDL 不进任何迁移记录、事后无从审计。调研过的 11 个同类系统(Celery / APScheduler / Alembic / Django contrib / Hangfire / Quartz.NET / dbt / Airbyte / Fivetran / Prefect / Airflow)里没有一个把它作为默认行为。
### 请先读这一条:三处破坏性变更
| # | 变更 | 影响与应对 |
|---|---|---|
| ① | **Postgres 侧不再自动补列**(缺省转为 manual 档) | 库升级带来新列时,旧表不会被自动 `ALTER`:库改为发**一条** warning 点名缺失的维度并附上可直接执行的 SQL,同时按现有列裁剪 `INSERT` 继续写入——**缺的那几列静默不落库**,直到有人执行那几条 SQL。要恢复旧行为设 `PGW_TELEMETRY_SCHEMA_MODE=auto`。SQLite 侧缺省不变(仍 auto),理由见下 |
| ② | 两个 recorder 新增 **keyword-only 必填**参数 `auto_migrate` | `SQLiteRecorder(db_path, *, auto_migrate)` 与 `PostgresRecorder(dsn, *, pool=None, auto_migrate)`;直接构造 recorder 的调用点必须补这个参数,不传即 `TypeError`。**故意不给默认值**:缺省规则只写在 config 一处,不与类签名漂移 |
| ③ | `GatewaySettings` 新增**必填**字段 `telemetry_auto_migrate: bool` | 只影响「构造函数全量注入」这条装配路(测试/高级用法);`from_env()` / `from_settings()` 的用户零改动。`telemetry_backend="none"` 时该字段在 `__post_init__` 归一为 `False` |
### 新增
- **`PGW_TELEMETRY_SCHEMA_MODE`(可选键,值域 `auto` / `manual`)**,**三态**:不设 = 按后端派生,显式设置 = 两侧都可覆盖。派生规则**有意不对称**——`postgres``manual`,`sqlite``auto`。理由:PG 侧是共享的生产表,有 DBA、有迁移工具、讲最小权限,DDL 的执行时机该由他们挑;SQLite 侧是下游自己的本地文件(典型是 `runs/*.db`),没有 DBA、没有迁移工具、没有第二个系统碰它,`ALTER` 是毫秒级元数据操作,要求"升级后手工跑一条 SQL"是给零运维场景强加运维步骤。
- **公共函数 `telemetry_schema_sql(backend) -> str`**(已进顶层 `__all__`):返回可直接粘进迁移文件的完整脚本——注释头 + `CREATE TABLE IF NOT EXISTS`(全量列)+ 各补列语句。PG 变体带 `ADD COLUMN IF NOT EXISTS`,整段**可重复执行**;SQLite 无该语法,以注释标明"仅当该列不存在时执行"。非法 `backend``ValueError`
- **manual 档的缺列告警**逐列点名并写明后果(「以下维度不会被记录: tenant_id, meta」),附上可直接执行的 ALTER,且**只在准备期发一次**,不逐行刷屏。只说"缺列"是不够的:静默丢维度的后果是多租户账目全归空串且无任何报错。
### 变更
- **Postgres 的写入去掉了冲突目标**:`ON CONFLICT (call_id) DO NOTHING``ON CONFLICT DO NOTHING`。普通表上语义**逐字等价**(表上只有主键这一个唯一约束),但带目标的版本要求恰好匹配 `(call_id)` 的唯一约束,而 PostgreSQL 要求分区表的唯一约束必须包含分区键——按 `created_at` 分区后主键变成 `(call_id, created_at)`,该语句会被 PG 直接拒收,且失败只逐行 warning,表现为分区部署下遥测全线静默丢数据。SQLite 的 `INSERT OR IGNORE` 本就无目标,未动。
- **manual 档按现有列裁剪 `INSERT`**。这不是可选增强而是关掉 `ALTER` 的前提:旧表缺列时若仍发全量 `INSERT`,每一行都会因未知列被拒 → 遥测彻底丢失,比自动补列更严重地违反「遥测必录」。列探测失败、或探测结果与库认识的列毫无交集时,保守回落全量列(与今天的行为一致)。
- **schema 常量收敛为单一事实源** `telemetry/schema.py`(内部模块):列序、两端 DDL、两端补列语句、`INSERT` 构造与缺列告警此前在两个 recorder 各存一份。收敛的理由是**正确性**而非整洁——打印给下游的 SQL 必须与库真正执行的 DDL 同源,多处各存一份必然漂移,而漂移的表现是"下游照打印的 SQL 建完表,库仍报缺列"。
### 不变
- **manual 档仍然建表**。issue 把建表列为现状描述而非指控(它已在 #9 收口为"PG 侧先 `to_regclass` 探测、表在就不发 DDL")。新建表没有既有数据、没有并发访问者,不存在锁队列与数据风险,而停掉它会让"零配置起步"这条路彻底断掉。
- **auto 档行为与从前逐字相同**,包括补列失败时**不裁剪**:该档承诺的是"把列补上",补不上就让缺列以逐行 warning 暴露;要降级写入请显式选 manual。
- 降级方向不变:缺列、补列失败、写入失败一律只 warning,绝不冒泡打断业务调用;列名与列序不变;错误面零变更。
### 库对下游数据库的承诺(Expand/Contract,本版成文)
以下五条此前已被实现满足,但从未写成承诺。本版起它们是**承诺**:新列**只增不删不改名**且一律追加在既有列之后;新列必**可空**或带**非易失常量默认值**(PG 11+ 补列不重写全表,SQLite 补列是元数据操作);`INSERT` **永远显式写出列名**;库**从不 `SELECT *`**、从不读回这张表的数据(库只写不读,连探测都只查 catalog);写入的**冲突处理不绑定具体约束**。
合起来它们保证:你可以自行给 `llm_calls` 加列、加索引、挂 RLS,乃至把它建成 `PARTITION BY RANGE (created_at)` 的分区表,库的探测、补列与写入都照常工作。完整说明见 README「遥测表 schema 与升级纪律」——那份随包分发,`research-wiki/` 不在 sdist 内。
### 升级提示
-`from_env()` / `from_settings()` 装配的下游**无需改代码**;Postgres 下游升级后建议执行一次 `python -c "import polygateway; print(polygateway.telemetry_schema_sql('postgres'))"` 的输出,把新列补齐(不补则新维度不落库,库会在首次写入前用一条 warning 点名)。
- 直接构造 `SQLiteRecorder` / `PostgresRecorder` 或直接构造 `GatewaySettings` 的调用点必须补上新参数/新字段,否则 `TypeError`
## 1.2.1(2026-08-18)
每次调用现在可以带上**租户标识与任意调用方自定义维度**,并逐条落进遥测表(issue #11)。`llm_calls` 存的是**完整正文**(`digest_messages` 只对多模态 `image_url` 做 sha256,纯文本原样透传),多租户下游的合同与标书全文因此混在同一张表里,而原先的 22 列**没有任何租户维度**——能区分来源的只有 `session_id` / `parent_call_id` 两个调用方自填、库内不校验的自由字符串。
+61 -2
View File
@@ -125,7 +125,7 @@ resp = await client.chat(
**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 审出还有多少行待归属)。
存储上 `tenant_id` 两端都是 `TEXT NOT NULL DEFAULT ''`,`meta` 在 Postgres 是 `JSONB`、在 SQLite 是 `TEXT`;老表要补上这两列(补列是否由库自动执行取决于 `PGW_TELEMETRY_SCHEMA_MODE`,见[遥测表 schema 与升级纪律](#遥测表-schema-与升级纪律)),**补列后老行读出是空串而非 NULL**(NULL 在任何 RLS policy 下都对所有人不可见,空串则可用一条 SQL 审出还有多少行待归属)。
**库只提供列,不启用 RLS、不建索引。** 要数据库层的强制隔离,以下 DDL 是**下游 DBA 的职责,库不会代劳**;不执行则 `tenant_id` 只是一个可查可过滤的普通列,没有任何数据库层强制。库不代劳的原因是 default-deny:启用 RLS 而没有匹配的 policy = 零行可写且静默不报错,会让非多租户部署的遥测全量写失败。
@@ -152,6 +152,64 @@ CREATE INDEX CONCURRENTLY idx_llm_calls_tenant_created
| 租户上下文必须在**显式事务内**用 `set_config('app.tenant_id', ..., true)` | asyncpg 默认 autocommit,单发 `SET LOCAL` 会当场失效,而 PG **只发 warning 不报错**;表现是 policy 永远拿不到租户 → fail-closed 到零行 |
| policy 必须同时写 `USING``WITH CHECK` | 只写前者则租户 A 读不到 B 的行,却**能插入标着 B 的行**——污染发生在写入侧,读侧查不出来 |
## 遥测表 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 **同源**(同一份常量),照它建完表,库探测到的列就是齐的:
```python
import polygateway
print(polygateway.telemetry_schema_sql("postgres")) # 或 "sqlite";非法值抛 ValueError
```
```bash
# 直接落成迁移文件:注释头 + 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 要求分区表唯一约束含分区键),库的探测、补列与写入照常工作 |
## 错误模型(四分类)
一切失败在 transport 层翻译为四类之一,治理行为由分类决定,业务侧不需要判断状态码:
@@ -195,6 +253,7 @@ CREATE INDEX CONCURRENTLY idx_llm_calls_tenant_created
| `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 与升级纪律](#遥测表-schema-与升级纪律) |
| `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)。
@@ -224,7 +283,7 @@ graph LR
| `telemetry/` | SQLite / Postgres 遥测后端 |
| `structured/` | 结构化输出策略 |
依赖纪律由 import-linter 机械化执法(`make lint`)。完整架构决策(D1-D14 含论证过程)见 [research-wiki/ARCHITECTURE.md](research-wiki/ARCHITECTURE.md)。
依赖纪律由 import-linter 机械化执法(`make lint`)。完整架构决策(D1-D15 含论证过程)见 [research-wiki/ARCHITECTURE.md](research-wiki/ARCHITECTURE.md)。
## 可靠性证据
+18 -1
View File
@@ -119,7 +119,7 @@ HTTP API → arq 队列 → worker 协程 脚本 → asyncio.gather 协
---
## 3. 架构决策记录(D1D14,含讨论过程与备选方案)
## 3. 架构决策记录(D1D15,含讨论过程与备选方案)
> 每条决策记录格式:**决策 / 背景与讨论 / 被否决的备选 / 影响**。这些决策已与人类逐条确认;推翻任何一条需要人类批准并修订本节。
@@ -248,6 +248,20 @@ HTTP API → arq 队列 → worker 协程 脚本 → asyncio.gather 协
**影响**: §5.2 `structured` 参数三档语义、§7.9 重写为阶梯、§6.1 ResultInvalid 行注 D14;缓存写入发生在阶梯通过之后(§7.5 "不固化坏结果"的执行点);反馈模板与策略升级细则留 M1 设计文档。
### D15 库对下游数据库只做 SELECT/INSERT + 可选 CREATE;改结构与删数据归下游(2026-08-19,issue #13)
**决策**: 遥测表 `llm_calls` 是**下游的表**,不是库的私有存储。库对它发出的语句只有三类——catalog 探测(PG `to_regclass` + `pg_attribute`,SQLite `PRAGMA table_info`)、显式列名的 `INSERT`、以及表不存在时的 `CREATE TABLE IF NOT EXISTS`;**改结构(`ALTER`)与删数据(`UPDATE`/`DELETE`/`TRUNCATE`/`DROP`)一律归下游**。`ALTER` 保留唯一一个受控出口:`PGW_TELEMETRY_SCHEMA_MODE=auto` 时给已存在的旧表补列,而该档在 PG 侧**不是缺省**(缺省按后端派生: sqlite→auto、postgres→manual)。配套五条 Expand/Contract 承诺:新列只增不删不改名且追加在既有列之后、新列必可空或带非易失常量默认值、`INSERT` 永远显式列名、库从不 `SELECT *` 也从不读回该表数据、写入的冲突处理不绑定具体约束。
**背景与讨论**: 补列此前没有任何开关,库一升级、下次调用即在下游生产库上发 DDL。issue #13 的三条指控成立: ① 与最小权限原则冲突;② 多进程/多版本共存时谁先补列是竞态;③ DDL 不进任何迁移记录,DBA 事后无从审计。量级判据是 `ALTER TABLE ADD COLUMN` 取 ACCESS EXCLUSIVE 锁,会排在长事务后阻塞该表其后的所有查询,而遥测是业务路径上的内联 `await`。调研的 11 个同类系统(Celery / APScheduler / Alembic / Django contrib / Hangfire / Quartz.NET / dbt / Airbyte / Fivetran / Prefect / Airflow)中**没有一个**把"库在下游库里自动 ALTER 出列"作为默认行为。
两条边界是讨论出来的、不是照抄先例: **① 缺省按后端不对称**(D-a,人类拍板)——issue 引用的全部先例语境都是共享的生产 PG,而本库的 SQLite 侧是下游自己的本地文件(没有 DBA、没有迁移工具、没有第二个系统碰它,`ALTER` 是毫秒级元数据操作),两侧统一 manual 会给零运维场景强加运维步骤;两侧有意不对称在本库已有先例(§7.8 的建表探测,issue #9)。**② manual 档不连 `CREATE TABLE` 一起停**——新建表没有既有数据与并发访问者,不存在锁队列与数据风险,停掉它会让"零配置起步"断掉(Celery 的先例同样是"自动建表 + 永不 ALTER")。**③ 关掉 `ALTER` 必须配套按现有列裁剪 `INSERT`**,否则旧表缺列时每行写入都被拒,是把自动补列换成静默全失能,比原问题更严重地违反「遥测必录」。
五条承诺本身是既有实现的**成文化**(零代码变更),但成文后才可被下游依赖——它同时是遥测保留期方案(issue #12)能成立的前提: 下游拿这份 schema 自己加 `PARTITION BY RANGE (created_at)` 建成分区表后,库的 `to_regclass` 探测、列探测与 `INSERT` 路由都照常工作。第五条(冲突处理不绑定约束)是审查带出的**新增**承诺,并伴随一处真实修复,见 §7.8。
**被否决的备选**: 两侧统一缺省 manual(语义最一致,但现有 SQLite 下游升级即需人工干预,而这些场景根本没有承接手工 SQL 的角色);保持 auto 缺省只加关闭档(默认状态仍是"库在下游生产表上发不受控 DDL",issue 的核心诉求未被满足);Celery 式"自动建表但永不 ALTER、无开关"(SQLite 场景纯净损失,且真想要自动补列的下游没有出路);APScheduler 4.x 式"schema 不认识就拒绝启动"(与「遥测初始化失败必须静默降级」的库铁律正面冲突,不可选)。
**影响**: §7.8 补列一节按档位重写;新增配置键 `PGW_TELEMETRY_SCHEMA_MODE`(§9)与公共函数 `telemetry_schema_sql`;两个 recorder 新增 keyword-only 必填参数 `auto_migrate``GatewaySettings` 新增必填字段 `telemetry_auto_migrate`(缺省规则只写在 config 一处,不与类签名漂移);五条承诺进 README(随包分发)。
---
## 4. 总体架构
@@ -495,6 +509,8 @@ flowchart TB
(`cached_prompt_tokens`/`model_reported` 为 2026-07-31 issue #3 新增,端口由 18 字段扩为 20;两个后端在初始化期对已存在的旧表幂等补列——`CREATE TABLE IF NOT EXISTS` 不会给旧表加列,不补则每行写入都被逐行 warning 丢弃。补列一律**先探测缺列再 ALTER**(`ADD COLUMN IF NOT EXISTS` 即使列已存在也先取 ACCESS EXCLUSIVE 锁,而遥测内联 await,锁共享审计表会拖垮业务调用),且**失败只逐行降级、绝不置结构性失能标志**。**建表同理(2026-08-07,issue #9)**: PG 对 schema 的 CREATE 权限检查早于 `IF NOT EXISTS` 的存在性判断(16.14 实测,只授表级 `SELECT, INSERT` 的角色写得进去却建不了表),故 PG 侧必须**先 `to_regclass` 探测、表在就不发 DDL**;SQLite 侧实测在解析期即短路(持排他锁/只读文件下该语句均通过),无同款风险,**有意不加探测**。由此把"结构性失能"的判据从「初始化时出过异常」收窄为「确定写不进去」——仅建池失败与"表确定不存在且建不出来"判死,探测/取连接失败只跳过本次并留待下次重试。新列在 DDL 里必须排在 `created_at` **之后**,与 `ALTER TABLE ADD COLUMN` 的追加位置一致,否则新建库与升级库的物理列序分叉)。链路: `session_id`/`parent_call_id` 由调用方传入贯穿(agent step → LLM call)。`messages` 落库前对多模态 part 先摘要(与缓存 key 共用同一摘要函数,§7.5)——Video-Tree 现状 base64 整段进 SQLite 导致 db 膨胀(`llm.py:330`),库内修复(2026-07-20,VT 迁移缺口 R12)。
**schema 单一事实源、档位与冲突目标(2026-08-19,issue #13,决策见 D15)**: 列序、两端 DDL、两端补列语句、`INSERT` 构造与缺列告警收敛进 `telemetry/schema.py`——此前在两个 recorder 各存一份,而公共函数 `telemetry_schema_sql` 打印给下游的 SQL 必须与库真正执行的 DDL **同源**,三份必然漂移,漂移的表现是"下游照打印的 SQL 建完表,库仍报缺列"。补列自此由 `PGW_TELEMETRY_SCHEMA_MODE` 控制(三态: 不设按后端派生 sqlite→auto / postgres→manual,显式设置两侧均可覆盖): manual 档一条 DDL 都不发,改为按探测到的现有列**裁剪 `INSERT`**(裁剪是关掉 ALTER 的前提,否则缺列旧表每行写入都被拒 = 遥测全失)并发**一条**点名缺列、附可执行 SQL 的 warning;auto 档行为不变,且补列失败时**不裁剪**(该档承诺"把列补上",补不上就让缺列以逐行 warning 暴露)。**库内执行的补列语句与打印给人的那份是两套文本**: 库内不用 `ADD COLUMN IF NOT EXISTS`(它即便列已存在也先取 ACCESS EXCLUSIVE 锁,故库侧一律先探测后 ALTER),打印的那份带,以保证下游可重复执行。同批把 PG 写入的 `ON CONFLICT (call_id) DO NOTHING` 改为**无冲突目标**的 `ON CONFLICT DO NOTHING`: 带目标的语句要求恰好匹配 `(call_id)` 的唯一约束,而 PG 要求分区表的唯一约束必须包含分区键——按 `created_at` 分区(issue #12)后主键变成 `(call_id, created_at)`,该语句被 PG 直接拒收,而写失败只逐行 warning,表现为分区部署下遥测全线静默丢数据;无目标版本在两种表形态上都合法,普通表上语义逐字等价(表上只有主键这一个唯一约束),SQLite 的 `INSERT OR IGNORE` 本就无目标。
- 后端: `SQLiteRecorder`(默认;WAL + busy_timeout、`INSERT OR IGNORE` 幂等、`asyncio.to_thread` 桥接、初始化/写入失败全降级不冒泡)与 `PostgresRecorder`
- **单一 helper 铁律**: 遥测调用点收敛为一个内部函数/上下文管理器;Video-Tree 与 GovDoc 各有 4-5 处逐字复制的 `record_llm_call(15 个参数)` 是本条的直接教训。
- 成本: `pricing.py` 维护 model → (input 单价, output 单价, **可选** cached_input 单价) 表,遥测时换算 `cost` 字段;查不到价格记 None 并 warning,**不阻塞调用**。缓存读取单价(2026-07-31,issue #3)只在配置了该档且本次有命中时启用,按 `(prompt - cached) × input + cached × cached_input` 分段计价;**未配该档绝不按经验折扣率猜**,退化为全额输入价(P5)。命中数超过输入总数时按总数夹取并 warning,不产生负成本。
@@ -558,6 +574,7 @@ src/polygateway/
- **per-scope 韧性配置(2026-07-20,CHS 迁移缺口 G4)**: 韧性参数支持按 scope 覆盖——`{SCOPE}__RETRY__MAX_ATTEMPTS` / `{SCOPE}__BREAKER__FAIL_THRESHOLD` / `{SCOPE}__BREAKER__COOLDOWN_S` / `{SCOPE}__BACKPRESSURE__STALL_WINDOW_S` / `{SCOPE}__SELECTOR` / `{SCOPE}__GLOBAL__MAX_CONCURRENCY|RPM|TPM`(CHS 现状: VLM 与 OCR 两 scope 参数各异)。平铺键(`LLM_*`)是单 scope 场景的简写;两者并存时 scope 键优先。
- **装配只有两条路**: `GatewayClient.from_env()`/`from_settings(settings)`(工厂,覆盖 90% 用户;补上三项目每次手写、GovDoc 缺失的"配置→client"一段)或构造函数全量依赖注入(测试/高级用户)。库内部任何组件**不得自读环境变量**(显式优于隐式)。
- 后端选择即配置: 如 `PGW_LIMITER_BACKEND=memory|redis``PGW_TELEMETRY_BACKEND=sqlite|postgres``PGW_QUOTA_FULL=wait|fail_fast`(命名待 M1 设计文档定稿)。
- **`PGW_TELEMETRY_SCHEMA_MODE=auto|manual`(2026-08-19,issue #13,D15)**: 可选键、**三态**——不设 = 按后端派生(sqlite→auto、postgres→manual),显式设置则两侧都可覆盖。派生只发生在 config 层一处,落到 `GatewaySettings.telemetry_auto_migrate`(无默认值,与既有全部字段一致;`telemetry_backend=none` 时无人消费,归一为 `False`),recorder 的 `auto_migrate` 是 keyword-only **必填**参数——关键行为参数不给默认值(P4),缺省规则也就不会与类签名漂移。
---
+2
View File
@@ -23,6 +23,7 @@ from polygateway.errors import (
from polygateway.ocr import OcrClient
from polygateway.pricing import ModelPrice, PricingTable
from polygateway.providers import DEFAULT_PROFILES, ProviderProfile, register_provider
from polygateway.telemetry.schema import telemetry_schema_sql
from polygateway.types import (
EmbeddingResponse,
LLMResponse,
@@ -64,4 +65,5 @@ __all__ = [
"__version__",
"gather_bounded",
"register_provider",
"telemetry_schema_sql",
]
+15 -5
View File
@@ -367,7 +367,9 @@ class RedisGate:
keys=[self._key(source_name)], args=[owner, self._probe_ttl_ms]
)
except RedisError as exc:
raise GovernanceBackendError(f"熔断后端 try_enter 失败: {exc}", scope=self._scope) from exc
raise GovernanceBackendError(
f"熔断后端 try_enter 失败: {exc}", scope=self._scope
) from exc
return self._decision(source_name, result)
async def record_success(
@@ -385,7 +387,9 @@ class RedisGate:
try:
result = await self._success_lua(keys=[self._key(entry.source_name)], args=args)
except RedisError as exc:
raise GovernanceBackendError(f"熔断后端 record_success 失败: {exc}", scope=self._scope) from exc
raise GovernanceBackendError(
f"熔断后端 record_success 失败: {exc}", scope=self._scope
) from exc
return self._update(result)
async def record_failure(
@@ -407,7 +411,9 @@ class RedisGate:
try:
result = await self._failure_lua(keys=[self._key(entry.source_name)], args=args)
except RedisError as exc:
raise GovernanceBackendError(f"熔断后端 record_failure 失败: {exc}", scope=self._scope) from exc
raise GovernanceBackendError(
f"熔断后端 record_failure 失败: {exc}", scope=self._scope
) from exc
return self._update(result)
async def release_probe(self, entry: GateDecision) -> GateUpdate:
@@ -419,7 +425,9 @@ class RedisGate:
keys=[self._key(entry.source_name)], args=[entry.epoch, entry.probe_owner]
)
except RedisError as exc:
raise GovernanceBackendError(f"熔断后端 release_probe 失败: {exc}", scope=self._scope) from exc
raise GovernanceBackendError(
f"熔断后端 release_probe 失败: {exc}", scope=self._scope
) from exc
return self._update(result)
async def retry_after_s(self, sources: tuple[str, ...]) -> float:
@@ -429,7 +437,9 @@ class RedisGate:
try:
result = await self._retry_after_lua(keys=[self._key(s) for s in sources])
except RedisError as exc:
raise GovernanceBackendError(f"熔断后端 retry_after_s 失败: {exc}", scope=self._scope) from exc
raise GovernanceBackendError(
f"熔断后端 retry_after_s 失败: {exc}", scope=self._scope
) from exc
return int(result) / 1000.0
async def aclose(self) -> None:
+15 -5
View File
@@ -247,7 +247,9 @@ class RedisLimiter:
],
)
except RedisError as exc:
raise GovernanceBackendError(f"限流后端 try_acquire 失败: {exc}", scope=self._scope) from exc
raise GovernanceBackendError(
f"限流后端 try_acquire 失败: {exc}", scope=self._scope
) from exc
if ok != 1:
return None
return _RedisPermit(self, source_key, lease_id, est_tokens, window)
@@ -265,7 +267,9 @@ class RedisLimiter:
try:
await self._release_lua(keys=[gl, sl], args=[lease_id])
except RedisError as exc:
raise GovernanceBackendError(f"限流后端 release 失败: {exc}", scope=self._scope) from exc
raise GovernanceBackendError(
f"限流后端 release 失败: {exc}", scope=self._scope
) from exc
async def _settle_tpm(self, source_key: str, delta: int, window: int) -> None:
wk = self._window_keys(source_key, window)
@@ -283,7 +287,9 @@ class RedisLimiter:
wk = self._window_keys(source_key, window)
res = await self._stats_lua(keys=[sl, wk["s_rpm"], wk["s_tpm"]])
except RedisError as exc:
raise GovernanceBackendError(f"限流后端 source_stats 失败: {exc}", scope=self._scope) from exc
raise GovernanceBackendError(
f"限流后端 source_stats 失败: {exc}", scope=self._scope
) from exc
return SourceStats(
inflight=int(res[0]),
rpm_used=max(0, int(res[1])),
@@ -295,14 +301,18 @@ class RedisLimiter:
try:
await self._progress_mark_lua(keys=[self._progress_key()], args=[_PROGRESS_TTL_S])
except RedisError as exc:
raise GovernanceBackendError(f"限流后端 mark_progress 失败: {exc}", scope=self._scope) from exc
raise GovernanceBackendError(
f"限流后端 mark_progress 失败: {exc}", scope=self._scope
) from exc
async def progress_age_s(self) -> float:
"""距上次全局成功的秒数;仅键缺失(-1)= 从未进展 → inf(CHS limiter.py:208)。"""
try:
res = await self._progress_age_lua(keys=[self._progress_key()])
except RedisError as exc:
raise GovernanceBackendError(f"限流后端 progress_age_s 失败: {exc}", scope=self._scope) from exc
raise GovernanceBackendError(
f"限流后端 progress_age_s 失败: {exc}", scope=self._scope
) from exc
return float("inf") if int(res) == -1 else int(res) / 1000.0
async def aclose(self) -> None:
+6 -2
View File
@@ -400,11 +400,15 @@ def _build_telemetry(settings: GatewaySettings) -> TelemetryRecorder | None:
from polygateway.telemetry.postgres import PostgresRecorder
assert settings.telemetry_pg_dsn is not None # 内部不变量: _validate_telemetry 已保证
return PostgresRecorder(settings.telemetry_pg_dsn)
return PostgresRecorder(
settings.telemetry_pg_dsn, auto_migrate=settings.telemetry_auto_migrate
)
from polygateway.telemetry.sqlite import SQLiteRecorder
assert settings.telemetry_sqlite_path is not None # 内部不变量: _validate_telemetry 已保证
return SQLiteRecorder(settings.telemetry_sqlite_path)
return SQLiteRecorder(
settings.telemetry_sqlite_path, auto_migrate=settings.telemetry_auto_migrate
)
def _build_structured(
+45
View File
@@ -55,6 +55,9 @@ _LIMITER_BACKENDS = frozenset({"memory", "redis"})
_BREAKER_BACKENDS = frozenset({"memory", "redis"})
_CACHE_BACKENDS = frozenset({"redis", "memory", "none"})
_TELEMETRY_BACKENDS = frozenset({"sqlite", "postgres", "none"})
# 遥测 schema 档位(issue #13): auto 允许 recorder 给旧表 ALTER 补列,manual 不发 DDL
_SCHEMA_MODES = frozenset({"auto", "manual"})
_SCHEMA_MODE_KEY = "PGW_TELEMETRY_SCHEMA_MODE"
_REDIS_DEPENDENT_BACKENDS = ("limiter_backend", "breaker_backend", "cache_backend")
# 背压默认(M1 仅 poll 生效;CHS _BACKOFF_S=0.05 同源)
_DEFAULT_STALL_WINDOW_S = 300.0
@@ -131,6 +134,11 @@ class GatewaySettings:
telemetry_backend: str
telemetry_sqlite_path: str | None
telemetry_pg_dsn: str | None
# 是否允许 recorder 给已存在的旧表自动 ALTER 补列(issue #13);env 的三态
# 派生只写在 `_load_schema_mode` 一处,不与 recorder 的类签名漂移。
# backend=none 时恒 False 这条跨字段不变量则由 `_validate_telemetry`
# 把关,对直接构造与 `dataclasses.replace` 同样生效
telemetry_auto_migrate: bool
redis_url: str | None
pricing_path: str | None
structured_max_retries: int
@@ -211,7 +219,14 @@ class GatewaySettings:
剥而不是拒: 两条装配路对同一 DSN 应产出同一结果。但不静默——`from_env`
那条路在 `_load_pg_dsn` 就剥干净了,能走到这里的只有手工构造的调用方,
他有权知道库动了他给的值。
`telemetry_auto_migrate` 同理归一化而非报错: backend=none 时根本没有
recorder 消费它,True 是个自相矛盾却无害的状态。`from_env` 那条路的派生
已经给出 False,归一化是为了直接构造与 `dataclasses.replace` 也一致——
不变量挂在构造期,才不用每加一个装配工厂就多一处要同步。
"""
if self.telemetry_backend == "none" and self.telemetry_auto_migrate:
object.__setattr__(self, "telemetry_auto_migrate", False)
if self.telemetry_backend == "sqlite" and not self.telemetry_sqlite_path:
raise ValueError("telemetry_backend=sqlite 时必须提供 telemetry_sqlite_path")
if self.telemetry_backend != "postgres":
@@ -434,6 +449,7 @@ def _load_pgw(env: Mapping[str, str]) -> dict[str, object]:
redis_url = env.get("REDIS_URL") or None
if "redis" in (limiter_backend, breaker_backend) and redis_url is None:
raise ValueError("缺关键配置: 限流/熔断后端取 redis 需设置 REDIS_URL")
auto_migrate = _load_schema_mode(env, telemetry_backend)
return {
"limiter_backend": limiter_backend,
"breaker_backend": breaker_backend,
@@ -444,6 +460,7 @@ def _load_pgw(env: Mapping[str, str]) -> dict[str, object]:
if telemetry_backend == "sqlite"
else None,
"telemetry_pg_dsn": _load_pg_dsn(env) if telemetry_backend == "postgres" else None,
"telemetry_auto_migrate": auto_migrate,
"redis_url": redis_url,
"pricing_path": env.get("PGW_PRICING_PATH") or None,
"structured_max_retries": _load_structured_retries(env),
@@ -451,6 +468,34 @@ def _load_pgw(env: Mapping[str, str]) -> dict[str, object]:
}
def _load_schema_mode(env: Mapping[str, str], telemetry_backend: str) -> bool:
"""把 `PGW_TELEMETRY_SCHEMA_MODE` 的三态解成 `telemetry_auto_migrate`(issue #13)。
三态: 键未设 → 按后端**不对称**派生;显式 auto/manual → 两侧都可覆盖。
不对称的理由是两个后端的风险量级不同: SQLite 是下游自己的本地文件(没有
DBA、没有迁移工具、没有第二个系统碰它),ALTER 是毫秒级元数据操作,要求
手工跑 SQL 是给零运维场景强加运维步骤;PG 是共享的生产表,ALTER 取
ACCESS EXCLUSIVE 锁会排在长事务后阻塞该表其后的所有查询,而遥测是业务
路径上的内联 await。
`_load_choice` 带 default,不能直接用来读这个键——default 会把"未设"
"设成默认值"抹平成同一种,三态就塌回两态,后端派生也就再没机会生效。故
先用 `_first` 探"设没设",确认设了才交给 `_load_choice` 做值域校验(错误
信息点出 env 键名这件事仍由它负责)。
Args:
env: 已合并的环境映射。
telemetry_backend: 已校验过值域的遥测后端名。
Returns:
recorder 是否获准给旧表自动 ALTER 补列;backend=none 时无人消费,
构造期守卫会再把它归一化为 False。
"""
if _first(env, _SCHEMA_MODE_KEY) is None:
return telemetry_backend == "sqlite"
return _load_choice(env, _SCHEMA_MODE_KEY, _SCHEMA_MODES, "auto") == "auto"
def _strip_dsn_driver(dsn: str) -> str:
"""剥 SQLAlchemy 风格的 `+driver` 后缀(asyncpg 不认);已干净的原样返回。"""
scheme, sep, rest = dsn.partition("://")
+101 -104
View File
@@ -21,56 +21,17 @@ from typing import TYPE_CHECKING
from loguru import logger
from polygateway.telemetry.schema import (
COLUMNS,
PG_BACKFILL,
PG_DDL,
insert_sql,
missing_columns_warning,
)
if TYPE_CHECKING:
import asyncpg
_DDL = """
CREATE TABLE IF NOT EXISTS llm_calls (
call_id TEXT PRIMARY KEY,
parent_call_id TEXT,
session_id TEXT,
model TEXT NOT NULL,
provider TEXT NOT NULL,
source_name TEXT NOT NULL,
messages TEXT NOT NULL,
response TEXT NOT NULL,
thinking TEXT NOT NULL DEFAULT '',
prompt_tokens INTEGER NOT NULL,
completion_tokens INTEGER NOT NULL,
usage_source TEXT NOT NULL,
latency_ms INTEGER NOT NULL,
ttft_ms DOUBLE PRECISION,
max_inter_token_ms DOUBLE PRECISION,
cache_hit BOOLEAN NOT NULL DEFAULT FALSE,
error TEXT,
cost DOUBLE PRECISION,
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
cached_prompt_tokens INTEGER,
model_reported TEXT,
sampling TEXT,
reasoning_tokens INTEGER,
tenant_id TEXT NOT NULL DEFAULT '',
meta JSONB NOT NULL DEFAULT '{}'::jsonb
);
"""
# 新列排在 created_at 之后: 与旧表 ALTER 追加的位置一致(见 sqlite.py 同款注释)
_BACKFILL = (
("cached_prompt_tokens", "ALTER TABLE llm_calls ADD COLUMN cached_prompt_tokens INTEGER"),
("model_reported", "ALTER TABLE llm_calls ADD COLUMN model_reported TEXT"),
("sampling", "ALTER TABLE llm_calls ADD COLUMN sampling TEXT"),
("reasoning_tokens", "ALTER TABLE llm_calls ADD COLUMN reasoning_tokens INTEGER"),
# 两个默认值都是非易失常量,PG 11+ 只改 catalog 不重写全表,故大表补列亦是秒级
(
"tenant_id",
"ALTER TABLE llm_calls ADD COLUMN tenant_id TEXT NOT NULL DEFAULT ''",
),
(
"meta",
"ALTER TABLE llm_calls ADD COLUMN meta JSONB NOT NULL DEFAULT '{}'::jsonb",
),
)
# 探测表是否存在;不需要任何权限,且与 INSERT 走同一套 search_path 解析
_TABLE_EXISTS = "SELECT to_regclass('llm_calls')"
@@ -80,44 +41,22 @@ _EXISTING_COLUMNS = (
"WHERE attrelid = to_regclass('llm_calls') AND attnum > 0 AND NOT attisdropped"
)
_COLUMNS = (
"call_id",
"parent_call_id",
"session_id",
"model",
"provider",
"source_name",
"messages",
"response",
"thinking",
"prompt_tokens",
"completion_tokens",
"usage_source",
"latency_ms",
"ttft_ms",
"max_inter_token_ms",
"cache_hit",
"error",
"cost",
"cached_prompt_tokens",
"model_reported",
"sampling",
"reasoning_tokens",
"tenant_id",
"meta",
)
_INSERT = (
f"INSERT INTO llm_calls ({', '.join(_COLUMNS)}) "
f"VALUES ({', '.join(f'${i + 1}' for i in range(len(_COLUMNS)))}) "
"ON CONFLICT (call_id) DO NOTHING"
)
class PostgresRecorder:
"""TelemetryRecorder 端口的 Postgres 实现;asyncpg 原生异步,无线程桥接。"""
def __init__(self, dsn: str, *, pool: asyncpg.Pool | None = None) -> None:
def __init__(self, dsn: str, *, pool: asyncpg.Pool | None = None, auto_migrate: bool) -> None:
"""记下装配参数(不连库);列与 INSERT 语句在首次准备期定型。
Args:
dsn: asyncpg 连接串(已剥驱动后缀)。
pool: 外部注入的池;注入方自己负责关闭。
auto_migrate: True 则给已存在的旧表自动补列;False(PG 侧的缺省档)
则一条 ALTER 都不发——`ALTER TABLE ADD COLUMN` 取 ACCESS EXCLUSIVE
锁,会排在长事务后阻塞该表其后所有查询,而遥测是业务路径上的内联
await。keyword-only **必填**: 缺省规则只写在 config 一处,不与本类
签名漂移(设计 D-c)。
"""
try:
import asyncpg # noqa: F401 - 仅探测 extra 是否安装
except ImportError as exc:
@@ -127,6 +66,10 @@ class PostgresRecorder:
self._dsn = dsn
self._pool: asyncpg.Pool | None = pool
self._external_pool = pool is not None
self._auto_migrate = auto_migrate
# 先按全量列定型: 准备期探测失败时保守沿用全量(今天的行为)
self._columns: tuple[str, ...] = COLUMNS
self._insert = insert_sql("postgres", COLUMNS)
self._schema_ready = False
self._failed = False # 结构性降级标志: 置位后所有写入短路
self._init_lock = asyncio.Lock()
@@ -169,7 +112,7 @@ class PostgresRecorder:
"""备好表并交回可用的池;瞬时失败只跳过本次,确定写不进去才判死。"""
try:
async with pool.acquire() as conn:
writable = await self._prepare_table(conn)
columns = await self._prepare_table(conn)
except asyncio.CancelledError:
raise
except Exception as exc:
@@ -177,14 +120,20 @@ class PostgresRecorder:
# 只跳过本次记录,下次调用重新准备
logger.warning("Postgres 遥测建表探测失败(跳过本条,下次重试): {}", exc)
return None
if not writable:
if columns is None:
self._failed = True
return None
# 写入列、语句与就绪标志必须**一起**生效: `_ensure_ready` 只看 `_schema_ready`
# 就绕开 `_init_lock` 直接返回池,先置就绪会开出"已就绪但语句还是旧的"的窗口
self._columns = columns
self._insert = insert_sql("postgres", columns)
self._schema_ready = True
return pool
async def _prepare_table(self, conn: object) -> bool:
"""备好 `llm_calls`;**表存在就绝不发 DDL**。返回 False 仅表示表确定不存在。
async def _prepare_table(self, conn: object) -> tuple[str, ...] | None:
"""备好 `llm_calls` 并返回本实例要写的列;**表存在就绝不发 DDL**。
返回 None 仅表示表确定不存在且建不出来(唯一允许判死的情形)。
`CREATE TABLE IF NOT EXISTS` 不能无条件发: PostgreSQL 对 schema 的
CREATE 权限检查**早于** `IF NOT EXISTS` 的存在性判断(PG 16.14 实测:
@@ -197,33 +146,77 @@ class PostgresRecorder:
"""
exists = await conn.fetchval(_TABLE_EXISTS) is not None # type: ignore[attr-defined]
if exists:
await self._backfill_columns(conn) # 旧表可能缺列;失败只逐行降级
return True
return await self._resolve_columns(conn) # 旧表可能缺列
try:
await conn.execute(_DDL) # type: ignore[attr-defined]
await conn.execute(PG_DDL) # type: ignore[attr-defined]
except asyncio.CancelledError:
raise
except Exception as exc:
logger.warning("Postgres 遥测建表失败(表不存在,记录无处可落): {}", exc)
return False
return True # 新建表列已齐全,无需再走补列
return None
return COLUMNS # 新建表列已齐全,无需再走补列
async def _backfill_columns(self, conn: object) -> None:
"""给已存在的旧表补新列(issue #3);**先探测再 ALTER,失败绝不置 `_failed`**
async def _resolve_columns(self, conn: object) -> tuple[str, ...]:
"""探测旧表现有列并定型写入列: auto 档先补齐,manual 档改为裁剪(issue #13)。
两条纪律各有实测理由:
① 不置 `_failed`: 应用账号只有 INSERT 权限时,`ALTER TABLE` 的 ownership
检查早于 `IF NOT EXISTS` 的存在性判断——列明明齐全也会失败。置位会让
整个 recorder 永久 no-op,与「补列失败只降级为逐行丢弃」的承诺相悖
(SQLite 侧同款守卫,两侧必须对称)。
② 先探测: `ADD COLUMN IF NOT EXISTS` 即便列已存在,也会**先取 ACCESS
EXCLUSIVE 锁**再判存在性(实测会被一个开着的读事务阻塞)。遥测是内联
await,让每个进程的首次写入都去抢共享审计表的排他锁,等于用记录基础设施
**先探测**的理由(两档共用): `ADD COLUMN IF NOT EXISTS` 即便列已存在,也会
**先取 ACCESS EXCLUSIVE 锁**再判存在性(实测会被一个开着的读事务阻塞)。遥测是
内联 await,让每个进程的首次写入都去抢共享审计表的排他锁,等于用记录基础设施
拖垮业务调用。探测走 ACCESS SHARE,稳态下一条 ALTER 都不会发。
探测失败保守沿用全量列(今天的行为): 猜不出真实列集合时,让写入照常尝试。
"""
try:
existing = {row["attname"] for row in await conn.fetch(_EXISTING_COLUMNS)} # type: ignore[attr-defined]
for column, statement in _BACKFILL:
except asyncio.CancelledError:
raise
except Exception as exc:
logger.warning("Postgres 遥测列探测失败(沿用全量列,写入将逐行降级): {}", exc)
return COLUMNS
if self._auto_migrate:
await self._backfill_columns(conn, existing)
return COLUMNS
return self._trim_columns(existing)
def _trim_columns(self, existing: set[str]) -> tuple[str, ...]:
"""manual 档: 按现有列裁剪写入列,并把缺列一次讲清楚。
裁剪是关掉 ALTER 的**前提**而非增强: 旧表缺列时仍发全量 INSERT,每一行
都会因未知列被拒 → 遥测彻底丢失,比自动 ALTER 更严重地违反"遥测必录"
探测结果与 `COLUMNS` 毫无交集时视同探测异常保守回落全量: 空列集拼不出合法
INSERT,`insert_sql` 会 ValueError,而 `_prepare_schema` 里那次调用在 try
**之外**,异常会顺着 `record_llm_call` 一路冒给业务调用方(遥测绝不冒泡)
——回落必须发生在把空列集交给它之前。
"""
effective = tuple(column for column in COLUMNS if column in existing)
if not effective:
logger.warning(
"Postgres 遥测表 llm_calls 没有任何本库认识的列(沿用全量列,写入将逐行降级);"
"现有列: {}",
sorted(existing),
)
return COLUMNS
missing = [column for column in COLUMNS if column not in existing]
if missing:
# 单参数传入: 补列 SQL 里带 `'{}'::jsonb` 字面量,拼进 format 模板会被当占位符
logger.warning(
"{}",
missing_columns_warning("postgres", missing, alien_table="call_id" not in existing),
)
return effective
async def _backfill_columns(self, conn: object, existing: set[str]) -> None:
"""auto 档: 给已存在的旧表补新列(issue #3);**失败绝不置 `_failed`**。
不置 `_failed` 的实测理由: 应用账号只有 INSERT 权限时,`ALTER TABLE` 的
ownership 检查早于 `IF NOT EXISTS` 的存在性判断——列明明齐全也会失败。置位会让
整个 recorder 永久 no-op,与「补列失败只降级为逐行丢弃」的承诺相悖
(SQLite 侧同款守卫,两侧必须对称)。补列失败后写入沿用全量列(今天的行为):
auto 档承诺的是"把列补上",补不上就让缺列以逐行 warning 暴露;要降级写入
请显式选 manual。
"""
try:
for column, statement in PG_BACKFILL:
if column not in existing:
await conn.execute(statement) # type: ignore[attr-defined]
except asyncio.CancelledError:
@@ -232,14 +225,18 @@ class PostgresRecorder:
logger.warning("Postgres 遥测补列失败(写入将逐行降级): {}", exc)
async def record_llm_call(self, **fields: object) -> None:
"""写一行遥测;单条失败逐条 warning 丢弃(两级降级之二),绝不冒泡。"""
"""写一行遥测;单条失败逐条 warning 丢弃(两级降级之二),绝不冒泡。
取值按 `self._columns`(manual 档可能已被裁剪),与 `self._insert` 的
占位符同序——两者必须一起改,分开改就是把值写进错位的列。
"""
pool = await self._ensure_ready()
if pool is None:
return
row = tuple(fields[col] for col in _COLUMNS)
row = tuple(fields[col] for col in self._columns)
try:
async with pool.acquire() as conn:
await conn.execute(_INSERT, *row)
await conn.execute(self._insert, *row)
except asyncio.CancelledError:
raise
except Exception as exc:
+300
View File
@@ -0,0 +1,300 @@
"""遥测表 `llm_calls` 的 schema 单一事实源: 列序、两端 DDL、补列语句与 INSERT 构造。
两个 recorder(`sqlite.py` / `postgres.py`)与公共函数 `telemetry_schema_sql` 共用本模块。
收敛的理由是**正确性**而非整洁: 打印给下游的 SQL 必须与库真正执行的 DDL 同源——常量在
多处各存一份必然漂移,而漂移的表现是"下游照打印的 SQL 建完表,库仍报缺列"
**`COLUMNS` 是 INSERT 字段序,不是物理列序**: 数据库自填的 `created_at` 不在其中(它带
`DEFAULT now()` / `datetime('now')`,库从不显式写它)。物理表列 = 24 个 INSERT 字段 +
`created_at` = 25;列数断言一律按物理列数写,两套口径混用是最易错处。
本模块只依赖标准库: `telemetry/` 与 `backends/`、`transports/`、`structured/` 同层且
互不依赖(import-linter 契约执法)。
"""
from __future__ import annotations
from typing import TYPE_CHECKING
if TYPE_CHECKING:
from collections.abc import Sequence
TABLE = "llm_calls"
# 支持的后端;`insert_sql` / `telemetry_schema_sql` 的取值域
_BACKENDS = ("sqlite", "postgres")
SQLITE_DDL = """
CREATE TABLE IF NOT EXISTS llm_calls (
call_id TEXT PRIMARY KEY,
parent_call_id TEXT,
session_id TEXT,
model TEXT NOT NULL,
provider TEXT NOT NULL,
source_name TEXT NOT NULL,
messages TEXT NOT NULL,
response TEXT NOT NULL,
thinking TEXT NOT NULL DEFAULT '',
prompt_tokens INTEGER NOT NULL,
completion_tokens INTEGER NOT NULL,
usage_source TEXT NOT NULL,
latency_ms INTEGER NOT NULL,
ttft_ms REAL,
max_inter_token_ms REAL,
cache_hit INTEGER NOT NULL DEFAULT 0,
error TEXT,
cost REAL,
created_at TEXT NOT NULL DEFAULT (datetime('now')),
cached_prompt_tokens INTEGER,
model_reported TEXT,
sampling TEXT,
reasoning_tokens INTEGER,
tenant_id TEXT NOT NULL DEFAULT '',
meta TEXT NOT NULL DEFAULT '{}'
);
"""
PG_DDL = """
CREATE TABLE IF NOT EXISTS llm_calls (
call_id TEXT PRIMARY KEY,
parent_call_id TEXT,
session_id TEXT,
model TEXT NOT NULL,
provider TEXT NOT NULL,
source_name TEXT NOT NULL,
messages TEXT NOT NULL,
response TEXT NOT NULL,
thinking TEXT NOT NULL DEFAULT '',
prompt_tokens INTEGER NOT NULL,
completion_tokens INTEGER NOT NULL,
usage_source TEXT NOT NULL,
latency_ms INTEGER NOT NULL,
ttft_ms DOUBLE PRECISION,
max_inter_token_ms DOUBLE PRECISION,
cache_hit BOOLEAN NOT NULL DEFAULT FALSE,
error TEXT,
cost DOUBLE PRECISION,
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
cached_prompt_tokens INTEGER,
model_reported TEXT,
sampling TEXT,
reasoning_tokens INTEGER,
tenant_id TEXT NOT NULL DEFAULT '',
meta JSONB NOT NULL DEFAULT '{}'::jsonb
);
"""
# 新列必须排在 created_at 之后: 旧表只能经 ALTER 追加到末尾,新建库若把它们
# 插在前面,两条路径的物理列序会分叉(列序断言测试无合规修法)。
SQLITE_BACKFILL = (
("cached_prompt_tokens", "INTEGER"),
("model_reported", "TEXT"),
("sampling", "TEXT"),
("reasoning_tokens", "INTEGER"),
# NOT NULL 补列必须带非 NULL 常量默认值,否则 SQLite 直接拒绝该 ALTER
# ("Cannot add a NOT NULL column with default value NULL"),补列全盘失败。
("tenant_id", "TEXT NOT NULL DEFAULT ''"),
("meta", "TEXT NOT NULL DEFAULT '{}'"),
)
# PG 补列的列定义。语句由此派生成两份文本(见下),使"库内执行的那份"与"打印给
# 下游的那份"的列集合与列定义**无法分叉**——本模块存在的全部理由就是不许漂移。
_PG_BACKFILL_DECLS = (
("cached_prompt_tokens", "INTEGER"),
("model_reported", "TEXT"),
("sampling", "TEXT"),
("reasoning_tokens", "INTEGER"),
# 两个默认值都是非易失常量,PG 11+ 只改 catalog 不重写全表,故大表补列亦是秒级
("tenant_id", "TEXT NOT NULL DEFAULT ''"),
("meta", "JSONB NOT NULL DEFAULT '{}'::jsonb"),
)
# 新列排在 created_at 之后: 与旧表 ALTER 追加的位置一致(见 SQLITE_BACKFILL 同款注释)。
# **库内执行的这份有意不带 `IF NOT EXISTS`**: PG 对它即便列已存在也会先取 ACCESS
# EXCLUSIVE 锁,而遥测是业务路径上的内联 await,故库侧一律"先探测后 ALTER"
# (postgres.py `_backfill_columns` 记有实测)。给人执行的那份见 `telemetry_schema_sql`。
PG_BACKFILL = tuple(
(column, f"ALTER TABLE {TABLE} ADD COLUMN {column} {decl}")
for column, decl in _PG_BACKFILL_DECLS
)
COLUMNS = (
"call_id",
"parent_call_id",
"session_id",
"model",
"provider",
"source_name",
"messages",
"response",
"thinking",
"prompt_tokens",
"completion_tokens",
"usage_source",
"latency_ms",
"ttft_ms",
"max_inter_token_ms",
"cache_hit",
"error",
"cost",
"cached_prompt_tokens",
"model_reported",
"sampling",
"reasoning_tokens",
"tenant_id",
"meta",
)
_COLUMN_SET = frozenset(COLUMNS)
def insert_sql(backend: str, columns: Sequence[str]) -> str:
"""按给定列构造 INSERT;列必须是 `COLUMNS` 的非空子集,否则 ValueError。
子集校验是**注入面的闸**: 列名来自数据库探测结果,不是常量,不校验就等于把外部
字符串拼进 SQL(占位符只保护值,保护不了列名)。空集同样来自探测结果,而
`INSERT INTO llm_calls () VALUES ()` 两端都语法非法——本函数自己拒,不把这个
不变量押在调用方身上。sqlite 用 `?`、postgres 用 `$n`,
两端的重复键处理都不绑定具体约束名(`INSERT OR IGNORE` / `ON CONFLICT`)。
**PG 的 `ON CONFLICT` 一律不带冲突目标,不得"顺手"补回 `(call_id)`**: PG 要求
分区表的唯一约束必须包含分区键,按 `created_at` 分区(issue #12 的保留期方案)后
主键变成 `(call_id, created_at)`,带目标的语句匹配不到任何约束,PG 直接拒收
("there is no unique or exclusion constraint matching the ON CONFLICT
specification"),而遥测写失败只逐行 warning——分区部署下会全线静默丢数据。
无目标版本在两种表形态上都合法,普通表上语义逐字等价(表上只有主键一个唯一约束)。
Args:
backend: `"sqlite"` 或 `"postgres"`。
columns: 要写入的列,顺序即占位符顺序(调用方须按同序取值)。
Returns:
完整的 INSERT 语句。
Raises:
ValueError: backend 不在取值域内,columns 为空,或含 `COLUMNS` 之外的列名。
"""
if backend not in _BACKENDS:
raise ValueError(f"未知遥测后端 {backend!r}: 只支持 {list(_BACKENDS)}")
selected = tuple(columns)
if not selected:
raise ValueError("遥测 INSERT 至少需要一列: 空列集合会拼出语法非法的 SQL")
unknown = [column for column in selected if column not in _COLUMN_SET]
if unknown:
raise ValueError(f"列名不在遥测 schema 内(拒绝拼进 SQL): {unknown}")
names = ", ".join(selected)
if backend == "sqlite":
placeholders = ", ".join("?" for _ in selected)
return f"INSERT OR IGNORE INTO {TABLE} ({names}) VALUES ({placeholders})"
placeholders = ", ".join(f"${i + 1}" for i in range(len(selected)))
return f"INSERT INTO {TABLE} ({names}) VALUES ({placeholders}) ON CONFLICT DO NOTHING"
# 缺列告警要打印的补列语句: 库内执行的那份怎么写,打印给人的就怎么写(同源不许漂移)。
# SQLite 侧常量只有列定义,故在此按 TABLE 拼成整条 ALTER;PG 侧常量本就是整条语句。
_ALTER_BY_BACKEND = {
"sqlite": {
column: f"ALTER TABLE {TABLE} ADD COLUMN {column} {decl}"
for column, decl in SQLITE_BACKFILL
},
"postgres": dict(PG_BACKFILL),
}
_BACKEND_LABELS = {"sqlite": "SQLite", "postgres": "Postgres"}
# PG 的 ALTER 取 ACCESS EXCLUSIVE 锁,执行时机得由 DBA 自己挑;SQLite 是下游本地文件,无此顾虑
_EXECUTION_NOTES = {"sqlite": "", "postgres": "(建议挑低峰,ALTER 取 ACCESS EXCLUSIVE 锁)"}
def missing_columns_warning(backend: str, missing: Sequence[str], *, alien_table: bool) -> str:
"""拼 manual 档的缺列告警: 逐列点名 + 讲清后果 + 给出可直接执行的 SQL。
只说"缺列"是不够的: 静默丢维度的后果是多租户账目全归空串且无任何报错,
看告警的人必须一眼看到丢的是哪几个维度、以及怎么补。
**住在本模块而不是两个 recorder 里**: 这条消息拼的是给人执行的 DDL,与库自己
执行的 ALTER 必须同源——本模块存在的全部理由就是不许这两者漂移。
Args:
backend: `"sqlite"` 或 `"postgres"`。
missing: 缺失的列名(按 `COLUMNS` 保序)。
alien_table: 连主键列 `call_id` 都没有——该表多半不是本库的 `llm_calls`。
Returns:
单条 warning 的完整文本(库只在准备期发一次,不逐行发)。
Raises:
ValueError: backend 不在取值域内。
"""
if backend not in _BACKENDS:
raise ValueError(f"未知遥测后端 {backend!r}: 只支持 {list(_BACKENDS)}")
alters = _ALTER_BY_BACKEND[backend]
selected = tuple(missing)
statements = [f"{alters[column]};" for column in selected if column in alters]
unknown = [column for column in selected if column not in alters]
if unknown:
# 这些列本库从未经 ALTER 补过(建表即有),给不出单条 ALTER,指向完整脚本
statements.append(
f"-- 另缺 {', '.join(unknown)};完整建表脚本见 "
f'polygateway.telemetry_schema_sql("{backend}")'
)
label = _BACKEND_LABELS[backend]
head = (
f"{label} 遥测表 {TABLE} 缺主键列 call_id,很可能不是本库的遥测表"
"(库不做二次判定,仍照常尝试写入)"
if alien_table
else f"{label} 遥测表 {TABLE} 缺列,且 auto_migrate=False(库不发任何 DDL)"
)
return (
f"{head};以下维度不会被记录: {', '.join(selected)}"
f"补列请自行执行{_EXECUTION_NOTES[backend]}:\n" + "\n".join(statements)
)
def telemetry_schema_sql(backend: str) -> str:
"""返回可直接粘进迁移文件的完整脚本(建表 + 各补列语句 + 注释)。
给不愿意让库在自己的生产表上发 DDL 的下游用: 输出与库运行时执行的 DDL 同源,
照它建完表,库探测到的列就是齐的。
**补列语句与库内执行的那份是两套文本,不是一份**: 这份给人执行,必须可重复执行,
故 PG 变体带 `ADD COLUMN IF NOT EXISTS`(它会先取 ACCESS EXCLUSIVE 锁,但执行时机
由 DBA 自己挑,锁风险可控);库内那份不带,靠先探测后 ALTER 规避锁。SQLite 没有
`ADD COLUMN IF NOT EXISTS` 语法,只能以注释交代"仅当该列不存在时执行"
Args:
backend: `"sqlite"` 或 `"postgres"`。
Returns:
含注释的完整 SQL 脚本。
Raises:
ValueError: backend 不在取值域内。
"""
if backend not in _BACKENDS:
raise ValueError(f"未知遥测后端 {backend!r}: 只支持 {list(_BACKENDS)}")
if backend == "sqlite":
ddl = SQLITE_DDL
notes = (
f"-- 旧表补列(库升级后新增的列)。SQLite 无 ADD COLUMN IF NOT EXISTS 语法,\n"
f"-- 以下每条**仅当该列不存在时执行**(先 PRAGMA table_info({TABLE}) 对照)。"
)
alters = [
f"ALTER TABLE {TABLE} ADD COLUMN {column} {decl};" for column, decl in SQLITE_BACKFILL
]
else:
ddl = PG_DDL
notes = (
"-- 旧表补列(库升级后新增的列)。带 IF NOT EXISTS,整段可重复执行;\n"
"-- 注意它即便列已存在也会先取 ACCESS EXCLUSIVE 锁,请挑低峰执行。"
)
alters = [
f"ALTER TABLE {TABLE} ADD COLUMN IF NOT EXISTS {column} {decl};"
for column, decl in _PG_BACKFILL_DECLS
]
header = (
f"-- PolyGateway 遥测表 {TABLE}({backend})\n"
f'-- 由 polygateway.telemetry_schema_sql("{backend}") 生成,与库运行时执行的 DDL 同源。\n'
"-- 新建库执行整段;已有旧表则建表语句自动跳过,只需关注下方补列语句。"
)
return "\n".join([header, "", ddl.strip(), "", notes, *alters, ""])
+75 -86
View File
@@ -22,117 +22,102 @@ from pathlib import Path
from loguru import logger
_DDL = """
CREATE TABLE IF NOT EXISTS llm_calls (
call_id TEXT PRIMARY KEY,
parent_call_id TEXT,
session_id TEXT,
model TEXT NOT NULL,
provider TEXT NOT NULL,
source_name TEXT NOT NULL,
messages TEXT NOT NULL,
response TEXT NOT NULL,
thinking TEXT NOT NULL DEFAULT '',
prompt_tokens INTEGER NOT NULL,
completion_tokens INTEGER NOT NULL,
usage_source TEXT NOT NULL,
latency_ms INTEGER NOT NULL,
ttft_ms REAL,
max_inter_token_ms REAL,
cache_hit INTEGER NOT NULL DEFAULT 0,
error TEXT,
cost REAL,
created_at TEXT NOT NULL DEFAULT (datetime('now')),
cached_prompt_tokens INTEGER,
model_reported TEXT,
sampling TEXT,
reasoning_tokens INTEGER,
tenant_id TEXT NOT NULL DEFAULT '',
meta TEXT NOT NULL DEFAULT '{}'
);
"""
# 新列必须排在 created_at 之后: 旧表只能经 ALTER 追加到末尾,新建库若把它们
# 插在前面,两条路径的物理列序会分叉(列序断言测试无合规修法)。
_BACKFILL_COLUMNS = (
("cached_prompt_tokens", "INTEGER"),
("model_reported", "TEXT"),
("sampling", "TEXT"),
("reasoning_tokens", "INTEGER"),
# NOT NULL 补列必须带非 NULL 常量默认值,否则 SQLite 直接拒绝该 ALTER
# ("Cannot add a NOT NULL column with default value NULL"),补列全盘失败。
("tenant_id", "TEXT NOT NULL DEFAULT ''"),
("meta", "TEXT NOT NULL DEFAULT '{}'"),
)
_COLUMNS = (
"call_id",
"parent_call_id",
"session_id",
"model",
"provider",
"source_name",
"messages",
"response",
"thinking",
"prompt_tokens",
"completion_tokens",
"usage_source",
"latency_ms",
"ttft_ms",
"max_inter_token_ms",
"cache_hit",
"error",
"cost",
"cached_prompt_tokens",
"model_reported",
"sampling",
"reasoning_tokens",
"tenant_id",
"meta",
)
_INSERT = (
f"INSERT OR IGNORE INTO llm_calls ({', '.join(_COLUMNS)}) "
f"VALUES ({', '.join('?' for _ in _COLUMNS)})"
from polygateway.telemetry.schema import (
COLUMNS,
SQLITE_BACKFILL,
SQLITE_DDL,
insert_sql,
missing_columns_warning,
)
class SQLiteRecorder:
"""TelemetryRecorder 端口的 SQLite 实现;初始化/写入失败全降级 warning。"""
def __init__(self, db_path: Path | str) -> None:
def __init__(self, db_path: Path | str, *, auto_migrate: bool) -> None:
"""建连接与表,并按探测到的列定型本实例的 INSERT 语句。
Args:
db_path: 库文件路径;父目录不存在会自动创建。
auto_migrate: True 则给已存在的旧表自动补列(SQLite 侧的缺省档:
下游本地文件,无 DBA 无迁移工具);False 则一条 ALTER 都不发,
改为按现有列裁剪写入。keyword-only **必填**: 缺省规则只写在
config 一处,不与本类签名漂移(设计 D-c)。
"""
self._auto_migrate = auto_migrate
self._lock = threading.Lock()
self._conn: sqlite3.Connection | None = None
# 先按全量列定型: 连接失败/探测失败时保守沿用全量(今天的行为)
self._columns: tuple[str, ...] = COLUMNS
self._insert = insert_sql("sqlite", COLUMNS)
try:
path = Path(db_path)
path.parent.mkdir(parents=True, exist_ok=True)
conn = sqlite3.connect(path, check_same_thread=False, timeout=10.0)
conn.execute("PRAGMA journal_mode=WAL")
conn.execute("PRAGMA busy_timeout=5000")
conn.execute(_DDL)
conn.execute(SQLITE_DDL)
conn.commit()
self._conn = conn
except (OSError, sqlite3.Error) as exc:
logger.warning("SQLite 遥测初始化失败,后续记录降级为 no-op: {}", exc)
self._backfill_columns()
self._prepare_columns()
def _backfill_columns(self) -> None:
"""给已存在的旧表补新列(issue #3);独立 try,失败只降级为逐行丢弃
def _prepare_columns(self) -> None:
"""探测现有列后定型写入: auto 档补齐缺列,manual 档改为裁剪写入(issue #13)。
必须放在 `self._conn` 赋值**之后**并先判空: 初始化失败时连接为 None,
无守卫的补列会抛 AttributeError 逃出 `__init__`,把"静默降级"变成崩溃。
补列失败也绝不清空 `self._conn`——那会让整个 recorder 永久 no-op,
比逐行丢弃严重得多。
无守卫的探测会抛 AttributeError 逃出 `__init__`,把"静默降级"变成崩溃。
探测失败保守沿用全量列(今天的行为): 猜不出真实列集合时,让写入照常尝试。
"""
if self._conn is None:
return
try:
existing = {row[1] for row in self._conn.execute("PRAGMA table_info(llm_calls)")}
except sqlite3.Error as exc:
logger.warning("SQLite 遥测列探测失败(写入将逐行降级): {}", exc)
logger.warning("SQLite 遥测列探测失败(沿用全量列,写入将逐行降级): {}", exc)
return
for column, decl in _BACKFILL_COLUMNS:
if self._auto_migrate:
self._backfill_columns(existing)
return
self._adopt_existing_columns(existing)
def _adopt_existing_columns(self, existing: set[str]) -> None:
"""manual 档: 不发任何 DDL,按现有列裁剪 INSERT,并把缺列一次讲清楚。
裁剪是关掉 ALTER 的**前提**而非增强: 旧表缺列时仍发全量 INSERT,每一行
都会因未知列被拒 → 遥测彻底丢失,比自动 ALTER 更严重地违反"遥测必录"
探测结果与 `COLUMNS` 毫无交集时视同探测异常保守回落全量: 空列集拼不出合法
INSERT,`insert_sql` 会 ValueError,而遥测构造期抛异常就是把"初始化失败静默
降级"的铁律破成崩溃——回落必须发生在把空列集交给它之前。
"""
effective = tuple(column for column in COLUMNS if column in existing)
if not effective:
logger.warning(
"SQLite 遥测表 llm_calls 没有任何本库认识的列(沿用全量列,写入将逐行降级);"
"现有列: {}",
sorted(existing),
)
return
self._columns = effective
self._insert = insert_sql("sqlite", effective)
missing = [column for column in COLUMNS if column not in existing]
if missing:
# 单参数传入: 补列 SQL 里带 `'{}'` 字面量,拼进 format 模板会被当占位符
logger.warning(
"{}",
missing_columns_warning("sqlite", missing, alien_table="call_id" not in existing),
)
def _backfill_columns(self, existing: set[str]) -> None:
"""auto 档: 给已存在的旧表补新列(issue #3);逐列独立 try,失败只降级为逐行丢弃。
补列失败绝不清空 `self._conn`——那会让整个 recorder 永久 no-op,
比逐行丢弃严重得多。失败后写入沿用全量列(今天的行为): auto 档承诺的是
"把列补上",补不上就让缺列以逐行 warning 暴露;要降级写入请显式选 manual。
"""
assert self._conn is not None # 内部不变量: 调用方已判空
for column, decl in SQLITE_BACKFILL:
if column in existing:
continue
# 逐列独立 try: 一列撞上 duplicate 不得让后面的列漏补
@@ -145,10 +130,14 @@ class SQLiteRecorder:
logger.warning("SQLite 遥测补列失败(写入将逐行降级): {}", exc)
async def record_llm_call(self, **fields: object) -> None:
"""写一行遥测;字段集合即 24 字段冻结签名(ports.TelemetryRecorder)。"""
"""写一行遥测;字段集合即 24 字段冻结签名(ports.TelemetryRecorder)。
取值按 `self._columns`(manual 档可能已被裁剪),与 `self._insert` 的
占位符同序——两者必须一起改,分开改就是把值写进错位的列。
"""
if self._conn is None:
return
row = tuple(fields[col] for col in _COLUMNS)
row = tuple(fields[col] for col in self._columns)
try:
await asyncio.to_thread(self._write, row)
except (OSError, sqlite3.Error) as exc:
@@ -157,7 +146,7 @@ class SQLiteRecorder:
def _write(self, row: tuple) -> None:
assert self._conn is not None # 内部不变量: 调用方已判空
with self._lock:
self._conn.execute(_INSERT, row)
self._conn.execute(self._insert, row)
self._conn.commit()
def close(self) -> None:
+6 -6
View File
@@ -119,7 +119,7 @@ class TestBreakerRecoveryFullChain:
class TestCancellationThroughStack:
async def test_cancel_mid_request_releases_and_records(self, tmp_path):
recorder = SQLiteRecorder(tmp_path / "t.db")
recorder = SQLiteRecorder(tmp_path / "t.db", auto_migrate=True)
entered = asyncio.Event()
async def hanging_handler(request):
@@ -144,7 +144,7 @@ class TestCancellationThroughStack:
class TestTelemetryAcrossPaths:
async def test_success_cache_hit_and_failure_rows(self, tmp_path):
recorder = SQLiteRecorder(tmp_path / "t.db")
recorder = SQLiteRecorder(tmp_path / "t.db", auto_migrate=True)
client = _full_client(lambda req: _sse(), telemetry=recorder, cache=InMemoryCache())
await client.chat([{"role": "user", "content": "hi"}]) # 成功(尝试行)
await client.chat([{"role": "user", "content": "hi"}]) # 缓存命中行
@@ -155,7 +155,7 @@ class TestTelemetryAcrossPaths:
assert hits == 1 and total == 2
async def test_transient_attempts_each_recorded(self, tmp_path):
recorder = SQLiteRecorder(tmp_path / "t.db")
recorder = SQLiteRecorder(tmp_path / "t.db", auto_migrate=True)
calls = {"n": 0}
def flaky(request):
@@ -193,7 +193,7 @@ class TestRejectionReasonIsQueryable:
)
async def test_rejected_call_leaves_the_reason_in_telemetry(self, tmp_path):
recorder = SQLiteRecorder(tmp_path / "t.db")
recorder = SQLiteRecorder(tmp_path / "t.db", auto_migrate=True)
client = _full_client(
lambda req: httpx.Response(400, content=self._BODY.encode()), telemetry=recorder
)
@@ -257,7 +257,7 @@ class TestSamplingThroughStack:
return _sse()
db = tmp_path / "t.db"
recorder = SQLiteRecorder(db)
recorder = SQLiteRecorder(db, auto_migrate=True)
client = _full_client(handler, telemetry=recorder)
await client.chat([{"role": "user", "content": "hi"}], overlay={"seed": 42})
recorder.close()
@@ -276,7 +276,7 @@ class TestSamplingThroughStack:
src = dataclasses.replace(_source(), extra_body={"temperature": 0})
db = tmp_path / "t.db"
recorder = SQLiteRecorder(db)
recorder = SQLiteRecorder(db, auto_migrate=True)
client = GatewayClient(
scope="llm",
sources=[src],
+301 -16
View File
@@ -14,12 +14,14 @@ import asyncio
import json
import os
import re
from datetime import UTC, datetime, timedelta
from uuid import uuid4
import pytest
from dotenv import dotenv_values
from polygateway.telemetry.postgres import PostgresRecorder
from polygateway.telemetry.schema import COLUMNS, telemetry_schema_sql
_EXPECTED_COLUMNS = [
"call_id",
@@ -88,8 +90,13 @@ async def dsn():
async def _record_minimal(
recorder: PostgresRecorder, call_id: str | None = None, **overrides
) -> None:
fields = {
) -> dict[str, object]:
"""记一行最小遥测,并**返回实际提交的字段**供调用方逐列比对回读结果。
返回值不是顺手加的: 逐列断言若在测试里另抄一份期望值,抄错的那一列会以
"库写错列位"的形态误报,而漏抄的列则悄悄不被验证。
"""
fields: dict[str, object] = {
"call_id": call_id if call_id is not None else _cid("c1"),
"parent_call_id": None,
"session_id": "sess-1",
@@ -118,6 +125,7 @@ async def _record_minimal(
}
fields.update(overrides)
await recorder.record_llm_call(**fields)
return fields
async def _fetch(dsn: str, sql: str, *args):
@@ -130,6 +138,17 @@ async def _fetch(dsn: str, sql: str, *args):
await conn.close()
async def _execute_script(dsn: str, sql: str) -> None:
"""整段执行多语句脚本(不带参数,走简单查询协议)——模拟下游把脚本贴进 psql。"""
import asyncpg
conn = await asyncpg.connect(dsn, timeout=10)
try:
await conn.execute(sql)
finally:
await conn.close()
_LEGACY_DDL = """
CREATE TABLE {schema}.llm_calls (
call_id TEXT PRIMARY KEY,
@@ -184,7 +203,7 @@ class TestObservabilityColumns:
"""issue #3: 两列写入可回读,且已存在的 18 列旧表会被自动补列。"""
async def test_values_round_trip(self, dsn):
recorder = PostgresRecorder(dsn)
recorder = PostgresRecorder(dsn, auto_migrate=True)
try:
await _record_minimal(recorder, call_id=_cid("hit"), cached_prompt_tokens=64)
await _record_minimal(recorder, call_id=_cid("zero"), cached_prompt_tokens=0)
@@ -212,7 +231,7 @@ class TestObservabilityColumns:
async def test_legacy_table_is_upgraded_in_place(self, legacy_schema):
"""18 列旧表不补列的话,每行写入都会被逐行 warning 丢弃(遥测静默全失)。"""
schema_dsn, schema = legacy_schema
recorder = PostgresRecorder(schema_dsn)
recorder = PostgresRecorder(schema_dsn, auto_migrate=True)
try:
await _record_minimal(
recorder, call_id=_cid("legacy"), cached_prompt_tokens=7, model_reported="m-real"
@@ -237,7 +256,7 @@ class TestObservabilityColumns:
class TestSchema:
async def test_schema_has_frozen_columns_in_order(self, dsn):
recorder = PostgresRecorder(dsn)
recorder = PostgresRecorder(dsn, auto_migrate=True)
try:
await _record_minimal(recorder)
rows = await _fetch(
@@ -250,7 +269,7 @@ class TestSchema:
await recorder.aclose()
async def test_call_id_idempotent(self, dsn):
recorder = PostgresRecorder(dsn)
recorder = PostgresRecorder(dsn, auto_migrate=True)
try:
await _record_minimal(recorder, call_id=_cid("dup"))
await _record_minimal(recorder, call_id=_cid("dup"), response="second")
@@ -262,7 +281,7 @@ class TestSchema:
await recorder.aclose()
async def test_concurrent_writes_all_land(self, dsn):
recorder = PostgresRecorder(dsn)
recorder = PostgresRecorder(dsn, auto_migrate=True)
try:
await asyncio.gather(
*(_record_minimal(recorder, call_id=_cid(f"c{i}")) for i in range(50))
@@ -280,14 +299,14 @@ class TestSchema:
class TestDegradation:
async def test_unreachable_server_degrades_silently(self):
"""结构性失败(建池不通)→ warning 一次后永久降级,业务零感知。"""
recorder = PostgresRecorder("postgresql://u:p@127.0.0.1:1/x")
recorder = PostgresRecorder("postgresql://u:p@127.0.0.1:1/x", auto_migrate=True)
await _record_minimal(recorder) # 不抛
await _record_minimal(recorder, call_id=_cid("c2")) # 已降级短路,同样不抛
await recorder.aclose()
async def test_row_failure_does_not_poison_later_rows(self, dsn):
"""运行时单条写失败(NUL 字节文本被 PG 拒)→ 丢该行,后续行照常落库。"""
recorder = PostgresRecorder(dsn)
recorder = PostgresRecorder(dsn, auto_migrate=True)
try:
await _record_minimal(recorder, call_id=_cid("bad"), response="nul\x00byte")
await _record_minimal(recorder, call_id=_cid("good"))
@@ -301,7 +320,7 @@ class TestDegradation:
await recorder.aclose()
async def test_aclose_idempotent(self, dsn):
recorder = PostgresRecorder(dsn)
recorder = PostgresRecorder(dsn, auto_migrate=True)
await _record_minimal(recorder)
await recorder.aclose()
await recorder.aclose()
@@ -320,7 +339,7 @@ async def least_privilege_dsn(dsn):
"""
import asyncpg
from polygateway.telemetry.postgres import _DDL
from polygateway.telemetry.schema import PG_DDL
name = f"pgwtest_lp_{uuid4().hex[:8]}"
admin = await asyncpg.connect(dsn, timeout=10)
@@ -332,7 +351,7 @@ async def least_privilege_dsn(dsn):
await admin.execute(f"CREATE ROLE {name} LOGIN PASSWORD '{_PROBE_PASSWORD}'")
await admin.execute(f"CREATE SCHEMA {name}")
await admin.execute(f"SET search_path = {name}")
await admin.execute(_DDL) # 表由**别的账号**建好,与现场一致
await admin.execute(PG_DDL) # 表由**别的账号**建好,与现场一致
await admin.execute(f"GRANT USAGE ON SCHEMA {name} TO {name}")
await admin.execute(f"GRANT SELECT, INSERT ON {name}.llm_calls TO {name}")
# 关键: 绝不 GRANT CREATE ON SCHEMA —— 缺的正是这一项
@@ -373,7 +392,7 @@ class TestLeastPrivilegeDeployment:
async def test_records_land_without_schema_create_privilege(self, least_privilege_dsn):
"""修复前: 建表被拒 → _failed → 整个进程一条不落(下游 150 次调用全丢)。"""
low_dsn, schema = least_privilege_dsn
recorder = PostgresRecorder(low_dsn)
recorder = PostgresRecorder(low_dsn, auto_migrate=True)
try:
await _record_minimal(recorder, call_id=_cid("lp1"))
await _record_minimal(recorder, call_id=_cid("lp2"), cost=1.5)
@@ -428,6 +447,15 @@ _PRE_TENANT_INSERT = (
)
# `_PRE_TENANT_DDL` 的物理列(23 个): 由 `_EXPECTED_COLUMNS` 去掉 issue #11 的两个新维度
# 派生而非另抄一份——两份常量必然漂移,而漂移的表现是"manual 档没补列"这条断言假绿。
# 去掉后的顺序与 DDL 逐字一致(tenant_id/meta 在 DDL 里本就排在末尾)。
_PRE_TENANT_COLUMNS = [c for c in _EXPECTED_COLUMNS if c not in ("tenant_id", "meta")]
# 回读要逐列比对的字段: 物理列去掉库从不显式写的 created_at,恰好 22 个
_PRE_TENANT_WRITTEN_COLUMNS = [c for c in _PRE_TENANT_COLUMNS if c != "created_at"]
def _search_path_dsn(dsn: str, schema: str) -> str:
sep = "&" if "?" in dsn else "?"
return f"{dsn}{sep}options=-csearch_path%3D{schema}"
@@ -533,7 +561,7 @@ class TestCallerDimensionsAcceptance:
async def test_fresh_schema_round_trips_the_dimensions(self, fresh_schema):
"""新建库: 列齐全,且维度值原样读回——只验列存在会漏掉写错列位的错。"""
fresh_dsn, schema = fresh_schema
recorder = PostgresRecorder(fresh_dsn)
recorder = PostgresRecorder(fresh_dsn, auto_migrate=True)
try:
await _record_minimal(
recorder, call_id=_cid("dim"), tenant_id="tenant-a", meta='{"batch": "b7"}'
@@ -567,7 +595,7 @@ class TestCallerDimensionsAcceptance:
审计出来,历史欠账是可见、可量化、可补录的。
"""
schema_dsn, schema = pre_tenant_schema
recorder = PostgresRecorder(schema_dsn)
recorder = PostgresRecorder(schema_dsn, auto_migrate=True)
try:
await _record_minimal(
recorder, call_id=_cid("new"), tenant_id="tenant-a", meta='{"k": 1}'
@@ -619,7 +647,7 @@ class TestCallerDimensionsAcceptance:
置 `_failed` 会让整个进程从此一条遥测都不写(比逐行丢弃严重得多),
且一旦 DBA 补上列也不会自愈——必须等重启。
"""
recorder = PostgresRecorder(least_privilege_pre_tenant_dsn)
recorder = PostgresRecorder(least_privilege_pre_tenant_dsn, auto_migrate=True)
try:
await _record_minimal(recorder, call_id=_cid("lpp1")) # 不得抛
assert recorder._failed is False
@@ -628,3 +656,260 @@ class TestCallerDimensionsAcceptance:
assert any("写入失败" in m for m in captured_warnings)
finally:
await recorder.aclose()
# issue #12 的目标表形态: 按 created_at 做 RANGE 分区(过期清理 DROP PARTITION 而非 DELETE)。
# PG 强制分区表的唯一约束必须包含分区键,故主键只能是 (call_id, created_at) ——
# 这正是带目标的 `ON CONFLICT (call_id)` 再也匹配不到约束的现场。
_PARTITIONED_DDL = """
CREATE TABLE {schema}.llm_calls (
call_id TEXT NOT NULL,
parent_call_id TEXT,
session_id TEXT,
model TEXT NOT NULL,
provider TEXT NOT NULL,
source_name TEXT NOT NULL,
messages TEXT NOT NULL,
response TEXT NOT NULL,
thinking TEXT NOT NULL DEFAULT '',
prompt_tokens INTEGER NOT NULL,
completion_tokens INTEGER NOT NULL,
usage_source TEXT NOT NULL,
latency_ms INTEGER NOT NULL,
ttft_ms DOUBLE PRECISION,
max_inter_token_ms DOUBLE PRECISION,
cache_hit BOOLEAN NOT NULL DEFAULT FALSE,
error TEXT,
cost DOUBLE PRECISION,
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
cached_prompt_tokens INTEGER,
model_reported TEXT,
sampling TEXT,
reasoning_tokens INTEGER,
tenant_id TEXT NOT NULL DEFAULT '',
meta JSONB NOT NULL DEFAULT '{{}}'::jsonb,
PRIMARY KEY (call_id, created_at)
) PARTITION BY RANGE (created_at)
"""
_PARTITION_DDL = (
"CREATE TABLE {schema}.llm_calls_current PARTITION OF {schema}.llm_calls "
"FOR VALUES FROM ('{start}') TO ('{end}')"
)
def _current_month_bounds() -> tuple[str, str]:
"""当前月的 [月初, 下月初) 边界字面量;分区键落在区间外会因找不到分区而写失败。"""
now = datetime.now(UTC)
start = now.replace(day=1, hour=0, minute=0, second=0, microsecond=0)
end = (start + timedelta(days=32)).replace(day=1)
fmt = "%Y-%m-%d %H:%M:%S%z"
return start.strftime(fmt), end.strftime(fmt)
@pytest.fixture
async def partitioned_schema(dsn):
"""自建临时 schema 里造一张按 created_at RANGE 分区的表 + 覆盖当前月的分区。
与 legacy_schema 同款隔离: 绝不碰共享的 public.llm_calls,teardown 只 DROP
自己建的 schema(CASCADE 连分区一并删)。
"""
import asyncpg
name = f"pgwtest_part_{uuid4().hex[:8]}"
start, end = _current_month_bounds()
conn = await asyncpg.connect(dsn, timeout=10)
try:
await conn.execute(f"CREATE SCHEMA {name}")
await conn.execute(_PARTITIONED_DDL.format(schema=name))
await conn.execute(_PARTITION_DDL.format(schema=name, start=start, end=end))
finally:
await conn.close()
yield _search_path_dsn(dsn, name), name
conn = await asyncpg.connect(dsn, timeout=10)
try:
await conn.execute(f"DROP SCHEMA {name} CASCADE")
finally:
await conn.close()
class TestConflictTargetFreeInsert:
"""issue #13: INSERT 不绑定冲突目标,普通表与分区表两种形态都写得进去。"""
async def test_plain_table_still_dedupes_by_call_id(self, fresh_schema, captured_warnings):
"""普通表上语义不变: 重复 call_id 仍只落一行,且不是被拒后丢弃。
表上只有主键这一个唯一约束,故无目标的 DO NOTHING 与 `(call_id)` 逐字等价;
断言"无写入失败 warning"是为了区分"冲突被忽略""整条被 PG 拒收"
"""
fresh_dsn, _ = fresh_schema
recorder = PostgresRecorder(fresh_dsn, auto_migrate=True)
try:
await _record_minimal(recorder, call_id=_cid("nodup"))
await _record_minimal(recorder, call_id=_cid("nodup"), response="second")
assert [m for m in captured_warnings if "写入失败" in m] == []
rows = await _fetch(
fresh_dsn, "SELECT response FROM llm_calls WHERE call_id = $1", _cid("nodup")
)
assert [r["response"] for r in rows] == ["ok"] # 首行胜出,写入幂等
finally:
await recorder.aclose()
async def test_partitioned_table_accepts_writes(self, partitioned_schema, captured_warnings):
"""分区表上写入成功且能读回——改动前这里必红。
带目标的 `ON CONFLICT (call_id)` 在主键为 `(call_id, created_at)` 的表上
匹配不到任何约束,PG 报 "there is no unique or exclusion constraint matching
the ON CONFLICT specification";该错误被逐行降级吞成 warning,于是分区部署下
遥测全线写不进去却一声不吭,只能靠"读不回来"暴露。
"""
part_dsn, _ = partitioned_schema
recorder = PostgresRecorder(part_dsn, auto_migrate=True)
try:
await _record_minimal(recorder, call_id=_cid("part"), tenant_id="tenant-p")
assert [m for m in captured_warnings if "写入失败" in m] == []
rows = await _fetch(
part_dsn,
"SELECT call_id, tenant_id FROM llm_calls WHERE call_id = $1",
_cid("part"),
)
assert [(r["call_id"], r["tenant_id"]) for r in rows] == [(_cid("part"), "tenant-p")]
finally:
await recorder.aclose()
class TestManualSchemaModeAcceptance:
"""issue #13 manual 档的真实实例验收: 旧表原样不动,写入照常,缺列只作提示。
manual 档的承诺是"库一条 DDL 都不发"——单元测试只能验"没调用 execute",
真表上才验得了"表结构确实没变"。两条用例分别覆盖有权补列却不补(纪律)与
无权补列(现场),后者正是 auto 档会刷出 `补列失败` warning 的那张表。
"""
async def test_manual_leaves_the_stale_table_untouched(
self, pre_tenant_schema, captured_warnings
):
"""22 字段旧表 + manual: 列一个不加,行照常落库,缺的两维度静默不写。
与 `test_pre_tenant_table_gains_columns_and_old_rows_stay_auditable` 恰成对照:
同一张表、同一份负载,只有 `auto_migrate` 不同,列数就必须是 23 与 25 之别。
"""
schema_dsn, schema = pre_tenant_schema
recorder = PostgresRecorder(schema_dsn, auto_migrate=False)
try:
recorded = await _record_minimal(
recorder, call_id=_cid("man"), tenant_id="tenant-a", meta='{"k": 1}'
)
cols = await _fetch(
schema_dsn,
"SELECT column_name FROM information_schema.columns "
"WHERE table_schema = $1 AND table_name = 'llm_calls' ORDER BY ordinal_position",
schema,
)
# 表结构逐字不动: 既没多出 tenant_id/meta,也没被顺手改了列序
assert [r["column_name"] for r in cols] == _PRE_TENANT_COLUMNS
names = ", ".join(_PRE_TENANT_WRITTEN_COLUMNS)
rows = await _fetch(
schema_dsn, f"SELECT {names} FROM llm_calls WHERE call_id = $1", _cid("man")
)
assert len(rows) == 1 # 裁剪后的 INSERT 真写进去了,不是被 PG 拒收
# 其余 22 列逐列与提交值相等: 少写两列最容易引发的错是剩下的值整体错位
assert dict(rows[0]) == {c: recorded[c] for c in _PRE_TENANT_WRITTEN_COLUMNS}
assert [m for m in captured_warnings if "写入失败" in m] == []
assert [m for m in captured_warnings if "补列失败" in m] == []
notices = [m for m in captured_warnings if "auto_migrate=False" in m]
assert len(notices) == 1 # 准备期一次讲清,不逐行刷屏
assert "以下维度不会被记录: tenant_id, meta" in notices[0]
finally:
await recorder.aclose()
async def test_manual_on_a_role_that_cannot_alter_emits_no_backfill_failure(
self, least_privilege_pre_tenant_dsn, captured_warnings
):
"""缺列旧表 + 只授 SELECT/INSERT 的角色 + manual: 补列失败的 warning 彻底消失。
auto 档在这张表上会刷出 `补列失败` 再刷 `写入失败`(见
`test_backfill_failure_degrades_per_row_not_wholesale`)——那是 issue #13 要
消灭的噪声。manual 档下 ALTER 压根不发,取而代之的是一条点名缺列并附可直接
执行的 ALTER 的提示,而遥测照常落库。
"""
recorder = PostgresRecorder(least_privilege_pre_tenant_dsn, auto_migrate=False)
try:
recorded = await _record_minimal(
recorder, call_id=_cid("manlp1"), tenant_id="tenant-b", meta='{"k": 2}'
)
await _record_minimal(recorder, call_id=_cid("manlp2"), cost=2.5)
assert [m for m in captured_warnings if "补列失败" in m] == []
assert [m for m in captured_warnings if "写入失败" in m] == []
assert recorder._failed is False
notices = [m for m in captured_warnings if "auto_migrate=False" in m]
assert len(notices) == 1 # 准备期一次,第二行不再重复
assert "以下维度不会被记录: tenant_id, meta" in notices[0]
# 提示里的 SQL 必须可直接粘贴执行,而不是只报个列名
assert (
"ALTER TABLE llm_calls ADD COLUMN tenant_id TEXT NOT NULL DEFAULT '';" in notices[0]
)
assert (
"ALTER TABLE llm_calls ADD COLUMN meta JSONB NOT NULL DEFAULT '{}'::jsonb;"
in notices[0]
)
# 该角色无权 ALTER,表必然还是旧形态: 缺的两列确实没被写
cols = await _fetch(
least_privilege_pre_tenant_dsn,
"SELECT column_name FROM information_schema.columns "
"WHERE table_schema = current_schema() AND table_name = 'llm_calls' "
"ORDER BY ordinal_position",
)
assert [r["column_name"] for r in cols] == _PRE_TENANT_COLUMNS
names = ", ".join(_PRE_TENANT_WRITTEN_COLUMNS)
rows = await _fetch(
least_privilege_pre_tenant_dsn,
f"SELECT {names} FROM llm_calls WHERE call_id LIKE $1 ORDER BY call_id",
f"{_RUN_PREFIX}-manlp%",
)
assert [r["call_id"] for r in rows] == [_cid("manlp1"), _cid("manlp2")]
assert dict(rows[0]) == {c: recorded[c] for c in _PRE_TENANT_WRITTEN_COLUMNS}
assert rows[1]["cost"] == 2.5
finally:
await recorder.aclose()
_PHYSICAL_COLUMNS_SQL = (
"SELECT column_name FROM information_schema.columns "
"WHERE table_schema = $1 AND table_name = 'llm_calls' ORDER BY ordinal_position"
)
class TestPublishedSchemaScript:
"""issue #13: README 叫下游执行的那份脚本,在真实实例上必须建得出、且可重复执行。
这份脚本是 `telemetry_schema_sql("postgres")` 的输出,manual 档下游拿它建表,
库随后靠列探测决定写哪些列——脚本与 `COLUMNS` 一旦漂移,表现是"照文档建完表,
库仍报缺列"。人工核对不构成回归保护: 改一次 README 或 DDL 就会悄悄失去它。
"""
async def test_script_builds_the_full_table_and_is_rerunnable(self, fresh_schema):
"""空 schema 里执行一遍建出全部物理列;再执行一遍不报错。
第二遍是 `ADD COLUMN IF NOT EXISTS` 的幂等性验收: 去掉 IF NOT EXISTS 后,
建表语句会被 `IF NOT EXISTS` 跳过而补列语句撞上 "column ... already exists",
整段脚本第二次执行即失败——而"可重复执行"正是这份脚本对下游的承诺。
"""
fresh_dsn, schema = fresh_schema
script = telemetry_schema_sql("postgres")
await _execute_script(fresh_dsn, script)
actual = [r["column_name"] for r in await _fetch(fresh_dsn, _PHYSICAL_COLUMNS_SQL, schema)]
# 物理列 = 24 个 INSERT 字段 + 库从不显式写的 created_at;对着库常量比,不另抄一份
assert set(actual) == set(COLUMNS) | {"created_at"}
# 列序也不许漂: 新列必须排在 created_at 之后,否则新建库与 ALTER 升级的列序分叉
assert actual == _EXPECTED_COLUMNS
await _execute_script(fresh_dsn, script) # 可重复执行: 第二遍不得抛
rerun = [r["column_name"] for r in await _fetch(fresh_dsn, _PHYSICAL_COLUMNS_SQL, schema)]
assert rerun == actual # 且第二遍没有偷偷改动表结构
+2 -6
View File
@@ -321,9 +321,7 @@ class TestStallBudget:
async def advance(_n):
clock.advance(_STALL)
mw = _mw(
[src], limiter, [], clock=clock, sleep=BoundedSleep(advance), transport=transport
)
mw = _mw([src], limiter, [], clock=clock, sleep=BoundedSleep(advance), transport=transport)
with pytest.raises(AllSourcesExhausted) as ei:
await mw(_REQ)
assert ei.value.reason == "stalled" # 不是 retry_exhausted: 429 确实没烧重试预算
@@ -541,9 +539,7 @@ class TestUnknownSourceIsAssemblyDefect:
"""
src = make_source("s1")
# 限流后端的源名单与治理循环拿到的源对不上 = 装配缺陷
limiter = InMemoryLimiter(
scope="llm", sources={"other": src}, global_limits=_NO_GLOBAL
)
limiter = InMemoryLimiter(scope="llm", sources={"other": src}, global_limits=_NO_GLOBAL)
gate = QuotaGate(limiter, scope="llm")
with pytest.raises(SourceNotConfiguredError) as ei:
await getattr(gate, method)(src)
+65
View File
@@ -331,6 +331,61 @@ class TestAssemblyGuards:
assert GatewaySettings.from_env("LLM", env=env_ok).backpressure.stall_window_s == 60.0
class TestTelemetrySchemaMode:
"""PGW_TELEMETRY_SCHEMA_MODE 三态(issue #13 设计 §4.1)。
键未设时按后端**不对称**派生: SQLite 是下游自己的本地文件(没有 DBA、
没有迁移工具、没有第二个系统碰它),补列是毫秒级元数据操作,故默认 auto;
PG 是共享生产表,ALTER 取 ACCESS EXCLUSIVE 锁会阻塞该表其后的所有查询,
而遥测是业务路径上的内联 await,故默认 manual。显式设置两侧都可覆盖——
"可覆盖"正是三态相对两态多出来的那一态,派生本身盖不住它。
"""
def _sqlite_env(self, **overrides):
return _env(
PGW_TELEMETRY_BACKEND="sqlite",
PGW_TELEMETRY_SQLITE_PATH="logs/telemetry.db",
**overrides,
)
def _pg_env(self, **overrides):
return _env(
PGW_TELEMETRY_BACKEND="postgres",
PGW_TELEMETRY_PG_DSN="postgresql://u:p@h:5432/polygateway",
**overrides,
)
def test_unset_key_derives_auto_for_sqlite(self):
s = GatewaySettings.from_env("LLM", env=self._sqlite_env())
assert s.telemetry_auto_migrate is True
def test_unset_key_derives_manual_for_postgres(self):
s = GatewaySettings.from_env("LLM", env=self._pg_env())
assert s.telemetry_auto_migrate is False
def test_unset_key_derives_manual_for_none_backend(self):
"""backend=none 无 recorder 消费该字段,派生结果必须是 False 而非 sqlite 那档。"""
s = GatewaySettings.from_env("LLM", env=_env())
assert s.telemetry_auto_migrate is False
def test_explicit_manual_overrides_sqlite_default(self):
s = GatewaySettings.from_env(
"LLM", env=self._sqlite_env(PGW_TELEMETRY_SCHEMA_MODE="manual")
)
assert s.telemetry_auto_migrate is False
def test_explicit_auto_overrides_postgres_default(self):
s = GatewaySettings.from_env("LLM", env=self._pg_env(PGW_TELEMETRY_SCHEMA_MODE="auto"))
assert s.telemetry_auto_migrate is True
def test_invalid_mode_rejected_naming_the_env_key(self):
"""报错须点出 env 键名: 这条路的调用方看得懂的是键名,不是字段名。"""
with pytest.raises(ValueError, match="PGW_TELEMETRY_SCHEMA_MODE"):
GatewaySettings.from_env(
"LLM", env=self._sqlite_env(PGW_TELEMETRY_SCHEMA_MODE="enabled")
)
class TestOcrSettings:
"""M3 OcrSettings(设计 §3.4): 复用 GatewaySettings,无 OCR 专用键。"""
@@ -555,6 +610,16 @@ class TestCrossFieldInvariants:
with pytest.raises(ValueError, match="telemetry_pg_dsn"):
dataclasses.replace(base, telemetry_backend="postgres")
def test_none_backend_forces_auto_migrate_off(self):
"""backend=none 时没有 recorder 消费该字段,True 是自相矛盾的状态(issue #13)。
env 路的派生已给出 False,但直接构造与 dataclasses.replace 这两条同等
官方的装配路仍能把 True 传进来——不变量归位到构造期,三条路才一致。
"""
base = self._base() # telemetry_backend="none"
replaced = dataclasses.replace(base, telemetry_auto_migrate=True)
assert replaced.telemetry_auto_migrate is False
# —— 标量域 ——
def test_negative_structured_retries_rejected(self):
+12
View File
@@ -25,3 +25,15 @@ def test_ocr_public_surface_exported():
):
assert hasattr(polygateway, name), name
assert name in polygateway.__all__, name
def test_telemetry_schema_sql_exported():
"""issue #13: manual 档下游需要主动索取"库要求的最小 schema"的顶层入口。
同时钉住公共面**只增这一个名字**: `missing_columns_warning` 是 recorder 内部
共用的文案构造函数,导出它等于多一份永久承诺(库承诺公共面只增不删)。
"""
assert "telemetry_schema_sql" in polygateway.__all__
assert callable(polygateway.telemetry_schema_sql)
assert "missing_columns_warning" not in polygateway.__all__
assert not hasattr(polygateway, "missing_columns_warning")
+386 -36
View File
@@ -115,33 +115,169 @@ async def _record_minimal(recorder, call_id="c1", **overrides):
await recorder.record_llm_call(**fields)
class TestBackendColumnParity:
"""两个后端的 `_COLUMNS` 必须逐字同名同序(issue #11)。
@pytest.fixture
def captured_warnings():
"""捕获库发出的 WARNING;loguru 不经标准 logging,pytest 的 caplog 抓不到。
emitter 只组装一份 `fields`,两个后端各自按自己的 `_COLUMNS` 取值;两份清单
一旦分叉,同一次调用在 SQLite 上写得进、在 PG 上抛 KeyError 被降级吞掉,
差异只在换后端时才暴露。列**序**同样断言: INSERT 用位置占位符,顺序错位
会把值写进错误的列而不报错。
名字避开裸 `warnings`: 那会遮蔽标准库模块名,本文件将来任何一次
`import warnings` 都会与它静默互相顶掉,而报错点离真因很远。
"""
from loguru import logger
messages: list[str] = []
sink_id = logger.add(messages.append, level="WARNING")
yield messages
logger.remove(sink_id)
# 搬迁前(1.2.1)两个 recorder 各自持有的 INSERT 常量原文,逐字冻结在此。
# 这两条字符串是"纯搬迁不改行为"的机械证据: 构造逻辑换了地方,产物必须一字不差。
_FROZEN_SQLITE_INSERT = (
"INSERT OR IGNORE INTO llm_calls (call_id, parent_call_id, session_id, model, provider, "
"source_name, messages, response, thinking, prompt_tokens, completion_tokens, usage_source, "
"latency_ms, ttft_ms, max_inter_token_ms, cache_hit, error, cost, cached_prompt_tokens, "
"model_reported, sampling, reasoning_tokens, tenant_id, meta) "
"VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)"
)
_FROZEN_PG_INSERT = (
"INSERT INTO llm_calls (call_id, parent_call_id, session_id, model, provider, source_name, "
"messages, response, thinking, prompt_tokens, completion_tokens, usage_source, latency_ms, "
"ttft_ms, max_inter_token_ms, cache_hit, error, cost, cached_prompt_tokens, model_reported, "
"sampling, reasoning_tokens, tenant_id, meta) "
"VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16, $17, $18, "
"$19, $20, $21, $22, $23, $24) "
# 无冲突目标(issue #13 Task 2): 带 `(call_id)` 的版本在按 created_at 分区、
# 主键为 (call_id, created_at) 的表上匹配不到约束,PG 直接拒收整条写入
"ON CONFLICT DO NOTHING"
)
def _first_occurrence_order(text: str, names: list[str]) -> list[str]:
"""按各列名在 text 中首次出现的位置排序,用于比对"列名出现顺序""""
found = [(text.index(name), name) for name in names if name in text]
return [name for _, name in sorted(found)]
class TestSchemaModule:
"""`telemetry/schema.py` 是 schema 单一事实源(issue #13 Task 1)。
库执行的 DDL 与打印给下游的 SQL 必须同源: 常量分散在两个 recorder 里各存一份时,
公共函数再写一份就是三份,漂移的表现是"下游照打印的 SQL 建完表,库仍报缺列"
"""
def test_columns_and_ddl_are_frozen(self):
"""列序与两端 DDL 逐字未变(搬迁不得改动任何一个字符)。"""
from polygateway.telemetry.schema import COLUMNS, PG_BACKFILL, PG_DDL, SQLITE_DDL
# COLUMNS 是 INSERT 字段序,不含数据库自填的 created_at
assert list(COLUMNS) == [c for c in _EXPECTED_COLUMNS if c != "created_at"]
assert len(COLUMNS) == 24
# 两端 DDL 的列出现顺序 == 物理列序(created_at 在第 19 位)
for ddl in (SQLITE_DDL, PG_DDL):
assert _first_occurrence_order(ddl, _EXPECTED_COLUMNS) == _EXPECTED_COLUMNS
assert "CREATE TABLE IF NOT EXISTS llm_calls" in SQLITE_DDL
assert "created_at TEXT NOT NULL DEFAULT (datetime('now'))" in SQLITE_DDL
assert "created_at TIMESTAMPTZ NOT NULL DEFAULT now()" in PG_DDL
assert "meta JSONB NOT NULL DEFAULT '{}'::jsonb" in PG_DDL
# 库内执行的补列语句不带 IF NOT EXISTS(它即便列已存在也先取 ACCESS EXCLUSIVE 锁)
assert PG_BACKFILL[0] == (
"cached_prompt_tokens",
"ALTER TABLE llm_calls ADD COLUMN cached_prompt_tokens INTEGER",
)
assert PG_BACKFILL[-1] == (
"meta",
"ALTER TABLE llm_calls ADD COLUMN meta JSONB NOT NULL DEFAULT '{}'::jsonb",
)
assert all("IF NOT EXISTS" not in stmt for _, stmt in PG_BACKFILL)
def test_insert_sql_reproduces_the_frozen_statements(self):
"""`insert_sql(backend, COLUMNS)` 与搬迁前的 `_INSERT` 一致(PG 侧去掉冲突目标)。"""
from polygateway.telemetry.schema import COLUMNS, insert_sql
assert insert_sql("sqlite", COLUMNS) == _FROZEN_SQLITE_INSERT
assert insert_sql("postgres", COLUMNS) == _FROZEN_PG_INSERT
# 裁剪列表按位置占位符重新编号,不留空洞
assert insert_sql("postgres", ["call_id", "model"]) == (
"INSERT INTO llm_calls (call_id, model) VALUES ($1, $2) ON CONFLICT DO NOTHING"
)
# 冲突目标不得被"顺手"补回: 分区表上它会让每一条遥测都被 PG 拒收
assert "ON CONFLICT (" not in insert_sql("postgres", COLUMNS)
def test_insert_sql_rejects_foreign_columns_and_backends(self):
"""列名来自数据库探测结果而非常量,子集校验是唯一的注入面闸门。"""
from polygateway.telemetry.schema import COLUMNS, insert_sql
with pytest.raises(ValueError, match="call_id_x"):
insert_sql("sqlite", ["call_id_x"])
with pytest.raises(ValueError):
insert_sql("sqlite", ["call_id", "meta); DROP TABLE llm_calls; --"])
with pytest.raises(ValueError, match="mysql"):
insert_sql("mysql", COLUMNS)
def test_insert_sql_rejects_an_empty_column_set(self):
"""空列集合两端都拼出语法非法的 SQL,构造器自己拒,不押在调用方的不变量上。
入参来自数据库探测结果: 探测到一张与本库毫无共同列的同名表,`effective`
就是空的。真放行会产出 `INSERT OR IGNORE INTO llm_calls () VALUES ()`,
错误要到执行时才由数据库报,离真因很远。
"""
from polygateway.telemetry.schema import insert_sql
for backend in ("sqlite", "postgres"):
with pytest.raises(ValueError, match="至少需要一列"):
insert_sql(backend, [])
def test_schema_sql_is_paste_ready_and_same_source(self):
"""打印给下游的脚本与库执行的 DDL 同源,且对人可重复执行。"""
from polygateway.telemetry.schema import PG_BACKFILL, SQLITE_BACKFILL, telemetry_schema_sql
pg = telemetry_schema_sql("postgres")
lite = telemetry_schema_sql("sqlite")
for script in (pg, lite):
# 24 个 INSERT 字段 + created_at 全在,且首次出现顺序与建表 DDL 一致
assert _first_occurrence_order(script, _EXPECTED_COLUMNS) == _EXPECTED_COLUMNS
assert "CREATE TABLE IF NOT EXISTS llm_calls" in script
# 人执行的那份必须幂等: PG 用 ADD COLUMN IF NOT EXISTS(与库内那份有意不同)
for column, _ in PG_BACKFILL:
assert f"ALTER TABLE llm_calls ADD COLUMN IF NOT EXISTS {column} " in pg
# SQLite 无该语法(写上去直接语法错误),只能以注释交代执行前提
lite_alters = [line for line in lite.splitlines() if line.startswith("ALTER TABLE")]
assert len(lite_alters) == len(SQLITE_BACKFILL)
assert all("IF NOT EXISTS" not in line for line in lite_alters)
for column, _ in SQLITE_BACKFILL:
assert f"ALTER TABLE llm_calls ADD COLUMN {column} " in lite
assert "不存在" in lite
with pytest.raises(ValueError, match="mysql"):
telemetry_schema_sql("mysql")
class TestBackendColumnParity:
"""两个后端的列清单必须逐字同名同序(issue #11)。
emitter 只组装一份 `fields`,两个后端各按自己的清单取值;两份清单一旦分叉,
同一次调用在 SQLite 上写得进、在 PG 上抛 KeyError 被降级吞掉,差异只在换后端时
才暴露。issue #13 起两端共用 `schema.COLUMNS`,故这里断言的是"共用"本身
(同一个对象则永远无从分叉),列**序**仍单独断言: INSERT 用位置占位符,
顺序错位会把值写进错误的列而不报错。
"""
def test_two_backends_agree_on_columns(self):
from polygateway.telemetry.postgres import _COLUMNS as PG_COLUMNS
from polygateway.telemetry.sqlite import _COLUMNS as SQLITE_COLUMNS
from polygateway.telemetry import postgres, sqlite
from polygateway.telemetry.schema import COLUMNS
assert SQLITE_COLUMNS == PG_COLUMNS
assert sqlite.COLUMNS is COLUMNS
assert postgres.COLUMNS is COLUMNS
def test_caller_dimensions_are_appended_last(self):
"""新列只能追加在末尾: 旧表经 ALTER 补列必落末尾,插在中间会让两条路径分叉。"""
from polygateway.telemetry.postgres import _COLUMNS as PG_COLUMNS
from polygateway.telemetry.sqlite import _COLUMNS as SQLITE_COLUMNS
from polygateway.telemetry.schema import COLUMNS
assert SQLITE_COLUMNS[-2:] == ("tenant_id", "meta")
assert PG_COLUMNS[-2:] == ("tenant_id", "meta")
assert COLUMNS[-2:] == ("tenant_id", "meta")
class TestSQLiteRecorder:
async def test_schema_has_frozen_columns(self, tmp_path):
recorder = SQLiteRecorder(tmp_path / "t.db")
recorder = SQLiteRecorder(tmp_path / "t.db", auto_migrate=True)
await _record_minimal(recorder)
recorder.close()
cols = [
@@ -150,7 +286,7 @@ class TestSQLiteRecorder:
assert cols == _EXPECTED_COLUMNS
async def test_call_id_idempotent(self, tmp_path):
recorder = SQLiteRecorder(tmp_path / "t.db")
recorder = SQLiteRecorder(tmp_path / "t.db", auto_migrate=True)
await _record_minimal(recorder, call_id="dup")
await _record_minimal(recorder, call_id="dup", response="second")
recorder.close()
@@ -162,7 +298,7 @@ class TestSQLiteRecorder:
assert rows == [("ok",)] # INSERT OR IGNORE: 第二次静默忽略
async def test_concurrent_writes_all_land(self, tmp_path):
recorder = SQLiteRecorder(tmp_path / "t.db")
recorder = SQLiteRecorder(tmp_path / "t.db", auto_migrate=True)
await asyncio.gather(*(_record_minimal(recorder, call_id=f"c{i}") for i in range(50)))
recorder.close()
(count,) = (
@@ -171,12 +307,12 @@ class TestSQLiteRecorder:
assert count == 50
async def test_unwritable_path_degrades_silently(self):
recorder = SQLiteRecorder(Path("/nonexistent-root/deep/t.db"))
recorder = SQLiteRecorder(Path("/nonexistent-root/deep/t.db"), auto_migrate=True)
await _record_minimal(recorder) # 不抛
recorder.close()
async def test_observability_columns_round_trip(self, tmp_path):
recorder = SQLiteRecorder(tmp_path / "t.db")
recorder = SQLiteRecorder(tmp_path / "t.db", auto_migrate=True)
await _record_minimal(recorder, call_id="c-hit", cached_prompt_tokens=64)
await _record_minimal(recorder, call_id="c-zero", cached_prompt_tokens=0)
await _record_minimal(recorder, call_id="c-none", model_reported="MiniMax-Text-01")
@@ -192,7 +328,7 @@ class TestSQLiteRecorder:
async def test_reasoning_tokens_column_round_trip(self, tmp_path):
"""issue #6: 7 / 0 / None 三种值各自如实落库,0 与 NULL 不得混同。"""
recorder = SQLiteRecorder(tmp_path / "t.db")
recorder = SQLiteRecorder(tmp_path / "t.db", auto_migrate=True)
await _record_minimal(recorder, call_id="r-some", reasoning_tokens=7)
await _record_minimal(recorder, call_id="r-zero", reasoning_tokens=0)
await _record_minimal(recorder, call_id="r-none", reasoning_tokens=None)
@@ -208,7 +344,7 @@ class TestSQLiteRecorder:
async def test_sampling_column_round_trips(self, tmp_path):
"""issue #4: 采样参数落库,否则事后无法证明某批数据跑在什么温度下。"""
recorder = SQLiteRecorder(tmp_path / "t.db")
recorder = SQLiteRecorder(tmp_path / "t.db", auto_migrate=True)
await _record_minimal(recorder, call_id="c-s", sampling='{"seed": 42, "temperature": 0}')
await _record_minimal(recorder, call_id="c-plain")
recorder.close()
@@ -255,7 +391,7 @@ class TestSQLiteColumnBackfill:
legacy.commit()
legacy.close()
recorder = SQLiteRecorder(db)
recorder = SQLiteRecorder(db, auto_migrate=True)
await _record_minimal(recorder, cached_prompt_tokens=7, model_reported="m-real")
recorder.close()
@@ -280,7 +416,7 @@ class TestSQLiteColumnBackfill:
conn.commit()
conn.close()
recorder = SQLiteRecorder(db) # 不得抛
recorder = SQLiteRecorder(db, auto_migrate=True) # 不得抛
assert recorder._conn is not None # 补列失败 ≠ recorder 失能(D1 纪律)
await _record_minimal(recorder) # 不得抛
recorder.close()
@@ -337,7 +473,7 @@ class TestSQLiteCallerDimensionsAcceptance:
async def test_fresh_db_round_trips_the_dimensions(self, tmp_path):
"""新建库: 列齐全,且维度值原样读回——只验列存在会漏掉写错列位的错。"""
db = tmp_path / "fresh.db"
recorder = SQLiteRecorder(db)
recorder = SQLiteRecorder(db, auto_migrate=True)
await _record_minimal(
recorder, call_id="c-dim", tenant_id="tenant-a", meta='{"batch": "b7"}'
)
@@ -363,7 +499,7 @@ class TestSQLiteCallerDimensionsAcceptance:
db = tmp_path / "pre_tenant.db"
_make_pre_tenant_db(db)
recorder = SQLiteRecorder(db)
recorder = SQLiteRecorder(db, auto_migrate=True)
await _record_minimal(recorder, call_id="new-row", tenant_id="tenant-a", meta='{"k": 1}')
recorder.close()
@@ -402,7 +538,7 @@ class TestSQLiteCallerDimensionsAcceptance:
# finally 还原权限位: 任一断言先失败时,不还原会让 tmp_path 清理连带报错,
# 把"某条断言失败"的真因盖成一个无关的 PermissionError
try:
recorder = SQLiteRecorder(db) # 不得抛
recorder = SQLiteRecorder(db, auto_migrate=True) # 不得抛
assert recorder._conn is not None # 补列失败 ≠ recorder 失能
await _record_minimal(recorder, call_id="doomed") # 只读库写不进,但不得抛
recorder.close()
@@ -415,6 +551,144 @@ class TestSQLiteCallerDimensionsAcceptance:
) # 补列确实没成功,用例不是在只读库上空转
class TestSQLiteSchemaMode:
"""issue #13: `auto_migrate` 两档——auto 保持自动补列,manual 只裁剪写入不发 DDL。
列数断言一律按**物理列数**写: 旧表 22 个 INSERT 字段 + `created_at` = 23,
补齐后 24 + `created_at` = 25。混用 INSERT 字段数与物理列数是本处最易错的地方。
"""
def _physical_columns(self, db: Path) -> list[str]:
conn = sqlite3.connect(db)
try:
return [r[1] for r in conn.execute("PRAGMA table_info(llm_calls)")]
finally:
conn.close()
async def test_manual_mode_trims_the_insert_instead_of_altering(
self, tmp_path, captured_warnings
):
"""manual + 22 字段旧表: 一条 ALTER 都不发,写入按现有列裁剪后照样落库。
裁剪是关掉 ALTER 的前提: 不裁剪的话每行 INSERT 都撞 `no column named
tenant_id` 而被整行丢弃——那是把自动补列换成静默全失能。
"""
db = tmp_path / "manual_legacy.db"
_make_pre_tenant_db(db)
recorder = SQLiteRecorder(db, auto_migrate=False)
await _record_minimal(recorder, call_id="new-row", tenant_id="tenant-a", meta='{"k": 1}')
recorder.close()
assert len(self._physical_columns(db)) == 23 # 未 ALTER: 物理列数原封不动
conn = sqlite3.connect(db)
assert conn.execute(
"SELECT response, model FROM llm_calls WHERE call_id = 'new-row'"
).fetchone() == ("ok", "m") # 裁剪后的列值仍对得上位
conn.close()
assert len(captured_warnings) == 1 # 缺列只讲一次,不逐行刷屏
message = captured_warnings[0]
assert "tenant_id" in message and "meta" in message # 逐列点名
assert "不会被记录" in message # 讲清后果
assert "ALTER TABLE" in message # 给出可直接执行的补列 SQL
async def test_auto_mode_still_upgrades_the_legacy_table(self, tmp_path):
"""auto + 同款旧表: 现状回归,补列后物理列数 23 → 25。"""
db = tmp_path / "auto_legacy.db"
_make_pre_tenant_db(db)
recorder = SQLiteRecorder(db, auto_migrate=True)
await _record_minimal(recorder, call_id="new-row", tenant_id="tenant-a")
recorder.close()
assert self._physical_columns(db) == _EXPECTED_COLUMNS
assert len(self._physical_columns(db)) == 25
async def test_manual_mode_still_creates_a_fresh_table(self, tmp_path):
"""manual 只管 ALTER,不管 CREATE: 全新库照建,25 个物理列齐全(设计 §4.2)。"""
db = tmp_path / "manual_fresh.db"
recorder = SQLiteRecorder(db, auto_migrate=False)
await _record_minimal(recorder, call_id="c-fresh", tenant_id="tenant-a")
recorder.close()
assert self._physical_columns(db) == _EXPECTED_COLUMNS
conn = sqlite3.connect(db)
assert (
conn.execute("SELECT tenant_id FROM llm_calls WHERE call_id = 'c-fresh'").fetchone()[0]
== "tenant-a"
)
conn.close()
async def test_table_without_call_id_escalates_the_wording(self, tmp_path, captured_warnings):
"""缺主键列 call_id = 该表压根不是本库的 llm_calls: 措辞升级,但库不做二次判定。"""
db = tmp_path / "alien.db"
conn = sqlite3.connect(db)
conn.execute("CREATE TABLE llm_calls (model TEXT, provider TEXT)")
conn.commit()
conn.close()
recorder = SQLiteRecorder(db, auto_migrate=False) # 不得抛
await _record_minimal(recorder) # 照常尝试写入
recorder.close()
message = "\n".join(captured_warnings)
assert "call_id" in message
assert "不是本库" in message
async def test_no_recognizable_column_falls_back_to_the_full_column_set(
self, tmp_path, captured_warnings
):
"""探测结果与 COLUMNS 毫无交集视同探测异常: 保守回落全量列。
`insert_sql` 自己拒空列集合(见 `test_insert_sql_rejects_an_empty_column_set`),
故这里回落不发生就不是"拼出空语句",而是 ValueError 逃出 `__init__` ——
遥测初始化失败必须静默降级,崩溃比丢维度严重得多。
"""
from polygateway.telemetry.schema import COLUMNS
db = tmp_path / "foreign.db"
conn = sqlite3.connect(db)
conn.execute("CREATE TABLE llm_calls (foo TEXT, bar TEXT)")
conn.commit()
conn.close()
recorder = SQLiteRecorder(db, auto_migrate=False) # 不得抛
assert recorder._columns == COLUMNS
await _record_minimal(recorder) # 写不进去,但只逐行 warning,不抛
recorder.close()
assert captured_warnings # 沉默地退化成空语句是最坏结果,必须有声
async def test_empty_probe_result_degrades_instead_of_raising(
self, tmp_path, captured_warnings
):
"""探测返回空集合时走回落,绝不让 `insert_sql` 的 ValueError 逃出去。
SQLite 建不出零列的表,故直接喂空探测结果调那条分支——它正是
`insert_sql` 拒空之后唯一可能把"静默降级"变成崩溃的入口。
"""
from polygateway.telemetry.schema import COLUMNS
db = tmp_path / "empty_probe.db"
recorder = SQLiteRecorder(db, auto_migrate=False)
recorder._adopt_existing_columns(set()) # 不得抛
assert recorder._columns == COLUMNS
await _record_minimal(recorder, call_id="c-after") # 写入照常
recorder.close()
conn = sqlite3.connect(db)
assert conn.execute(
"SELECT response FROM llm_calls WHERE call_id = 'c-after'"
).fetchone() == ("ok",)
conn.close()
assert [m for m in captured_warnings if "没有任何本库认识的列" in m] # 只有 warning
async def test_auto_migrate_is_required_keyword_only(self, tmp_path):
"""关键行为参数不给默认值(P4): 缺省规则只写在 config 一处,不与类签名漂移。"""
with pytest.raises(TypeError):
SQLiteRecorder(tmp_path / "t.db") # type: ignore[call-arg]
class _FakePgConn:
"""记录执行过的语句;可让 ALTER/CREATE/探测抛错以模拟权限不足与抖动。
@@ -493,7 +767,9 @@ class TestPostgresBackfillDiscipline:
def _recorder(self, conn):
from polygateway.telemetry.postgres import PostgresRecorder
return PostgresRecorder("postgresql://u:p@h:5432/polygateway", pool=_FakePgPool(conn))
return PostgresRecorder(
"postgresql://u:p@h:5432/polygateway", pool=_FakePgPool(conn), auto_migrate=True
)
async def test_alter_failure_does_not_disable_the_recorder(self):
"""ALTER 失败(如账号只有 INSERT 权限)不得置 _failed —— 那会让遥测全灭。"""
@@ -516,10 +792,10 @@ class TestPostgresBackfillDiscipline:
async def test_missing_columns_are_added_once(self):
conn = _FakePgConn(self._LEGACY)
await _record_minimal(self._recorder(conn))
from polygateway.telemetry.postgres import _BACKFILL
from polygateway.telemetry.schema import PG_BACKFILL
altered = [s for s in conn.statements if s.startswith("ALTER TABLE")]
assert len(altered) == len(_BACKFILL) # 旧表缺全部补列,故一列一条 ALTER
assert len(altered) == len(PG_BACKFILL) # 旧表缺全部补列,故一列一条 ALTER
assert all("IF NOT EXISTS" not in s for s in altered) # 探测已确认缺列,无需再判
@@ -547,7 +823,9 @@ class TestPostgresTableProbe:
def _recorder(self, conn):
from polygateway.telemetry.postgres import PostgresRecorder
return PostgresRecorder("postgresql://u:p@h:5432/polygateway", pool=_FakePgPool(conn))
return PostgresRecorder(
"postgresql://u:p@h:5432/polygateway", pool=_FakePgPool(conn), auto_migrate=True
)
def _created(self, conn):
return [s for s in conn.statements if s.lstrip().startswith("CREATE TABLE")]
@@ -595,6 +873,79 @@ class TestPostgresTableProbe:
assert [s for s in conn.statements if s.startswith("INSERT INTO llm_calls")]
class TestPostgresSchemaMode:
"""issue #13: PG 侧 manual 档一条 ALTER 都不发,改按现有列裁剪 INSERT。
真实 PG 的验收在 `tests/integration/test_postgres_telemetry.py`;这里用 fake 连接
锁住"发了哪些语句",无 DSN 环境下集成用例被 skip 时仍有回归保护。
"""
_LEGACY = ["call_id", "cost", "created_at"]
def _recorder(self, conn, *, auto_migrate):
from polygateway.telemetry.postgres import PostgresRecorder
return PostgresRecorder(
"postgresql://u:p@h:5432/polygateway", pool=_FakePgPool(conn), auto_migrate=auto_migrate
)
async def test_manual_mode_trims_the_insert_instead_of_altering(self, captured_warnings):
conn = _FakePgConn(self._LEGACY)
recorder = self._recorder(conn, auto_migrate=False)
await _record_minimal(recorder)
assert not [s for s in conn.statements if s.startswith("ALTER TABLE")]
assert "INSERT INTO llm_calls (call_id, cost) VALUES ($1, $2) ON CONFLICT DO NOTHING" in (
conn.statements
)
assert recorder._columns == ("call_id", "cost")
message = "\n".join(captured_warnings)
assert "tenant_id" in message and "meta" in message # 逐列点名
assert "不会被记录" in message # 讲清后果
assert "ALTER TABLE llm_calls ADD COLUMN tenant_id" in message # 可直接执行的 SQL
async def test_manual_mode_still_creates_a_missing_table(self):
"""manual 只管 ALTER 不管 CREATE: 新建表列已齐全,写入照发全量列。"""
from polygateway.telemetry.schema import COLUMNS
conn = _FakePgConn([])
recorder = self._recorder(conn, auto_migrate=False)
await _record_minimal(recorder)
assert [s for s in conn.statements if s.lstrip().startswith("CREATE TABLE")]
assert recorder._columns == COLUMNS
async def test_no_recognizable_column_falls_back_to_the_full_column_set(
self, captured_warnings
):
"""PG 侧同款回落(SQLite 侧对称用例见 TestSQLiteSchemaMode)。
表存在(`to_regclass` 非空)但列与 `COLUMNS` 毫无交集: 裁剪结果为空,
必须回落全量而不是把空列集交给 `insert_sql`——`_prepare_schema` 里那次
调用在 try 之外,ValueError 会顺着 `record_llm_call` 冒给业务调用方。
"""
from polygateway.telemetry.schema import COLUMNS
conn = _FakePgConn(["foo", "bar"])
recorder = self._recorder(conn, auto_migrate=False)
await _record_minimal(recorder) # 不得抛
assert recorder._columns == COLUMNS
assert not [s for s in conn.statements if s.startswith("ALTER TABLE")]
assert [m for m in captured_warnings if "没有任何本库认识的列" in m]
async def test_auto_mode_still_backfills(self):
"""auto 档现状回归: 缺列照补,补完写全量列。"""
from polygateway.telemetry.schema import COLUMNS, PG_BACKFILL
conn = _FakePgConn(self._LEGACY)
recorder = self._recorder(conn, auto_migrate=True)
await _record_minimal(recorder)
assert len([s for s in conn.statements if s.startswith("ALTER TABLE")]) == len(PG_BACKFILL)
assert recorder._columns == COLUMNS
class _MemoryRecorder:
def __init__(self):
self.rows = []
@@ -604,16 +955,15 @@ class _MemoryRecorder:
class TestEmitterRecorderContract:
"""emitter 的实参键集合必须与两个后端的 _COLUMNS 完全一致(issue #3)。
"""emitter 的实参键集合必须与后端的 `schema.COLUMNS` 完全一致(issue #3)。
两个后端的 `row = tuple(fields[col] for col in _COLUMNS)` 都在 try **之外**,
两个后端的 `row = tuple(fields[col] for col in COLUMNS)` 都在 try **之外**,
emitter 漏传一个键就抛 KeyError,被 `_record` 的 except Exception 吞成 warning
→ 遥测静默全丢。而 8 个 `**fields` 形态的 fake 一个都拦不住,故显式断言。
"""
async def test_emitter_supplies_exactly_the_backend_columns(self):
from polygateway.telemetry.postgres import _COLUMNS as PG_COLUMNS
from polygateway.telemetry.sqlite import _COLUMNS as SQLITE_COLUMNS
from polygateway.telemetry.schema import COLUMNS
rec = _MemoryRecorder()
await TelemetryEmitter(rec).emit_attempt(
@@ -624,11 +974,11 @@ class TestEmitterRecorderContract:
response=_resp(),
error=None,
)
assert set(rec.rows[0]) == set(SQLITE_COLUMNS) == set(PG_COLUMNS)
assert set(rec.rows[0]) == set(COLUMNS)
@pytest.mark.parametrize("emit", ["attempt", "cache_hit", "terminal_failure"])
async def test_every_entry_point_supplies_the_same_keys(self, emit):
from polygateway.telemetry.sqlite import _COLUMNS as SQLITE_COLUMNS
from polygateway.telemetry.schema import COLUMNS
rec = _MemoryRecorder()
emitter = TelemetryEmitter(rec)
@@ -647,7 +997,7 @@ class TestEmitterRecorderContract:
await emitter.emit_terminal_failure(
request=_REQ, call_id="c", latency_ms=1, error="dead"
)
assert set(rec.rows[0]) == set(SQLITE_COLUMNS)
assert set(rec.rows[0]) == set(COLUMNS)
class TestEmitterObservabilityFields:
+1 -1
View File
@@ -127,7 +127,7 @@ async def _worker_async(args: argparse.Namespace, worker_idx: int) -> None:
env = _merged_env()
run_id = args.run_id
telemetry_path = _ROOT / f"data/soak/telemetry_{run_id}_{worker_idx}.db"
recorder = SQLiteRecorder(telemetry_path)
recorder = SQLiteRecorder(telemetry_path, auto_migrate=True)
if args.scenario == "P7":
from polygateway.ocr import OcrClient