Files
Codex 82e2c91f42
Publish documentation / publish (push) Successful in 1m27s
feat: implement memory and evidence administration with guided repairs
Add PostgreSQL-backed memory, editable evidence with source review and activation, and human-approved archive repairs across the harness, API, and UI. Include migrations, deployment support, regression coverage, and validation documentation.

Refresh permissions from validated session roles so existing administrator logins can access newly deployed archive management features.
2026-09-10 10:31:34 +02:00

473 lines
21 KiB
Python

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