"""
artifact_store — Concrete ArtifactStore and SchemaValidator implementations.
Extracted from rca_reasoning_orchestrator.py. The parent module re-exports
both classes for backward-compatible imports.
"""
from __future__ import annotations
import json
import os
import tempfile
from pathlib import Path
from typing import Any, Dict, List
[docs]
JsonDict = Dict[str, Any]
[docs]
class NoOpSchemaValidator:
"""Permissive :class:`SchemaValidator` that only enforces JSON-object shape.
Implements all three validator styles but performs no schema checks — useful
as a default when schema validation is not wired. Every artifact is reported
as valid provided it is a JSON object.
"""
[docs]
def validate(self, artifact_name: str, payload: JsonDict) -> None:
"""Legacy style: raise ``TypeError`` unless *payload* is a JSON object."""
if not isinstance(payload, dict):
raise TypeError(f"{artifact_name} must be a JSON object")
[docs]
def validate_artifact(self, artifact_name: str, payload: JsonDict) -> JsonDict:
"""Per-artifact style: shape-check *payload* and return an ``ok`` report."""
self.validate(artifact_name, payload)
return {
"ok": True,
"issues": [],
"artifact": artifact_name,
"mode": "noop",
}
[docs]
def validate_run_bundle(self, **kwargs: Any) -> JsonDict:
"""Bundle style: shape-check each non-None keyword artifact and return ``ok``."""
for artifact_name, payload in kwargs.items():
if payload is None:
continue
self.validate(artifact_name, payload)
return {
"ok": True,
"issues": [],
"artifact": "bundle",
"mode": "noop",
}
[docs]
class FileArtifactStore:
""":class:`ArtifactStore` that persists artifacts as JSON files on disk.
Each artifact is written atomically to ``<root_dir>/<run_id>/<name>.json``
so a reader never observes a partial write.
"""
def __init__(self, root_dir: str | Path):
[docs]
self.root_dir = Path(root_dir)
[docs]
def save(self, run_id: str, artifact_name: str, payload: JsonDict) -> str:
"""Persist a single-object artifact and return its file path."""
return self._write_atomic(run_id, artifact_name, payload)
[docs]
def save_list(self, run_id: str, artifact_name: str, payload: List[JsonDict]) -> str:
"""Persist a list-valued artifact and return its file path."""
return self._write_atomic(run_id, artifact_name, payload)
[docs]
def load(self, run_id: str, artifact_name: str) -> Any:
"""Load and return a previously saved artifact, or None if absent."""
path = self.root_dir / run_id / f"{artifact_name}.json"
if not path.exists():
return None
return json.loads(path.read_text(encoding="utf-8"))
[docs]
def is_run_complete(self, run_id: str) -> bool:
"""Return True only if run_status.json exists and run_complete is True.
External callers (notebooks, replay scripts) should check this before
loading or re-validating artifacts from a run directory. A run that
crashed mid-pipeline will have run_complete=False (or no run_status.json
at all), and its artifacts should not be treated as authoritative.
"""
status = self.load(run_id, "run_status")
return bool((status or {}).get("run_complete", False))
[docs]
def _write_atomic(self, run_id: str, artifact_name: str, payload: Any) -> str:
"""Write payload to <run_dir>/<artifact_name>.json via temp-file + rename.
The rename is atomic on POSIX and near-atomic on Windows (py3.3+).
A reader can never observe a partial write.
"""
run_dir = self.root_dir / run_id
run_dir.mkdir(parents=True, exist_ok=True)
target = run_dir / f"{artifact_name}.json"
text = json.dumps(payload, indent=2, default=str)
fd, tmp_path = tempfile.mkstemp(dir=run_dir, prefix=f".{artifact_name}_", suffix=".tmp")
try:
with os.fdopen(fd, "w", encoding="utf-8") as fh:
fh.write(text)
Path(tmp_path).replace(target)
except Exception:
try:
os.unlink(tmp_path)
except OSError:
pass
raise
return str(target)