docs: record the telemetry pool semantics and ownership rule
ARCHITECTURE 7.8 gains the pool resource semantics, the two-sentence failure verdict and the two config keys; the ownership rule lands in a new 4.5 because it is a cross-subsystem discipline, not a telemetry convention. CHANGELOG leads with the three items downstream must read first: the 3.12 floor, the connection count going from 10 per client to on demand, and aclose no longer closing injected components.
This commit is contained in:
@@ -325,6 +325,32 @@ flowchart TB
|
||||
6. **开路/全源耗尽**: `CircuitOpenError` / `AllSourcesExhausted` → 按配置 wait(等待恢复,含 stall 判定)或 fail-fast 上抛。
|
||||
7. **任意时刻取消**: `CancelledError` 穿透所有层;in-flight permit 与连接在 finally 释放。
|
||||
|
||||
### 4.5 资源所有权纪律: 谁建的谁关,注入的一律不碰(2026-08-24,issue #15)
|
||||
|
||||
这是**跨子系统的通用纪律**,不是遥测的局部约定。它被写下来的直接原因是: 库对"谁建的、谁负责关"从来没有统一说法,于是同一个根因在三个地方长出三种形态——
|
||||
|
||||
| 形态 | 位置(修复前) | 性质 |
|
||||
|---|---|---|
|
||||
| `GatewayClient.aclose()` 无条件关掉**注入的** telemetry,共享 recorder 被第一个关闭的 client 弄死(`embedding.py`/`ocr.py` 各有一份逐字复制) | `client.py:271-273` | 越权 |
|
||||
| `RedisCache.aclose()` 无条件关掉**注入的** redis 客户端 | `redis_cache.py:43` | 越权 |
|
||||
| `_build_limiter`/`_build_breaker` **自建**的 redis 客户端从来没人关(`aclose` 压根不持有 limiter/breaker 的引用) | `client.py:263-280` | 泄漏 |
|
||||
| 对照组: `RedisLimiter._owns_client` 的纪律**一直是对的** | `limiter.py:185-191, 318-322` | 正确先例 |
|
||||
|
||||
纪律把已有的那个正确先例推广为全库唯一说法,分两层落地:
|
||||
|
||||
| 层 | 所有权归属 | 落法 |
|
||||
|---|---|---|
|
||||
| 组件**内部**自建的连接(limiter/breaker/cache 的 redis 客户端) | 组件自己 | 组件的 `aclose` 自查 `_owns_client`;调用方无条件调用即安全 |
|
||||
| client **自建**的整个组件(transport / recorder / limiter / breaker / cache) | client | 工厂构造后置 `_owns_*` 私有属性,`aclose` 只关自建的;三处复制的 `getattr(..., "aclose")` 鸭子探测收敛为一个内部 helper(同时探测 `aclose`/`close`,SQLite recorder 只有同步 `close()`) |
|
||||
|
||||
三条实现细则各自都是"少写一条就等于纪律不成立":
|
||||
|
||||
1. **默认必须是"不拥有"**。`__init__` 是全量注入路径,经它传入的一切组件一律 `_owns_* = False`,只有三个工厂在真正自建时置 True。默认若反过来,直接构造路径下共享 transport 仍会被第一个 client 关掉。
|
||||
2. **判定一律用 `is None` / `is not None`,不用 `or`**。工厂里 `limiter or _build_limiter(...)` 这种写法在注入一个 falsy 后端时会走自建分支,而所有权标志按 `is None` 判成 False——两者一漂移就等于又造了一个 `aclose` 越权。这是所有权判定能成立的**必要条件**,不是风格偏好。
|
||||
3. **三个 client(chat/embedding/ocr)必须逐一持有 limiter/breaker 引用并各自被测试钉一次**。收敛成 helper 之后仍要三处各钉一次,否则下次有人把逻辑复制回去无人发现;`GatewayClient` 此前把 limiter/breaker 交给 `RetryMW` 后自己不留引用,`aclose` 因此触达不到自建的 redis 客户端,泄漏就是这么来的。
|
||||
|
||||
公共 API 面零变化(`_owns_*` 是私有属性)。**对下游的可见后果**只有一条,且必须显式声明: `aclose()` 不再关闭注入进来的组件,若有下游依赖了"注入后由 client 代关",升级后需自己关。
|
||||
|
||||
---
|
||||
|
||||
## 5. 核心类型
|
||||
@@ -525,6 +551,38 @@ flowchart TB
|
||||
- **单一 helper 铁律**: 遥测调用点收敛为一个内部函数/上下文管理器;Video-Tree 与 GovDoc 各有 4-5 处逐字复制的 `record_llm_call(15 个参数)` 是本条的直接教训。
|
||||
- 成本: `pricing.py` 维护 model → (input 单价, output 单价, **可选** cached_input 单价) 表,遥测时换算 `cost` 字段;查不到价格记 None 并 warning,**不阻塞调用**。缓存读取单价(2026-07-31,issue #3)只在配置了该档且本次有命中时启用,按 `(prompt - cached) × input + cached × cached_input` 分段计价;**未配该档绝不按经验折扣率猜**,退化为全额输入价(P5)。命中数超过输入总数时按总数夹取并 warning,不产生负成本。
|
||||
|
||||
**遥测池的资源语义(2026-08-24,issue #15)**: `PostgresRecorder` 此前 `create_pool(dsn, timeout=10)` 继承 asyncpg 默认的 `min_size=max_size=10`,而 asyncpg 的 `min_size` 语义是"**预连接**"不是"下限"(`pool.py:457` 的 `if self._minsize:`)——建池是一次全有全无的重资源动作: 拿不到 10 条就抛异常。这让遥测成为全库唯一预占资源的组件(httpx transport 与三个 redis 后端全是按需建连),也就成了共享实例余量紧张时**必然第一个倒下**的一环,而它承担的恰恰是最不该悄悄失败的职责。改为 `create_pool(dsn, min_size=0, max_size=<PGW_TELEMETRY_PG_POOL_MAX>, timeout=<预算>, command_timeout=<预算>)`,三条随之确立:
|
||||
|
||||
| 语义 | 内容 |
|
||||
|---|---|
|
||||
| 建池零成本 | `min_size=0` 时 `_initialize` 只造 holder 对象、**一条连接都不连**(实测 0.000s,指向不可达端口也照样成功)。稳态占用由"每 client 常驻 10 条"变为"实际并发,闲时 0";真实 PG 实测: 建 recorder 后 0 → 一次写入后 1 → 20 行并发后 4(= `pool_max`)→ `aclose` 后 0 |
|
||||
| 只暴露 `max_size` | `min_size` **有意不给配置项**: 它唯一的作用是把上面那个脆点装回来,换取的只是首次写入省下 ≈390ms 建连。库没有理由提供一个只会伤人的旋钮(P1+P5)。`max_size` 则必须暴露——继承第三方默认值等于库对自己的资源占用不表态(P4) |
|
||||
| 写入有硬预算 | 整次写入(准备 + acquire + execute)由 `asyncio.timeout(PGW_TELEMETRY_PG_WRITE_TIMEOUT_S)` 包一层,超时按行级丢弃。把"遥测绝不拖垮业务"从"靠各处 timeout 参数凑"升级为一条可陈述、可测试的保证 |
|
||||
|
||||
两处实现纪律,都是"看起来完成了、其实资源还挂着"的形态,必须写下来否则会被改回去: ① **不得用 `async with pool.acquire(...)`**——`Pool.release()` 是 `await asyncio.shield(ch.release(timeout))` 且默认复用 acquire 记录的 `ch._timeout`(asyncpg `pool.py:886-889, 930-937`),外层预算到期时 cancel 在 `execute` 处抛出,异常传播中执行的那个 shielded release **会正常等到完成**,业务路径真实上界变成 ≈ 2 × 预算;故改为显式 `acquire(timeout=剩余预算)` + `finally: release(con, timeout=1s)`,释放超时即 `con.terminate()`,承诺精确化为"主写入尝试 ≤ 预算,释放路径独立有界"。② **`aclose()` 必须有界且终局**: `Pool.close()` 会 `await` 每个 holder 的 `wait_until_released()`,in-flight 未释放时无限等、60 秒只发一条 warning(`pool.py:939-948, 961-972`),故走 `asyncio.wait_for` + 超时 `terminate()`;同时置 `_closed`,此后写入短路且**不复活**——原实现关完池后下一次写入会拿 DSN 悄悄自建一个新池,注入方以为自己管着全部连接、实际早已不是(issue #15 实施期发现,是下面所有权根因的又一处表现)。
|
||||
|
||||
**遥测失败的三分判据(2026-08-24,issue #15)**: 判死判据此前挂在"**哪一步**失败"(`_open_pool` 失败即永久判死),而那一步里同时藏着两类性质完全不同的失败——DSN 非法(进程内不可能改变)与 `too many clients` / 网络抖动(外部状态,随时可能好)。判据改挂"失败是**什么性质**",两句话说完:
|
||||
|
||||
1. **致命 = 失败原因完全在进程内部且不可变**;其余一切失败都可能被外部修好,故一律带冷却重试。
|
||||
2. **行级 vs 环境级看"失败与这一行的数据有没有关系"**: 只与本行数据有关(换一行可能成功)= 行级;与数据无关、每一行都会同样失败 = 环境级。
|
||||
|
||||
| 档 | 覆盖(按 SQLSTATE 分类而非异常类白名单——SQLSTATE 是 PG 标准,不随 asyncpg 版本漂移) | 处置 |
|
||||
|---|---|---|
|
||||
| 配置级致命 | `ClientConfigurationError`(DSN 不可解析);`create_pool` 抛的 `ValueError`/`TypeError` | 永久 no-op + 一条 **error**(人配错了,不是 warning) |
|
||||
| 环境级不可用 | SQLSTATE 类 `08`/`53`(含 53300 too many connections)/`57`/`28`/`3D`,具体码 `42501`(无权限)/`42P01`(表不存在);`OSError`/`ConnectionError`/`TimeoutError`/其余 `InterfaceError`;表确定不存在且建不出来 | **冷却降级**(内部常量 60s,不给配置项——无部署差异理由),到期放行**一次**重新准备,成功即恢复 |
|
||||
| 行级拒绝 | 其余 `PostgresError`(`22`/`23` 等数据与约束类),以及**具名例外 `42703`(缺列)** | 逐条 warning 丢弃,不降级,接入节流复述 |
|
||||
|
||||
三点必须一起记住,否则后来人会把判据改回去: ① **致命档窄到只剩 DSN 一类是有意的**——认证失败、库不存在、表建不出来一律归环境级,因为 DBA 改完密码/建完表就该自动恢复,而永久失能是最坏结局,只留给"重试在任何时刻都不可能成功"的情形;②**`42703` 是唯一具名例外**,按第 2 句它本该是环境级(缺列时每行都失败),归行级是因为 issue #13 定下了优先级更高的承诺——manual 档缺列时按现有列裁剪 `INSERT` 继续写、缺列以逐行 warning 暴露,即"部分列写进去了"这件事本身有价值,不该被冷却掉;新增例外必须同款论证。③ **认不出的失败一律归最轻档(行级)**,这个保守缺省在建池路径上是安全的,理由是 `min_size=0` 让建池不触库(实测 0.000s),"下次调用重试建池"本身**零成本**——原实现注释担心的"每次重试内联吞一次 connect 超时"在新语义下不再成立。
|
||||
|
||||
**降级的可见性与可编程性(2026-08-24,issue #15)**: 铁律里"遥测后端挂 → 静默降级"的"静默"指的是**不向调用方冒泡**,不是"没有日志、没有状态"。此前它被实现成了后者——全程只有一条 warning,长跑进程里等同于消失(issue 是人工比对"日志里的完成里程碑条数 vs `llm_calls` 行数"才发现的,期间 19 次调用一行未落);SQLite 侧更糟,初始化失败后写入直接 `return`,连 warning 都没有。"遥测必录"铁律的实质要求是: **库做不到必录时,必须持续、可编程地让下游知道**。落法是 `telemetry/status.py` 的 `TelemetryStatusTracker`——两个 recorder 共用、不含任何后端知识(只接受"降级了/恢复了/丢了一行"三个事实),进入与恢复各一条日志,降级期间按行数(100 行)与时间(300s)双阈值节流复述,`snapshot()` 给只读 `TelemetryStatus`(`degraded`/`fatal`/`reason`/`degraded_for_s`/`dropped_rows`/`retry_after_s`),经三个 client 的 `telemetry_status` 属性出口。三条设计约束:
|
||||
|
||||
- **不叫 `health`**: 该词在 `ports.py` 已被 `OcrTransport.check_health`(源探活)与 `SourceSelector.health(source_name) -> float`(成功率 EWMA)占用两次,库内 `health` 一律指"源的健康度";这里描述的是"这个 recorder 现在能不能写、为什么不能、丢了多少",是状态不是评分(P2)。
|
||||
- **不并入 `TelemetryRecorder` 主 Protocol**,新起**独立**端口 `TelemetryStatusProvider`: 前者是 `@runtime_checkable`,而 runtime 检查按属性存在性做——加一个成员会让所有只实现 `record_llm_call` 的对象**当场不再是** `TelemetryRecorder`,库内与下游的同款 `isinstance` 断言升级即断。client 侧取值经**一处** `isinstance` 判定,不重演 `aclose` 那种三处复制的鸭子类型。
|
||||
- **`TelemetryStatus` 进顶层 `__all__`**(与 `SourceStats` 不同): 后者是端口内部快照、下游不消费,而本类型是 `client.telemetry_status` 的返回类型,下游要拿它做类型标注与对账——"顶层导出即公共 API 面"的约定要求它出现在那里。端口 `TelemetryStatusProvider` 则不导出(库外无实现者,导出即多一份永久承诺)。
|
||||
- **SQLite 侧只做可见性**,不做 lazy 化与冷却重连: 它的失败模式(本地目录不可写、文件损坏)在装配期就暴露给下游,不是"跑到一半悄悄断",永久降级在那里语义基本正确。这个不对称是已知且有理由的;tracker 与快照两侧共用,将来要对称时接口已就位。
|
||||
|
||||
**资源所有权在遥测侧的落点**: 通用纪律见 §4.5。对遥测的直接后果是 §7.7 R5 那条"共享必须显式注入"第一次真正可用——`PostgresRecorder(dsn, pool=<外部池>)` 与"多个 client 注入同一个 recorder"都不再被第一个 `aclose()` 弄死,issue #15 提的"共享池"方向由此以显式注入形态自然成立,不需要任何隐式全局注册表(那会违反"纯 asyncio 中立: 无全局状态、无模块级单例")。
|
||||
|
||||
### 7.9 结构化输出阶梯(D14)
|
||||
|
||||
| 级 | 内容 | 成本 |
|
||||
@@ -587,6 +645,7 @@ src/polygateway/
|
||||
- **`{SCOPE}__CIRCUIT_OPEN=fail_fast|wait`(2026-08-19,issue #14)**: 熔断全拒时的处置,与 `{SCOPE}__QUOTA_FULL` 同形同族(上一条"后端选择即配置"里记的 `PGW_QUOTA_FULL` 是 M1 定稿前的暂拟名,实际落地为 scope 键 `{SCOPE}__QUOTA_FULL`)。缺省 **fail_fast** = 存量下游的控制流逐字不变;**单源 scope 应显式配 `wait`**。两键值域相同但语义不同故分列: 配额满是"排队等自己的份额"(必然轮到),熔断开路是"等这个源恢复"(未必恢复),调用方可能想要"配额满就等、源坏了就立刻失败"。落到 `GatewaySettings.circuit_open`(无默认值,与既有全部字段一致),校验收敛在唯一消费者 `SourceAdmission` 一处——三个客户端构造函数此前各带一份 `quota_full` 校验,再加一键就是八处复制。
|
||||
- **`PGW_TELEMETRY_SCHEMA_MODE=auto|manual`(2026-08-19,issue #13,D15)**: 可选键、**三态**——不设 = 按后端派生(sqlite→auto、postgres→manual),显式设置则两侧都可覆盖。派生只发生在 config 层一处,落到 `GatewaySettings.telemetry_auto_migrate`(无默认值,与既有全部字段一致;`telemetry_backend=none` 时无人消费,归一为 `False`),recorder 的 `auto_migrate` 是 keyword-only **必填**参数——关键行为参数不给默认值(P4),缺省规则也就不会与类签名漂移。
|
||||
- **`PGW_TELEMETRY_TEXT_CAP`(2026-08-19,issue #12)**: 可选正整数键、**二态**——不设 = 不截断(缺省)。与相邻的 `SCHEMA_MODE` 不同,这里"未设"本身就是最终答案,没有需要按后端派生的第二种缺省。落到 `GatewaySettings.telemetry_text_cap: int | None`(同样无默认值),`TelemetryEmitter.text_cap` 是 keyword-only 必填参数。值域(`> 0`)在 settings 与 emitter **两处**校验: 前者只管 env 一条路,而"构造函数全量注入"是库承诺的另一条公共装配路,`text_cap=0` 会让每条正文只剩一个省略标记(P5 不得静默)。
|
||||
- **`PGW_TELEMETRY_PG_POOL_MAX` / `PGW_TELEMETRY_PG_WRITE_TIMEOUT_S`(2026-08-24,issue #15)**: 两个可选键,**缺省 4 与 5.0**——与相邻三个遥测键不同,这两个有真正的默认值而不是"无默认值的必填字段",因为它们回答的是"库该占多少资源",而库对此**必须有一个可陈述的表态**(不表态就等于继承第三方默认值,那正是 issue 的病根,见 §7.8);缺省写在 config 一处,`PostgresRecorder` 的 `pool_max`/`write_timeout_s` 是 keyword-only **必填**参数(与 `auto_migrate` 同一纪律: 缺省规则不与类签名漂移)。值域校验(`pool_max >= 1`、`write_timeout_s > 0`)落 `GatewaySettings._validate_telemetry`,与 `telemetry_text_cap` 同一先例覆盖**三条装配路**(直接构造 / `dataclasses.replace` / env),报错文本同时点字段名与 env 键名。两键都带 `PG` 前缀与 `PGW_TELEMETRY_PG_DSN` 对齐: SQLite 侧的等价物(`busy_timeout=5000`)本次不动,这个不对称是已知且有理由的(§7.8 末)。**冷却期 60s 有意不给键**——无部署差异理由(P1 YAGNI)。`pool_max` 的调参口径必须按实测折算而非按 `pool_max / RTT` 估算: 跨内网 RTT ≈ 123ms 的实验室 PG 上 `pool_max=4` 实测约 **15.6 行/秒**(50 行并发批 3.2s),一次 `INSERT` 的实际往返比一次 `SELECT 1` 重一倍。
|
||||
|
||||
---
|
||||
|
||||
|
||||
@@ -2,7 +2,7 @@
|
||||
|
||||
- **issue**: #15(共享 PostgreSQL 实例,`max_connections=100`,多 worker × 多 scope 部署)
|
||||
- **核查基准**: HEAD 1.2.4。issue 按 1.1.2 运行环境提交并已自行复核 1.2.4,本文逐条重核**全部成立**: `postgres.py:100`(建池不传 min/max)、`postgres.py:106`(建池失败即永久判死)、`postgres.py:246-251`(`aclose` 不清 `_failed`)、`client.py:405-420`(每个 client 各 new 一个 recorder)。`min_size`/`max_size` 在整个包内**一次都没出现过**。
|
||||
- **状态**: 人类已确认方案与全部四组改动 + 缺省值(2026-08-24);**Codex 已审,7 条全部处置完毕(§9)**;待人类审 → 实施
|
||||
- **状态**: **已实施**(2026-08-24,分支 `feat/issue-15-telemetry-pool-lifecycle`,T0–T7 见实现计划末尾的提交表)。人类已确认方案与全部四组改动 + 缺省值;**Codex 已审,7 条全部处置完毕(§9)**;实施期的三处修订以 §10 标注
|
||||
- **实测环境**: asyncpg 0.31.0;真实实验室 PG(`polygateway` 专用库,跨内网 RTT ≈ 123ms)
|
||||
|
||||
## 1. 问题的真实形状
|
||||
@@ -23,6 +23,8 @@ asyncpg `pool.py:457` 是 `if self._minsize:` ——为 0 时 `_initialize` 只
|
||||
|
||||
**关键推论**: `min_size=0` 不只是"调小",它把建池从一次全有全无的重资源动作变成**零成本、不触库**的动作。这一步走出去,后面三层的性质全变。
|
||||
|
||||
**实施后在同一台真实实验室 PG 上的复测(T6,按唯一 `application_name` 过滤 `pg_stat_activity`)**,是全套证据里最直观的一条: 修复前建完 recorder 即 **10** 条连接;修复后 **0**(建 recorder)→ **1**(一次写入)→ **4**(20 行并发,恰为 `pool_max`)→ **0**(`aclose` 后)。四个数字逐一对应上表的四行推论。
|
||||
|
||||
### 1.2 缺陷一: 库对自己的资源占用从未表态——而这是全库唯一一处
|
||||
|
||||
`create_pool(self._dsn, timeout=10)` 继承第三方默认值(P4/P5: 默认参数掩盖关键逻辑)。横向扫过库内每一处外部资源:
|
||||
@@ -61,6 +63,8 @@ issue 建议"让指向同一 DSN 的多个 recorder 共享一个池"。这条路
|
||||
| `_build_limiter`/`_build_breaker` **自建**的 redis 客户端从来没人关(`aclose` 压根不碰 limiter/breaker) | `client.py:263-280` | **泄漏** |
|
||||
| 对照组: `RedisLimiter._owns_client` 纪律**是对的** | `limiter.py:185-191, 318-322` | 正确先例 |
|
||||
|
||||
**实施期挖出的第四个现象(T4,本设计原稿未预见)**: 注入外部池时,`aclose()` 之后的下一次写入会拿 DSN **偷偷自建一个池**——注入方以为自己管着全部连接,实际早已不是。它与上表三条同一根因(库不区分"这个资源是谁的"),只是表现在**关闭之后**而非关闭当时,故原稿按"谁关谁的"扫一遍时没看见。修法归入 §3.2 第 4 点的"关了就是关了": 置 `_closed` 后写入短路且不复活。
|
||||
|
||||
三个现象一个根因: **库对"谁建的、谁负责关"没有统一纪律**。ARCH §7.7 R5 规定"共享必须显式注入",但显式注入这条正道今天是坏的,下游只能退回"每 client 各占一份"——缺陷一的放大器由此长在架构里,而不是长在某个默认值里。
|
||||
|
||||
## 2. 备选方案与否决理由
|
||||
@@ -111,7 +115,7 @@ issue 建议"让指向同一 DSN 的多个 recorder 共享一个池"。这条路
|
||||
4. **`aclose()` 的语义钉死为"关了就是关了"**: 置 `_closed`,此后写入短路且**不复活**。今天"关完还能自己重建池"的灰色状态取消。issue 提的"`aclose` 不清 `_failed`"由冷却机制解决,不由 `aclose` 解决——恢复是运行时行为,不是关闭动作的副作用。
|
||||
**关闭动作本身也必须有界(Codex 审查,已核实)**: `Pool.close()` 会 `await` 每个 holder 的 `wait_until_released()`,in-flight 未释放时**无限等**,60 秒只发一条 warning(`pool.py:939-948, 961-972`);asyncpg 自己的 docstring 就写着"advisable to use `asyncio.wait_for` to set a timeout"。故 `aclose()` 走 `asyncio.wait_for(pool.close(), ...)`,超时后 `pool.terminate()`,外部取消照常穿透——否则"遥测不得拖垮业务"在收尾路径上开了个口子。
|
||||
|
||||
分类函数是全库唯一一处 PG 失败分类,作 `postgres.py` 模块级私有函数(与 recorder 同文件、只服务 PG;不新起文件避免碎片化)。
|
||||
分类函数是全库唯一一处 PG 失败分类,作 `postgres.py` 模块级私有函数(与 recorder 同文件、只服务 PG;不新起文件避免碎片化)。**认不出的失败归最轻档(行级)**是它的保守缺省,而这个缺省在**建池路径**上安全的理由比"最轻档代价最小"更强(实施期核实,T5): `min_size=0` 让建池不触库(实测 0.000s),所以"归行级 = 下次调用再重试一次建池"本身**零成本**——`postgres.py:104-105` 那条注释担心的"每次重试内联吞一次 connect 超时"是 `min_size=10` 语义下的顾虑,在新语义下**不成立**。这是 §1.1 那个关键推论的又一处红利: 地基一换,原本需要小心处理的保守缺省变成了白拿。
|
||||
|
||||
### 3.3 C 组 · 降级可见 + 可编程
|
||||
|
||||
@@ -146,6 +150,7 @@ issue 建议"让指向同一 DSN 的多个 recorder 共享一个池"。这条路
|
||||
- **默认必须是"不拥有"**: `__init__` 是全量注入路径(`client.py:126-149`),经它传入的一切组件一律视为**外部所有**(`_owns_* = False`),只有三个工厂在 `or _build_*` / `if telemetry is not None else _build_telemetry` 真正自建时才置 True。原稿只写了"工厂置位"没写死这条默认,Codex 据此指出直接构造路径下共享 transport 仍会被第一个 client 关掉——那是实现走偏的后果,但默认值本就该在设计里定死,故补。
|
||||
- 三处复制的 `getattr(..., "aclose")` 收敛为一个内部 helper;所有权修正必须三处一致,复制就是下一个 bug 的种子。
|
||||
- `aclose` 补关 limiter/breaker——修掉现存泄漏。这需要**三个 client 都新持引用**: 今天 `GatewayClient.__init__` 把 limiter/breaker 交给 `RetryMW` 后自己不留引用(`client.py:133-134`),embedding/ocr 同样(`embedding.py:497-498`、`ocr.py:510-511` 自建、`embedding.py:452-461`、`ocr.py:462-467` 的 `aclose` 触达不到)。内存后端无 `aclose`,helper 探测后跳过。
|
||||
- **判定一律用 `is None` / `is not None`,不得用 `or`**(实施期补,T1): 工厂里 `limiter or _build_limiter(...)` 这种写法在注入一个 falsy 后端时会走自建分支,而所有权标志按 `is None` 判成 False——两者一漂移就等于又造了一个 `aclose` 越权。这是所有权判定能成立的**必要条件**,不是风格偏好,故写进设计而非留在代码里。
|
||||
- **零公共 API 面变化**: `_owns_*` 是私有属性,由工厂置位。
|
||||
- 有了 D 组,issue 的"共享池"方向以**显式注入**形态自然成立(`PostgresRecorder(dsn, pool=...)` 已支持且不关外部池),无需任何隐式全局。
|
||||
|
||||
@@ -153,7 +158,7 @@ issue 建议"让指向同一 DSN 的多个 recorder 共享一个池"。这条路
|
||||
|
||||
| 键 | 字段 | 缺省 | 依据 |
|
||||
|---|---|---|---|
|
||||
| `PGW_TELEMETRY_PG_POOL_MAX` | `telemetry_pg_pool_max: int` | **4** | 稳态吞吐 ≈ `max_size / RTT` = 4/0.123 ≈ **32 行/秒**,覆盖单 client 数十并发;闲时占 0,不构成常驻负担。issue 现场 4 client × 4 = 峰值 16、稳态趋近 0(今天是 40 条常驻) |
|
||||
| `PGW_TELEMETRY_PG_POOL_MAX` | `telemetry_pg_pool_max: int` | **4** | 稳态吞吐**实测约 15.6 行/秒**(见 §6 的口径更正;原稿按 `max_size / RTT` 估的 32 行/秒偏乐观一倍),覆盖单 client 十余并发;闲时占 0,不构成常驻负担。issue 现场 4 client × 4 = 峰值 16、稳态趋近 0(今天是 40 条常驻) |
|
||||
| `PGW_TELEMETRY_PG_WRITE_TIMEOUT_S` | `telemetry_pg_write_timeout_s: float` | **5.0** | 实测稳态 123ms、首次含建连 513ms;5s 宽松且**有界**。同时用作 connect / acquire / 整次写入硬上界 |
|
||||
| (无键) | 冷却期 | 60s,**内部常量** | 无部署差异理由(P1 YAGNI) |
|
||||
|
||||
@@ -202,10 +207,10 @@ issue 建议"让指向同一 DSN 的多个 recorder 共享一个池"。这条路
|
||||
| 维度 | 结论 |
|
||||
|---|---|
|
||||
| 首次写入延迟 | `min_size=0` 把 ≈390ms 建连从"装配期"挪到"首次写入"。稳态无差异(实测 123ms);空闲超 `max_inactive_connection_lifetime`(asyncpg 缺省 300s,不暴露)后再付一次。相对一次秒级 LLM 调用可忽略 |
|
||||
| 突发排队(**热池稳态**) | 业务并发 > `pool_max` 时遥测写入排队。64 行同时到达、32 行/秒 → 最坏约 2s,在 5s 预算内;超出即丢行(铁律"丢一条 < 拖垮调用") |
|
||||
| 突发排队(**热池稳态**) | 业务并发 > `pool_max` 时遥测写入排队。按下一格更正后的实测口径(15.6 行/秒): 50 行同时到达 → 实测 3.2s,在 5s 预算内但**余量只剩约 1.5 倍**(原稿按 32 行/秒估算时以为余量有 3 倍);超出即丢行(铁律"丢一条 < 拖垮调用") |
|
||||
| 突发排队(**冷启动/空闲后**) | 上一格的算术只在"schema 已就绪且连接已热"时成立。空闲超回收期后连接归 0,第一波要重新建连(实测 ≈390ms),且首次准备被 `_init_lock`(`postgres.py:75, 83-91`)串行保护——冷启动的最坏延迟不是 `64 / 32 ≈ 2s`。Codex 审查指出原稿这段易被读成两种情形通用,故拆开写。冷启动上界仍由硬预算封顶,超出即丢行 |
|
||||
| Python 版本 | **本条取舍已消解**(人类决策,2026-08-24): 最低版本提到 **3.12**(`requires-python = ">=3.12"`、ruff `target-version = "py312"`),3.11.0/3.11.1 的 `uncancel` 缺陷不再在支持范围内,`asyncio.timeout` 可直接用,不必退回 `wait_for`。代价见 §7 |
|
||||
| `pool_max` 的调参口径 | 须进文档: 期望吞吐 ≈ `pool_max / RTT`。共享一个 recorder 给多 client 时并发汇聚,应相应放大 |
|
||||
| `pool_max` 的调参口径(**实施期更正,T3 实测**) | 原稿的 `期望吞吐 ≈ pool_max / RTT`(4/0.123 ≈ 32 行/秒)**偏乐观一倍**: T3 实测 50 行并发批耗时 **3.2s**,即约 **15.6 行/秒**、每条连接约 4 行/秒——一次 `INSERT` 的实际往返比一次 `SELECT 1`(RTT 的测法)重。取舍方向不变(超预算丢行 < 拖垮业务),但 `.env.example` 与 README 的调参口径**必须写实测数字**,否则下游按错公式放大,以为 `pool_max=8` 能到 64 行/秒(实为约 31)。共享一个 recorder 给多 client 时并发在此汇聚,应按 client 数相应放大 |
|
||||
| 冷却期的丢数 | 降级 60s 期间的行**确实丢了**,只是可见、可计数、且到期自动恢复。这是"遥测降级不得拖垮业务"的既有方向(ARCH 降级方向铁律),本设计不改方向,只改**可恢复性与可见性** |
|
||||
| `42703` 缺列的持续逐行重试 | 缺列时每次调用付一次 acquire+execute(≈123ms 内联)且逐行 warning,不进冷却。**这是判据的唯一具名例外**(§3.2 第 2 点),由 issue #13 的"缺列须逐行暴露"承诺定死;代价由节流复述抵消。`42501`/`42P01` 原稿同归此格,经 Codex 审查已改判环境级 |
|
||||
| 快照计数的线程安全 | `dropped_rows` 是单事件循环内的 int 自增。库不承诺跨线程共享同一 recorder("纯 asyncio 中立"),最坏是计数不准,不会崩 |
|
||||
@@ -260,3 +265,13 @@ CHANGELOG 有三处需"请先读这一条"待遇:
|
||||
**本轮自查另补两条 Codex 未发现的**: ① `health` 一词在 `ports.py` 已被 `check_health`(`:90`)与 `health(source_name) -> float`(`:234`)占用两次,故快照改名 `TelemetryStatus`(§3.3);② `asyncio.timeout` 是 3.11 新增而 `requires-python = ">=3.11"`,3.11.0/3.11.1 的 `uncancel` 有已知缺陷,实施时须在"抬最低版本"与"改用 `wait_for`"之间选一(§6)。
|
||||
|
||||
Codex 的取消穿透实测与本会话结论一致(外部 `task.cancel()` 在 `asyncio.timeout` 内冒出的是 `CancelledError` 而非 `TimeoutError`),两处独立验证互为佐证。
|
||||
|
||||
## 10. 实施期修订(2026-08-24,T0–T7 执行中发现)
|
||||
|
||||
设计经人类审后实施,过程中三处需要回改设计本身——都不是措辞问题,而是"原稿的事实基础不够"。逐条落回正文而非只记在这里,以免后来人读正文时踩同一个坑。
|
||||
|
||||
| # | 修订 | 落点 |
|
||||
|---|---|---|
|
||||
| 1 | **吞吐算术偏乐观一倍**。原稿按 `pool_max / RTT` 估 32 行/秒,T3 实测 50 行并发批 3.2s(≈15.6 行/秒)——`INSERT` 的实际往返比测 RTT 用的 `SELECT 1` 重。方向不变,但下游调参必须拿实测数字 | §3.5 表、§6 两格 |
|
||||
| 2 | **原稿未预见的一处真 bug**: 注入外部池时 `aclose()` 之后的下一次写入会拿 DSN 偷偷自建一个池。与 §1.5 三条同根因,只是表现在关闭之后,T4 修掉 | §1.5 |
|
||||
| 3 | **两条论证被补强**: ①"认不出的失败归行级"这个保守缺省在建池路径上安全,理由是 `min_size=0` 让重试建池零成本(T5);②所有权判定必须用 `is not None` 而非 `or`,否则注入 falsy 后端时自建分支与所有权标志漂移(T1) | §3.2 末、§3.4 |
|
||||
|
||||
@@ -3,6 +3,7 @@
|
||||
- **设计**: `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 步的事)
|
||||
- **状态**: **已实施**(2026-08-24)。T0–T7 全部提交完成,提交表见文末;合并前的三道门(`pytest -m slow`、独立 verifier、整分支审查)见「完成判据」)
|
||||
|
||||
## 目标
|
||||
|
||||
@@ -131,9 +132,9 @@ def _classify_failure(exc: BaseException) -> str: ...
|
||||
|
||||
### 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 可达性都会改变它
|
||||
- [x] 从 `main` 建分支 `feat/issue-15-telemetry-pool-lifecycle`
|
||||
- [x] 把工作区现有改动分两次提交: ① `chore: 最低 Python 提到 3.12 并改用 PEP 695 泛型语法`(`pyproject.toml`/`README.md`/`CLAUDE.md`/`client.py`/`streaming.py`);② `docs: issue #15 设计文档与 wiki 登记`(`research-wiki/`)
|
||||
- [x] 记录基线用例计数(执行时实测;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` → 分支名正确。
|
||||
|
||||
@@ -151,6 +152,8 @@ def _classify_failure(exc: BaseException) -> str: ...
|
||||
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` 判定。
|
||||
|
||||
**执行留痕(T1)**: 工厂里既有的 `limiter or _build_limiter(...)` 一律改成了 `is not None` 判定。理由是注入一个 **falsy** 后端时 `or` 会走自建分支,而所有权标志按 `is None` 判成 False——两者一漂移就等于又造了一个 `aclose` 越权。这不是风格偏好,是所有权判定能成立的**必要条件**,已回写设计 §3.4。
|
||||
|
||||
**同时修掉的现存泄漏**: `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`,探测后跳过。
|
||||
@@ -164,7 +167,7 @@ def _classify_failure(exc: BaseException) -> str: ...
|
||||
|
||||
**验证**: `pytest tests/unit/test_client.py tests/unit/test_embedding.py tests/unit/test_ocr_client.py -q` → PASS;`make check` 绿;全套件绿。
|
||||
|
||||
- [ ] 提交: `fix: 统一资源所有权纪律(谁建的谁关),修 aclose 越权与 redis 客户端泄漏`
|
||||
- [x] 提交: `fix: 统一资源所有权纪律(谁建的谁关),修 aclose 越权与 redis 客户端泄漏`
|
||||
|
||||
---
|
||||
|
||||
@@ -195,7 +198,7 @@ def _classify_failure(exc: BaseException) -> str: ...
|
||||
|
||||
**验证**: `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 + 只读快照 + 节流日志)`
|
||||
- [x] 提交: `feat: 遥测降级升格为一等状态(共用 tracker + 只读快照 + 节流日志)`
|
||||
|
||||
---
|
||||
|
||||
@@ -228,7 +231,7 @@ def _classify_failure(exc: BaseException) -> str: ...
|
||||
|
||||
**验证**: `pytest tests/unit/test_config.py tests/unit/test_telemetry.py -q` → PASS;`make check` 绿;全套件绿。
|
||||
|
||||
- [ ] 提交: `feat: 遥测池显式声明资源占用(min_size=0/max_size 可配)并给写入硬预算`
|
||||
- [x] 提交: `feat: 遥测池显式声明资源占用(min_size=0/max_size 可配)并给写入硬预算`
|
||||
|
||||
---
|
||||
|
||||
@@ -249,7 +252,7 @@ def _classify_failure(exc: BaseException) -> str: ...
|
||||
|
||||
**验证**: `pytest tests/unit/test_telemetry.py -q` → PASS;全套件绿。
|
||||
|
||||
- [ ] 提交: `fix: 遥测池关闭有界化(wait_for + terminate),关闭后不再复活`
|
||||
- [x] 提交: `fix: 遥测池关闭有界化(wait_for + terminate),关闭后不再复活`
|
||||
|
||||
---
|
||||
|
||||
@@ -280,7 +283,7 @@ def _classify_failure(exc: BaseException) -> str: ...
|
||||
|
||||
**验证**: `pytest tests/unit/test_telemetry.py -q` → PASS;`radon cc src/polygateway/telemetry/postgres.py -n C -s` → 无输出;全套件绿。
|
||||
|
||||
- [ ] 提交: `fix: 遥测失败按性质三分,永久判死收窄到 DSN 不可解析,其余带冷却自愈`
|
||||
- [x] 提交: `fix: 遥测失败按性质三分,永久判死收窄到 DSN 不可解析,其余带冷却自愈`
|
||||
|
||||
---
|
||||
|
||||
@@ -305,7 +308,7 @@ def _classify_failure(exc: BaseException) -> str: ...
|
||||
|
||||
**验证**: `pytest tests/integration/test_postgres_telemetry.py -q` → PASS(或无 DSN 时全 skip);全套件绿。
|
||||
|
||||
- [ ] 提交: `test: 真实 PG 验证遥测池不预连接与降级自愈`
|
||||
- [x] 提交: `test: 真实 PG 验证遥测池不预连接与降级自愈`
|
||||
|
||||
---
|
||||
|
||||
@@ -321,13 +324,13 @@ def _classify_failure(exc: BaseException) -> str: ...
|
||||
|
||||
**验证**: `make check` 绿;人工通读 `.env.example` 两键注释,确认调参口径可执行。
|
||||
|
||||
- [ ] 提交: `docs: 遥测池资源语义、失败判据与所有权纪律成文`
|
||||
- [x] 提交: `docs: 遥测池资源语义、失败判据与所有权纪律成文`
|
||||
|
||||
---
|
||||
|
||||
## 完成判据(合并前)
|
||||
|
||||
- [ ] T0-T7 全部提交完成,每次提交都过了提交门(ruff + radon + 全套件)
|
||||
- [x] T0-T7 全部提交完成,每次提交都过了提交门(ruff + radon + 全套件)
|
||||
- [ ] `pytest -m slow` 单独跑过一次(发布清单第 4 步;本次改动触及遥测写入路径,e2e 与 Redis 时间语义变体必须实测)
|
||||
- [ ] 派**全新上下文**的 verifier subagent 独立验证(`verification-before-completion`,里程碑级/合并前 MANDATORY)
|
||||
- [ ] 整分支审查(`requesting-code-review`,合并前 MANDATORY)
|
||||
@@ -346,3 +349,28 @@ def _classify_failure(exc: BaseException) -> str: ...
|
||||
|
||||
Codex 给的 `_failed` 分布数字(源码 5 处 / 测试断言 8 处)与本地实测(源码 7 / unit 6 / integration 6)不一致,以实测为准——它漏了 docstring 里那两处,而那两处恰恰是**必须改**的(留着就是指向已删字段的说明)。
|
||||
|
||||
## 实际提交(2026-08-24,分支 `feat/issue-15-telemetry-pool-lifecycle`)
|
||||
|
||||
| 任务 | hash | message 首行 |
|
||||
|---|---|---|
|
||||
| T0 ① | `157a27f` | `chore: require python 3.12 and adopt PEP 695 type parameters` |
|
||||
| T0 ② | `e7caa50` | `docs: plan the telemetry pool lifecycle rework for issue 15` |
|
||||
| T1 | `e69ca4c` | `fix: make every client close what it built and nothing else` |
|
||||
| T2 | `f958138` | `feat: make telemetry degradation a first-class state` |
|
||||
| T3 | `84c2cc1` | `feat: make the telemetry pool declare what it costs` |
|
||||
| T4 | `bc071c6` | `fix: make closing the telemetry pool bounded and final` |
|
||||
| T5 | `eef2fdc` | `fix: judge telemetry failures by nature, not by step` |
|
||||
| T6 | `bfeda5b` | `test: prove on real PG that the pool never preconnects` |
|
||||
| T7 ⓪ | `69a5b5f` | `test: pin the cooldown assertion to a fake clock`(T5 留下的一处间歇红: 快照里的 `retry_after_s` 是时间差,却用真实时钟断言 60.0) |
|
||||
| T7 ① | `7834d75` | `feat: export TelemetryStatus from the package root` |
|
||||
| T7 ② | 本文件所在的这次提交 | `docs: record the telemetry pool semantics and ownership rule` |
|
||||
|
||||
T7 分两次提交是因为它含一处**公共 API 面**改动(`TelemetryStatus` 进顶层 `__all__`,决策见下),与纯文档的回滚粒度不同。
|
||||
|
||||
**T7 执行期追加的决策与发现**(计划原稿只列了四项文档任务):
|
||||
|
||||
| # | 内容 | 落点 |
|
||||
|---|---|---|
|
||||
| 1 | `TelemetryStatus` 进 `polygateway.__all__`。issue #15 的核心诉求之一是下游能**编程对账**,而 `client.telemetry_status` 的返回类型若不能从顶层 import,下游做类型标注就得深入 `polygateway.types`——与"顶层导出即公共 API 面"的约定冲突。T2 参照的 `SourceStats` 先例**不适用**: 那是端口内部快照、下游不消费。端口 `TelemetryStatusProvider` 仍不导出 | `__init__.py`、`tests/unit/test_package.py`、ARCH §7.8 |
|
||||
| 2 | 吞吐算术更正为实测值(15.6 行/秒),`.env.example` / README 的调参口径按实测写 | 设计 §3.5/§6/§10 |
|
||||
| 3 | "重试建池已零成本"这条红利与"关闭后偷偷复活"这个 bug 分别补进设计 §3.2 / §1.5 | 设计 §10 |
|
||||
|
||||
Reference in New Issue
Block a user