Files
2026-09-04 14:58:42 +08:00

79 lines
2.4 KiB
Python

"""可恢复运行的原子清单持久化。"""
from __future__ import annotations
import json
from pathlib import Path
from typing import Any
from scripts.dynamic_compile.fast.storage import atomic_write_json
def read_json_object(path: Path) -> dict[str, Any] | None:
try:
value = json.loads(path.read_text(encoding="utf-8"))
except (OSError, json.JSONDecodeError):
return None
return value if isinstance(value, dict) else None
class RunManifest:
def __init__(
self,
path: Path,
*,
harness: str,
model: str,
external_model: str,
task_dir: Path,
):
self.path = path
inputs = {
"harness": harness,
"model": model,
"external_model": external_model,
"task": str(task_dir),
}
if path.is_file():
existing = read_json_object(path)
if existing is None:
raise ValueError(f"invalid run manifest: {path}")
if existing.get("schema_version") != "1.0":
raise ValueError(f"unsupported run manifest schema: {path}")
if existing.get("inputs") != inputs:
raise ValueError(
f"run directory belongs to different pipeline inputs: {path.parent}"
)
stages = existing.get("stages")
if not isinstance(stages, dict):
raise ValueError(f"invalid stages in run manifest: {path}")
self.value = existing
self.value["status"] = "running"
self.value.pop("error", None)
self.save()
return
self.value: dict[str, Any] = {
"schema_version": "1.0",
"status": "running",
"inputs": inputs,
"stages": {},
}
self.save()
def save(self) -> None:
atomic_write_json(self.path, self.value)
def stage(self, name: str, status: str, **details: Any) -> None:
self.value["stages"][name] = {"status": status, **details}
self.save()
def complete(self, final_skill: Path) -> None:
self.value["status"] = "complete"
self.value["final_skill"] = str(final_skill)
self.save()
def fail(self, error: Exception) -> None:
self.value["status"] = "failed"
self.value["error"] = f"{type(error).__name__}: {error}"
self.save()