fix: judge telemetry failures by nature, not by step
The pool exhaustion in issue #15 was fatal only because min_size=10 forced a transient error to surface at pool creation, and that step was hardcoded to permanent death. Step is the wrong axis: it conflates "the DSN cannot be parsed" with "someone else holds all the connections right now". Failures are now classified by two rules. Fatal means the cause lies entirely inside this process and cannot change, which only the construction-time DSN satisfies. Everything else splits on whether the failure has anything to do with this row's data: row-level failures drop one row and keep trying, environment-level failures cool down for 60s and then get exactly one retry, so a restarted database or a DBA creating the table heals on its own. 42703 (missing column) is the single named exception and stays row-level even though every row fails alike: issue #13 promised that the manual mode trims the INSERT and exposes drift per row, and that promise outranks the rule. Any future exception owes the same argument. The _failed boolean is gone; the tracker is the only degradation state, because two copies of the same fact drift apart. Closing stays outside that state: it is the caller's own decision, not an anomaly to recover from, so the snapshot reports it through dropped_rows and the drop reason instead of raising the degraded flag on every clean shutdown.
This commit is contained in:
@@ -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 会补一条新的)。
|
||||
|
||||
@@ -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]
|
||||
|
||||
+214
-14
@@ -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:
|
||||
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
|
||||
|
||||
Reference in New Issue
Block a user