Source code for src.dackar.RCA.orchestrators.artifact_store

"""
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)