"""Persistent curated files, immutable consolidation candidates and explicit activation. The working tree is primary data. Core consumers use only active_snapshot(); an editor save or failed activation never switches that pointer. No Git/network/DWH I/O. """ from __future__ import annotations import fcntl import hashlib import json import os import tempfile from collections.abc import Callable from contextlib import contextmanager from pathlib import Path import yaml from .canonical import ( MAX_CURATED_FILE_BYTES, CuratedEvidence, EvidenceProvenance, ManualEvidenceProvenance, dump_curated_markdown, parse_curated_markdown, ) class ArchiveConflict(ValueError): """The curator must reconcile a concurrent change before replacing it.""" def _digest(value: bytes) -> str: return hashlib.sha256(value).hexdigest() def _content(unit: CuratedEvidence) -> str: return _digest( json.dumps( unit.model_dump(mode="json", exclude={"schema_version", "provenance"}), sort_keys=True, ensure_ascii=False, ).encode() ) def _atomic(path: Path, data: str) -> None: path.parent.mkdir(parents=True, exist_ok=True) fd, temporary = tempfile.mkstemp(prefix=".write-", dir=path.parent) try: with os.fdopen(fd, "w") as handle: handle.write(data) handle.flush() os.fsync(handle.fileno()) os.replace(temporary, path) finally: if os.path.exists(temporary): os.unlink(temporary) class LocalEvidenceArchive: def __init__(self, workspace_root: Path): self.root = workspace_root.resolve() self.evidence = self.root / "evidence" self.metadata = self.evidence / ".local" if self.evidence.is_symlink() or self.metadata.is_symlink(): raise ValueError("The local Evidence archive must use persistent regular directories") @contextmanager def operation(self): if any( path.is_symlink() for path in ( self.evidence, self.metadata, self.metadata / "snapshots", self.metadata / "state.yaml", ) ): raise ValueError("Evidence archive metadata must not use symlinks") self.root.mkdir(parents=True, exist_ok=True) lock = self.root / ".evidence-archive.lock" if lock.is_symlink(): raise ValueError("Evidence lock must not be a symlink") with lock.open("a") as handle: fcntl.flock(handle, fcntl.LOCK_EX) try: yield finally: fcntl.flock(handle, fcntl.LOCK_UN) def _state(self): path = self.metadata / "state.yaml" if not path.exists(): return { "schema_version": 1, "active": None, "pending": None, "baseline": None, "deleted_ids": [], "suppressed_sources": [], } value = yaml.safe_load(path.read_text()) if not isinstance(value, dict) or value.get("schema_version") != 1: raise ValueError("Unsupported local Evidence archive state") return value def _write_state(self, state): _atomic(self.metadata / "state.yaml", yaml.safe_dump(state, sort_keys=True)) def _snapshot(self, revision: str) -> Path: if ( not isinstance(revision, str) or len(revision) != 64 or any(c not in "0123456789abcdef" for c in revision) ): raise ValueError("Invalid Evidence snapshot identity") path = self.metadata / "snapshots" / revision if not path.is_dir() or path.is_symlink(): raise ValueError("Consolidated Evidence snapshot is missing") return path def active_snapshot(self) -> Path | None: """Only this immutable source is eligible for core consumption.""" with self.operation(): revision = self._state()["active"] return self._snapshot(revision) if revision else None def validate(self): """Check the working files without modifying them or changing active content.""" with self.operation(): units = self._units(self._files(self.evidence)) state = self._state() revision = state["pending"] or state["active"] or state.get("baseline") previous = ( self._units(self._files(self._snapshot(revision)), allow_review=True) if revision else {} ) for identity, (_, unit) in units.items(): old = previous.get(identity) if isinstance(unit.provenance, EvidenceProvenance) and ( old is None or _content(old[1]) == _content(unit) ): self._validate_document_source(unit.provenance) return len(units) def _files(self, root: Path): result = {} curated = root / "curated" if curated.is_symlink(): raise ValueError("Curated directory must not be a symlink") if not curated.is_dir(): raise ValueError(f"{curated}: curated archive is unavailable; absence is not deletion") for path in sorted(curated.rglob("*")): if path.is_symlink(): raise ValueError(f"{path}: symlinks are not supported in the curated archive") if path.suffix != ".md" or path.name.upper().startswith("README"): continue if path.stat().st_size > MAX_CURATED_FILE_BYTES: raise ValueError(f"{path}: Evidence file exceeds the size limit") result[path.relative_to(root).as_posix()] = path.read_bytes() return result def _units(self, files, *, allow_review=False): units = {} for relative, data in files.items(): try: unit = parse_curated_markdown(data.decode("utf-8"), path=Path(relative)) if unit.schema_version != 4: raise ValueError("Run evidence migrate before consolidating legacy units") if unit.id in units: raise ValueError(f"Duplicate Evidence identity {unit.id}") if unit.review_items and not allow_review: raise ValueError("Resolve review items before consolidation") units[unit.id] = (relative, unit) except ValueError as error: raise ValueError(f"{relative}: {error}") from error return units def initialize(self): """Capture migrated/refined content before edits, without activating review items.""" with self.operation(): state = self._state() if state.get("baseline") or state["active"] or state["pending"]: return files = self._files(self.evidence) self._units(files, allow_review=True) source = self.evidence / "source" if source.is_symlink(): raise ValueError("Source symlinks are not supported") if source.exists(): for path in sorted(source.rglob("*")): if path.is_symlink(): raise ValueError("Source symlinks are not supported") if path.is_file(): if path.stat().st_size > MAX_CURATED_FILE_BYTES: raise ValueError(f"{path}: source exceeds the size limit") files[path.relative_to(self.evidence).as_posix()] = path.read_bytes() revision = _digest(b"".join(k.encode() + b"\0" + v for k, v in sorted(files.items()))) snapshots = self.metadata / "snapshots" snapshots.mkdir(parents=True, exist_ok=True) destination = snapshots / revision if not destination.exists(): with tempfile.TemporaryDirectory(prefix=".baseline-", dir=snapshots) as tmp: candidate = Path(tmp) / "snapshot" (candidate / "curated").mkdir(parents=True) for relative, data in files.items(): path = candidate / relative path.parent.mkdir(parents=True, exist_ok=True) path.write_bytes(data) os.replace(candidate, destination) state["baseline"] = revision self._write_state(state) def consolidate(self, *, actor: str, activate: Callable[[Path], None] | None = None): """Persist one valid candidate; optionally activate it through the existing index stage. A missing callback deliberately reports pending_activation, never success. A retry of unchanged files reuses the same candidate after an index failure. """ if not actor.strip(): raise ValueError("A curator identity is required") with self.operation(): return self._consolidate(actor, activate) def _consolidate(self, actor, activate): state = self._state() self._finish_import(state) self._finish_normalization(state) original_files = self._files(self.evidence) units = self._units(original_files) previous_revision = state["pending"] or state["active"] or state.get("baseline") previous_root = self._snapshot(previous_revision) if previous_revision else None previous = ( self._units(self._files(previous_root), allow_review=True) if previous_root else {} ) deleted = set(state["deleted_ids"]) | (previous.keys() - units.keys()) suppressed = set(state["suppressed_sources"]) for identity in previous.keys() - units.keys(): provenance = previous[identity][1].provenance source = ( provenance.original if isinstance(provenance, ManualEvidenceProvenance) else provenance ) if source: suppressed.add(source.source_file) files = {} for identity, (relative, unit) in units.items(): old = previous.get(identity) changed = old is not None and _content(old[1]) != _content(unit) approved = state.get("approved_imports", {}).get(identity) == _digest( dump_curated_markdown(unit).encode() ) if approved: pass # An explicit source decision authorized this exact content and provenance. elif changed or ( isinstance(unit.provenance, ManualEvidenceProvenance) and (old is None or unit.provenance.declared_by == "local curator") ): origin = old[1].provenance if old else unit.provenance origin = origin.original if isinstance(origin, ManualEvidenceProvenance) else origin unit = unit.model_copy( update={ "provenance": ManualEvidenceProvenance(declared_by=actor, original=origin) } ) elif old and unit.provenance != old[1].provenance: raise ArchiveConflict( f"{relative}: provenance is managed; edit the content instead" ) if isinstance(unit.provenance, EvidenceProvenance): self._validate_document_source(unit.provenance) files[relative] = dump_curated_markdown(unit).encode() origin = ( unit.provenance.original if isinstance(unit.provenance, ManualEvidenceProvenance) else unit.provenance ) if origin: candidates = [previous_root / origin.source_file] if previous_root else [] candidates.append(self.evidence / origin.source_file) from .authoring import normalize_source_text for source in candidates: if source.is_file() and not source.is_symlink(): if source.stat().st_size > MAX_CURATED_FILE_BYTES: raise ValueError(f"{source}: source exceeds the size limit") raw = source.read_bytes() if ( "sha256:" + _digest(normalize_source_text(raw.decode()).encode()) == origin.source_sha256 ): if origin.source_file in files and files[origin.source_file] != raw: raise ArchiveConflict( "Different source revisions require explicit source resolution" ) files[origin.source_file] = raw break else: raise ValueError( f"{origin.source_file}: the recorded original document is unavailable" ) # Record current declarations and original document lineage distinctly. manifest = { "schema_version": 1, "units": { identity: { "file": relative, "content_hash": _content(parse_curated_markdown(files[relative].decode())), "provenance": parse_curated_markdown( files[relative].decode() ).provenance.model_dump(mode="json"), } for identity, (relative, _) in sorted(units.items()) }, "deleted_ids": sorted(deleted), "suppressed_sources": sorted(suppressed), } files["local-manifest.yaml"] = yaml.safe_dump( manifest, allow_unicode=True, sort_keys=True ).encode() revision = _digest( b"".join(path.encode() + b"\0" + data + b"\0" for path, data in sorted(files.items())) ) snapshots = self.metadata / "snapshots" snapshots.mkdir(parents=True, exist_ok=True) destination = snapshots / revision if not destination.exists(): with tempfile.TemporaryDirectory(prefix=".candidate-", dir=snapshots) as tmp: candidate = Path(tmp) / "snapshot" (candidate / "curated").mkdir(parents=True) for relative, data in files.items(): path = candidate / relative path.parent.mkdir(parents=True, exist_ok=True) path.write_bytes(data) os.replace(candidate, destination) if self._files(self.evidence) != original_files: raise ArchiveConflict( "Evidence files changed during consolidation; retry with the current files" ) # The candidate exists first. Pending state makes every subsequent interruption recoverable. state.update( pending=revision, deleted_ids=sorted(deleted), suppressed_sources=sorted(suppressed), normalization={relative: _digest(data) for relative, data in original_files.items()}, ) state.pop("approved_imports", None) self._write_state(state) self._finish_normalization(state) if activate is None: return { "status": "pending_activation", "revision": revision, "snapshot": str(destination), } activate(destination) state.update(active=revision, pending=None) self._write_state(state) return { "status": "active", "revision": revision, "snapshot": str(destination), "units": len(units), "deleted": len(previous.keys() - units.keys()), } def _finish_import(self, state): """Replay an explicit source decision, rejecting intervening external edits.""" writes = state.get("import_writes") if writes is None: return for relative, change in writes.items(): path = self.evidence / relative if not relative.startswith(("curated/", "source/acquired/")) or ".." in Path(relative).parts: raise ValueError("Invalid import destination") if any(p.is_symlink() for p in [path, *path.parents] if p != self.root.parent): raise ValueError("Import destinations must not use symlinks") current = _digest(path.read_bytes()) if path.exists() else None after = _digest(change["after"].encode()) if change["after"] is not None else None if current not in (change["before"], after): raise ArchiveConflict("Evidence changed during a source decision; restore or review the file") for relative, change in writes.items(): path = self.evidence / relative if change["after"] is None: path.unlink(missing_ok=True) else: _atomic(path, change["after"]) del state["import_writes"] self._write_state(state) def _finish_normalization(self, state): """Replay interrupted managed writes only where the user's bytes are unchanged.""" if "normalization" not in state: return snapshot = self._snapshot(state["pending"]) current = self._files(self.evidence) for relative, original_hash in state["normalization"].items(): if relative in current and _digest(current[relative]) == original_hash: _atomic(self.evidence / relative, (snapshot / relative).read_text()) _atomic( self.evidence / "local-manifest.yaml", (snapshot / "local-manifest.yaml").read_text() ) del state["normalization"] self._write_state(state) def _validate_document_source(self, provenance): from .authoring import _normalize, normalize_source_text path = self.evidence / provenance.source_file if ( path.is_symlink() or not path.is_file() or not path.resolve().is_relative_to(self.evidence.resolve()) ): raise ValueError(f"{provenance.source_file}: source document is unavailable") if path.stat().st_size > MAX_CURATED_FILE_BYTES: raise ValueError(f"{provenance.source_file}: source exceeds the size limit") source = normalize_source_text(path.read_text()) if "sha256:" + _digest(source.encode()) != provenance.source_sha256: raise ArchiveConflict( f"{provenance.source_file}: source changed; explicit source refresh is required" ) if any(_normalize(excerpt) not in source for excerpt in provenance.supporting_excerpts): raise ValueError(f"{provenance.source_file}: supporting excerpt is missing") def save(self, value: CuratedEvidence, *, expected_revision: str | None, actor: str): """Explicit workflow correction with optimistic concurrency; activation is a separate boundary.""" if not actor.strip(): raise ValueError("A curator identity is required") with self.operation(): return self._save(value, expected_revision=expected_revision, actor=actor) def _save(self, value, *, expected_revision, actor): """Save while the caller holds operation(), including workflow receipt recovery.""" self._finish_normalization(self._state()) units = self._units(self._files(self.evidence)) existing = units.get(value.id) if (existing is None) != (expected_revision is None): raise ArchiveConflict("Evidence was created or removed since review") if existing and _content(existing[1]) != expected_revision: raise ArchiveConflict("Evidence changed since review") relative = ( existing[0] if existing else f"curated/{value.kind}/{value.id.removeprefix('evidence:')}.md" ) if existing and existing[1].kind != value.kind: raise ValueError("An update cannot change the Evidence kind") if existing: value = value.model_copy(update={"provenance": existing[1].provenance}) _atomic(self.evidence / relative, dump_curated_markdown(value)) return self._consolidate(actor, None) def get(self, identity: str): with self.operation(): relative, unit = self._units(self._files(self.evidence))[identity] return {"unit": unit, "revision": _content(unit), "path": str(self.evidence / relative)} def remove(self, identity: str, *, expected_revision: str, actor: str): if not actor.strip(): raise ValueError("A curator identity is required") with self.operation(): self._finish_normalization(self._state()) relative, unit = self._units(self._files(self.evidence))[identity] if _content(unit) != expected_revision: raise ArchiveConflict("Evidence changed since review") (self.evidence / relative).unlink() return self._consolidate(actor, None)