feat(stores): 补上易失存储实现,stores 拆成四个文件
design 0003 的否决方案一节承诺过「显式命名的、明确不提供恢复的内存实现」,那个实现一直 没写,migrations/dissect.md 还专门登记着「不要照那句话去找一个不存在的类」。 **名字取 VolatileRunStore 不取「不提供恢复」**:一个真的读不回自己写过的东西的存储过不了 自家的契约套件(套件第一条要的就是「写进去的意图读得回来」),而这一层的准入标准就是那套 套件。一个过不了自家准入标准的实现不该存在。它不提供的是跨进程恢复,Volatile 说的正是这 件事。 桶里存编码后的载荷、读的时候才解码,两个方向的别名都堵上:调用方写完再改自己手里那个 dict 改不到日志,读回来的日志被就地改动也污染不了存储本身。落盘那个实现每次重新解析文件,天然 如此,这个实现靠同一条路径对齐它。**写入只编码不解码**——落盘那边写的时候只做 json.dumps, 一条字段类型不对的记录写得进去、读的时候才炸,两边现在一致。 两张平行的表(记录类→标签、标签→解码器)合成一张三元组再派生视图。加上易失实现要用的第三 张视图之后,三张表之间那个谁也不检查的一致性要求就不可能被违反了。
This commit is contained in:
+13
-304
@@ -1,310 +1,19 @@
|
||||
"""存储接缝的第一个实现:一次运行一个文件,一行一条记录,逐行追加。
|
||||
"""存储接缝的两个实现:一个逐行追加进本地文件,一个只留在进程内存里。
|
||||
|
||||
**这个模块公开,但不进 `polyloop/__init__.py`**,必须显式 import(`0003` 决策八第 9 条)。
|
||||
|
||||
它满足 `research-wiki/design/0003-public-api-shape.md` 决策四那张写入序列表与
|
||||
`0005-storage-atomicity-and-record-fields.md` 决策五那条前缀持久性要求;文件布局、坏行怎么算、
|
||||
`fsync` 在哪几处,定在 `research-wiki/design/0011-jsonl-run-store.md`。
|
||||
|
||||
**`record` 是这一层的保留键。** 每行是「一个类型标签加那条记录的全部字段」,而
|
||||
`polyloop.serialization` 编出来的载荷只有记录类自己的字段——标签是这里加的。五个记录类现在
|
||||
都没有叫 `record` 的字段,将来也不许加:加了的话编码出来的键会和标签撞,而撞的表现是解码时
|
||||
把一条记录读成另一种。
|
||||
**这一层公开,但不进 `polyloop/__init__.py`**,必须显式 import(`0003` 决策八第 9 条)。存储
|
||||
是必填的装配项,库不给默认值:顺手提供一个默认实现等于替所有下游选了日志落在哪儿,而那是
|
||||
每个项目自己的运维决定。调用方要么从这里挑一个,要么自己写一个。
|
||||
|
||||
**关系数据库那一种形态本库不提供。** 表结构、事务边界、连接管理都在项目那边,库替它写一个
|
||||
通用实现只会写出一个谁都不合用的。存储接缝的意义就是让它自己实现,而 `tests/contract/` 是
|
||||
它的准入标准。
|
||||
通用实现只会写出一个谁都不合用的。存储接缝的意义就是让它自己实现,而契约套件
|
||||
(`polyloop.testing`,见 `research-wiki/design/0014-contract-suite-distribution.md` 决策一)
|
||||
是它的准入标准:那套用例只说行为、不碰形态,跑全绿就算合格。
|
||||
|
||||
**这一层也不提供「把日志转成别的格式」「按时间范围查」这类操作。** 存储接缝上多一个方法,就是
|
||||
给每一个下游实现多加一份永久要求,而这些事拿 `read_log` 读回来的日志在库外面做即可。
|
||||
"""
|
||||
|
||||
import asyncio
|
||||
import json
|
||||
import os
|
||||
import re
|
||||
from collections.abc import Mapping
|
||||
from pathlib import Path
|
||||
from polyloop.stores._jsonl import RECORD_KEY, JsonlRunStore
|
||||
from polyloop.stores._volatile import VolatileRunStore
|
||||
|
||||
from polyloop.ports import RunLog
|
||||
from polyloop.serialization import (
|
||||
DecodeError,
|
||||
decode_intent,
|
||||
decode_model_call_result,
|
||||
decode_run_finished,
|
||||
decode_run_started,
|
||||
decode_step_completed,
|
||||
encode,
|
||||
)
|
||||
from polyloop.types import Intent, ModelCallResult, RunFinished, RunStarted, StepCompleted
|
||||
|
||||
#: 每行那个类型标签的键名。见模块 docstring:它是保留键。
|
||||
RECORD_KEY = "record"
|
||||
|
||||
#: 运行标识同时是文件名,所以它必须是一个安全的文件名。
|
||||
#:
|
||||
#: 不转义也不哈希:那样文件名就不再等于运行标识,而按运行标识去目录里找文件是最自然的用法。
|
||||
#: 这条校验挡住的不只是可读性——运行标识是调用方给的不透明字符串,里面出现 `../` 的话,
|
||||
#: 写文件会跑到目录外面去。
|
||||
_SAFE_RUN_ID = re.compile(r"[A-Za-z0-9._-]+")
|
||||
|
||||
_TAGS: Mapping[type, str] = {
|
||||
RunStarted: "run_started",
|
||||
Intent: "intent",
|
||||
ModelCallResult: "model_call_result",
|
||||
StepCompleted: "step_completed",
|
||||
RunFinished: "run_finished",
|
||||
}
|
||||
|
||||
_DECODERS = {
|
||||
"run_started": decode_run_started,
|
||||
"intent": decode_intent,
|
||||
"model_call_result": decode_model_call_result,
|
||||
"step_completed": decode_step_completed,
|
||||
"run_finished": decode_run_finished,
|
||||
}
|
||||
|
||||
|
||||
class JsonlRunStore:
|
||||
"""把一次运行的日志逐行追加进 `<目录>/<运行标识>.jsonl`。
|
||||
|
||||
**一次运行一个文件,不是一个大文件加一列运行标识。** 大文件上「读回某一次运行的整份日志」
|
||||
要扫全文,而那件事在每次开工前都会做一遍;更要命的是两次并发运行会往同一个文件追加,
|
||||
前缀持久性就从「同一文件的追加序」退化成「两条交错的序」。
|
||||
|
||||
它满足 `polyloop.ports.RunStore`,但不显式继承那个 Protocol:结构化子类型不需要继承。
|
||||
"""
|
||||
|
||||
__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
|
||||
|
||||
def __repr__(self) -> str:
|
||||
return f"JsonlRunStore(directory={str(self._directory)!r})"
|
||||
|
||||
def parameters(self) -> Mapping[str, str]:
|
||||
"""上报可复现参数。
|
||||
|
||||
**目录不进快照。** 它是这份日志本身所在的地方——目录要是不一样,根本读不到这份日志、
|
||||
也就走不到比对那一步。把它记进去只会在换一台机器、挂载点变了的时候报出一次假的漂移,
|
||||
而那次续跑其实完全正常。
|
||||
"""
|
||||
return {"kind": "jsonl"}
|
||||
|
||||
# -- 写 ------------------------------------------------------------------
|
||||
|
||||
async def write_run_started(self, record: RunStarted) -> None:
|
||||
"""写运行开始记录。**独占创建**:文件已存在就直接失败。
|
||||
|
||||
驱动入口在开工前会先读一次日志判断这个标识有没有用过,但那是先读后写,两个进程同时
|
||||
读到空、同时开始写的窗口它挡不住。独占创建把那个窗口关掉,代价是一个标志位。
|
||||
|
||||
两个进程同时跑同一个运行标识的后果很具体:交错的记录序会让恢复读到同一步的两条意图,
|
||||
判成「日志被并发写过」,于是这次运行从此续不了——而两边的模型调用都已经花过钱了。
|
||||
"""
|
||||
await self._append(record, fsync=True, exclusive=True)
|
||||
|
||||
async def write_intent(self, record: Intent) -> None:
|
||||
"""写一条意图。**耐久屏障**:必须落盘才能往下走。
|
||||
|
||||
意图必须在副作用之前就持久,这是意图日志的全部意义。
|
||||
"""
|
||||
await self._append(record, fsync=True)
|
||||
|
||||
async def write_model_call_result(self, record: ModelCallResult) -> None:
|
||||
"""不做 `fsync`:后面紧跟的不是副作用。
|
||||
|
||||
它靠前缀持久性兜——同一个文件的追加写,下一次 `fsync`(那必定是一条意图,或者运行
|
||||
结束)会把它一起刷下去。所以「结果还没落盘、动作意图落了盘」这个状态在这份实现上
|
||||
不可能出现,而那正是 `0005` 决策五要防的。
|
||||
"""
|
||||
await self._append(record)
|
||||
|
||||
async def write_step_completed(self, record: StepCompleted) -> None:
|
||||
"""动作结果与步记录一次原子落地。
|
||||
|
||||
它们本来就是同一个记录类的两个字段,所以「一次原子写」在这份实现上就是**一行**:
|
||||
一行要么完整地在文件里,要么是被丢掉的撕裂尾行,没有中间态。
|
||||
"""
|
||||
await self._append(record)
|
||||
|
||||
async def write_run_finished(self, record: RunFinished) -> None:
|
||||
"""写结束标记并 `fsync`。丢了的话这次运行看起来还能续,而它已经跑完了。"""
|
||||
await self._append(record, fsync=True)
|
||||
|
||||
# -- 读 ------------------------------------------------------------------
|
||||
|
||||
async def read_log(self, run_id: str) -> RunLog:
|
||||
"""读回整份日志。文件不存在时返回空日志,不抛异常。
|
||||
|
||||
驱动入口靠这条判断「这个标识是不是已经有日志了」。抛异常的话那个判断就得写成捕获
|
||||
异常,而用捕获异常做流程控制会把真正的存储故障一起吞掉——于是「磁盘挂了」会被读成
|
||||
「这是一次全新的运行」,然后覆盖式地重跑一遍。
|
||||
"""
|
||||
path = self._path(run_id)
|
||||
if not path.exists():
|
||||
return RunLog()
|
||||
raw = await asyncio.to_thread(path.read_bytes)
|
||||
return _parse(raw, run_id)
|
||||
|
||||
# -- 内部 ----------------------------------------------------------------
|
||||
|
||||
def _path(self, run_id: str) -> Path:
|
||||
if not _SAFE_RUN_ID.fullmatch(run_id) or run_id.startswith("."):
|
||||
raise ValueError(
|
||||
f"运行标识 {run_id!r} 不能直接当文件名。这份实现要求它只含字母、数字、点、"
|
||||
"下划线与连字符,且不以点开头——它同时是文件名,而按标识去目录里找文件是最"
|
||||
"自然的用法"
|
||||
)
|
||||
return self._directory / f"{run_id}.jsonl"
|
||||
|
||||
async def _append(
|
||||
self,
|
||||
record: RunStarted | Intent | ModelCallResult | StepCompleted | RunFinished,
|
||||
*,
|
||||
fsync: bool = False,
|
||||
exclusive: bool = False,
|
||||
) -> None:
|
||||
tag = _TAGS[type(record)]
|
||||
line = json.dumps({RECORD_KEY: tag, **encode(record)}, ensure_ascii=False) + "\n"
|
||||
path = self._path(record.run_id)
|
||||
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`、关闭。
|
||||
|
||||
**不长期持有文件句柄。** 持有要为每个运行标识维护一份状态,而那份状态在并发下就是共享
|
||||
可变状态;打开的成本相对一次 `fsync` 可以忽略,一次 `fsync` 相对一次模型调用又可以忽略。
|
||||
"""
|
||||
path.parent.mkdir(parents=True, exist_ok=True)
|
||||
flags = os.O_WRONLY | os.O_APPEND | os.O_CREAT
|
||||
if exclusive:
|
||||
flags |= os.O_EXCL
|
||||
descriptor = os.open(path, flags, 0o644)
|
||||
try:
|
||||
self._write_all(descriptor, payload)
|
||||
if fsync:
|
||||
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:
|
||||
"""把这些字节全部写进去。
|
||||
|
||||
**这是那个可注入的故障点。** 契约套件验不了原子写的「一起不可见」那一半(要在写入中途
|
||||
杀进程,而它跑在一个进程里),`0011` 定的做法是在这一层留一个可替换的内部函数,单元测试
|
||||
把它换成「写一半就抛异常」。它不出现在任何接缝签名上,换一个存储实现就没有它。
|
||||
|
||||
循环是因为 `os.write` 允许短写。短写留下的半行正是撕裂尾行,读那边会丢掉它。
|
||||
"""
|
||||
written = 0
|
||||
while written < len(payload):
|
||||
written += os.write(descriptor, payload[written:])
|
||||
|
||||
|
||||
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
|
||||
|
||||
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
|
||||
record = _decode_line(chunk, run_id=run_id, number=number)
|
||||
if isinstance(record, RunStarted):
|
||||
started = record
|
||||
elif isinstance(record, Intent):
|
||||
intents.append(record)
|
||||
elif isinstance(record, ModelCallResult):
|
||||
model_results.append(record)
|
||||
elif isinstance(record, StepCompleted):
|
||||
steps.append(record)
|
||||
else:
|
||||
finished = record
|
||||
|
||||
return RunLog(
|
||||
started=started,
|
||||
intents=tuple(intents),
|
||||
model_results=tuple(model_results),
|
||||
steps=tuple(steps),
|
||||
finished=finished,
|
||||
)
|
||||
|
||||
|
||||
def _decode_line(chunk: bytes, *, run_id: str, number: int) -> object:
|
||||
"""解一条被换行终结的行。解不开就是损坏,直接报错。
|
||||
|
||||
这里不再有「解不开就当撕裂尾行」那条路——撕裂由有没有换行判定,进不到这个函数。
|
||||
"""
|
||||
where = f"运行 {run_id!r} 的日志第 {number} 行"
|
||||
try:
|
||||
payload = json.loads(chunk)
|
||||
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:
|
||||
raise DecodeError(f"{where}没有 {RECORD_KEY!r} 标签,这份文件不是本库写的")
|
||||
tag = payload[RECORD_KEY]
|
||||
decoder = _DECODERS.get(tag)
|
||||
if decoder is None:
|
||||
raise DecodeError(f"{where}的记录类型 {tag!r} 认不得,这份文件不是本库写的")
|
||||
return decoder(payload)
|
||||
|
||||
|
||||
__all__ = ["RECORD_KEY", "JsonlRunStore"]
|
||||
__all__ = ["RECORD_KEY", "JsonlRunStore", "VolatileRunStore"]
|
||||
|
||||
Reference in New Issue
Block a user