fix(stores): 修 Codex 对抗审查报的五条,其中两条同一根因

最实的一条:读取端只要一段能解析成 JSON 就收下,没检查它后面有没有换行。而短写完全可能
正好写完整个 JSON 对象、只差那个换行——那次写从来没被确认过(调用方的 await 还没返回),
按契约就是「没发生」,但它会被当成一条有效的动作意图读回来,恢复据此判成「状态未知」并可能
重放,而那个动作一定没执行过(调用方是在写意图返回之后才去执行的)。

判据改成「这一行有没有被换行终结」,不是「能不能解析」。同一个改动顺带修掉第三条:一行完整
终结的坏行(比如被外部追加的 {})此前会被当成撕裂尾行吞掉,读成「少了一条记录但看起来完整」
的日志;现在终结过的行解不开就是损坏,直接报错。

其余三条:
- 新建日志文件不 fsync 父目录。os.fsync(fd) 刷的是文件内容,刷不到「这个目录里多了一个
  文件」这条目录项;掉电后内容可能在而文件不存在,read_log 走「文件不存在」返回空日志,
  驱动入口判成全新运行,一次已经花过钱的运行静默没了留痕。只在新建时刷。
- 同一运行标识上的并发写会交错:一条记录可能由不止一次 os.write 写完,而 O_APPEND 只保证
  每次 write 的追加位置原子,保证不了一条逻辑行整体原子。按运行标识加锁串起来(不同运行
  照样并行),跨进程那一半仍靠独占创建挡。有一条用短写逼出那个窗口的测试。
- 往返测试的 TOTAL_WRITES 是硬编码,而且漏写结束标记它发现不了(恢复会把最后一步之后那次
  停止判定重演一遍,得出同样结果)。加一条把十次写的记录类型序列整个钉死的测试。
This commit is contained in:
2026-08-10 03:52:35 -04:00
parent 7f6066701e
commit aeb575e0f7
4 changed files with 299 additions and 34 deletions
+61 -28
View File
@@ -72,10 +72,17 @@ class JsonlRunStore:
它满足 `polyloop.ports.RunStore`,但不显式继承那个 Protocol:结构化子类型不需要继承。
"""
__slots__ = ("_directory", "_write_all")
__slots__ = ("_directory", "_locks", "_write_all")
def __init__(self, *, directory: Path | str) -> None:
self._directory = Path(directory)
#: 每个运行标识一把锁,把同一个文件上的写串起来。
#:
#: 一条记录可能由不止一次 `os.write` 写完(`os.write` 允许短写),而 `O_APPEND` 只保证
#: 每一次 `os.write` 的追加位置原子,保证不了「一条逻辑行整体原子」。两个协程同时往同一
#: 个文件写时,一次短写会让两条记录交错成一段谁也解不开的字节。锁把这件事挡在进程内;
#: 跨进程那一半靠运行开始记录的独占创建挡(见 `write_run_started`)。
self._locks: dict[str, asyncio.Lock] = {}
#: **可注入的故障点**,见 `_write_all_bytes` 的 docstring。做成实例属性而不是方法,
#: 是因为 `__slots__` 让方法替换不掉,而替换它正是那条测试唯一的做法。
self._write_all = _write_all_bytes
@@ -169,11 +176,13 @@ class JsonlRunStore:
tag = _TAGS[type(record)]
line = json.dumps({RECORD_KEY: tag, **encode(record)}, ensure_ascii=False) + "\n"
path = self._path(record.run_id)
# 写入与 `fsync` 都是阻塞调用,而 `fsync` 在忙盘上可以到几十毫秒。直接在事件循环里做
# 会把同一个循环上所有并发运行一起卡住。
await asyncio.to_thread(
self._write_line, path, line.encode("utf-8"), fsync=fsync, exclusive=exclusive
)
lock = self._locks.setdefault(record.run_id, asyncio.Lock())
async with lock:
# 写入与 `fsync` 都是阻塞调用,而 `fsync` 在忙盘上可以到几十毫秒。直接在事件循环里
# 做会把同一个循环上所有并发运行一起卡住。锁按运行标识分,所以不同运行照样并行。
await asyncio.to_thread(
self._write_line, path, line.encode("utf-8"), fsync=fsync, exclusive=exclusive
)
def _write_line(self, path: Path, payload: bytes, *, fsync: bool, exclusive: bool) -> None:
"""打开、追加、按需 `fsync`、关闭。
@@ -192,6 +201,24 @@ class JsonlRunStore:
os.fsync(descriptor)
finally:
os.close(descriptor)
if exclusive:
_fsync_directory(path.parent)
def _fsync_directory(directory: Path) -> None:
"""把新建文件的目录项刷下去。
`os.fsync(fd)` 刷的是那个文件的内容,刷不到「这个目录里多了一个文件」这条目录项。掉电之后
内容可能在、而文件根本不存在——那时 `read_log` 走「文件不存在」返回空日志,驱动入口据此
判成一次全新的运行,于是一次已经开始过、可能已经花过钱的运行静默没了留痕。
只在新建文件时做:往已有文件追加不改目录项。
"""
descriptor = os.open(directory, os.O_RDONLY)
try:
os.fsync(descriptor)
finally:
os.close(descriptor)
def _write_all_bytes(descriptor: int, payload: bytes) -> None:
@@ -211,30 +238,33 @@ def _write_all_bytes(descriptor: int, payload: bytes) -> None:
def _parse(raw: bytes, run_id: str) -> RunLog:
"""把一份文件内容还原成日志。
**第一条解不开的行就是日志的结尾**,它后面还有内容就不是撕裂而是损坏。追加写只在末尾产生
撕裂;中间出现读不了的字节意味着别的东西动过这个文件,那时跳过那一行接着读会拼出一份少了
几条记录、看起来却完整的日志,而恢复会照它做判断
**判据是「这一行有没有被换行终结」,不是「它能不能解析」。** 一次写入是先写整行再由调用方
等到它返回,所以文件末尾那段没有换行的字节对应的那次写**从来没有被确认过**——按契约它就是
没发生,丢掉它正是「要么都可见、要么都不可见」的落地方式
照「能不能解析」判会漏掉一个很具体的场景:短写正好写完了整个 JSON 对象、只差最后那个换行。
那段字节解得开,于是一条从没被确认的动作意图被当成有效记录读回来,恢复据此判成「状态未知」
并可能重放——而那个动作其实一定没执行过,因为调用方是在写意图返回之后才去执行的。
**被换行终结的行必须解得开**,解不开就是损坏,直接报错。追加写只在末尾产生撕裂;一条完整
终结的行读不了,说明别的东西动过这个文件,那时跳过它接着读会拼出一份少了几条记录、看起来
却完整的日志,而恢复会照它做判断。
"""
started: RunStarted | None = None
intents: list[Intent] = []
model_results: list[ModelCallResult] = []
steps: list[StepCompleted] = []
finished: RunFinished | None = None
torn_at: int | None = None
for number, chunk in enumerate(raw.split(b"\n"), start=1):
chunks = raw.split(b"\n")
# 文件以换行结尾时最后一段是空的;不以换行结尾说明最后那次写没写完。
terminated = chunks[:-1] if chunks and chunks[-1].strip() else chunks
for number, chunk in enumerate(terminated, start=1):
if not chunk.strip():
# 空行不携带记录,也不是撕裂的证据。
continue
if torn_at is not None:
raise DecodeError(
f"运行 {run_id!r} 的日志第 {torn_at} 行读不了,而第 {number} 行还有内容。"
"追加写只在末尾产生撕裂,中间读不了说明这个文件被别的东西动过"
)
record = _decode_line(chunk)
if record is None:
torn_at = number
continue
record = _decode_line(chunk, run_id=run_id, number=number)
if isinstance(record, RunStarted):
started = record
elif isinstance(record, Intent):
@@ -255,22 +285,25 @@ def _parse(raw: bytes, run_id: str) -> RunLog:
)
def _decode_line(chunk: bytes) -> object | None:
"""解一行。解不开返回 `None`——调用方据此判断它是不是撕裂的尾行
def _decode_line(chunk: bytes, *, run_id: str, number: int) -> object:
"""解一条被换行终结的行。解不开就是损坏,直接报错
**认不得的类型标签不算撕裂,直接报错。** 一行完整的 JSON 带着一个我们不认识的标签,说明
这份日志是别的版本或者别的东西写的,不是被杀在写一半。
这里不再有「解不开就当撕裂尾行」那条路——撕裂由有没有换行判定,进不到这个函数。
"""
where = f"运行 {run_id!r} 的日志第 {number}"
try:
payload = json.loads(chunk)
except (UnicodeDecodeError, json.JSONDecodeError):
return None
except (UnicodeDecodeError, json.JSONDecodeError) as exc:
raise DecodeError(
f"{where}读不了({type(exc).__name__})。它是被换行终结的完整一行,"
"说明这个文件被别的东西动过——追加写只在末尾产生撕裂"
) from exc
if not isinstance(payload, dict) or RECORD_KEY not in payload:
return None
raise DecodeError(f"{where}没有 {RECORD_KEY!r} 标签,这份文件不是本库写的")
tag = payload[RECORD_KEY]
decoder = _DECODERS.get(tag)
if decoder is None:
raise DecodeError(f"日志里出现认不得的记录类型 {tag!r},这份文件不是本库写的")
raise DecodeError(f"{where}的记录类型 {tag!r} 认不得,这份文件不是本库写的")
return decoder(payload)