diff --git a/src/polygateway/telemetry/postgres.py b/src/polygateway/telemetry/postgres.py index 37947e7..8a32589 100644 --- a/src/polygateway/telemetry/postgres.py +++ b/src/polygateway/telemetry/postgres.py @@ -1,17 +1,21 @@ -"""Postgres 遥测后端(M2 设计 §5): asyncpg lazy 池 + 两级降级。 +"""Postgres 遥测后端(M2 设计 §5): asyncpg lazy 池 + 按失败性质三分的降级。 参考仓无先例(三项目遥测全 SQLite);asyncpg 工程写法取 GovDoc `taskrun/postgres_store.py`($n 占位、`ON CONFLICT DO NOTHING`),但其 -"失败冒泡"方向按遥测铁律**有意反转**: -① 结构性失败 → warning 一次后永久降级(所有写入短路); -② 运行时单条写失败 → 逐条 warning 丢弃,不降级不重试(连接抖动由 -asyncpg 池自恢复;避免浸泡开头一次抖动导致后续全程失遥测)。 +"失败冒泡"方向按遥测铁律**有意反转**: 遥测失败一律不冒泡,只降级。 构造不连库(lazy),24 列 schema 与 SQLite 版同名同序。 -**"结构性"的判据是「确定写不进去」,不是「初始化时出过错」**(issue #9): -只有建池失败(重试要在业务路径上内联吞掉 connect 超时)与"表确定不存在 -且建不出来"(后续 INSERT 必然全败)才判死;探测失败、补列失败、取连接 -失败一律只 warning,让写入照常尝试或下次调用重试。 +**降级档位挂在"失败是什么性质",不挂"哪一步失败"**(issue #15)。挂步骤是 +issue 的病灶: `min_size=10` 把"连接耗尽"这种瞬时错误逼到建池那一步,于是它被 +一刀切成了永久判死,整进程从此一条遥测都不落,只有重启能恢复。判据两句: + +1. **致命 = 失败原因完全在进程内部且不可变**。DSN 是构造期定死的字符串,是唯一 + 满足这条的东西;认证失败、库不存在、表建不出来一律不算——DBA 改完就该好。 +2. **行级 vs 环境级看失败与"这一行的数据"有没有关系**: 只与本行数据有关(换一行 + 可能成功)= 行级,逐条丢弃;与数据无关、每一行都会同样失败 = 环境级,进冷却。 + +见 `_classify_failure`(全库唯一一处 PG 失败分类)与 `_handle_failure`(三个降级点 +唯一一处处置)。 """ from __future__ import annotations @@ -60,6 +64,71 @@ _EXISTING_COLUMNS = ( "WHERE attrelid = to_regclass('llm_calls') AND attnum > 0 AND NOT attisdropped" ) +# 环境级降级的冷却期(issue #15)。**不暴露配置**: 它的取值只影响"多久重试一次" +# 这个内部节奏,任何取值都不改变对外承诺(降级可见、可自愈、有界成本),给出旋钮 +# 只会多一个下游要理解却调不对的东西(设计 §3.5) +_DEGRADE_COOLDOWN_S = 60.0 + +_FATAL = "fatal" +"""配置级致命: 原因完全在进程内部且不可变 → 永久 no-op + 一条 error。""" + +_UNAVAILABLE = "unavailable" +"""环境级不可用: 与本行数据无关、每行都会同样失败 → 冷却降级,到期重试一次。""" + +_ROW = "row" +"""行级拒绝: 只与本行数据有关 → 逐条 warning 丢弃,不降级。""" + +# 环境级的 SQLSTATE 类(前两位): 08 连接、53 资源不足(含 53300 too many +# connections)、57 管理干预、28 认证、3D 库不存在。共同点是"与这一行的数据无关, +# 换一行照样失败",且都能被外部修好 +_UNAVAILABLE_SQLSTATE_CLASSES = frozenset({"08", "53", "57", "28", "3D"}) + +# 类 42 整体归行级(见 `_classify_failure` 的默认档),但这两个码与本行数据无关: +# 42501 = 账号被收走 INSERT 权限,42P01 = 表被迁走/删掉。它们是持续性的环境状态, +# 按类归行级会让每次 LLM 调用都内联付一次往返、刷一条 warning,且永不自愈 +_UNAVAILABLE_SQLSTATES = frozenset({"42501", "42P01"}) + +# **判据的唯一具名例外**(issue #13 的更高优先级承诺): 42703 = 缺列。按判据第 2 句 +# 它本该是环境级(缺列时每一行都失败),归行级是因为 manual 档会按现有列裁剪 INSERT +# 继续写——"部分列写进去了 + 缺列逐行 warning 暴露"本身有价值,是下游发现 schema +# 漂移的唯一信号,不该被冷却掉。**新增例外必须同款论证**: 说清它为什么值得违反判据 +_ROW_SQLSTATES = frozenset({"42703"}) + + +def _classify_failure(exc: BaseException) -> str: + """按**失败的性质**分档(全库唯一一处 PG 失败分类);判据见模块 docstring。 + + 分类只认 SQLSTATE 与异常类型,不认"在哪一步失败"——后者正是 issue #15 的病灶。 + SQLSTATE 而非 asyncpg 异常类白名单: 前者是 PG 标准,不随驱动版本漂移。 + + **认不出来的失败一律给最轻的一档**(`_ROW`): 升档(冷却 60s)要有依据,没依据就 + 宁可每次调用多付一次内联往返,也不拿 60 秒的遥测去赌一个猜测。issue #9 定下的 + "探测抖动只跳过本次、下次重试"正是靠这条默认保住的。 + """ + if isinstance(exc, ValueError | TypeError): + # DSN 不可解析(实测: 端口写成非数字 → 裸 ValueError;scheme 不对 → + # ClientConfigurationError,它本身就是 ValueError 子类)与建池参数非法。 + # 这些是构造期就定死的进程内部事实,重试在任何时刻都不可能成功 + return _FATAL + sqlstate = getattr(exc, "sqlstate", None) + if isinstance(sqlstate, str): + if sqlstate in _ROW_SQLSTATES: + return _ROW + if sqlstate[:2] in _UNAVAILABLE_SQLSTATE_CLASSES or sqlstate in _UNAVAILABLE_SQLSTATES: + return _UNAVAILABLE + # 其余 PostgresError(22 数据异常、23 约束冲突等)都是这一行的数据问题 + return _ROW + # 没有 SQLSTATE = 话还没说到 PG 就断了: OSError(含 ConnectionError 与 + # TimeoutError)与 asyncpg 自己的 InterfaceError,都与本行数据无关 + return _UNAVAILABLE if isinstance(exc, OSError | _interface_error()) else _ROW + + +def _interface_error() -> type[BaseException]: + """asyncpg 的 `InterfaceError` 类型;延迟取用以免模块导入期硬依赖 extra。""" + import asyncpg + + return asyncpg.InterfaceError + class PostgresRecorder: """TelemetryRecorder 端口的 Postgres 实现;asyncpg 原生异步,无线程桥接。""" @@ -105,10 +174,9 @@ class PostgresRecorder: self._columns: tuple[str, ...] = COLUMNS self._insert = insert_sql("postgres", COLUMNS) self._schema_ready = False - self._failed = False # 结构性降级标志: 置位后所有写入短路 self._closed = False # 关了就是关了: 置位后写入短路且**不重建池** - # 降级的可编程出口与节流日志;`_failed` 与它并存是 issue #15 的过渡态, - # 判据改造(冷却自愈)落地时状态收归 tracker 一处 + # 降级状态**只此一份**: 是否短路写入、多久重试一次、下游查到什么, + # 全由 tracker 回答。两份状态(曾经的 `_failed` 布尔 + tracker)必然漂移 self._status = TelemetryStatusTracker(backend="postgres", now=now) self._init_lock = asyncio.Lock() @@ -118,18 +186,22 @@ class PostgresRecorder: return self._status.snapshot() async def _ensure_ready(self) -> asyncpg.Pool | None: - """lazy 建池+备表;判死只认「确定写不进去」(issue #9),其余失败都留活路。 + """lazy 建池+备表;降级期间**零成本短路**,冷却到期放行一次重新准备。 + + `should_retry()` 是纯时间比较,不触库: 降级期间的调用因此既不内联吞 + connect 超时(`postgres.py` 老注释担心的正是这个),也不需要重启进程—— + 成本变成"每 60s 一次、上界一个写入预算",有界且可解释。 `_closed` 在锁内**必须复查**: 等锁期间发生的 `aclose` 否则会被这次 等待"绕过",等到锁时照旧建出一个没人负责关的池(注入档更隐蔽—— 注入方以为自己管着全部连接,实际早已不是)。 """ - if self._closed or self._failed: + if self._closed or not self._status.should_retry(): return None if self._schema_ready: return self._pool async with self._init_lock: - if self._closed or self._failed: + if self._closed or not self._status.should_retry(): return None if self._schema_ready: return self._pool @@ -139,7 +211,7 @@ class PostgresRecorder: return await self._prepare_schema(pool) async def _open_pool(self) -> asyncpg.Pool | None: - """建池;失败即永久降级(唯一一处「无条件判死」)。 + """建池;失败按性质分档处置(见 `_handle_failure`),不再一律判死。 **池的资源占用由本库显式声明**(issue #15): `min_size=0` 的语义是"不预 连接"(asyncpg `pool.py:457` 为 0 时只造 holder 对象,一条连接都不连), @@ -163,10 +235,7 @@ class PostgresRecorder: except asyncio.CancelledError: raise except Exception as exc: - # 池建不出来 = 确定写不进去;且每次调用重试都要内联吞掉 connect - # 超时,而遥测是业务路径上的 await —— 此处必须永久降级 - self._failed = True - self._status.enter_degraded(f"建池失败: {exc}", fatal=True, cooldown_s=None) + self._handle_failure(exc, stage="建池") return None return self._pool @@ -185,14 +254,15 @@ class PostgresRecorder: except asyncio.CancelledError: raise except Exception as exc: - # 池已在手,取连接/探测失败多为瞬时抖动: 不判死也不标就绪, - # 只跳过本次记录,下次调用重新准备 - logger.warning("Postgres 遥测建表探测失败(跳过本条,下次重试): {}", exc) + self._handle_failure(exc, stage="建表探测") return None if columns is None: - self._failed = True + # 表确定不存在且建不出来: 与本行数据无关(每行都会同样失败)且能被 + # 外部修好(DBA 建了表就该自愈)—— 判据第 2 句下的环境级 self._status.enter_degraded( - "表 llm_calls 不存在且建不出来(记录无处可落)", fatal=True, cooldown_s=None + "表 llm_calls 不存在且建不出来(记录无处可落)", + fatal=False, + cooldown_s=_DEGRADE_COOLDOWN_S, ) return None # 写入列、语句与就绪标志必须**一起**生效: `_ensure_ready` 只看 `_schema_ready` @@ -205,7 +275,7 @@ class PostgresRecorder: async def _prepare_table(self, conn: object) -> tuple[str, ...] | None: """备好 `llm_calls` 并返回本实例要写的列;**表存在就绝不发 DDL**。 - 返回 None 仅表示表确定不存在且建不出来(唯一允许判死的情形)。 + 返回 None 仅表示表确定不存在且建不出来(调用方据此进环境级冷却降级)。 `CREATE TABLE IF NOT EXISTS` 不能无条件发: PostgreSQL 对 schema 的 CREATE 权限检查**早于** `IF NOT EXISTS` 的存在性判断(PG 16.14 实测: @@ -278,11 +348,11 @@ class PostgresRecorder: return effective async def _backfill_columns(self, conn: object, existing: set[str]) -> None: - """auto 档: 给已存在的旧表补新列(issue #3);**失败绝不置 `_failed`**。 + """auto 档: 给已存在的旧表补新列(issue #3);**失败绝不让 recorder 降级**。 - 不置 `_failed` 的实测理由: 应用账号只有 INSERT 权限时,`ALTER TABLE` 的 - ownership 检查早于 `IF NOT EXISTS` 的存在性判断——列明明齐全也会失败。置位会让 - 整个 recorder 永久 no-op,与「补列失败只降级为逐行丢弃」的承诺相悖 + 不降级的实测理由: 应用账号只有 INSERT 权限时,`ALTER TABLE` 的 + ownership 检查早于 `IF NOT EXISTS` 的存在性判断——列明明齐全也会失败。降级会让 + 整个 recorder 停写(环境级还要停满一个冷却期),与「补列失败只降级为逐行丢弃」的承诺相悖 (SQLite 侧同款守卫,两侧必须对称)。补列失败后写入沿用全量列(今天的行为): auto 档承诺的是"把列补上",补不上就让缺列以逐行 warning 暴露;要降级写入 请显式选 manual。 @@ -296,6 +366,49 @@ class PostgresRecorder: except Exception as exc: logger.warning("Postgres 遥测补列失败(写入将逐行降级): {}", exc) + def _handle_failure(self, exc: BaseException, *, stage: str) -> None: + """按分档处置一次遥测失败;三个降级点(建池/建表探测/写入)共用这一处。 + + 收敛成一处不只是去重: 三处各写一遍处置,就是三处各自漂移一次判据的机会, + 而判据漂移正是 issue #15 的病灶(注释写着"确定写不进去",代码做的是别的事)。 + + `stage` 只进日志文案,**不参与分档**——挂步骤分档正是要被拆掉的错法。 + """ + verdict = _classify_failure(exc) + if verdict == _FATAL: + # error 而非 warning: 这是人配错了,且本进程内不会自愈,运维要看见 + logger.error( + "Postgres 遥测{}失败: 配置有误,本进程内不会自愈(请修正 DSN 后重启): {}", + stage, + exc, + ) + self._status.enter_degraded( + f"{stage}失败(配置有误): {exc}", fatal=True, cooldown_s=None + ) + elif verdict == _UNAVAILABLE: + # 每一行都会同样失败 → 冷却期内不再内联重试;`_schema_ready` 一并作废, + # 到期那次要重新走准备(表被删/权限被收回都得靠重新准备才能发现已修好) + self._schema_ready = False + self._status.enter_degraded( + f"{stage}失败: {exc}", fatal=False, cooldown_s=_DEGRADE_COOLDOWN_S + ) + else: + # 行级不进降级: 换一行可能就成了。逐条出声是 issue #13 的承诺 + # (缺列靠这条 warning 暴露 schema 漂移),不因刷屏而节流掉 + logger.warning("Postgres 遥测{}失败(丢弃该行,下次调用照常重试): {}", stage, exc) + + def _drop_reason(self) -> str: + """写不进去时说清是**哪一种**写不进去: 关了 / 降级中 / 本次没准备好。 + + 三者的处置完全不同(一个是调用方自己关了却还在写、一个等自愈、一个下次 + 就会重试),混成一句话会让对账的人分不清该等还是该修。 + """ + if self._closed: + return "遥测已关闭" + if self._status.snapshot().degraded: + return "遥测降级中" + return "后端本次未准备好(下次调用重试)" + async def record_llm_call(self, **fields: object) -> None: """写一行遥测;整次写入受硬预算约束,失败逐条丢弃(两级降级之二),绝不冒泡。 @@ -319,8 +432,10 @@ class PostgresRecorder: ) self._status.record_drop("写入超预算") except Exception as exc: - # 遥测铁律: 丢一条 < 拖垮调用;仅记 warning(非 pass),池自恢复 - logger.warning("Postgres 遥测写入失败(丢弃该行): {}", exc) + # 遥测铁律: 丢一条 < 拖垮调用。这一行无论如何都没了,区别只在于 + # **下一行还试不试**——那由失败的性质决定,不由这里决定 + self._handle_failure(exc, stage="写入") + self._status.record_drop("写入失败") async def _write_row(self, fields: dict[str, object]) -> None: """预算内的写入本体: 准备 → 取连接 → 执行 → 归还。 @@ -330,10 +445,8 @@ class PostgresRecorder: """ pool = await self._ensure_ready() if pool is None: - # 降级期间静默 return 就是 issue #15 的破口: 丢行必须计数且节流出声。 - # 关闭后的丢行同样要计数,但原因不是降级——两者的处置完全不同 - # (一个等自愈,一个是调用方自己关了却还在写) - self._status.record_drop("遥测已关闭" if self._closed else "遥测已降级") + # 降级期间静默 return 就是 issue #15 的破口: 丢行必须计数且节流出声 + self._status.record_drop(self._drop_reason()) return row = tuple(fields[col] for col in self._columns) conn = await pool.acquire(timeout=self._write_timeout_s) @@ -341,6 +454,9 @@ class PostgresRecorder: await conn.execute(self._insert, *row) finally: await self._release(pool, conn) + # **恢复的唯一权威证据是一次真正写成功**(未降级时是廉价 no-op)。放在这里 + # 而不是准备期: 准备通过不代表写得进去(权限只到 SELECT 时正是如此) + self._status.recover() async def _release(self, pool: asyncpg.Pool, conn: object) -> None: """归还连接;归还路径独立有界,失败即断开(下次 acquire 会补一条新的)。 diff --git a/tests/integration/test_postgres_telemetry.py b/tests/integration/test_postgres_telemetry.py index 49c8551..68b6d29 100644 --- a/tests/integration/test_postgres_telemetry.py +++ b/tests/integration/test_postgres_telemetry.py @@ -413,13 +413,13 @@ class TestLeastPrivilegeDeployment: await conn.close() async def test_records_land_without_schema_create_privilege(self, least_privilege_dsn): - """修复前: 建表被拒 → _failed → 整个进程一条不落(下游 150 次调用全丢)。""" + """修复前: 建表被拒 → 整体判死 → 整个进程一条不落(下游 150 次调用全丢)。""" low_dsn, schema = least_privilege_dsn recorder = _recorder(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) - assert recorder._failed is False # 判死开关不得被建表权限触发 + assert recorder.telemetry_status.degraded is False # 建表权限不得触发降级 rows = await _fetch( low_dsn, "SELECT call_id, cost FROM llm_calls WHERE call_id LIKE $1 ORDER BY call_id", @@ -665,15 +665,15 @@ class TestCallerDimensionsAcceptance: async def test_backfill_failure_degrades_per_row_not_wholesale( self, least_privilege_pre_tenant_dsn, captured_warnings ): - """补列失败的降级方向: 记 warning、不置 `_failed`、后续 INSERT 仍照发。 + """补列失败的降级方向: 记 warning、不整体降级、后续 INSERT 仍照发。 - 置 `_failed` 会让整个进程从此一条遥测都不写(比逐行丢弃严重得多), - 且一旦 DBA 补上列也不会自愈——必须等重启。 + 整体降级会让整个进程停写(比逐行丢弃严重得多),而缺列(SQLSTATE 42703) + 是判据的唯一具名例外: 必须逐行暴露,好让下游看见 schema 漂移(issue #13)。 """ recorder = _recorder(least_privilege_pre_tenant_dsn, auto_migrate=True) try: await _record_minimal(recorder, call_id=_cid("lpp1")) # 不得抛 - assert recorder._failed is False + assert recorder.telemetry_status.degraded is False assert any("补列失败" in m for m in captured_warnings) # 缺列的表上 INSERT 必然失败;逐行 warning 正是"INSERT 照发了"的证据 assert any("写入失败" in m for m in captured_warnings) @@ -867,7 +867,7 @@ class TestManualSchemaModeAcceptance: 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 + assert recorder.telemetry_status.degraded 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] diff --git a/tests/unit/test_telemetry.py b/tests/unit/test_telemetry.py index 84a722e..d6256c3 100644 --- a/tests/unit/test_telemetry.py +++ b/tests/unit/test_telemetry.py @@ -729,6 +729,7 @@ class _FakePgConn: probe_errors: int = 0, hang_insert: bool = False, fail_terminate: bool = False, + insert_error: BaseException | None = None, ): self.existing = existing self.fail_alter = fail_alter @@ -737,6 +738,8 @@ class _FakePgConn: # 只挂 INSERT: 准备期照常完成,挂住的才是业务路径上那次内联 await self.hang_insert = hang_insert self.fail_terminate = fail_terminate + # INSERT 阶段抛出的真实 PG 异常(带 SQLSTATE),用来钉失败三分的边界 + self.insert_error = insert_error self.terminated = False self.statements: list[str] = [] @@ -747,8 +750,11 @@ class _FakePgConn: async def execute(self, sql, *args): self.statements.append(sql) - if sql.startswith("INSERT INTO") and self.hang_insert: - await asyncio.sleep(3600) + if sql.startswith("INSERT INTO"): + if self.hang_insert: + await asyncio.sleep(3600) + if self.insert_error is not None: + raise self.insert_error if sql.startswith("ALTER TABLE") and self.fail_alter: raise RuntimeError("must be owner of table llm_calls") if sql.lstrip().startswith("CREATE TABLE"): @@ -784,8 +790,12 @@ class _FakePgPool: hang_acquire: bool = False, fail_release: bool = False, hang_close: bool = False, + acquire_error: BaseException | None = None, ): self._conn = conn + # 取连接阶段抛出的真实异常: 连接耗尽/DSN 非法都在这一步现形 + # (`min_size=0` 之后建池不触库,实测 create_pool 连 DSN 都不解析) + self.acquire_error = acquire_error self.hang_acquire = hang_acquire self.fail_release = fail_release # 模拟 asyncpg 的 `Pool.close()` 在 in-flight 连接未归还时**无限等** @@ -802,6 +812,8 @@ class _FakePgPool: self.acquire_timeouts.append(timeout) if self.hang_acquire: await asyncio.sleep(3600) + if self.acquire_error is not None: + raise self.acquire_error self.acquired += 1 return self._conn @@ -849,11 +861,11 @@ class TestPostgresBackfillDiscipline: ) async def test_alter_failure_does_not_disable_the_recorder(self): - """ALTER 失败(如账号只有 INSERT 权限)不得置 _failed —— 那会让遥测全灭。""" + """ALTER 失败(如账号只有 INSERT 权限)不得让 recorder 降级 —— 那会让遥测全灭。""" conn = _FakePgConn(self._LEGACY, fail_alter=True) recorder = self._recorder(conn) await _record_minimal(recorder) # 不得抛 - assert recorder._failed is False + assert recorder.telemetry_status.degraded is False assert any(s.startswith("INSERT INTO llm_calls") for s in conn.statements) async def test_no_alter_when_columns_already_exist(self): @@ -922,7 +934,7 @@ class TestPostgresTableProbe: conn = _FakePgConn(self._CURRENT, fail_create=True) recorder = self._recorder(conn) await _record_minimal(recorder) # 不得抛 - assert recorder._failed is False + assert recorder.telemetry_status.degraded 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): @@ -932,15 +944,16 @@ class TestPostgresTableProbe: 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 recorder.telemetry_status.degraded 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。""" + async def test_create_failure_on_missing_table_enters_cooldown(self): + """表确定不存在且建不出来 = 环境级(DBA 建了表就该好): 冷却降级,不判死。""" conn = _FakePgConn([], fail_create=True) recorder = self._recorder(conn) await _record_minimal(recorder) # 不得抛 - assert recorder._failed is True + status = recorder.telemetry_status + assert status.degraded is True and status.fatal is False 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): @@ -948,7 +961,8 @@ class TestPostgresTableProbe: conn = _FakePgConn(self._CURRENT, probe_errors=1) recorder = self._recorder(conn) await _record_minimal(recorder, call_id="first") # 不得抛 - assert recorder._failed is False + # 不认识的失败不给升级: 探测抖动只丢本行,绝不进 60s 冷却(issue #9) + assert recorder.telemetry_status.degraded 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")] @@ -1914,7 +1928,7 @@ class TestSQLiteStatusVisibility: class TestPostgresStatusVisibility: - """PG 侧的判死本任务不改判据,只让它经 tracker 变得可见(计划 T2)。""" + """PG 侧的降级必须能被下游查到(计划 T2 建立可见性,T5 改判据)。""" def _recorder(self, conn): from polygateway.telemetry.postgres import PostgresRecorder @@ -1928,13 +1942,15 @@ class TestPostgresStatusVisibility: ) async def test_unusable_table_shows_up_in_the_status(self, captured_warnings): - """表确定不存在且建不出来 = 既有的判死档;现在它要能被下游查到。""" + """表确定不存在且建不出来: 降级可见,且是**可自愈**的环境级而非永久判死。""" + from polygateway.telemetry.postgres import _DEGRADE_COOLDOWN_S + recorder = self._recorder(_FakePgConn([], fail_create=True)) await _record_minimal(recorder) status = recorder.telemetry_status - assert status.degraded is True and status.fatal is True - assert status.dropped_rows == 1 # 判死那一次调用本身也丢了一行 - assert recorder._failed is True # 过渡期两份状态并存(T5 收掉 `_failed`) + assert status.degraded is True and status.fatal is False # 环境级: 建了表就该自愈 + assert status.dropped_rows == 1 # 降级那一次调用本身也丢了一行 + assert status.retry_after_s == pytest.approx(_DEGRADE_COOLDOWN_S) async def test_healthy_recorder_is_not_degraded(self): recorder = self._recorder(_FakePgConn(list(_EXPECTED_COLUMNS))) @@ -2115,6 +2131,29 @@ class TestPostgresCloseIsBounded: assert pool.acquired == 2 # 准备期 + 首次写入;关闭后一次都没有 assert recorder.telemetry_status.dropped_rows == dropped_before + 1 + async def test_close_is_not_degradation(self, monkeypatch, captured_warnings): + """**关闭 ≠ 降级**(T4 留下的语义问题,T5 收口)。 + + `degraded` 的含义是"后端本该可写却写不进去,库正在设法恢复"。关闭是调用方 + 自己的决定,没有异常、也按设计不会自愈——把它记成降级,等于让每一次正常 + 收尾都发一次降级信号,下游"degraded 就告警"的规则会被每次退出打穿。 + 关闭后真正要对账的是"还有多少行没落地",那由 `dropped_rows` 与逐条原因 + 承担,不必污染 `degraded`。 + """ + pool = _FakePgPool(_FakePgConn(list(_EXPECTED_COLUMNS))) + recorder, _ = self._self_built(monkeypatch, pool) + await _record_minimal(recorder) + await recorder.aclose() + captured_warnings.clear() + + await _record_minimal(recorder, call_id="c2") + + status = recorder.telemetry_status + assert status.degraded is False and status.fatal is False + assert status.reason is None and status.retry_after_s is None + assert status.dropped_rows == 1 # 丢了多少行照样可对账 + assert any("遥测已关闭" in m for m in captured_warnings) # 且分得清是哪一种丢 + async def test_aclose_is_idempotent(self, monkeypatch): pool = _FakePgPool(_FakePgConn(list(_EXPECTED_COLUMNS))) recorder, _ = self._self_built(monkeypatch, pool) @@ -2180,3 +2219,164 @@ class TestPostgresReleaseDegradation: conn = _FakePgConn(list(_EXPECTED_COLUMNS), fail_terminate=True) await _record_minimal(_pg_recorder(pool=_FakePgPool(conn, fail_release=True))) assert any("断开失败" in m for m in captured_warnings) + + +class TestPostgresFailureClassification: + """issue #15 B 组: 判死判据从"哪一步失败"改为"失败是什么性质"(设计 §3.2)。 + + 判据两句: ①**致命 = 失败原因完全在进程内部且不可变**;②**行级 vs 环境级看 + 失败与这一行的数据有没有关系**。三档边界两侧各钉一次——按 SQLSTATE 前两位 + 一刀切正是本 issue 之前的错法,回归会当场红。 + """ + + def _self_built(self, monkeypatch, *, outcomes, clock): + """走**自建池**那条路;`outcomes` 逐次消费,元素是异常就抛出。 + + 必须自建而非注入: 注入档走 `_external_pool=True` 分支,完全绕过建池, + 而 issue 现场的失败恰恰发生在建池那一步。 + """ + import asyncpg + + created: list[str] = [] + + async def fake_create_pool(dsn, **kwargs): + created.append(dsn) + outcome = outcomes[min(len(created) - 1, len(outcomes) - 1)] + if isinstance(outcome, BaseException): + raise outcome + return outcome + + monkeypatch.setattr(asyncpg, "create_pool", fake_create_pool) + return _pg_recorder(now=clock), created + + async def test_pool_exhaustion_degrades_with_cooldown_and_self_heals(self, monkeypatch): + """**issue 场景直接回归**: 53300 落在准备期,过去 = 整进程永久失遥测。 + + `too many clients` 是外部状态,别人还连接就该好——它永远不满足"原因完全 + 在进程内部且不可变",故绝不许判死,只许冷却重试。 + """ + import asyncpg + + from polygateway.telemetry.postgres import _DEGRADE_COOLDOWN_S + + clock = _FakeClock() + conn = _FakePgConn(list(_EXPECTED_COLUMNS)) + recorder, created = self._self_built( + monkeypatch, + outcomes=[ + asyncpg.exceptions.TooManyConnectionsError("sorry, too many clients already"), + _FakePgPool(conn), + ], + clock=clock, + ) + + await _record_minimal(recorder, call_id="c1") # 不得抛 + status = recorder.telemetry_status + assert status.degraded is True + assert status.fatal is False # ← 现状在这里判死,整进程从此一条不落 + assert status.retry_after_s == pytest.approx(_DEGRADE_COOLDOWN_S) + assert not conn.statements + + await _record_minimal(recorder, call_id="c2") # 冷却期内零成本短路 + assert len(created) == 1 # 不再内联吞一次 connect 超时 + + clock.advance(_DEGRADE_COOLDOWN_S) + await _record_minimal(recorder, call_id="c3") + + assert len(created) == 2 # 到期放行一次重新准备 + assert [s for s in conn.statements if s.startswith("INSERT INTO llm_calls")] + assert recorder.telemetry_status.degraded is False # 自愈,无需重启进程 + assert recorder.telemetry_status.dropped_rows == 2 # 降级期间那两行确实丢了 + + async def test_unparseable_dsn_is_fatal_and_costs_nothing_afterwards(self, captured_warnings): + """DSN 是构造期定死的字符串: 唯一"进程内不可能变好"的东西,故唯一的致命档。""" + import asyncpg + + clock = _FakeClock() + pool = _FakePgPool( + _FakePgConn(list(_EXPECTED_COLUMNS)), + acquire_error=asyncpg.exceptions.ClientConfigurationError("invalid DSN: bad scheme"), + ) + recorder = _pg_recorder(pool=pool, now=clock) + + await _record_minimal(recorder, call_id="c1") # 不得抛 + status = recorder.telemetry_status + assert status.degraded is True and status.fatal is True + assert status.retry_after_s is None # 本进程内不会自愈 + assert any("重启" in m for m in captured_warnings) # 恢复条件必须写在日志里 + + clock.advance(1_000_000.0) + attempts = len(pool.acquire_timeouts) + await _record_minimal(recorder, call_id="c2") + assert len(pool.acquire_timeouts) == attempts # 此后零成本短路,不再触库 + + @pytest.mark.parametrize("error_name", ["InsufficientPrivilegeError", "UndefinedTableError"]) + async def test_environment_level_sqlstates_enter_cooldown(self, error_name): + """`42501`/`42P01` 与这一行的数据无关(每一行都会同样失败)→ 环境级。 + + 它们与 `42703` 同属 SQLSTATE `42` 类却分属两档: 判据看的是"失败与这一行的 + 数据有没有关系",不是前两位。 + """ + import asyncpg + + from polygateway.telemetry.postgres import _DEGRADE_COOLDOWN_S + + clock = _FakeClock() + conn = _FakePgConn( + list(_EXPECTED_COLUMNS), insert_error=getattr(asyncpg.exceptions, error_name)("boom") + ) + recorder = _pg_recorder(pool=_FakePgPool(conn), now=clock) + + await _record_minimal(recorder, call_id="c1") + status = recorder.telemetry_status + assert status.degraded is True and status.fatal is False + assert status.retry_after_s == pytest.approx(_DEGRADE_COOLDOWN_S) + + inserts = len([s for s in conn.statements if s.startswith("INSERT INTO llm_calls")]) + await _record_minimal(recorder, call_id="c2") + # 冷却期内不再每行内联付一次往返;权限/建表修好后由冷却到期自动恢复 + assert len([s for s in conn.statements if s.startswith("INSERT INTO llm_calls")]) == inserts + + async def test_missing_column_stays_row_level(self, captured_warnings): + """`42703` 是判据的**唯一具名例外**,由 issue #13 定死: 缺列要逐行暴露。 + + 按判据第 2 句它本该是环境级(缺列时每行都失败),归行级是因为 manual 档 + 会裁剪 INSERT 继续写,"部分列写进去了 + 逐行 warning"本身有价值, + 不该被冷却掉——下游正是靠这条 warning 发现 schema 漂移的。 + """ + import asyncpg + + conn = _FakePgConn( + list(_EXPECTED_COLUMNS), + insert_error=asyncpg.exceptions.UndefinedColumnError('column "meta" does not exist'), + ) + recorder = _pg_recorder(pool=_FakePgPool(conn)) + + await _record_minimal(recorder, call_id="c1") + await _record_minimal(recorder, call_id="c2") + + status = recorder.telemetry_status + assert status.degraded is False # 不进冷却 + assert status.dropped_rows == 2 + # 每一行都照发 INSERT,每一行都出声: schema 漂移必须持续可见 + assert len([s for s in conn.statements if s.startswith("INSERT INTO llm_calls")]) == 2 + # 逐行那条不节流(与累计计数那条区分开): 每丢一行都要出声 + assert len([m for m in captured_warnings if "写入失败(丢弃该行" in m]) == 2 + + async def test_uncreatable_table_recovers_once_the_dba_creates_it(self): + """表建不出来是环境级: DBA 建完表,冷却到期就该自己好,不必重启进程。""" + from polygateway.telemetry.postgres import _DEGRADE_COOLDOWN_S + + clock = _FakeClock() + conn = _FakePgConn([], fail_create=True) + recorder = _pg_recorder(pool=_FakePgPool(conn), now=clock) + + await _record_minimal(recorder, call_id="c1") + assert recorder.telemetry_status.degraded is True + + conn.existing = list(_EXPECTED_COLUMNS) # DBA 手工建了表 + clock.advance(_DEGRADE_COOLDOWN_S) + await _record_minimal(recorder, call_id="c2") + + assert [s for s in conn.statements if s.startswith("INSERT INTO llm_calls")] + assert recorder.telemetry_status.degraded is False