From cc01d5ed6283419febb773fc5bd5358f3d95185e Mon Sep 17 00:00:00 2001 From: iomgaa Date: Thu, 16 Jul 2026 08:53:50 -0400 Subject: [PATCH] fix: address Codex review of telemetry fix (mkdir degrade, close, degrade tests) --- adapters/telemetry.py | 23 +++++++-- ...-07-16-telemetry-concurrency-fix-design.md | 8 ++++ tests/unit/test_telemetry.py | 47 ++++++++++++++++--- 3 files changed, 69 insertions(+), 9 deletions(-) diff --git a/adapters/telemetry.py b/adapters/telemetry.py index aea6d91..8e49d95 100644 --- a/adapters/telemetry.py +++ b/adapters/telemetry.py @@ -1,7 +1,11 @@ """SQLite 遥测记录器 — TelemetryRecorder Protocol 的生产实现。 -通过 asyncio.to_thread 将 SQLite 同步写入桥接到异步接口, -确保事件循环不被阻塞。表在首次写入时懒初始化。 +通过 asyncio.to_thread 将 SQLite 同步写入桥接到异步接口,确保事件循环不被阻塞。 +构造时建单持久连接 + 建表(对齐 app/harness/log.py:HarnessLog 的并发写模式), +写入经进程内 threading.Lock 串行化,消除多连接并发写的 database is locked。 + +零丢失保证范围 = 单进程、单 recorder 实例(当前 main.py / video_split_cli 均单实例 +注入)。同进程多个 recorder 指向同一 db 会退回跨连接竞争——本实现不支持该场景。 """ from __future__ import annotations @@ -68,6 +72,8 @@ class SQLiteTelemetryRecorder: self._db_path = db_path self._lock = threading.Lock() self._conn: sqlite3.Connection | None = None + # mkdir / connect / PRAGMA / 建表统一纳入降级边界:任一失败(OSError 含 + # PermissionError、sqlite3.Error)都降级为 self._conn=None,绝不冒泡拖垮初始化。 try: Path(db_path).parent.mkdir(parents=True, exist_ok=True) conn = sqlite3.connect(str(db_path), check_same_thread=False) @@ -76,10 +82,21 @@ class SQLiteTelemetryRecorder: conn.execute(self._CREATE_TABLE_SQL) conn.commit() self._conn = conn - except sqlite3.Error as exc: + except (OSError, sqlite3.Error) as exc: logger.warning("遥测连接初始化失败(已降级,后续写入丢弃): {}", exc) self._conn = None + def close(self) -> None: + """幂等关闭持久连接(对齐 HarnessLog;进程退出前可选调以释放 fd)。 + + 不调也不丢数据——每次 _write 已 commit 落 WAL,进程退出 OS 回收 fd、 + WAL 已提交内容下次打开自动 checkpoint 恢复。 + """ + with self._lock: + if self._conn is not None: + self._conn.close() + self._conn = None + def _write( self, *, diff --git a/research-wiki/designs/2026-07-16-telemetry-concurrency-fix-design.md b/research-wiki/designs/2026-07-16-telemetry-concurrency-fix-design.md index cff61d0..d5d2a62 100644 --- a/research-wiki/designs/2026-07-16-telemetry-concurrency-fix-design.md +++ b/research-wiki/designs/2026-07-16-telemetry-concurrency-fix-design.md @@ -65,6 +65,14 @@ telemetry **不能**复用 HarnessLog 实例,三条硬隔离: | **续跑** | 不适用(遥测无状态;进程退出 WAL 自动恢复已 commit 的) | | **原子性** | 单条 insert+commit 原子,无半写 | +**零丢失保证的适用范围(Codex 审明确)**: +- **单进程、单 recorder 实例**:锁与连接是实例字段,串行化只在同一实例内成立。同进程多个 recorder 指向同一 db 会退回跨连接竞争——当前 `main.py` / `video_split_cli` 均单实例注入,不踩;本实现不支持多实例同库(YAGNI,若未来需要再引 class-level registry)。 +- **唯一 call_id**:`INSERT OR IGNORE` 下重复 call_id 是**预期忽略**(幂等),不计作丢失。 + +**降级边界(Codex 审加固)**:`__init__` 的 mkdir / connect / PRAGMA / 建表统一纳入 `except (OSError, sqlite3.Error)` 降级(`self._conn=None`),任一失败都不冒泡拖垮初始化;`_write` 遇 `self._conn is None` 或 execute 抛错均降级 warning。守住"遥测失败绝不拖垮 LLM 调用"哲学。 + +**生命周期**:补幂等 `close()`(对齐 HarnessLog)供进程退出前可选调释放 fd;不调也不丢数据(WAL 已 commit)。telemetry 是长生命周期单例,无 context-manager 场景,故 close 为可选而非强制。 + ## 7. 测试 - **并发写不锁死**(核心):多线程/多协程并发调 `record_llm_call`(如 32 并发 × N 条),断言全部落库、零 `database is locked`、零丢失(行数 == 写入数)。这是复现 bug 的真实场景测试。 diff --git a/tests/unit/test_telemetry.py b/tests/unit/test_telemetry.py index edcb405..ddbd3e0 100644 --- a/tests/unit/test_telemetry.py +++ b/tests/unit/test_telemetry.py @@ -125,11 +125,43 @@ async def test_duplicate_call_id_does_not_raise(recorder, db_path): @pytest.mark.asyncio -async def test_db_error_does_not_propagate(tmp_path): - """SQLite 写入失败时 record_llm_call 应静默降级,不抛异常。""" - bad_recorder = SQLiteTelemetryRecorder(db_path=tmp_path / "nonexistent_dir" / "bad.db") - kwargs = _make_call_kwargs() - await bad_recorder.record_llm_call(**kwargs) +async def test_init_db_error_does_not_propagate(db_path, monkeypatch): + """构造期 connect 失败应降级(self._conn=None),不冒泡拖垮初始化。""" + + def _boom_connect(*args, **kwargs): + raise sqlite3.OperationalError("unable to open database file") + + monkeypatch.setattr(sqlite3, "connect", _boom_connect) + recorder = SQLiteTelemetryRecorder(db_path=db_path) # 不抛 + assert recorder._conn is None + await recorder.record_llm_call(**_make_call_kwargs()) # 写也降级不抛 + + +@pytest.mark.asyncio +async def test_init_mkdir_error_does_not_propagate(db_path, monkeypatch): + """构造期 mkdir 失败(PermissionError/OSError)也应降级,不冒泡。""" + + def _boom_mkdir(*args, **kwargs): + raise PermissionError("cannot create dir") + + monkeypatch.setattr("adapters.telemetry.Path.mkdir", _boom_mkdir) + recorder = SQLiteTelemetryRecorder(db_path=db_path) # 不抛 + assert recorder._conn is None + + +@pytest.mark.asyncio +async def test_write_db_error_does_not_propagate(recorder): + """写入期 execute 失败应静默降级,不抛异常(遥测失败绝不拖垮 LLM 调用)。""" + + class _BoomConn: + def execute(self, *args, **kwargs): + raise sqlite3.OperationalError("disk I/O error") + + def commit(self): + pass + + recorder._conn = _BoomConn() # 连接存在但写抛错,验证 _write 的 except 降级 + await recorder.record_llm_call(**_make_call_kwargs()) # 降级不抛 @pytest.mark.asyncio @@ -173,17 +205,20 @@ def test_uses_single_persistent_connection(db_path, monkeypatch): 每次写新建连接是并发锁竞争根源(多连接争 SQLite 写锁,撑爆 busy_timeout); 单连接 + 进程内 Lock 串行化把并发控制拉到进程内,消除 SQLite 层锁竞争。 """ - connect_calls = {"n": 0} + connect_calls = {"n": 0, "kwargs": None} real_connect = sqlite3.connect def _counting_connect(*args, **kwargs): connect_calls["n"] += 1 + connect_calls["kwargs"] = kwargs return real_connect(*args, **kwargs) monkeypatch.setattr(sqlite3, "connect", _counting_connect) recorder = SQLiteTelemetryRecorder(db_path=db_path) after_init = connect_calls["n"] + # 跨线程共享连接必须 check_same_thread=False(串行性由 self._lock 保证) + assert connect_calls["kwargs"].get("check_same_thread") is False for _ in range(10): recorder._write(**_make_call_kwargs()) after_writes = connect_calls["n"]