From 2e028d38f2226486412ca3c03a6b8d0b023a0a37 Mon Sep 17 00:00:00 2001 From: iomgaa Date: Fri, 7 Aug 2026 11:21:33 -0400 Subject: [PATCH] fix: probe for the telemetry table before creating it PostgreSQL checks the schema CREATE privilege before the IF NOT EXISTS existence test, so an account with only table-level INSERT was denied on CREATE TABLE IF NOT EXISTS even though the table was right there and writable. The denial set _failed and the whole recorder went no-op for the process lifetime, silently: 150+ calls downstream lost their latency, token and cost rows with nothing but one warning to show for it. The probe is the direct fix. The larger fix is the criterion: structural degradation now means "provably cannot write" (pool creation failed, or the table is absent and cannot be created), not "something threw during init" -- a probe or acquire failure just skips the row and retries on the next call. SQLite stays as it is on purpose. Measured: it short-circuits the statement at parse time, so it passes even under another connection's EXCLUSIVE lock or on a read-only file. A probe there would buy nothing; the docstring now says so to keep symmetry-minded future edits away. --- research-wiki/ARCHITECTURE.md | 2 +- .../designs/issue9-telemetry-ddl-probe.md | 62 +++++++++++ research-wiki/index.md | 5 +- research-wiki/log.md | 3 + src/polygateway/telemetry/postgres.py | 103 ++++++++++++++---- src/polygateway/telemetry/sqlite.py | 8 ++ tests/integration/test_postgres_telemetry.py | 86 +++++++++++++++ tests/unit/test_telemetry.py | 98 ++++++++++++++++- 8 files changed, 340 insertions(+), 27 deletions(-) create mode 100644 research-wiki/designs/issue9-telemetry-ddl-probe.md diff --git a/research-wiki/ARCHITECTURE.md b/research-wiki/ARCHITECTURE.md index bb8dbac..5168f23 100644 --- a/research-wiki/ARCHITECTURE.md +++ b/research-wiki/ARCHITECTURE.md @@ -475,7 +475,7 @@ flowchart TB **`sampling` 列(2026-07-31,issue #4,端口 20 → 21)**: 列语义 = 「调用方采样意图 ⊎ 生效源 `extra_body`」的 canonical JSON,空则 NULL。**不含**结构化注入的 `response_format`——列名是采样参数,schema 不是,且数 KB schema 逐行落库会让审计表无谓膨胀。三个 emit 入口口径必须各自定死,否则同一列在不同行含义不同: `emit_attempt`(RetryMW 调用,**唯一**有生效源者)并上 `source.extra_body`;`emit_cache_hit` / `emit_terminal_failure`(TelemetryMW 最外层调用)无 source 可言,只记调用级——与 `model`/`source_name` 在终态行置空是同一先例,且缓存命中行无损(`sampling` 已进缓存 key,能命中即意味调用级参数与历史那次逐字相同)。三者统一读 `request.sampling` 而非 `request.overlay`(后者在 RetryMW 处已被结构化注入污染、在 TelemetryMW 处未被污染,直接用必然三行分叉)。OCR/embedding 路径因决策 G 剥离 `extra_body`,该列恒 NULL。 -(`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,锁共享审计表会拖垮业务调用),且**失败只逐行降级、绝不置结构性失能标志**。新列在 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)。 +(`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)。 - 后端: `SQLiteRecorder`(默认;WAL + busy_timeout、`INSERT OR IGNORE` 幂等、`asyncio.to_thread` 桥接、初始化/写入失败全降级不冒泡)与 `PostgresRecorder`。 - **单一 helper 铁律**: 遥测调用点收敛为一个内部函数/上下文管理器;Video-Tree 与 GovDoc 各有 4-5 处逐字复制的 `record_llm_call(15 个参数)` 是本条的直接教训。 diff --git a/research-wiki/designs/issue9-telemetry-ddl-probe.md b/research-wiki/designs/issue9-telemetry-ddl-probe.md new file mode 100644 index 0000000..ff425f4 --- /dev/null +++ b/research-wiki/designs/issue9-telemetry-ddl-probe.md @@ -0,0 +1,62 @@ +--- +type: design +node_id: design:issue9-telemetry-ddl-probe +title: "建表前先探测,判死只认「确定写不进去」" +date: 2026-08-07 +--- + +# 建表前先探测,判死只认「确定写不进去」 + +**来源**: Gitea issue #9(CHSAnalyzer3 现场)|**范围**: `telemetry/postgres.py` 单模块,无独立 plan(小改动自判)|**相关**: [[design:response-observability-fields]](issue #3 修的是同一个坑的另一半) + +## 问题 + +应用账号有表级 `INSERT`、表也已存在,但没有 schema 的 `CREATE` 权限时,`_ensure_ready()` 的 `CREATE TABLE IF NOT EXISTS` 被拒 → `_failed = True` → **整个进程遥测永久 no-op**。业务调用一切正常,只留一行 warning,从外部完全看不出异常;下游 CHSAnalyzer3 首次端到端跑的 150+ 次调用数据因此全丢且无法补回。 + +## 根因 + +**PostgreSQL 对 schema 的 CREATE 权限检查早于 `IF NOT EXISTS` 的存在性判断**(`RangeVarGetAndCheckCreationNamespace()` 先 aclcheck 后查 relid)。这与 issue #3 里 `ALTER TABLE` 的 ownership 检查早于 `IF NOT EXISTS` 是同一类问题——当时只修了补列那一半,建表这一半原样留着,于是同一账号形态下"补列失败只丢一行日志接着干活,建表失败却把整个 recorder 判死"。 + +**实测(PostgreSQL 16.14,临时角色只授 `SELECT, INSERT ON llm_calls`)**: + +| 语句 | 结果 | +|---|---| +| `SELECT to_regclass('llm_calls')` | 非 NULL(表就在那儿) | +| `CREATE TABLE IF NOT EXISTS llm_calls (...)` | **被拒 InsufficientPrivilegeError: permission denied for schema** | +| `INSERT INTO llm_calls ...` | 通过 | +| `ALTER TABLE ... ADD COLUMN IF NOT EXISTS` | 被拒 must be owner(即 issue #3 那条) | + +## 选定方案 + +两条,第二条才是治本的那条: + +1. **表存在就绝不发 DDL**。探测走 `to_regclass`(不需要任何权限,且与 `INSERT` 走同一套 search_path 解析——裸 `CREATE TABLE` 落在首个**可建**的 schema,可能与写入命中的不是同一张表,故探测优先反而更准)。表不存在才建;新建表列已齐全,顺带跳过补列。 +2. **"结构性失能"的判据从「初始化时出过异常」收窄为「确定写不进去」**: + +| 情形 | 处置 | 理由 | +|---|---|---| +| 建池失败 | 永久 no-op | 重试要在业务调用路径上内联吞掉 connect 超时 | +| 表存在 | 不发 DDL,只补列(失败仅 warning) | 本 issue 的直接修复 | +| 表不存在 → 建表成功 | 就绪,跳过补列 | 新建表列已齐 | +| 表不存在 → 建表失败 | 永久 no-op | 后续 INSERT 必然全败,重试无意义、日志纯噪音 | +| 探测/取连接失败 | 只跳过本条,下次调用重试 | 瞬时抖动,判死代价远大于多一次往返 | + +**SQLite 侧有意不对称**:实测其对已存在的表在**解析期**就把 `CREATE TABLE IF NOT EXISTS` 短路掉——另一连接持 `BEGIN EXCLUSIVE`、或文件 `chmod 444` 时该语句均通过(同条件下 `INSERT` 与新表名建表分别报 database is locked / readonly database),既不抢写锁也不检查可写性。故 PG 侧的坑在此不存在,加探测零收益。**需要对称的是保证(表存在就不该因建表失败而失能),不是代码**;结论已钉进 `sqlite.py` 模块 docstring,防止后人为"对称"加回来。 + +## 被否决的备选 + +| 备选 | 否决理由 | +|---|---| +| 只加探测,`_failed` 语义不动(issue 原方案) | 治标。初始化瞬间的 DB 抖动、一次 `pool.acquire` 失败、search_path 配错仍会让整个进程永久失遥测——同一个开关,换个触发口 | +| 除建池外一律不判死 | 方向最统一,但表真的不存在时每次调用都发一条注定失败的 INSERT + 一条 warning(150 次调用 = 150 行噪音),而这种情形是**可确定判定**的,没必要留活路 | +| 捕获 `InsufficientPrivilegeError` 特判放行 | 按异常类型打补丁,漏一种错误码就复发;探测是把"该不该发这条 DDL"判断在前,与错误面无关 | +| SQLite 侧同步加探测 | 实测证明零收益,属为对称而对称的 gold-plating | + +## 遗留 + +**SQLite 的窄缝**:表不存在 + 构造瞬间库被排他锁(多进程共库)→ `__init__` 里的建表失败 → recorder 永久失能。修它要把 SQLite 也改成 lazy 重试结构,超出本 issue 范围,记此备查。 + +## 测试证据 + +- 单测 `TestPostgresTableProbe`(5 例,fake conn):表存在不发 DDL / DDL 被拒仍照常 INSERT 且 `_failed` 不置位 / 表缺失则建表且不补列 / 表缺失且建不出来才判死 / 探测失败下次重试。 +- 集成 `TestLeastPrivilegeDeployment`(真实 PG,临时 schema + 临时角色,teardown 删净):先钉死"该角色确实建不了表"这条库外事实,再验两行记录照常落库。**修复前该用例复现 issue 原文那行 warning 并失败**。 diff --git a/research-wiki/index.md b/research-wiki/index.md index 2cb4641..47f7e61 100644 --- a/research-wiki/index.md +++ b/research-wiki/index.md @@ -1,8 +1,8 @@ # Research Wiki 索引 -> 自动生成,更新时间:2026-08-06 14:58 UTC +> 自动生成,更新时间:2026-08-07 15:11 UTC -## design (25) +## design (26) - [2026-07-20-m1-core-design](designs/2026-07-20-m1-core-design.md) `design:2026-07-20-m1-core-design` - [2026-07-20-m2-distributed-design](designs/2026-07-20-m2-distributed-design.md) `design:2026-07-20-m2-distributed-design` - [2026-07-21-m25-resilience-design](designs/2026-07-21-m25-resilience-design.md) `design:2026-07-21-m25-resilience-design` @@ -25,6 +25,7 @@ - [M4 迁移验证设计(GovDoc→CHS,发 v1.0)](designs/m4-migration.md) `design:m4-migration` - [stall 判定改为非生产性等待口径](designs/issue8-stall-budget.md) `design:issue8-stall-budget` - [响应可观测字段扩展(Issue #3)](designs/response-observability-fields.md) `design:response-observability-fields` +- [建表前先探测,判死只认「确定写不进去」](designs/issue9-telemetry-ddl-probe.md) `design:issue9-telemetry-ddl-probe` - [推理开关能力建模与 reasoning_tokens 采集(issue #5 + #6)](designs/2026-08-02-thinking-capability-design.md) `design:2026-08-02-thinking-capability-design` - [治理后端故障归位为 scope 级不可用(Issue #7)](designs/governance-backend-error.md) `design:governance-backend-error` - [采样参数透传设计(issue #4)](designs/sampling-params.md) `design:sampling-params` diff --git a/research-wiki/log.md b/research-wiki/log.md index 08d31e2..2b8a376 100644 --- a/research-wiki/log.md +++ b/research-wiki/log.md @@ -94,3 +94,6 @@ - [2026-08-06 14:58 UTC] 新增 plan: issue #8 实施计划: stall 非生产性等待口径 (plan:issue8-stall-budget-plan) - [2026-08-06 14:58 UTC] 新增边: plan:issue8-stall-budget-plan --implements--> design:issue8-stall-budget - [2026-08-06 14:58 UTC] 重建索引: 61 篇页面 +- [2026-08-07 15:20 UTC] 新增 design: 建表前先探测,判死只认「确定写不进去」 (design:issue9-telemetry-ddl-probe) +- [2026-08-07 15:20 UTC] 重建索引: 62 篇页面 +- [2026-08-07 15:11 UTC] 重建索引: 62 篇页面 diff --git a/src/polygateway/telemetry/postgres.py b/src/polygateway/telemetry/postgres.py index a195aa5..24e7986 100644 --- a/src/polygateway/telemetry/postgres.py +++ b/src/polygateway/telemetry/postgres.py @@ -1,12 +1,17 @@ """Postgres 遥测后端(M2 设计 §5): asyncpg lazy 池 + 两级降级。 参考仓无先例(三项目遥测全 SQLite);asyncpg 工程写法取 GovDoc -`taskrun/postgres_store.py`($n 占位、`CREATE TABLE IF NOT EXISTS`、 -`ON CONFLICT DO NOTHING`),但其"失败冒泡"方向按遥测铁律**有意反转**: -① 结构性失败(建池/建表)→ warning 一次后永久降级(池置 None 短路); +`taskrun/postgres_store.py`($n 占位、`ON CONFLICT DO NOTHING`),但其 +"失败冒泡"方向按遥测铁律**有意反转**: +① 结构性失败 → warning 一次后永久降级(所有写入短路); ② 运行时单条写失败 → 逐条 warning 丢弃,不降级不重试(连接抖动由 asyncpg 池自恢复;避免浸泡开头一次抖动导致后续全程失遥测)。 -构造不连库(lazy),20 列 schema 与 SQLite 版同名同序。 +构造不连库(lazy),22 列 schema 与 SQLite 版同名同序。 + +**"结构性"的判据是「确定写不进去」,不是「初始化时出过错」**(issue #9): +只有建池失败(重试要在业务路径上内联吞掉 connect 超时)与"表确定不存在 +且建不出来"(后续 INSERT 必然全败)才判死;探测失败、补列失败、取连接 +失败一律只 warning,让写入照常尝试或下次调用重试。 """ from __future__ import annotations @@ -55,6 +60,9 @@ _BACKFILL = ( ("reasoning_tokens", "ALTER TABLE llm_calls ADD COLUMN reasoning_tokens INTEGER"), ) +# 探测表是否存在;不需要任何权限,且与 INSERT 走同一套 search_path 解析 +_TABLE_EXISTS = "SELECT to_regclass('llm_calls')" + # 探测现有列;尊重 search_path(to_regclass 按当前 search_path 解析) _EXISTING_COLUMNS = ( "SELECT attname FROM pg_attribute " @@ -111,30 +119,81 @@ class PostgresRecorder: self._init_lock = asyncio.Lock() async def _ensure_ready(self) -> asyncpg.Pool | None: - """lazy 建池+建表;结构性失败 warning 一次后永久降级(设计 §5 两级之一)。""" + """lazy 建池+备表;判死只认「确定写不进去」(issue #9),其余失败都留活路。""" if self._failed: return None if self._schema_ready: return self._pool async with self._init_lock: - if self._failed or self._schema_ready: - return None if self._failed else self._pool - try: - if self._pool is None: - import asyncpg - - self._pool = await asyncpg.create_pool(self._dsn, timeout=10) - async with self._pool.acquire() as conn: - await conn.execute(_DDL) - await self._backfill_columns(conn) - self._schema_ready = True - return self._pool - except asyncio.CancelledError: - raise - except Exception as exc: - self._failed = True - logger.warning("Postgres 遥测初始化失败,后续记录降级为 no-op: {}", exc) + if self._failed: return None + if self._schema_ready: + return self._pool + pool = await self._open_pool() + if pool is None: + return None + return await self._prepare_schema(pool) + + async def _open_pool(self) -> asyncpg.Pool | None: + """建池;失败即永久降级(唯一一处「无条件判死」)。""" + if self._pool is not None: + return self._pool + try: + import asyncpg + + self._pool = await asyncpg.create_pool(self._dsn, timeout=10) + except asyncio.CancelledError: + raise + except Exception as exc: + # 池建不出来 = 确定写不进去;且每次调用重试都要内联吞掉 connect + # 超时,而遥测是业务路径上的 await —— 此处必须永久降级 + self._failed = True + logger.warning("Postgres 遥测建池失败,后续记录降级为 no-op: {}", exc) + return None + return self._pool + + async def _prepare_schema(self, pool: asyncpg.Pool) -> asyncpg.Pool | None: + """备好表并交回可用的池;瞬时失败只跳过本次,确定写不进去才判死。""" + try: + async with pool.acquire() as conn: + writable = await self._prepare_table(conn) + except asyncio.CancelledError: + raise + except Exception as exc: + # 池已在手,取连接/探测失败多为瞬时抖动: 不判死也不标就绪, + # 只跳过本次记录,下次调用重新准备 + logger.warning("Postgres 遥测建表探测失败(跳过本条,下次重试): {}", exc) + return None + if not writable: + self._failed = True + return None + self._schema_ready = True + return pool + + async def _prepare_table(self, conn: object) -> bool: + """备好 `llm_calls`;**表存在就绝不发 DDL**。返回 False 仅表示表确定不存在。 + + `CREATE TABLE IF NOT EXISTS` 不能无条件发: PostgreSQL 对 schema 的 + CREATE 权限检查**早于** `IF NOT EXISTS` 的存在性判断(PG 16.14 实测: + 只授 `SELECT, INSERT ON llm_calls` 的角色,表明明在、也写得进去,这一句 + 照样被拒 `permission denied for schema`)。这与 `_backfill_columns` 撞的 + 是同一类问题(issue #3/#9),故守卫也必须同款: 先探测,后 DDL。 + 探测走 `to_regclass`,不需要任何权限,且与 INSERT 的 search_path 解析 + 口径一致——比裸 DDL 更准(裸 `CREATE TABLE` 落在首个**可建**的 schema, + 可能与 INSERT 命中的不是同一张表)。 + """ + exists = await conn.fetchval(_TABLE_EXISTS) is not None # type: ignore[attr-defined] + if exists: + await self._backfill_columns(conn) # 旧表可能缺列;失败只逐行降级 + return True + try: + await conn.execute(_DDL) # type: ignore[attr-defined] + except asyncio.CancelledError: + raise + except Exception as exc: + logger.warning("Postgres 遥测建表失败(表不存在,记录无处可落): {}", exc) + return False + return True # 新建表列已齐全,无需再走补列 async def _backfill_columns(self, conn: object) -> None: """给已存在的旧表补新列(issue #3);**先探测再 ALTER,失败绝不置 `_failed`**。 diff --git a/src/polygateway/telemetry/sqlite.py b/src/polygateway/telemetry/sqlite.py index b8622c6..8483c68 100644 --- a/src/polygateway/telemetry/sqlite.py +++ b/src/polygateway/telemetry/sqlite.py @@ -3,6 +3,14 @@ 蓝本 VT `adapters/telemetry.py`: 构造期建连接与表,失败降级为 no-op (记录基础设施不得拖垮业务调用);`INSERT OR IGNORE` 幂等(call_id 主键); 写入经 threading.Lock 串行化后由 `asyncio.to_thread` 执行,不阻塞事件循环。 + +**这里不做 postgres.py 那样的建表前探测,是实测后的有意不对称**(issue #9): +SQLite 对已存在的表在**解析期**就把 `CREATE TABLE IF NOT EXISTS` 短路掉, +既不抢写锁也不检查可写性——实测同一时刻另一连接持 `BEGIN EXCLUSIVE`、或 +文件 `chmod 444`,该语句均通过,而同条件下的 `INSERT` 与新表名建表分别报 +database is locked / readonly database。故 PG 侧"权限检查早于存在性判断" +的坑在此不存在,加探测零收益。**别为了代码对称把它加回来**;需要对称的是 +保证(表存在就不该因建表失败而失能),这一条两侧都已满足。 """ from __future__ import annotations diff --git a/tests/integration/test_postgres_telemetry.py b/tests/integration/test_postgres_telemetry.py index 11382ad..d1af681 100644 --- a/tests/integration/test_postgres_telemetry.py +++ b/tests/integration/test_postgres_telemetry.py @@ -13,6 +13,7 @@ from __future__ import annotations import asyncio import json import os +import re from uuid import uuid4 import pytest @@ -299,3 +300,88 @@ class TestDegradation: await _record_minimal(recorder) await recorder.aclose() await recorder.aclose() + + +_PROBE_PASSWORD = "pgw_issue9_probe" # 临时角色,teardown 删除;非任何真实凭据 + + +@pytest.fixture +async def least_privilege_dsn(dsn): + """临时 schema + 临时角色: 只授表级 SELECT/INSERT,**不授 schema CREATE**。 + + 这是 issue #9 的现场——最小权限部署的标准形态。fixture 建的一切 + (schema、表、角色)都在 teardown 里删净,共享的 public.llm_calls 不受影响; + 连不上或无权建角色(非超级用户)时 skip,不让 CI 假绿。 + """ + import asyncpg + + from polygateway.telemetry.postgres import _DDL + + name = f"pgwtest_lp_{uuid4().hex[:8]}" + admin = await asyncpg.connect(dsn, timeout=10) + try: + if not await admin.fetchval( + "SELECT rolcreaterole OR rolsuper FROM pg_roles WHERE rolname = current_user" + ): + pytest.skip("当前账号无权建临时角色,跳过最小权限用例") + 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(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 —— 缺的正是这一项 + finally: + await admin.close() + low = re.sub(r"//[^@/]+@", f"//{name}:{_PROBE_PASSWORD}@", dsn, count=1) + sep = "&" if "?" in low else "?" + yield f"{low}{sep}options=-csearch_path%3D{name}", name + admin = await asyncpg.connect(dsn, timeout=10) + try: + await admin.execute(f"DROP SCHEMA IF EXISTS {name} CASCADE") + await admin.execute(f"DROP OWNED BY {name}") + await admin.execute(f"DROP ROLE IF EXISTS {name}") + finally: + await admin.close() + + +class TestLeastPrivilegeDeployment: + """issue #9: 只有表级写权限的账号,遥测必须照常落库而不是整体判死。""" + + async def test_create_table_if_not_exists_is_denied_for_this_role(self, least_privilege_dsn): + """库外事实先钉死: 表存在、写得进去,DDL 仍被拒——PG 的权限检查早于 IF NOT EXISTS。 + + 修复依赖的是这条 PG 语义;若某天它变了,这里先红,而不是让下面那条 + 用例悄悄变成"永远通过"的空断言。 + """ + import asyncpg + + low_dsn, _ = least_privilege_dsn + conn = await asyncpg.connect(low_dsn, timeout=10) + try: + assert await conn.fetchval("SELECT to_regclass('llm_calls')") is not None + with pytest.raises(asyncpg.exceptions.InsufficientPrivilegeError): + await conn.execute("CREATE TABLE IF NOT EXISTS llm_calls (call_id TEXT)") + finally: + await conn.close() + + 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) + try: + await _record_minimal(recorder, call_id=_cid("lp1")) + await _record_minimal(recorder, call_id=_cid("lp2"), cost=1.5) + assert recorder._failed is False # 判死开关不得被建表权限触发 + rows = await _fetch( + low_dsn, + "SELECT call_id, cost FROM llm_calls WHERE call_id LIKE $1 ORDER BY call_id", + f"{_RUN_PREFIX}-lp%", + ) + assert [(r["call_id"], r["cost"]) for r in rows] == [ + (_cid("lp1"), None), + (_cid("lp2"), 1.5), + ] + assert schema # teardown 会连表带角色删净 + finally: + await recorder.aclose() diff --git a/tests/unit/test_telemetry.py b/tests/unit/test_telemetry.py index 5a2df3c..c332777 100644 --- a/tests/unit/test_telemetry.py +++ b/tests/unit/test_telemetry.py @@ -257,17 +257,41 @@ class TestSQLiteColumnBackfill: class _FakePgConn: - """记录执行过的语句;可让 ALTER 抛错以模拟权限不足。""" + """记录执行过的语句;可让 ALTER/CREATE/探测抛错以模拟权限不足与抖动。 - def __init__(self, existing: list[str], *, fail_alter: bool = False): + `existing` 为空列表即表示**表不存在**(与真实 PG 一致: `to_regclass` 为 NULL + 时列探测必然零行),故 `fetchval` 与 `fetch` 共用同一份事实。 + """ + + def __init__( + self, + existing: list[str], + *, + fail_alter: bool = False, + fail_create: bool = False, + probe_errors: int = 0, + ): self.existing = existing self.fail_alter = fail_alter + self.fail_create = fail_create + self.probe_errors = probe_errors self.statements: list[str] = [] async def execute(self, sql, *args): self.statements.append(sql) if sql.startswith("ALTER TABLE") and self.fail_alter: raise RuntimeError("must be owner of table llm_calls") + if sql.lstrip().startswith("CREATE TABLE"): + if self.fail_create: + raise RuntimeError("permission denied for schema public") + self.existing = list(_EXPECTED_COLUMNS) + + async def fetchval(self, sql, *args): + self.statements.append(sql) + if self.probe_errors > 0: + self.probe_errors -= 1 + raise RuntimeError("connection was closed in the middle of operation") + return "llm_calls" if self.existing else None async def fetch(self, sql, *args): self.statements.append(sql) @@ -338,6 +362,76 @@ class TestPostgresBackfillDiscipline: assert all("IF NOT EXISTS" not in s for s in altered) # 探测已确认缺列,无需再判 +class TestPostgresTableProbe: + """建表必须先探测,且"判死"只认"确定写不进去"(issue #9)。 + + 实测(PostgreSQL 16.14,只有表级 SELECT/INSERT 的角色): `CREATE TABLE IF NOT + EXISTS` 被拒 permission denied for schema,而同一连接的 `INSERT` 通过—— + PG 对 schema 的 CREATE 权限检查早于 `IF NOT EXISTS` 的存在性判断。无条件发 + DDL 会让这类最小权限部署的整个进程静默失遥测。 + """ + + _CURRENT = [ + "call_id", + "cost", + "created_at", + "cached_prompt_tokens", + "model_reported", + "sampling", + "reasoning_tokens", + ] + + def _recorder(self, conn): + from polygateway.telemetry.postgres import PostgresRecorder + + return PostgresRecorder("postgresql://u:p@h:5432/polygateway", pool=_FakePgPool(conn)) + + def _created(self, conn): + return [s for s in conn.statements if s.lstrip().startswith("CREATE TABLE")] + + async def test_existing_table_is_never_recreated(self): + """表已存在就一条 DDL 都不发——这是权限被拒的唯一根治办法。""" + conn = _FakePgConn(self._CURRENT) + await _record_minimal(self._recorder(conn)) + assert not self._created(conn) + + async def test_create_denied_on_existing_table_keeps_recording(self): + """就算 DDL 仍被发出并被拒,表存在时也不得判死整个 recorder。""" + conn = _FakePgConn(self._CURRENT, fail_create=True) + recorder = self._recorder(conn) + await _record_minimal(recorder) # 不得抛 + assert recorder._failed is False + assert any(s.startswith("INSERT INTO llm_calls") for s in conn.statements) + + async def test_missing_table_is_created_and_not_backfilled(self): + """表不存在→建表;新建表列已齐全,不得再发补列 ALTER。""" + conn = _FakePgConn([]) + recorder = self._recorder(conn) + await _record_minimal(recorder) + assert len(self._created(conn)) == 1 + assert not [s for s in conn.statements if s.startswith("ALTER TABLE")] + assert recorder._failed is False + assert any(s.startswith("INSERT INTO llm_calls") for s in conn.statements) + + async def test_create_failure_on_missing_table_degrades_to_noop(self): + """表确定不存在且建不出来 = 确定写不进去: 此时才允许永久 no-op。""" + conn = _FakePgConn([], fail_create=True) + recorder = self._recorder(conn) + await _record_minimal(recorder) # 不得抛 + assert recorder._failed is True + assert not [s for s in conn.statements if s.startswith("INSERT INTO llm_calls")] + + async def test_probe_failure_is_transient_not_terminal(self): + """探测失败多为连接抖动: 跳过本次,下次调用必须重试,绝不永久判死。""" + conn = _FakePgConn(self._CURRENT, probe_errors=1) + recorder = self._recorder(conn) + await _record_minimal(recorder, call_id="first") # 不得抛 + assert recorder._failed is False + assert not [s for s in conn.statements if s.startswith("INSERT INTO llm_calls")] + await _record_minimal(recorder, call_id="second") + assert [s for s in conn.statements if s.startswith("INSERT INTO llm_calls")] + + class _MemoryRecorder: def __init__(self): self.rows = []