Files
ThothII/harness/tht/evidence/imports.py
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

236 lines
13 KiB
Python

"""Explicit acquisition and durable source comparisons; never implicit runtime refresh."""
import base64
import json
from .authoring import (
RestructureRequest,
_allocate_evidence_id,
_candidate_to_evidence,
normalize_source_text,
)
from .canonical import (
CuratedEvidence,
EvidenceProvenance,
ManualEvidenceProvenance,
dump_curated_markdown,
)
from .corpus.normalize import _decode
from .local_archive import ArchiveConflict, LocalEvidenceArchive, _atomic, _digest
MAX_DOCUMENTS = 200
MAX_TOTAL_BYTES = 100 * 1024 * 1024
def _origin(unit):
return unit.provenance.original if isinstance(unit.provenance, ManualEvidenceProvenance) else unit.provenance
def _records(archive):
path = archive.metadata / "sources.json"
if path.is_symlink():
raise ValueError("Source metadata must not use symlinks")
return json.loads(path.read_text()) if path.exists() else {}
def _write_records(archive, records):
_atomic(archive.metadata / "sources.json", json.dumps(records, ensure_ascii=False, sort_keys=True))
def reviews(archive):
return [{k: v for k, v in row.items() if k not in {"expected", "text", "writes", "approved"}}
for row in _records(archive).values()]
def acquisition_sources(cfg):
"""Local drafts and original files, plus configured read-only remote connectors."""
from .adapters import FilesystemEvidenceSource
from .sources import build_sources
root = cfg.evidence.local_archive_root / "evidence"
# Canonical acquired versions are immutable lineage, never new input documents.
patterns = [str(p.relative_to(root)) for folder in ("incoming", "source")
for p in sorted((root / folder).rglob("*.md"))
if not p.is_relative_to(root / "source/acquired")]
result = [FilesystemEvidenceSource(root, patterns=patterns)] if patterns else []
# Filesystem descriptors select the installation's local authoring tree after E2.
remote = cfg.evidence.model_copy(update={"source_root": None,
"sources": [s for s in cfg.evidence.sources if s.type != "filesystem"]})
result.extend(build_sources(remote, acquisition=True))
return result
def refresh(archive: LocalEvidenceArchive, sources, restructurer):
"""Acquire everything successfully before recording proposals. Missing is never deletion."""
with archive.operation():
state = archive._state()
if state.get("pending") or state.get("import_writes"):
raise ArchiveConflict("Complete the pending consolidation before refreshing sources")
records = _records(archive)
if any(r["status"] == "applying" for r in records.values()):
raise ArchiveConflict("Retry the pending source decision before refreshing")
files = archive._files(archive.evidence)
units = archive._units(files, allow_review=True)
documents, total = {}, 0
for adapter in sources:
for item in adapter.discover():
document = adapter.acquire(item)
total += len(document.content)
if len(documents) >= MAX_DOCUMENTS or total > MAX_TOTAL_BYTES:
raise ValueError("Source refresh exceeds the local acquisition limit")
relative = item.metadata.get("relative_path") if item.uri.startswith("file:") else None
identity = f"local:{relative}" if relative else item.uri
key = _digest(identity.encode())
if key in documents:
raise ValueError("Duplicate acquisition identity")
documents[key] = (document, relative)
reserved = set(units) | set(state["deleted_ids"])
for record in records.values():
reserved.update(u["id"] for u in record.get("proposed", []))
changed, unchanged = 0, 0
acquired = {}
for key, (document, relative) in documents.items():
text = normalize_source_text(_decode(document))
sha = "sha256:" + _digest(text.encode())
old = records.get(key)
stale = old and old["status"] == "review" and any(
p not in files or _digest(files[p]) != h for p, h in old["expected"].items())
if old and old["sha256"] == sha and not stale:
old["availability"] = "available"
unchanged += 1
continue
current = {i: pair for i, pair in units.items()
if (old and i in old["unit_ids"]) or
(_origin(pair[1]) and _origin(pair[1]).source_file == relative)}
# Seed imported E2 document identity without asking the model to recurate unchanged text.
if old is None and current and all(_origin(u).source_sha256 == sha for _, u in current.values()):
records[key] = {"id": key, "uri": document.source.uri, "legacy_file": relative,
"sha256": sha, "unit_ids": sorted(current), "status": "accepted",
"availability": "available", "revision": sha[7:], "proposed": []}
unchanged += 1
continue
source_file = f"source/acquired/{key}/{sha[7:]}.md"
request = RestructureRequest(source_file=source_file, source_sha256=sha,
normalized_text=text, previous_units=tuple(u for _, u in current.values()))
proposed = []
seen = set()
suppressed = relative in state["suppressed_sources"] or any(
p.startswith(f"source/acquired/{key}/") for p in state["suppressed_sources"]) or (old and (
old.get("suppressed", False) or any(i in state["deleted_ids"] for i in old["unit_ids"])))
for candidate in restructurer.restructure(request):
identity = candidate.existing_id
if identity in state["deleted_ids"]:
continue
if identity is not None and identity not in current:
raise ValueError("Source proposal refers to an unrelated Evidence identity")
if identity is None:
if suppressed:
continue # A model-created identifier cannot bypass a curated deletion.
identity = _allocate_evidence_id(candidate.title, reserved)
reserved.add(identity)
if identity in seen:
raise ValueError("Source proposal repeats an Evidence identity")
seen.add(identity)
unit = _candidate_to_evidence(candidate, identity, source_file, sha)
dump_curated_markdown(unit) # Refuse an uneditable proposal before saving any review.
if any(excerpt not in text for excerpt in unit.provenance.supporting_excerpts):
raise ValueError("Source proposal contains an excerpt absent from the acquired document")
proposed.append(unit.model_dump(mode="json"))
expected = {p: _digest(files[p]) for p, _ in current.values()}
row = {"id": key, "uri": document.source.uri, "legacy_file": relative,
"sha256": sha, "source_file": source_file, "text": text,
"status": "review", "availability": "available", "suppressed": bool(suppressed),
"unit_ids": sorted(current), "expected": expected,
"current": [u.model_dump(mode="json") for _, u in current.values()], "proposed": proposed,
"removed_ids": sorted(set(current) - seen)}
row["revision"] = _digest(json.dumps(row, sort_keys=True).encode())
records[key] = row
acquired[key] = {"source": document.source.model_dump(mode="json"),
"media_type": document.media_type, "raw_base64": base64.b64encode(document.content).decode()}
changed += 1
for key, row in records.items():
if key not in documents:
row["availability"] = "missing"
# No writes above: an access/model failure preserves every previous review and active unit.
for key, value in acquired.items():
directory = archive.metadata / "acquisitions" / key
if any(p.is_symlink() for p in [directory, directory.parent]):
raise ValueError("Acquisition metadata must not use symlinks")
_atomic(directory / f"{records[key]['sha256'][7:]}.json", json.dumps(value, ensure_ascii=False))
_write_records(archive, records)
return {"status": "succeeded", "counts": {"changed": changed, "unchanged": unchanged,
"review": sum(r["status"] == "review" for r in records.values())}}
def decide(archive, *, source_id, revision, decision, actor, activate):
if decision not in {"keep", "replace"} or not actor.strip():
raise ValueError("Choose keep or replace and supply a curator")
with archive.operation():
records = _records(archive)
row = records.get(source_id)
if row is None or row["revision"] != revision:
raise ArchiveConflict("Source comparison changed; refresh the page")
if row["status"] == "applying":
if row["decision"] != decision:
raise ArchiveConflict("Retry the saved source decision before changing it")
state = archive._state()
state.update(import_writes=row["writes"], approved_imports=row["approved"])
archive._write_state(state)
result = archive._consolidate(actor, activate)
else:
if row["status"] != "review":
raise ArchiveConflict("This source comparison was already decided")
state = archive._state()
if state.get("pending") or state.get("import_writes"):
raise ArchiveConflict("Complete the pending consolidation first")
files = archive._files(archive.evidence)
units = archive._units(files, allow_review=True)
if any(_digest(files[p]) != h if p in files else True for p, h in row["expected"].items()):
raise ArchiveConflict("Curated files changed since source review; refresh the source comparison")
selected = [CuratedEvidence.model_validate(v) for v in row["proposed"]] if decision == "replace" else [units[i][1] for i in row["unit_ids"]]
writes, approved = {}, {}
def write(path, content):
current = archive.evidence / path
writes[path] = {"before": _digest(current.read_bytes()) if current.exists() else None, "after": content}
if decision == "replace":
write(row["source_file"], row["text"])
for identity in row["unit_ids"]:
write(units[identity][0], None)
for value in selected:
if value.review_items:
raise ValueError("The proposal needs review; correct the input draft and refresh, or keep local content")
if decision == "keep":
origin = _origin(value)
if origin:
old = archive._snapshot(state["active"] or state["baseline"])
candidate = old / origin.source_file
if not candidate.is_file():
candidate = archive.evidence / origin.source_file
text = normalize_source_text(candidate.read_text())
if "sha256:" + _digest(text.encode()) != origin.source_sha256:
raise ValueError("The original source version is unavailable")
path = f"source/acquired/{source_id}/{origin.source_sha256[7:]}.md"
write(path, text)
origin = EvidenceProvenance(source_file=path, source_sha256=origin.source_sha256,
supporting_excerpts=origin.supporting_excerpts)
value = value.model_copy(update={"provenance": ManualEvidenceProvenance(declared_by=actor, original=origin)})
if value.id in units and value.id not in row["unit_ids"]:
raise ArchiveConflict("A proposed Evidence identity was created elsewhere")
path = units[value.id][0] if value.id in units and units[value.id][1].kind == value.kind else f"curated/{value.kind}/{value.id[9:]}.md"
if path in files and value.id not in units:
raise ArchiveConflict("The proposed file path is occupied")
write(path, dump_curated_markdown(value))
approved[value.id] = _digest(dump_curated_markdown(value).encode())
# Journal before applying files; retry never silently clobbers an external edit.
row.update(status="applying", decision=decision, decided_by=actor,
next_unit_ids=[u.id for u in selected], writes=writes, approved=approved)
_write_records(archive, records)
state.update(import_writes=writes, approved_imports=approved)
archive._write_state(state)
result = archive._consolidate(actor, activate)
if result["status"] != "active":
raise ValueError("Source decision was saved but not activated; retry")
row.update(status="accepted" if decision == "replace" else "kept", unit_ids=row["next_unit_ids"])
_write_records(archive, records)
return {"status": "succeeded", "counts": {"units": result["units"]}}