"""Postgres 遥测后端(M2 设计 §5): asyncpg lazy 池 + 按失败性质三分的降级。 参考仓无先例(三项目遥测全 SQLite);asyncpg 工程写法取 GovDoc `taskrun/postgres_store.py`($n 占位、`ON CONFLICT DO NOTHING`),但其 "失败冒泡"方向按遥测铁律**有意反转**: 遥测失败一律不冒泡,只降级。 构造不连库(lazy),24 列 schema 与 SQLite 版同名同序。 **降级档位挂在"失败是什么性质",不挂"哪一步失败"**(issue #15)。挂步骤是 issue 的病灶: `min_size=10` 把"连接耗尽"这种瞬时错误逼到建池那一步,于是它被 一刀切成了永久判死,整进程从此一条遥测都不落,只有重启能恢复。判据两句: 1. **致命 = 失败原因完全在进程内部且不可变**。DSN 是构造期定死的字符串,是唯一 满足这条的东西;认证失败、库不存在、表建不出来一律不算——DBA 改完就该好。 2. **行级 vs 环境级看失败与"这一行的数据"有没有关系**: 只与本行数据有关(换一行 可能成功)= 行级,逐条丢弃;与数据无关、每一行都会同样失败 = 环境级,进冷却。 见 `_classify_failure`(全库唯一一处 PG 失败分类)与 `_handle_failure`(三个降级点 唯一一处处置)。 """ from __future__ import annotations import asyncio import time from typing import TYPE_CHECKING from loguru import logger from polygateway.telemetry.schema import ( COLUMNS, PG_BACKFILL, PG_DDL, insert_sql, missing_columns_warning, ) from polygateway.telemetry.status import TelemetryStatusTracker if TYPE_CHECKING: from collections.abc import Callable import asyncpg from polygateway.types import TelemetryStatus # 探测表是否存在;不需要任何权限,且与 INSERT 走同一套 search_path 解析 _TABLE_EXISTS = "SELECT to_regclass('llm_calls')" # 归还连接的独立上限(issue #15)。**不**复用写入预算: 写入预算已经花在 # acquire+execute 上,归还再给它一个同样大的额度,等于允许业务路径上的一次遥测 # 写入吃掉 2 倍预算。归还是本地动作(reset 一次往返),1 秒足够;超时即断开, # asyncpg 会在下次 acquire 时补一条新连接 _RELEASE_TIMEOUT_S = 1.0 # 关闭池的独立上限(issue #15)。**不**复用写入预算: 关闭跑在收尾路径而非业务 # 路径上,给它一个略宽的固定额度即可,但必须**有界**——asyncpg 的 # `Pool.close()` 会 await 每个 holder 的 `wait_until_released()`,in-flight # 连接不归还就无限等(`pool.py:939-948, 961-972`,60s 只发一条 warning), # 其 docstring 自己写着 "advisable to use asyncio.wait_for to set a timeout" _CLOSE_TIMEOUT_S = 5.0 # 探测现有列;尊重 search_path(to_regclass 按当前 search_path 解析) _EXISTING_COLUMNS = ( "SELECT attname FROM pg_attribute " "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 原生异步,无线程桥接。""" def __init__( self, dsn: str, *, pool: asyncpg.Pool | None = None, auto_migrate: bool, pool_max: int, write_timeout_s: float, now: Callable[[], float] = time.monotonic, ) -> None: """记下装配参数(不连库);列与 INSERT 语句在首次准备期定型。 Args: dsn: asyncpg 连接串(已剥驱动后缀)。 pool: 外部注入的池;注入方自己负责关闭。 auto_migrate: True 则给已存在的旧表自动补列;False(PG 侧的缺省档) 则一条 ALTER 都不发——`ALTER TABLE ADD COLUMN` 取 ACCESS EXCLUSIVE 锁,会排在长事务后阻塞该表其后所有查询,而遥测是业务路径上的内联 await。keyword-only **必填**: 缺省规则只写在 config 一处,不与本类 签名漂移(设计 D-c)。 pool_max: 自建池的连接数上限(issue #15)。稳态吞吐**按实测折算,不要按 `pool_max / RTT` 估**(那会乐观一倍): RTT ≈ 123ms 上 `pool_max=4` 实测约 15.6 行/秒(50 行并发批 3.2s)。与 `auto_migrate` 同一纪律: 必填,缺省只写在 config 一处。 write_timeout_s: 单次写入的硬预算,同时用作 connect 与 acquire 的上限。 now: 单调时钟,注入给降级 tracker(测试可推进冷却与节流窗口)。 """ try: import asyncpg # noqa: F401 - 仅探测 extra 是否安装 except ImportError as exc: raise ImportError( "Postgres 遥测未启用: 安装 pip install 'polygateway[postgres]' 后重试" ) from exc self._dsn = dsn self._pool: asyncpg.Pool | None = pool self._external_pool = pool is not None self._auto_migrate = auto_migrate self._pool_max = pool_max self._write_timeout_s = write_timeout_s # 先按全量列定型: 准备期探测失败时保守沿用全量(今天的行为) self._columns: tuple[str, ...] = COLUMNS self._insert = insert_sql("postgres", COLUMNS) self._schema_ready = False self._closed = False # 关了就是关了: 置位后写入短路且**不重建池** # 降级状态**只此一份**: 是否短路写入、多久重试一次、下游查到什么, # 全由 tracker 回答。两份状态(曾经的 `_failed` 布尔 + tracker)必然漂移 self._status = TelemetryStatusTracker(backend="postgres", now=now) self._init_lock = asyncio.Lock() @property def telemetry_status(self) -> TelemetryStatus: """当前可写状态快照(ports.TelemetryStatusProvider)。""" return self._status.snapshot() async def _ensure_ready(self) -> asyncpg.Pool | None: """lazy 建池+备表;降级期间**零成本短路**,冷却到期放行一次重新准备。 `should_retry()` 是纯时间比较,不触库: 降级期间的调用因此既不内联吞 connect 超时(`postgres.py` 老注释担心的正是这个),也不需要重启进程—— 成本变成"每 60s 一次、上界一个写入预算",有界且可解释。 `_closed` 在锁内**必须复查**: 等锁期间发生的 `aclose` 否则会被这次 等待"绕过",等到锁时照旧建出一个没人负责关的池(注入档更隐蔽—— 注入方以为自己管着全部连接,实际早已不是)。 """ 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 not self._status.should_retry(): 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: """建池;失败按性质分档处置(见 `_handle_failure`),不再一律判死。 **池的资源占用由本库显式声明**(issue #15): `min_size=0` 的语义是"不预 连接"(asyncpg `pool.py:457` 为 0 时只造 holder 对象,一条连接都不连), 建池因此从"要么拿到 10 条、要么失败"的重资源动作变成零成本、不触库的 动作;连接失败自然落到 acquire 那条本来就正确的"丢一行、池自恢复"路径。 `max_size` 是库对自己占用的表态——继承第三方默认值等于不表态(P4),而 那正是共享实例余量紧张时先倒下的原因。 """ if self._pool is not None: return self._pool try: import asyncpg self._pool = await asyncpg.create_pool( self._dsn, min_size=0, max_size=self._pool_max, timeout=self._write_timeout_s, command_timeout=self._write_timeout_s, ) except asyncio.CancelledError: raise except Exception as exc: self._handle_failure(exc, stage="建池") return None return self._pool async def _prepare_schema(self, pool: asyncpg.Pool) -> asyncpg.Pool | None: """备好表并交回可用的池;瞬时失败只跳过本次,确定写不进去才判死。 取连接走显式 acquire/release(理由见 `_release`): 准备期同样跑在调用方的 写入预算里,`async with` 那条路的归还会把真实上界撑到 ≈2× 预算。 """ try: conn = await pool.acquire(timeout=self._write_timeout_s) try: columns = await self._prepare_table(conn) finally: await self._release(pool, conn) except asyncio.CancelledError: raise except Exception as exc: self._handle_failure(exc, stage="建表探测") return None if columns is None: # 表确定不存在且建不出来: 与本行数据无关(每行都会同样失败)且能被 # 外部修好(DBA 建了表就该自愈)—— 判据第 2 句下的环境级 self._status.enter_degraded( "表 llm_calls 不存在且建不出来(记录无处可落)", fatal=False, cooldown_s=_DEGRADE_COOLDOWN_S, ) 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) -> tuple[str, ...] | None: """备好 `llm_calls` 并返回本实例要写的列;**表存在就绝不发 DDL**。 返回 None 仅表示表确定不存在且建不出来(调用方据此进环境级冷却降级)。 `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: return await self._resolve_columns(conn) # 旧表可能缺列 try: await conn.execute(PG_DDL) # type: ignore[attr-defined] except asyncio.CancelledError: raise except Exception as exc: logger.warning("Postgres 遥测建表失败(表不存在,记录无处可落): {}", exc) return None return COLUMNS # 新建表列已齐全,无需再走补列 async def _resolve_columns(self, conn: object) -> tuple[str, ...]: """探测旧表现有列并定型写入列: auto 档先补齐,manual 档改为裁剪(issue #13)。 **先探测**的理由(两档共用): `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] 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);**失败绝不让 recorder 降级**。 不降级的实测理由: 应用账号只有 INSERT 权限时,`ALTER TABLE` 的 ownership 检查早于 `IF NOT EXISTS` 的存在性判断——列明明齐全也会失败。降级会让 整个 recorder 停写(环境级还要停满一个冷却期),与「补列失败只降级为逐行丢弃」的承诺相悖 (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: raise 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: 级别由 tracker 按 `fatal` 决定(致命档发 # error——人配错了,本进程内不会自愈)。此处复制一条只会让同一个事实出 # 两条语义重复的日志,并给"级别"这个决策造出第二个源头 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: """写一行遥测;整次写入受硬预算约束,失败逐条丢弃(两级降级之二),绝不冒泡。 **硬预算**(issue #15): 准备 + 取连接 + 执行合计不得超过 `write_timeout_s`。 这把"遥测绝不拖垮业务"从"靠各处 timeout 参数凑"变成一条可陈述、可测试的 保证——此前 `pool.acquire()` 无超时(asyncpg 缺省 `timeout=None` = 无限等), 池满时会无限期挂在业务路径上。 外部取消照常穿透: `asyncio.timeout` 只把**自己**触发的 cancel 转成 TimeoutError,故 `CancelledError` 分支必须排在最前且原样 re-raise(铁律)。 """ try: async with asyncio.timeout(self._write_timeout_s): await self._write_row(fields) except asyncio.CancelledError: raise except TimeoutError: logger.warning( "Postgres 遥测写入超预算 {}s(丢弃该行);后端慢不得拖垮业务调用", self._write_timeout_s, ) self._status.record_drop("写入超预算") except Exception as exc: # 遥测铁律: 丢一条 < 拖垮调用。这一行无论如何都没了,区别只在于 # **下一行还试不试**——那由失败的性质决定,不由这里决定 self._handle_failure(exc, stage="写入") self._status.record_drop("写入失败") async def _write_row(self, fields: dict[str, object]) -> None: """预算内的写入本体: 准备 → 取连接 → 执行 → 归还。 取值按 `self._columns`(manual 档可能已被裁剪),与 `self._insert` 的 占位符同序——两者必须一起改,分开改就是把值写进错位的列。 """ pool = await self._ensure_ready() if pool is None: # 降级期间静默 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) try: 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 会补一条新的)。 **不用 `async with pool.acquire()`**(设计 §3.1,已核实): asyncpg 的 `Pool.release` 是 `await asyncio.shield(ch.release(timeout))`,且那个 timeout 默认复用 acquire 时记录的 `ch._timeout`(`pool.py:886-889, 930-937`)。写入预算到期时 cancel 在 execute 处抛出,异常传播中执行 `__aexit__`,此时没有新的 cancel 投递——那次 shielded release 会**正常 等到完成**,业务路径的真实上界因此变成 ≈2 × 预算。显式归还才能给它一个 独立的小上限,承诺才精确成立: 主写入尝试 ≤ 预算,归还路径独立有界。 """ try: await pool.release(conn, timeout=_RELEASE_TIMEOUT_S) except asyncio.CancelledError: raise except Exception as exc: # 含 TimeoutError: 归还超时与归还出错的处置相同——断开而不是留一条 # 状态不明的连接在池里(asyncpg 的 reset 失败路径也是这么做的) logger.warning("Postgres 遥测连接归还失败(强制断开): {}", exc) self._terminate(conn, label="连接") @staticmethod def _terminate(target: object, *, label: str) -> None: """强制断开一条连接或整个池;断开本身再失败也只记 warning(遥测绝不冒泡)。 `label` 必填(不给默认值): 两个调用点的诊断价值全在"拆的是哪一层", 默认值只会让其中一处悄悄报错成另一处。 """ try: target.terminate() # type: ignore[attr-defined] except asyncio.CancelledError: raise except Exception as exc: logger.warning("Postgres 遥测{}断开失败(交给上层自行回收): {}", label, exc) async def aclose(self) -> None: """幂等关闭自建池;**关了就是关了**,此后写入短路且不复活。注入的池归注入方管理。 取消"关完还能自己重建池"的灰色状态(设计 §3.2 第 4 点): 关闭是所有权的 终结,而恢复是运行时行为(冷却重试),不该是关闭动作的副作用。 **关闭动作本身也有界**: `Pool.close()` 会 await 每个 holder 的 `wait_until_released()`,in-flight 连接不归还就无限等——收尾路径上照样是 "遥测拖垮业务"。超时即 `terminate()` 强拆: 关闭已在进行,留着一个关不掉的池 既不会自愈也没人再来收。外部取消照常穿透(铁律),不当成一次关闭超时。 """ self._closed = True pool, self._pool = self._pool, None self._schema_ready = False if pool is None or self._external_pool: return try: await asyncio.wait_for(pool.close(), timeout=_CLOSE_TIMEOUT_S) except asyncio.CancelledError: raise except Exception as exc: # 含 TimeoutError: 关不掉与关出错的处置相同——强拆 logger.warning("Postgres 遥测池关闭失败(强制断开): {}", exc) self._terminate(pool, label="池")