docs: plan the telemetry pool lifecycle rework for issue 15
The design traces the incident to four stacked defects rather than one bad default: the pool is the only resource in the library that pre-allocates, the kill switch keys off which step failed instead of what failed, the degraded state can neither recover nor be observed, and the ownership rules make the sanctioned sharing path unusable. The plan sequences the tracker ahead of the pool and failure work so every commit stays green, and records two facts the implementer needs up front: the pool-construction path has zero test coverage today, and the commit gate runs the full suite plus a complexity ceiling.
This commit is contained in:
@@ -0,0 +1,348 @@
|
||||
# 实现计划: 遥测连接池的资源语义与生命周期(issue #15)
|
||||
|
||||
- **设计**: `research-wiki/designs/2026-08-24-issue15-telemetry-pool-lifecycle-design.md`(已过 Codex 审 + 人类审)
|
||||
- **涉及技术**: Python 3.12(PEP 695 已就位)、asyncpg 0.31 连接池、`asyncio.timeout`、PG SQLSTATE、frozen dataclass、`@runtime_checkable` Protocol、pytest(含真实 PG 的 integration)
|
||||
- **版号**: 1.3.0(人类已定;**本计划不 bump 版本号**,那是发布清单第 3 步的事)
|
||||
|
||||
## 目标
|
||||
|
||||
让遥测池的资源占用与真实负载挂钩,把"建池失败 → 整进程永久失遥测"这条路彻底拆掉,并让任何降级都可恢复、可见、可编程。
|
||||
|
||||
## 方案概述
|
||||
|
||||
`min_size=0` 让建池变成零成本动作(实测不触库),连接失败自动落到 `acquire` 那条本来就正确的"丢一行、池自恢复"路径;判死判据从"哪一步失败"改为"失败是什么性质",永久档窄到只剩"DSN 不可解析",其余一律 60s 冷却重试;降级状态升格为共用的一等对象(节流日志 + 只读快照);顺带把"谁建的谁关"统一为全库纪律,让 ARCH §7.7 R5 的显式共享真正可用。
|
||||
|
||||
## 保真校验适用性
|
||||
|
||||
**不适用**。遥测后端无参考实现蓝本(ARCHITECTURE.md §7.8 明记"参考仓无先例: 三项目遥测全 SQLite"),本计划不涉及 `reference/` 迁移。但有两条**同等强度的既有承诺**不得被本次改动破坏,各任务已挂检查点:
|
||||
|
||||
1. issue #13 的"manual 档缺列时裁剪 INSERT 继续写、逐行 warning 暴露"(T5 的 `42703` 例外);
|
||||
2. issue #9 的"表存在就绝不发 DDL"(`to_regclass` 先探测,T4/T5 不得碰这段控制流)。
|
||||
|
||||
## 起点状态(执行前必读)
|
||||
|
||||
- **工作区有未提交改动且在 `main` 上**: Python 3.12 迁移已执行完毕(`pyproject.toml` `requires-python`/`target-version`、`README.md` 两处、`CLAUDE.md` 技术栈、`client.py` 与 `streaming.py` 的 UP047 三处改 PEP 695),conda 环境已重建为 3.12.13 并补装 `build`/`twine`。**T0 的第一件事就是把它们落到分支上**。
|
||||
- **建池路径今天零测试覆盖**: 全 `tests/` 目录对 `create_pool` 与 `_open_pool` 的引用数为 **0**(执行前可自行复核)。现有 PG 用例一律经 `pool=_FakePgPool(...)` 注入,走的是 `_external_pool=True` 分支,**从不经过建池**。这正是 `min_size=10` 潜伏至今的原因,也意味着 T3 要新建这一路的第一个用例。
|
||||
|
||||
## 提交门(每个提交点都受此约束)
|
||||
|
||||
`.claude/scripts/hooks/pre-commit-guard.sh` 在检测到 `git commit` 时**阻塞式**执行: `ruff check src/`(任何问题即阻塞)、`radon cc src -n C`(圈复杂度 ≥ C 即阻塞)、`pytest tests/ --tb=line -q`(任一红即阻塞)。文件 > 200 行只是 warning,不阻塞。
|
||||
|
||||
两条由此而来的硬约束:
|
||||
|
||||
- **不得留红态跨提交**——任务边界必须切在"全绿"处,不能把一个行为拆成"改实现"和"改测试"两次提交。
|
||||
- **圈复杂度是真实风险**: `record_llm_call` 本次要同时接入硬预算、失败分类与 tracker。一旦逼近 C 就必须抽私有方法,**这不算计划外重构**,是提交门的硬要求。
|
||||
|
||||
## 文件结构
|
||||
|
||||
| 文件 | 动作 | 职责 |
|
||||
|---|---|---|
|
||||
| `src/polygateway/types.py` | 改 | 新增 `TelemetryStatus` frozen dataclass(与 `SourceStats` 同一先例) |
|
||||
| `src/polygateway/ports.py` | 改 | 新增**独立** `TelemetryStatusProvider` Protocol;`TelemetryRecorder` **一字不动** |
|
||||
| `src/polygateway/telemetry/status.py` | **新建** | `TelemetryStatusTracker`: 降级状态机 + 节流日志 + 快照。两个 recorder 共用,不含任何后端知识 |
|
||||
| `src/polygateway/telemetry/postgres.py` | 改 | 池语义、硬预算、失败三分、冷却降级、有界 `aclose`、接入 tracker |
|
||||
| `src/polygateway/telemetry/sqlite.py` | 改 | **仅**接入 tracker(补上今天缺失的降级 warning);不做 lazy 化与冷却 |
|
||||
| `src/polygateway/config.py` | 改 | 两个新键的加载与校验 |
|
||||
| `src/polygateway/client.py` | 改 | 所有权纪律 + `aclose` helper + `telemetry_status` 出口 |
|
||||
| `src/polygateway/embedding.py`、`ocr.py` | 改 | 同款所有权与出口(三处必须一致) |
|
||||
| `src/polygateway/backends/redis_cache.py` | 改 | 补 `_owns_client` 纪律 |
|
||||
| `tests/unit/test_telemetry.py` | 改 | `_FakePgPool` 改造 + 池语义/预算/分类/冷却/tracker 用例 |
|
||||
| `tests/unit/test_client.py` | 改 | 所有权层用例(三个 client 各钉一次) |
|
||||
| `tests/unit/test_config.py` | 改 | 两个新键的三条装配路 |
|
||||
| `tests/integration/test_postgres_telemetry.py` | 改 | 真实 PG: 连接数计数、降级恢复 |
|
||||
| `.env.example`、`README.md`、`CHANGELOG.md`、`research-wiki/ARCHITECTURE.md` | 改 | 配置面、能力表、发布说明、架构决策成文 |
|
||||
|
||||
## 关键接口(跨任务消费,此处定死)
|
||||
|
||||
`types.py` 新增(T2 建立,T4/T5/T6 消费):
|
||||
|
||||
```python
|
||||
@dataclass(frozen=True)
|
||||
class TelemetryStatus:
|
||||
"""遥测后端的可写状态快照;degraded 期间下游可据此对账(issue #15)。"""
|
||||
degraded: bool
|
||||
fatal: bool # True = 本进程内不可恢复(仅 DSN 不可解析一类)
|
||||
reason: str | None # 降级原因;未降级为 None
|
||||
degraded_for_s: float | None # 已降级时长;未降级为 None
|
||||
dropped_rows: int # 累计丢弃行数(进程生命周期内单调不减)
|
||||
retry_after_s: float | None # 距下次重新准备;fatal 或未降级为 None
|
||||
```
|
||||
|
||||
`ports.py` 新增(T2 建立)——**独立于 `TelemetryRecorder`**,理由见设计 §3.3:
|
||||
|
||||
```python
|
||||
@runtime_checkable
|
||||
class TelemetryStatusProvider(Protocol):
|
||||
"""可自述可写状态的遥测后端;与 TelemetryRecorder 分开是为了不破坏后者的
|
||||
runtime_checkable 语义(加成员会让只实现 record_llm_call 的对象当场不满足协议)。"""
|
||||
|
||||
@property
|
||||
def telemetry_status(self) -> TelemetryStatus: ...
|
||||
```
|
||||
|
||||
`telemetry/status.py` 新增(T2 建立,T4/T5 消费)。`now` 注入以便测试推进假时钟:
|
||||
|
||||
```python
|
||||
class TelemetryStatusTracker:
|
||||
def __init__(self, *, backend: str, now: Callable[[], float] = time.monotonic) -> None: ...
|
||||
def enter_degraded(self, reason: str, *, fatal: bool, cooldown_s: float | None) -> None: ...
|
||||
def recover(self) -> None: ...
|
||||
def record_drop(self, reason: str) -> None: ...
|
||||
def should_retry(self) -> bool: ... # fatal→False;冷却未到→False;到期→True
|
||||
def snapshot(self) -> TelemetryStatus: ...
|
||||
```
|
||||
|
||||
`PostgresRecorder.__init__` 新签名(T3 落地;`pool_max`/`write_timeout_s` keyword-only **必填**,与 `auto_migrate` 同一纪律——缺省只写在 config 一处):
|
||||
|
||||
```python
|
||||
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: ...
|
||||
```
|
||||
|
||||
`GatewaySettings` 新字段与 env 键(T3 落地):
|
||||
|
||||
| 字段 | env 键 | 缺省 | 校验(落 `_validate_telemetry`) |
|
||||
|---|---|---|---|
|
||||
| `telemetry_pg_pool_max: int` | `PGW_TELEMETRY_PG_POOL_MAX` | 4 | `>= 1`,否则 ValueError 点出字段名与键名 |
|
||||
| `telemetry_pg_write_timeout_s: float` | `PGW_TELEMETRY_PG_WRITE_TIMEOUT_S` | 5.0 | `> 0`,同上 |
|
||||
|
||||
失败三分(T5 落地,`postgres.py` 模块级私有函数,全库唯一一处 PG 失败分类):
|
||||
|
||||
```python
|
||||
_FATAL = "fatal" # 配置级致命 → 永久 no-op + 一条 error
|
||||
_UNAVAILABLE = "unavailable" # 环境级 → 60s 冷却降级
|
||||
_ROW = "row" # 行级 → 逐条 warning 丢弃
|
||||
|
||||
def _classify_failure(exc: BaseException) -> str: ...
|
||||
```
|
||||
|
||||
判据(设计 §3.2,两句): ①致命 = 原因完全在进程内部且不可变;②行级 vs 环境级看失败与**这一行的数据**有没有关系。落到具体码:
|
||||
|
||||
| 归档 | 覆盖 |
|
||||
|---|---|
|
||||
| `_FATAL` | `asyncpg.ClientConfigurationError`;`create_pool` 抛的 `ValueError`/`TypeError` |
|
||||
| `_UNAVAILABLE` | SQLSTATE 前两位 ∈ {`08`,`53`,`57`,`28`,`3D`} + 具体码 `42501`、`42P01`;`OSError`/`ConnectionError`/`TimeoutError`/其余 `InterfaceError` |
|
||||
| `_ROW` | 其余 `PostgresError`(`22`/`23` 等)+ **具名例外 `42703`**(缺列,由 issue #13 承诺定死) |
|
||||
|
||||
冷却期为模块级常量 `_DEGRADE_COOLDOWN_S = 60.0`(不暴露配置,设计 §3.5)。
|
||||
|
||||
## 任务清单
|
||||
|
||||
### T0 — 分支与基线(把已完成的 3.12 迁移落盘)
|
||||
|
||||
- [ ] 从 `main` 建分支 `feat/issue-15-telemetry-pool-lifecycle`
|
||||
- [ ] 把工作区现有改动分两次提交: ① `chore: 最低 Python 提到 3.12 并改用 PEP 695 泛型语法`(`pyproject.toml`/`README.md`/`CLAUDE.md`/`client.py`/`streaming.py`);② `docs: issue #15 设计文档与 wiki 登记`(`research-wiki/`)
|
||||
- [ ] 记录基线用例计数(执行时实测;2026-08-24 本机为 **973 passed / 23 skipped / 45 deselected**,覆盖率 94%)。该数只作**同环境**参照,不作硬验收——`addopts = "-m 'not slow'"` 与 Redis/PG 可达性都会改变它
|
||||
|
||||
**验证**: `make check` 全绿;`/home/iomgaa/miniconda3/envs/PolyGateway/bin/python -m pytest tests/ -q` → 全 PASS;`git rev-parse --abbrev-ref HEAD` → 分支名正确。
|
||||
|
||||
> **不要用 `make lint` 做验证**——它带 `--fix` 会自动改文件(`Makefile:11`),只读验证用 `make check`。
|
||||
> **不要用 `conda run ... pytest` 取统计数字**——实测其输出缓冲会把结尾的 `N passed` 与覆盖率整段吞掉,只剩 exit code(2026-08-24 踩过)。用环境解释器绝对路径直跑。
|
||||
|
||||
---
|
||||
|
||||
### T1 — D 组: 资源所有权纪律统一(独立回滚点)
|
||||
|
||||
**动**: `src/polygateway/client.py`、`embedding.py`、`ocr.py`、`backends/redis_cache.py`;测试 `tests/unit/test_client.py`。
|
||||
|
||||
**要实现的行为**: 全库唯一纪律 —— **谁建的谁关,注入的一律不碰**。分两层落:
|
||||
|
||||
1. **组件内部自建的连接**归组件自己: `RedisCache` 补 `_owns_client`(构造注入 → False;`from_url` → True),`aclose` 自查后再关。这是照抄 `backends/redis/limiter.py:185-191, 318-322` 的既有正确先例,`backends/redis/breaker.py:437-441` 同款。
|
||||
2. **client 自建的整个组件**归 client: 三个 client 各持 `_owns_transport/_owns_telemetry/_owns_cache/_owns_limiter/_owns_breaker`,**默认全 False**(`__init__` 是全量注入路径,经它传入的一切都是外部的),只有三个工厂在真正自建时置 True。工厂里 `transport` 恒自建(三处工厂都没有 transport 注入参数),`limiter`/`breaker`/`cache`/`telemetry` 按 `xxx is None` 判定。
|
||||
|
||||
**同时修掉的现存泄漏**: `GatewayClient.__init__` 今天把 limiter/breaker 交给 `RetryMW` 构造(`client.py:156-176`)后自己不留引用(`self._transport`/`_telemetry`/`_cache` 都存了,唯独这两个没存,见 `client.py:203-206`),`aclose` 因此**触达不到**自建的 redis 客户端。三个 client 都要新持 `self._limiter`/`self._breaker` 引用(仅为关闭)。embedding/ocr 的自建点在 `embedding.py:497-498`、`ocr.py:510-511`。
|
||||
|
||||
**收敛**: 三处复制的 `getattr(..., "aclose")` 探测(`client.py:268-280`、`embedding.py:452-461`、`ocr.py:462-467`)收敛为**一个**内部 helper。SQLite recorder 只有同步 `close()`,helper 须同时探测 `aclose`/`close`(今天 `client.py:274-277` 已有这个分支,embedding/ocr 也有,收敛后行为不变)。内存后端无 `aclose`,探测后跳过。
|
||||
|
||||
**测试要求**(先失败后通过): 假 recorder/transport/limiter/breaker/cache 各记 close 次数。
|
||||
- 注入的组件 `aclose` 后 close 次数 **0**;自建的为 **1**(工厂路径);
|
||||
- 自建 redis limiter/breaker 被关(**泄漏钉子**,今天必红);
|
||||
- 注入给 `RedisCache` 的客户端不被关;
|
||||
- **三个 client 逐一覆盖**——收敛成 helper 之后仍须三处各钉一次,否则下次有人把逻辑复制回去无人发现;
|
||||
- `aclose` 幂等(连调两次不重复关)。
|
||||
|
||||
**验证**: `pytest tests/unit/test_client.py tests/unit/test_embedding.py tests/unit/test_ocr_client.py -q` → PASS;`make check` 绿;全套件绿。
|
||||
|
||||
- [ ] 提交: `fix: 统一资源所有权纪律(谁建的谁关),修 aclose 越权与 redis 客户端泄漏`
|
||||
|
||||
---
|
||||
|
||||
### T2 — C 组基础设施: 状态快照 + tracker + 出口
|
||||
|
||||
**动**: `src/polygateway/types.py`、`ports.py`、**新建** `telemetry/status.py`、`telemetry/postgres.py`、`telemetry/sqlite.py`、`client.py`、`embedding.py`、`ocr.py`;测试 `tests/unit/test_telemetry.py`、`test_ports.py`、`test_client.py`。
|
||||
|
||||
**为什么排在 A/B 组之前**: T3/T5 的所有降级点都要向 tracker 报告。先建 tracker 则那两步直接写成最终形态,反之要返工一遍日志代码。
|
||||
|
||||
**要实现的行为**:
|
||||
1. `TelemetryStatus` 与 `TelemetryStatusProvider` 按上文"关键接口"定死。**`TelemetryRecorder` 一字不动**。
|
||||
2. `TelemetryStatusTracker` 状态机: `enter_degraded` 打一条 warning(含原因与恢复条件: 冷却剩余秒数,或 fatal 时写明"需改配置并重启");降级期间 `record_drop` **节流复述**(按丢弃行数与时间双阈值,阈值为模块常量);`recover` 打一条 info 并报告"期间丢弃 N 行";`should_retry` 是纯查询(fatal → False,冷却未到 → False)。
|
||||
3. 两个 recorder 各持一个 tracker,把**今天已有的**降级点接上去: PG 的建池失败与判死、SQLite 的初始化失败。**SQLite 侧同时补上今天缺失的那条 warning**——`sqlite.py:138-139` 初始化失败后写入直接 `return`,连一条日志都没有。
|
||||
4. 出口 `telemetry_status` 属性加到三个 client,取值经**一处** `isinstance(self._telemetry, TelemetryStatusProvider)` 判定,不满足或无遥测则返回 `None`。
|
||||
|
||||
**本任务不改任何失败判据**: PG 侧仍是"建池失败即永久判死",只是这次判死会经 tracker 变得可见。判据在 T5 改。这样本任务的行为变更面收敛为"日志更可见 + 多一个只读出口"。
|
||||
|
||||
**过渡期状态并存(有意,且必须在 T5 收掉)**: 本任务结束时 PG 侧的 `_failed` 布尔与 tracker 的 fatal 状态**并存**——判死点两边都写。这是为了让 T2 能独立全绿提交,不是最终形态;T5 删除 `_failed`,状态收归 tracker 一处。两份状态只允许存活这一个任务的跨度,拖久了必然漂移。
|
||||
|
||||
**契约检查点**: `tests/unit/test_ports.py:137,141` 的 `isinstance(_DummyRecorder(), TelemetryRecorder)` 断言必须**保持绿**——它是"没把状态并进主 Protocol"这条决策的机械化执法点,新增用例不得替代它。
|
||||
|
||||
**测试要求**(先失败后通过):
|
||||
- tracker 状态机六字段逐个钉: 未降级 → `degraded=False` 且三个可空字段为 None;进入降级 → `reason`/`retry_after_s` 正确;假时钟推进 → `degraded_for_s` 增长、`retry_after_s` 递减到 0;`recover` → 回到未降级且 `dropped_rows` **不清零**(进程生命周期内单调不减);
|
||||
- 节流复述: 连续 N 次 `record_drop` 只产生 M 条 warning(loguru sink 捕获断言),且 N 与 M 的关系由常量决定而非硬编码数字;
|
||||
- fatal 档: `should_retry()` 恒 False,`retry_after_s` 为 None;
|
||||
- SQLite 初始化失败(指向不可写目录)→ 有 warning **且** `telemetry_status.degraded is True`(今天必红,连 warning 都没有);
|
||||
- 三个 client 的 `telemetry_status`: 无遥测 → None;注入不实现该 Protocol 的假 recorder → None(不得抛 AttributeError);内置 recorder → 返回快照。
|
||||
|
||||
**验证**: `pytest tests/unit/test_telemetry.py tests/unit/test_ports.py tests/unit/test_client.py -q` → PASS;`lint-imports` 绿(新文件 `telemetry/status.py` 在实现层,只许依赖 `types`/`ports`/标准库,**不得**被 `transports`/`backends` import);全套件绿。
|
||||
|
||||
- [ ] 提交: `feat: 遥测降级升格为一等状态(共用 tracker + 只读快照 + 节流日志)`
|
||||
|
||||
---
|
||||
|
||||
### T3 — A 组: 池语义与两个新配置键
|
||||
|
||||
**动**: `src/polygateway/config.py`、`client.py`(`_build_telemetry`)、`telemetry/postgres.py`;测试 `tests/unit/test_config.py`、`test_telemetry.py`、**`tests/integration/test_postgres_telemetry.py`**。
|
||||
|
||||
> **本任务必须一次改完全部 20 处 `PostgresRecorder(` 构造点**(Codex 审查,已实测复核): `src/polygateway/client.py` 1 处 + `tests/unit/test_telemetry.py` 3 处 + **`tests/integration/test_postgres_telemetry.py` 16 处**。新签名的 `pool_max`/`write_timeout_s` 是 keyword-only **必填**,漏一处就 `TypeError`,而提交门跑的是**全套件**——集成测试那 16 处不能拖到 T6,否则 T3 根本提交不了。这是 `auto_migrate` 当初(issue #13)踩过的同一形态: 必填 keyword-only 的代价就是所有构造点同批改。
|
||||
|
||||
**要实现的行为**:
|
||||
1. 两个新配置键按"关键接口"那张表落地: `_load_pgw` 里读取(模板照 `config.py:524-543` 的 `_load_text_cap`),值域校验落 `_validate_telemetry`(与 `telemetry_text_cap` 同一先例,**一次覆盖直接构造 / `dataclasses.replace` / env 三条路**),报错文本同时点字段名与 env 键名。`_build_telemetry`(`client.py:405-420`)把两个值透传给 recorder。
|
||||
2. 建池改为 `create_pool(dsn, min_size=0, max_size=pool_max, timeout=write_timeout_s, command_timeout=write_timeout_s)`。
|
||||
3. **两处** `acquire` 都改为**显式** acquire/release,**不得**用 `async with pool.acquire(...)`——`_prepare_schema`(`postgres.py:114`)与 `record_llm_call`(`postgres.py:238`)。准备期同样在预算内、同样吃 shielded release 那一刀,只改一处等于留了半个坑:
|
||||
- `con = await pool.acquire(timeout=剩余预算)`;
|
||||
- `finally: await pool.release(con, timeout=<小的独立上限>)`,释放超时则 `con.terminate()`;
|
||||
- 整次写入(准备 + acquire + execute)由 `asyncio.timeout(write_timeout_s)` 包一层。
|
||||
|
||||
**理由(设计 §3.1,已核实)**: `Pool.release()` 是 `await asyncio.shield(ch.release(timeout))` 且默认复用 acquire 记录的 `ch._timeout`(asyncpg `pool.py:886-889, 930-937`)。外层预算到期时 cancel 在 `execute` 处抛出,异常传播中执行 `__aexit__`,此时没有新的 cancel 投递,那个 shielded release 会**正常等到完成**——用 `async with` 的真实上界是 ≈ 2 × 预算。
|
||||
|
||||
**必须同步改造 `_FakePgPool`**(`tests/unit/test_telemetry.py:751`): 它今天的 `acquire()` **无参**且只返回一个 `_Ctx` 异步上下文管理器,没有 `release`。改造为接受 `timeout=` 并提供 `release(con, timeout=)`,同时记录 acquire/release 的配对次数(T3 与 T5 的用例都要用)。不改造则全部 PG 用例当场红。
|
||||
|
||||
**取消穿透的实现纪律**(铁律): 降级路径(节流日志、tracker 更新、release 收尾)一律不得 `except CancelledError` 而不 re-raise;`except TimeoutError` 必须排在 `except Exception` 之前;严禁裸 `except BaseException`。既有 `postgres.py:101-102` 的 `except asyncio.CancelledError: raise` 写法是对的,延续它。
|
||||
|
||||
**测试要求**(先失败后通过。注意: 建池路径**今天零覆盖**,这里要建立第一个用例):
|
||||
- **主回归钉子**: monkeypatch `asyncpg.create_pool`,断言实参 `min_size == 0` 且 `max_size == 配置值`。这一条防的是回归到继承第三方默认值,是本 issue 的核心;
|
||||
- 配置键三条装配路: env 路读取正确、缺省为 4 / 5.0、直接构造与 `replace` 同样被校验拦住(`pool_max=0`、`write_timeout_s=0` 各一条,断言报错文本含字段名与键名);
|
||||
- 硬预算: 假 pool 的 acquire 挂住 → 丢一行且耗时 ≤ 预算(用假时钟或极小预算,**不要**在用例里真睡 5 秒);
|
||||
- **release 不泄漏**(Codex 审查钉子): `execute` 被预算取消后,断言 `_FakePgPool` 记录的 acquire/release 次数**配对**;
|
||||
- 外部 `CancelledError` 在预算内**不**被吞成 `TimeoutError`(直接钉铁律)。
|
||||
|
||||
**验证**: `pytest tests/unit/test_config.py tests/unit/test_telemetry.py -q` → PASS;`make check` 绿;全套件绿。
|
||||
|
||||
- [ ] 提交: `feat: 遥测池显式声明资源占用(min_size=0/max_size 可配)并给写入硬预算`
|
||||
|
||||
---
|
||||
|
||||
### T4 — B 组之一: 有界关闭
|
||||
|
||||
**动**: `src/polygateway/telemetry/postgres.py`;测试 `tests/unit/test_telemetry.py`。
|
||||
|
||||
**为什么单列一个任务**: 它与 T5 的失败判据无关,但同属"收尾路径的隐性无界等待",且能独立验证。合进 T5 会让那次提交同时动判据与关闭两件事,回滚粒度变粗。
|
||||
|
||||
**要实现的行为**: `aclose()` 语义钉死为"关了就是关了"——置 `_closed`,此后写入短路且**不复活**(取消今天"关完还能自己重建池"的灰色状态);关闭动作本身走 `asyncio.wait_for(pool.close(), timeout=...)`,超时后 `pool.terminate()`,外部取消照常穿透。
|
||||
|
||||
**理由(已核实)**: `Pool.close()` 会 `await` 每个 holder 的 `wait_until_released()`,in-flight 未释放时**无限等**,60 秒只发一条 warning(asyncpg `pool.py:939-948, 961-972`);asyncpg 自己的 docstring 就写着 "advisable to use `asyncio.wait_for` to set a timeout"。
|
||||
|
||||
**测试要求**(先失败后通过):
|
||||
- 假 holder 永不 release → `aclose()` 在超时后走 `terminate()` 返回,**不无限挂**(今天必红/挂死,用例须自带超时保护);
|
||||
- `aclose` 后再 `record_llm_call` → 直接短路,**不重建池**(断言 `create_pool` 未被再次调用);
|
||||
- `aclose` 幂等;注入的外部池仍**不**被关(`_external_pool` 既有纪律不得破)。
|
||||
|
||||
**验证**: `pytest tests/unit/test_telemetry.py -q` → PASS;全套件绿。
|
||||
|
||||
- [ ] 提交: `fix: 遥测池关闭有界化(wait_for + terminate),关闭后不再复活`
|
||||
|
||||
---
|
||||
|
||||
### T5 — B 组之二: 失败三分与冷却降级(本 issue 的核心)
|
||||
|
||||
**动**: `src/polygateway/telemetry/postgres.py`;测试 `tests/unit/test_telemetry.py`。
|
||||
|
||||
**要实现的行为**:
|
||||
1. 新增模块级 `_classify_failure`(按"关键接口"的三档表),全库唯一一处 PG 失败分类。
|
||||
2. 三个降级点改为按分类处置: `_open_pool`、`_prepare_schema`/`_prepare_table`、`record_llm_call`。
|
||||
- `_FATAL` → 永久 no-op + 一条 **error**(不是 warning: 这是人配错了),经 tracker 置 `fatal=True`;
|
||||
- `_UNAVAILABLE` → `tracker.enter_degraded(cooldown_s=_DEGRADE_COOLDOWN_S)`,此后 `_ensure_ready` 开头零成本短路(只比较时间戳,不触库),到期 `should_retry()` 放行**一次**重新准备,成功即 `tracker.recover()`;
|
||||
- `_ROW` → 逐条 warning 丢弃 + `tracker.record_drop()`,不降级。
|
||||
3. **删除 `_failed` 这个布尔**,状态收归 tracker 一处(否则两份状态必然漂移)。实测引用分布(执行时可自行复核): `src/polygateway/telemetry/postgres.py` **7 处**(74/79/84/106/124 是代码,209/211 在 `_backfill_columns` 的 docstring 里——**文档也要改**,否则留下指向已删字段的说明)、`tests/unit/test_telemetry.py` **6 处**、`tests/integration/test_postgres_telemetry.py` **6 处**,测试侧一并改为读 `telemetry_status` 快照。
|
||||
4. 判据的两条既有承诺不得破:
|
||||
- **`42703` 仍走 `_ROW`**(issue #13: manual 档缺列时裁剪 INSERT 继续写、逐行暴露)。这是判据的**唯一具名例外**,代码里必须有注释写明它是例外及理由;
|
||||
- **`_prepare_table` 的 `to_regclass` 先探测、表在就不发 DDL** 这段控制流(`postgres.py:147-157`)一行不动(issue #9)。
|
||||
|
||||
**圈复杂度检查点**: 本任务是三个降级点同时改,`record_llm_call` 与 `_ensure_ready` 最容易触到 radon 的 C 档而被提交门阻塞。逼近就抽私有方法(如 `_handle_failure(exc, *, stage)` 收敛三处处置)——这是提交门的硬要求,不算计划外重构。
|
||||
|
||||
**测试要求**(先失败后通过,分档逐个钉):
|
||||
- **issue 场景直接回归**: 建池阶段抛 `TooManyConnectionsError`(53300)→ **不** fatal、进冷却降级 → 假时钟推进 60s → 下次调用自动恢复并成功写入。今天这一条必红(现状是永久判死);
|
||||
- `ClientConfigurationError` → fatal + 一条 error + 此后零成本短路(断言不再调 `acquire`);
|
||||
- **分档边界两侧各钉一次**: `42501`/`42P01` → 进冷却降级;`42703` → 行级丢弃且**不**进降级;
|
||||
- `_prepare_table` 建表失败(表确定不存在)→ 冷却降级(不再是永久判死),DBA 建表后自动恢复;
|
||||
- 探测失败(既有 `probe_errors` 路径)仍只跳过本次、下次重试,**不**降级(issue #9 既有行为不得回归);
|
||||
- 全部现有 PG 用例保持绿(它们钉的是 issue #3/#9/#13 的承诺)。
|
||||
|
||||
**验证**: `pytest tests/unit/test_telemetry.py -q` → PASS;`radon cc src/polygateway/telemetry/postgres.py -n C -s` → 无输出;全套件绿。
|
||||
|
||||
- [ ] 提交: `fix: 遥测失败按性质三分,永久判死收窄到 DSN 不可解析,其余带冷却自愈`
|
||||
|
||||
---
|
||||
|
||||
### T6 — 真实 PG 集成验证
|
||||
|
||||
**动**: `tests/integration/test_postgres_telemetry.py`。
|
||||
|
||||
**纪律(该文件既有,不得破)**: `llm_calls` 是与真实批跑共享的表,**严禁 DROP/TRUNCATE**;以 run 级 `call_id` 前缀隔离,teardown 只删自己的行;DSN 缺失则 skip;不标 `slow`(与该文件既有用例一致)。
|
||||
|
||||
**要实现的行为(用例)**:
|
||||
1. **issue 的直接回归钉子**: 建 recorder 后本池连接数为 **0**,一次写入后 **≤1**,稳态 ≤ `pool_max`。
|
||||
2. 降级与恢复走**不可达 DSN** 的 recorder 验证(连接被拒 → 降级 → 假时钟/短冷却后重试),**不去动共享实例的 `max_connections`**。
|
||||
|
||||
**计数必须按唯一 `application_name` 过滤**,该实例被多项目共用,按库名或用户名计数会被别人的连接污染——那样的用例是**设计上就会间歇红**的信号污染源(CLAUDE.md §4.6)。
|
||||
|
||||
**怎么设这个 tag(Codex 指出原稿这里无法执行,已实测给出解法)**: recorder 的构造签名**没有** `server_settings`/`connect_kwargs` 入口,原稿那句"经 `server_settings=` 建池"落不了地。解法是走 **DSN 查询参数**——给 recorder 一个 `f"{dsn}?application_name={run级唯一值}"`,其余一切不变。
|
||||
|
||||
- 已实测(2026-08-24,真实实验室 PG): `create_pool(dsn + "?application_name=pgwtest-abc123", min_size=0, ...)` 后 `SHOW application_name` 返回该值,`pg_stat_activity` 按它过滤得连接数 1,`pool.close()` 后归零。
|
||||
- **不要**改用"测试自建池后以 `pool=` 注入": 那会走 `_external_pool=True` 分支、**完全绕过被测的建池路径**,而本任务要验的恰恰是自建池不预连接。
|
||||
- **不要**为此给 recorder 加 `server_settings` 入口: 纯测试便利不值得扩公共 API(P1)。
|
||||
- 注意 `config.py` 的 `_strip_dsn_driver` 只动 scheme 的 `+driver` 后缀,不碰查询参数;且集成测试直接构造 recorder、不经 config,两条路都不受影响。
|
||||
|
||||
**验证**: `pytest tests/integration/test_postgres_telemetry.py -q` → PASS(或无 DSN 时全 skip);全套件绿。
|
||||
|
||||
- [ ] 提交: `test: 真实 PG 验证遥测池不预连接与降级自愈`
|
||||
|
||||
---
|
||||
|
||||
### T7 — 文档、配置面与发布说明
|
||||
|
||||
**动**: `.env.example`、`README.md`、`CHANGELOG.md`、`research-wiki/ARCHITECTURE.md`。
|
||||
|
||||
**要实现的行为**:
|
||||
1. `.env.example`: 两个新键写在 `PGW_TELEMETRY_PG_DSN` 之后,沿用该文件既有的"键 + 缩进注释块讲清为什么"风格。`pool_max` 必须给**调参口径**: 期望吞吐 ≈ `pool_max / RTT`(附本次实测: 跨内网 RTT ≈ 123ms,`pool_max=4` ≈ 32 行/秒),并写明"共享一个 recorder 给多 client 时并发汇聚,应相应放大"。
|
||||
2. `README.md`: 配置表加两键;能力表反映"遥测降级可恢复 + 可查询状态";**核对安装命令里的版本约束**(发布清单第 1 步的老账: `==1.2.*` 这类极易漏改)。
|
||||
3. `ARCHITECTURE.md` §7.8 增补三条: 遥测池的资源语义(为何 `min_size=0`、为何不暴露 `min_size`)、失败三分的**两句判据**、**资源所有权纪律**(后者应作为跨子系统的通用纪律成文,而非遥测局部约定);§9 登记两个新键。
|
||||
4. `CHANGELOG.md`: 记在"未发布"下,三处"请先读这一条": ①最低 Python 提到 3.12(**唯一会让下游装不上**的变更);②遥测常驻连接从 `10 × client 数` 变按需(监控曲线会突变);③`aclose` 不再关闭注入的组件。
|
||||
|
||||
**验证**: `make check` 绿;人工通读 `.env.example` 两键注释,确认调参口径可执行。
|
||||
|
||||
- [ ] 提交: `docs: 遥测池资源语义、失败判据与所有权纪律成文`
|
||||
|
||||
---
|
||||
|
||||
## 完成判据(合并前)
|
||||
|
||||
- [ ] T0-T7 全部提交完成,每次提交都过了提交门(ruff + radon + 全套件)
|
||||
- [ ] `pytest -m slow` 单独跑过一次(发布清单第 4 步;本次改动触及遥测写入路径,e2e 与 Redis 时间语义变体必须实测)
|
||||
- [ ] 派**全新上下文**的 verifier subagent 独立验证(`verification-before-completion`,里程碑级/合并前 MANDATORY)
|
||||
- [ ] 整分支审查(`requesting-code-review`,合并前 MANDATORY)
|
||||
- [ ] 设计文档 §5 的每一条测试要求都能指到一个具体用例(逐条对照,不是"大致覆盖")
|
||||
|
||||
## 审查留痕(Codex,2026-08-24)
|
||||
|
||||
**Status: Issues Found → 2 条阻断级均已修订,2 条 Recommendation 采纳 1 条。**
|
||||
|
||||
| # | 结论 | 落点 |
|
||||
|---|---|---|
|
||||
| 1 | **采纳(阻断)**。新签名的 `pool_max`/`write_timeout_s` 是必填 keyword-only,而 `PostgresRecorder(` 共 **20 处**构造点,其中 **16 处在集成测试**。原稿 T3 只列了两个单元测试文件,漏掉的那 16 处会让 T3 的提交门(跑全套件)当场红 | T3 "动"一节 |
|
||||
| 2 | **采纳(阻断),并给出比建议更好的解法**。原稿 T6 写"经 `server_settings=` 建池"设唯一 `application_name`,但 recorder 签名根本没有这个入口,零上下文执行者会卡死。Codex 给的两条出路(注入外部池 / 加 recorder 入口)都有代价——前者绕过被测的建池路径,后者为测试便利扩公共 API。**实测发现第三条**: `?application_name=<tag>` 走 DSN 查询参数,asyncpg 认、PG 侧生效、关池后计数归零,**零 API 改动且真实覆盖建池路径** | T6 计数一节 |
|
||||
| 3 | **采纳(建议)**。`_failed` 计数原稿写"测试 13 处"不准。实测: 源码 7 处(**含 2 处在 docstring 里**,文档也要改)、unit 6 处、integration 6 处 | T5 第 3 点 |
|
||||
| 4 | 无需动作。Codex 复核确认了计划的两条硬断言: 建池路径零覆盖(`rg create_pool\|_open_pool tests` 无匹配)、`_FakePgPool` 定义于 `:751-765` 且只经三个 helper 注入(故改造类本身即可覆盖既有假池用例) | — |
|
||||
|
||||
Codex 给的 `_failed` 分布数字(源码 5 处 / 测试断言 8 处)与本地实测(源码 7 / unit 6 / integration 6)不一致,以实测为准——它漏了 docstring 里那两处,而那两处恰恰是**必须改**的(留着就是指向已删字段的说明)。
|
||||
|
||||
@@ -0,0 +1,18 @@
|
||||
---
|
||||
type: plan
|
||||
node_id: plan:plan-issue15-telemetry-pool-lifecycle
|
||||
title: "实现计划: 遥测连接池的资源语义与生命周期(issue #15)"
|
||||
date: 2026-08-24
|
||||
---
|
||||
|
||||
# 实现计划: 遥测连接池的资源语义与生命周期(issue #15)
|
||||
|
||||
正文: `2026-08-24-issue15-telemetry-pool-lifecycle.md`(326 行)。实现 [[design:issue15-telemetry-pool-lifecycle]]。
|
||||
|
||||
- **八个任务**: T0 分支与基线(把已完成的 Python 3.12 迁移落盘)→ T1 D 组所有权纪律(独立回滚点)→ T2 C 组 tracker 与状态快照 → T3 A 组池语义与两个新配置键 → T4 有界关闭 → T5 B 组失败三分与冷却降级(核心)→ T6 真实 PG 集成验证 → T7 文档与发布说明。
|
||||
- **顺序的关键理由**: tracker(T2)排在池语义(T3)与失败判据(T5)**之前**——后两步的每个降级点都要向 tracker 报告,反过来做要把日志代码返工一遍。代价是 T2 结束时 `_failed` 与 tracker 状态**临时并存**(为了让 T2 能独立全绿提交),T5 必须收掉,两份状态只允许存活一个任务的跨度。
|
||||
- **执行前必读的两条事实**: ① Python 3.12 迁移的改动**还在 main 的工作区未提交**(T0 第一件事就是落到分支);② **建池路径今天零测试覆盖**——全 `tests/` 对 `create_pool`/`_open_pool` 的引用数为 0,现有 PG 用例一律经 `pool=_FakePgPool(...)` 注入、走 `_external_pool=True` 分支从不建池。这正是 `min_size=10` 潜伏至今的原因,也意味着 T3 要建这一路的**第一个**用例。
|
||||
- **提交门是任务边界的实际约束**: `.claude/scripts/hooks/pre-commit-guard.sh` 对每次 `git commit` 阻塞式跑 ruff + radon(圈复杂度 ≥C 即拦)+ 全套件。由此两条硬约束: 不得留红态跨提交(不能把一个行为拆成"改实现"和"改测试"两次);T5 同时改三个降级点,`record_llm_call` 逼近 C 时必须抽私有方法——这不算计划外重构,是提交门的硬要求。
|
||||
- **两条既有承诺挂了检查点,不得被本次改动破坏**: [[design:issue13-schema-mode]] 的"manual 档缺列时裁剪 INSERT 继续写、逐行暴露"(故 `42703` 是失败分类的唯一具名例外)、[[design:issue9-telemetry-ddl-probe]] 的"表存在就绝不发 DDL"(`to_regclass` 探测那段控制流一行不动)。
|
||||
- **保真校验不适用**: 遥测后端无 `reference/` 蓝本(ARCH §7.8 明记"参考仓无先例: 三项目遥测全 SQLite")。
|
||||
|
||||
Reference in New Issue
Block a user