fix: address Codex review of telemetry fix (mkdir degrade, close, degrade tests)

This commit is contained in:
2026-07-16 08:53:50 -04:00
parent 065c8ae1b9
commit cc01d5ed62
3 changed files with 69 additions and 9 deletions
+20 -3
View File
@@ -1,7 +1,11 @@
"""SQLite 遥测记录器 — TelemetryRecorder Protocol 的生产实现。 """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 from __future__ import annotations
@@ -68,6 +72,8 @@ class SQLiteTelemetryRecorder:
self._db_path = db_path self._db_path = db_path
self._lock = threading.Lock() self._lock = threading.Lock()
self._conn: sqlite3.Connection | None = None self._conn: sqlite3.Connection | None = None
# mkdir / connect / PRAGMA / 建表统一纳入降级边界:任一失败(OSError 含
# PermissionError、sqlite3.Error)都降级为 self._conn=None,绝不冒泡拖垮初始化。
try: try:
Path(db_path).parent.mkdir(parents=True, exist_ok=True) Path(db_path).parent.mkdir(parents=True, exist_ok=True)
conn = sqlite3.connect(str(db_path), check_same_thread=False) conn = sqlite3.connect(str(db_path), check_same_thread=False)
@@ -76,10 +82,21 @@ class SQLiteTelemetryRecorder:
conn.execute(self._CREATE_TABLE_SQL) conn.execute(self._CREATE_TABLE_SQL)
conn.commit() conn.commit()
self._conn = conn self._conn = conn
except sqlite3.Error as exc: except (OSError, sqlite3.Error) as exc:
logger.warning("遥测连接初始化失败(已降级,后续写入丢弃): {}", exc) logger.warning("遥测连接初始化失败(已降级,后续写入丢弃): {}", exc)
self._conn = None 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( def _write(
self, self,
*, *,
@@ -65,6 +65,14 @@ telemetry **不能**复用 HarnessLog 实例,三条硬隔离:
| **续跑** | 不适用(遥测无状态;进程退出 WAL 自动恢复已 commit 的) | | **续跑** | 不适用(遥测无状态;进程退出 WAL 自动恢复已 commit 的) |
| **原子性** | 单条 insert+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. 测试 ## 7. 测试
- **并发写不锁死**(核心):多线程/多协程并发调 `record_llm_call`(如 32 并发 × N 条),断言全部落库、零 `database is locked`、零丢失(行数 == 写入数)。这是复现 bug 的真实场景测试。 - **并发写不锁死**(核心):多线程/多协程并发调 `record_llm_call`(如 32 并发 × N 条),断言全部落库、零 `database is locked`、零丢失(行数 == 写入数)。这是复现 bug 的真实场景测试。
+41 -6
View File
@@ -125,11 +125,43 @@ async def test_duplicate_call_id_does_not_raise(recorder, db_path):
@pytest.mark.asyncio @pytest.mark.asyncio
async def test_db_error_does_not_propagate(tmp_path): async def test_init_db_error_does_not_propagate(db_path, monkeypatch):
"""SQLite 写入失败时 record_llm_call 应静默降级,不抛异常""" """构造期 connect 失败应降级(self._conn=None),不冒泡拖垮初始化"""
bad_recorder = SQLiteTelemetryRecorder(db_path=tmp_path / "nonexistent_dir" / "bad.db")
kwargs = _make_call_kwargs() def _boom_connect(*args, **kwargs):
await bad_recorder.record_llm_call(**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 @pytest.mark.asyncio
@@ -173,17 +205,20 @@ def test_uses_single_persistent_connection(db_path, monkeypatch):
每次写新建连接是并发锁竞争根源(多连接争 SQLite 写锁,撑爆 busy_timeout); 每次写新建连接是并发锁竞争根源(多连接争 SQLite 写锁,撑爆 busy_timeout);
单连接 + 进程内 Lock 串行化把并发控制拉到进程内,消除 SQLite 层锁竞争。 单连接 + 进程内 Lock 串行化把并发控制拉到进程内,消除 SQLite 层锁竞争。
""" """
connect_calls = {"n": 0} connect_calls = {"n": 0, "kwargs": None}
real_connect = sqlite3.connect real_connect = sqlite3.connect
def _counting_connect(*args, **kwargs): def _counting_connect(*args, **kwargs):
connect_calls["n"] += 1 connect_calls["n"] += 1
connect_calls["kwargs"] = kwargs
return real_connect(*args, **kwargs) return real_connect(*args, **kwargs)
monkeypatch.setattr(sqlite3, "connect", _counting_connect) monkeypatch.setattr(sqlite3, "connect", _counting_connect)
recorder = SQLiteTelemetryRecorder(db_path=db_path) recorder = SQLiteTelemetryRecorder(db_path=db_path)
after_init = connect_calls["n"] 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): for _ in range(10):
recorder._write(**_make_call_kwargs()) recorder._write(**_make_call_kwargs())
after_writes = connect_calls["n"] after_writes = connect_calls["n"]