109 lines
4.0 KiB
Python
109 lines
4.0 KiB
Python
"""框架运行工件的确定性写入、校验与重放工具。"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import hashlib
|
||
import json
|
||
import os
|
||
import tempfile
|
||
from pathlib import Path
|
||
from typing import Any, Iterable, Mapping
|
||
|
||
|
||
class ArtifactError(ValueError):
|
||
"""运行工件不完整、格式非法或序列不连续。"""
|
||
|
||
|
||
def payload_sha256(value: Any) -> str:
|
||
"""按稳定 JSON 计算 payload 摘要。"""
|
||
|
||
raw = json.dumps(value, ensure_ascii=False, sort_keys=True, separators=(",", ":"))
|
||
return "sha256:" + hashlib.sha256(raw.encode("utf-8")).hexdigest()
|
||
|
||
|
||
def _validate_event(event: Mapping[str, Any], expected_seq: int) -> dict[str, Any]:
|
||
if not isinstance(event, Mapping):
|
||
raise ArtifactError("事件必须是对象")
|
||
source_seq = event.get("sourceSeq")
|
||
if source_seq != expected_seq:
|
||
raise ArtifactError(
|
||
f"事件序列不连续:期望 sourceSeq={expected_seq},实际 {source_seq!r}"
|
||
)
|
||
if not event.get("framework") or not event.get("kind") or not event.get("phase"):
|
||
raise ArtifactError("事件缺少 framework/kind/phase")
|
||
return dict(event)
|
||
|
||
|
||
def write_jsonl_atomic(path: str | Path, events: Iterable[Mapping[str, Any]]) -> str:
|
||
"""原子发布完整 JSONL 工件,返回文件摘要;发布前校验序列连续。"""
|
||
|
||
target = Path(path)
|
||
target.parent.mkdir(parents=True, exist_ok=True)
|
||
lines: list[str] = []
|
||
for expected_seq, event in enumerate(events, start=1):
|
||
validated = _validate_event(event, expected_seq)
|
||
lines.append(json.dumps(validated, ensure_ascii=False, sort_keys=True) + "\n")
|
||
if not lines:
|
||
raise ArtifactError("不能发布空运行工件")
|
||
content = "".join(lines).encode("utf-8")
|
||
descriptor, temporary_name = tempfile.mkstemp(
|
||
prefix=f".{target.name}.", dir=str(target.parent)
|
||
)
|
||
temporary = Path(temporary_name)
|
||
try:
|
||
with os.fdopen(descriptor, "wb") as stream:
|
||
stream.write(content)
|
||
stream.flush()
|
||
os.fsync(stream.fileno())
|
||
os.replace(temporary, target)
|
||
finally:
|
||
temporary.unlink(missing_ok=True)
|
||
return "sha256:" + hashlib.sha256(content).hexdigest()
|
||
|
||
|
||
def read_jsonl(path: str | Path, *, allow_seq_gap: bool = False) -> list[dict[str, Any]]:
|
||
"""读取并校验 JSONL;默认拒绝缺失或重复 sourceSeq。"""
|
||
|
||
source = Path(path)
|
||
try:
|
||
lines = source.read_text(encoding="utf-8").splitlines()
|
||
except (OSError, UnicodeError) as exc:
|
||
raise ArtifactError(f"运行工件不可读: {source}") from exc
|
||
events: list[dict[str, Any]] = []
|
||
expected = 1
|
||
for line_number, line in enumerate(lines, start=1):
|
||
if not line.strip():
|
||
continue
|
||
try:
|
||
raw = json.loads(line)
|
||
except json.JSONDecodeError as exc:
|
||
raise ArtifactError(f"第 {line_number} 行不是合法 JSON") from exc
|
||
if not isinstance(raw, Mapping):
|
||
raise ArtifactError(f"第 {line_number} 行不是事件对象")
|
||
if allow_seq_gap:
|
||
source_seq = raw.get("sourceSeq")
|
||
if not isinstance(source_seq, int) or source_seq < expected:
|
||
raise ArtifactError(f"第 {line_number} 行 sourceSeq 非法或重复")
|
||
expected = source_seq + 1
|
||
events.append(dict(raw))
|
||
else:
|
||
events.append(_validate_event(raw, expected))
|
||
expected += 1
|
||
if not events:
|
||
raise ArtifactError("运行工件为空")
|
||
return events
|
||
|
||
|
||
def detect_seq_gaps(events: Iterable[Mapping[str, Any]]) -> list[tuple[int, int]]:
|
||
"""返回排序后发现的 sourceSeq 缺口,不修改输入。"""
|
||
|
||
values = sorted(int(event["sourceSeq"]) for event in events)
|
||
gaps: list[tuple[int, int]] = []
|
||
for left, right in zip(values, values[1:]):
|
||
if right > left + 1:
|
||
gaps.append((left + 1, right - 1))
|
||
return gaps
|
||
|
||
|
||
__all__ = ["ArtifactError", "detect_seq_gaps", "payload_sha256", "read_jsonl", "write_jsonl_atomic"]
|