"""Artifact board: pending drafts become downstream inputs only after validation.""" from __future__ import annotations from typing import Any from uuid import uuid4 from deerflow.persistence.report_collaboration.base import ReportCollaborationStore from deerflow.persistence.report_collaboration.codec import iso_now class ArtifactBoard: """Persistence wrapper. QualityGate (RC-BE-010) will sit in front of ``validate``.""" def __init__(self, store: ReportCollaborationStore) -> None: self._store = store async def submit( self, *, session_id: str, run_id: str, node_id: str, attempt: int, producer_node_run_id: str, artifact_type: str, content: dict[str, Any], requirement_revision: int, summary: str | None = None, validation_status: str = "pending", ) -> dict[str, Any]: return await self._store.save_artifact( { "id": f"rca_{uuid4().hex[:16]}", "session_id": session_id, "run_id": run_id, "node_id": node_id, "attempt": attempt, "producer_node_run_id": producer_node_run_id, "artifact_type": artifact_type, "schema_version": int(content.get("schema_version") or 1), "requirement_revision": requirement_revision, "content": content, "summary": summary, "validation_status": validation_status, "superseded_by": None, "created_at": iso_now(), } ) async def mark(self, artifact: dict[str, Any], *, validation_status: str, superseded_by: str | None = None) -> dict[str, Any]: payload = dict(artifact) payload["validation_status"] = validation_status if superseded_by is not None: payload["superseded_by"] = superseded_by return await self._store.save_artifact(payload) async def supersede_previous(self, *, run_id: str, node_id: str, new_artifact_id: str, attempt: int | None = None) -> list[dict[str, Any]]: """Hide older attempts of the same node. Same-attempt siblings stay visible.""" updated: list[dict[str, Any]] = [] for item in await self._store.list_artifacts_full(run_id): if item.get("node_id") != node_id or item.get("id") == new_artifact_id: continue if item.get("superseded_by"): continue if attempt is not None and int(item.get("attempt") or 0) >= attempt: continue updated.append(await self.mark(item, validation_status="superseded", superseded_by=new_artifact_id)) return updated async def validated_for(self, run_id: str, artifact_ids: list[str] | None = None) -> list[dict[str, Any]]: rows = await self._store.list_artifacts_full(run_id) out = [ item for item in rows if item.get("validation_status") == "validated" and not item.get("superseded_by") ] if artifact_ids is None: return out allowed = set(artifact_ids) return [item for item in out if item["id"] in allowed]