diff --git a/.env.example b/.env.example index 14b7c58..08ca2b9 100644 --- a/.env.example +++ b/.env.example @@ -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 # 可选: {"": {"input_per_1m": x, "output_per_1m": y}};缺省 cost 恒 None # # 可选第三档 "cached_input_per_1m": z —— 供应商 prompt cache 命中部分的单价; diff --git a/CHANGELOG.md b/CHANGELOG.md index ad662b0..21638a8 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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` 两个调用方自填、库内不校验的自由字符串。 diff --git a/README.md b/README.md index 306168c..90eaf0d 100644 --- a/README.md +++ b/README.md @@ -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)。 ## 可靠性证据 diff --git a/research-wiki/ARCHITECTURE.md b/research-wiki/ARCHITECTURE.md index f321858..3d4653b 100644 --- a/research-wiki/ARCHITECTURE.md +++ b/research-wiki/ARCHITECTURE.md @@ -119,7 +119,7 @@ HTTP API → arq 队列 → worker 协程 脚本 → asyncio.gather 协 --- -## 3. 架构决策记录(D1–D14,含讨论过程与备选方案) +## 3. 架构决策记录(D1–D15,含讨论过程与备选方案) > 每条决策记录格式:**决策 / 背景与讨论 / 被否决的备选 / 影响**。这些决策已与人类逐条确认;推翻任何一条需要人类批准并修订本节。 @@ -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),缺省规则也就不会与类签名漂移。 --- diff --git a/src/polygateway/__init__.py b/src/polygateway/__init__.py index a7bf20f..5a3cf9f 100644 --- a/src/polygateway/__init__.py +++ b/src/polygateway/__init__.py @@ -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", ] diff --git a/src/polygateway/backends/redis/breaker.py b/src/polygateway/backends/redis/breaker.py index 3fe0de8..dd196bd 100644 --- a/src/polygateway/backends/redis/breaker.py +++ b/src/polygateway/backends/redis/breaker.py @@ -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: diff --git a/src/polygateway/backends/redis/limiter.py b/src/polygateway/backends/redis/limiter.py index 011e925..cb5d516 100644 --- a/src/polygateway/backends/redis/limiter.py +++ b/src/polygateway/backends/redis/limiter.py @@ -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: diff --git a/src/polygateway/client.py b/src/polygateway/client.py index e1413b4..f8352a6 100644 --- a/src/polygateway/client.py +++ b/src/polygateway/client.py @@ -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( diff --git a/src/polygateway/config.py b/src/polygateway/config.py index 3c1031d..7ab8e8b 100644 --- a/src/polygateway/config.py +++ b/src/polygateway/config.py @@ -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("://") diff --git a/src/polygateway/telemetry/postgres.py b/src/polygateway/telemetry/postgres.py index 3f35d96..101f240 100644 --- a/src/polygateway/telemetry/postgres.py +++ b/src/polygateway/telemetry/postgres.py @@ -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,让每个进程的首次写入都去抢共享审计表的排他锁,等于用记录基础设施 - 拖垮业务调用。探测走 ACCESS SHARE,稳态下一条 ALTER 都不会发。 + **先探测**的理由(两档共用): `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: diff --git a/src/polygateway/telemetry/schema.py b/src/polygateway/telemetry/schema.py new file mode 100644 index 0000000..dd8f34e --- /dev/null +++ b/src/polygateway/telemetry/schema.py @@ -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, ""]) diff --git a/src/polygateway/telemetry/sqlite.py b/src/polygateway/telemetry/sqlite.py index 4cb5a64..9b68ab8 100644 --- a/src/polygateway/telemetry/sqlite.py +++ b/src/polygateway/telemetry/sqlite.py @@ -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: diff --git a/tests/integration/test_governance_stack.py b/tests/integration/test_governance_stack.py index acffbcb..2b140f7 100644 --- a/tests/integration/test_governance_stack.py +++ b/tests/integration/test_governance_stack.py @@ -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], diff --git a/tests/integration/test_postgres_telemetry.py b/tests/integration/test_postgres_telemetry.py index 5850c2f..dd427cd 100644 --- a/tests/integration/test_postgres_telemetry.py +++ b/tests/integration/test_postgres_telemetry.py @@ -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 # 且第二遍没有偷偷改动表结构 diff --git a/tests/unit/test_backpressure.py b/tests/unit/test_backpressure.py index 774e904..a8d1d1a 100644 --- a/tests/unit/test_backpressure.py +++ b/tests/unit/test_backpressure.py @@ -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) diff --git a/tests/unit/test_config.py b/tests/unit/test_config.py index 0532dfc..0bf1d27 100644 --- a/tests/unit/test_config.py +++ b/tests/unit/test_config.py @@ -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): diff --git a/tests/unit/test_package.py b/tests/unit/test_package.py index 2335bdc..fd33bcf 100644 --- a/tests/unit/test_package.py +++ b/tests/unit/test_package.py @@ -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") diff --git a/tests/unit/test_telemetry.py b/tests/unit/test_telemetry.py index afec6b2..a04fb0d 100644 --- a/tests/unit/test_telemetry.py +++ b/tests/unit/test_telemetry.py @@ -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: diff --git a/tools/soak/run_soak.py b/tools/soak/run_soak.py index 53018e1..a4358db 100644 --- a/tools/soak/run_soak.py +++ b/tools/soak/run_soak.py @@ -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