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

261 lines
14 KiB
Python

"""Closed session repair choices, durable receipts and authoritative archive activation.
Memory commits its receipt with the card. Evidence records the choice before writing
files and recovers by comparing the approved result, never by repeating a stale write.
"""
import json
from pathlib import Path
from typing import Literal
from pydantic import BaseModel, ConfigDict, Field
from sqlalchemy import text
from tht.evidence.canonical import CuratedEvidence
from tht.evidence.local_archive import LocalEvidenceArchive, _content, _digest
from tht.memory.models import CardInput, MemoryConflict, MemoryNotFound
from tht.memory.review import digest
from tht.phase import current_phase, effective_decisions
class RepairOption(BaseModel):
model_config = ConfigDict(extra="forbid", str_strip_whitespace=True)
id: str = Field(pattern=r"^[a-zA-Z0-9_-]{1,80}$")
label: str = Field(min_length=1, max_length=1000)
archive: Literal["memory", "evidence"]
target_id: str = Field(min_length=1, max_length=100)
revision: str = Field(min_length=1, max_length=100)
content: dict
class RepairProposal(BaseModel):
model_config = ConfigDict(extra="forbid", str_strip_whitespace=True)
reason: str = Field(min_length=1, max_length=10000)
options: list[RepairOption] = Field(min_length=1, max_length=5)
def _context(snapshot):
return digest({"phase": current_phase(snapshot), "question": snapshot.artifacts.get("question"),
"decisions": [d.model_dump(mode="json") for d in effective_decisions(snapshot)]})
def _authorize(service, snapshot):
service._session(snapshot)
if snapshot.manifest.status in {"finalized", "archived"}:
raise MemoryConflict("Archive repair requires an open session")
def _archive(cfg, workspace):
root = cfg.evidence.local_archive_root if cfg.evidence else None
if not root or Path(root).name != workspace:
raise MemoryConflict("Local Evidence is unavailable in this workspace")
return LocalEvidenceArchive(root)
def _params(repo, snapshot, repair_id):
return {"w": repo.workspace_id, "s": snapshot.manifest.id, "r": repair_id}
def target(service, snapshot, cfg, archive, identity):
"""Read complete current content for a proposal, bound to the source session."""
_authorize(service, snapshot)
if archive == "memory":
value = service.repository.get(identity)
return {"revision": value.revision, "content": CardInput.model_validate(
value.model_dump(include=set(CardInput.model_fields))).model_dump(mode="json")}
if archive != "evidence":
raise ValueError("Unknown archive")
try:
value = _archive(cfg, service.repository.workspace_id).get(identity)
except KeyError:
raise MemoryNotFound("Evidence was not found in this workspace") from None
return {"revision": value["revision"], "content": value["unit"].model_dump(mode="json")}
def _read(repo, snapshot, repair_id):
with repo.transaction() as connection:
value = connection.execute(text("SELECT data FROM thoth_memory.archive_repairs "
"WHERE workspace_id=:w AND session_id=:s AND repair_id=:r"),
_params(repo, snapshot, repair_id)).scalar_one_or_none()
if value is None:
raise MemoryNotFound("Repair was not found in this session and workspace")
return value
def _write(repo, snapshot, repair_id, value):
with repo.transaction() as connection:
connection.execute(text("INSERT INTO thoth_memory.archive_repairs "
"(workspace_id,session_id,repair_id,data) VALUES (:w,:s,:r,CAST(:data AS jsonb)) "
"ON CONFLICT (workspace_id,session_id,repair_id) DO UPDATE SET data=EXCLUDED.data"),
{**_params(repo, snapshot, repair_id), "data": json.dumps(value)})
def prepare(service, snapshot, cfg, proposal: RepairProposal):
_authorize(service, snapshot)
if (len({o.id for o in proposal.options}) != len(proposal.options)
or any(o.id in {"reject", "continue"} for o in proposal.options)):
raise ValueError("Repair choices must have distinct identities")
items = []
for option in proposal.options:
item = option.model_dump(mode="json")
if option.archive == "memory":
before = service.repository.get(option.target_id)
if before.revision != option.revision:
raise MemoryConflict("Memory changed; prepare new choices")
item["before"] = before.model_dump(mode="json")
item["content"] = CardInput.model_validate(option.content).model_dump(mode="json")
else:
archive = _archive(cfg, service.repository.workspace_id)
with archive.operation():
state = archive._state()
if not state["active"] or state["pending"] or state.get("import_writes"):
raise MemoryConflict("Consolidate Evidence before proposing a session repair")
files = archive._files(archive.evidence)
if files != archive._files(archive._snapshot(state["active"])):
raise MemoryConflict("Consolidate external Evidence edits before session repair")
existing = archive._units(files).get(option.target_id)
if existing is None:
raise MemoryNotFound("Evidence was not found in this workspace")
relative, before = existing
if _content(before) != option.revision:
raise MemoryConflict("Evidence changed; prepare new choices")
value = CuratedEvidence.model_validate(option.content)
if (value.schema_version != 4 or value.id != before.id
or value.kind != before.kind or value.review_items):
raise ValueError("Repair must preserve Evidence identity/kind and resolve review items")
# Curators change knowledge, not the source history supplied by the archive.
value = value.model_copy(update={"provenance": before.provenance})
item.update(content=value.model_dump(mode="json"),
before=before.model_dump(mode="json"), file=relative,
other_files={p: _digest(v) for p, v in files.items() if p != relative})
items.append(item)
value = {"reason": proposal.reason, "options": items, "context": _context(snapshot),
"choice": None, "status": "proposed", "saved": False, "indexed": False}
repair_id = digest({"workspace": service.repository.workspace_id,
"session": snapshot.manifest.id, **value})
with service.repository.operation() as repo:
try:
_read(repo, snapshot, repair_id)
except MemoryNotFound:
_write(repo, snapshot, repair_id, value)
return show(service, snapshot, cfg, repair_id)
def show(service, snapshot, cfg, repair_id):
_authorize(service, snapshot)
value = _read(service.repository, snapshot, repair_id)
# A completed receipt describes history; current eligibility is checked afresh.
if value.get("saved"):
option = next(o for o in value["options"] if o["id"] == value["choice"])
try:
if option["archive"] == "memory":
current = service.repository.get(option["target_id"])
same = current.revision == value["saved_revision"]
indexed = same and current.indexed
else:
archive = _archive(cfg, service.repository.workspace_id)
with archive.operation():
units = archive._units(archive._files(archive.evidence))
current = units[option["target_id"]][1]
same = _content(current) == _content(
CuratedEvidence.model_validate(option["content"]))
state = archive._state()
active = archive._units(archive._files(archive._snapshot(state["active"]))) \
if state["active"] else {}
indexed = same and option["target_id"] in active and \
active[option["target_id"]][1] == current
value.update(indexed=bool(indexed), status="superseded" if not same else
"active" if indexed else "pending_activation")
except (MemoryNotFound, KeyError, ValueError):
value.update(indexed=False, status="superseded")
return {**value, "repair_id": repair_id, "can_apply": service.principal.is_admin}
def list_repairs(service, snapshot):
_authorize(service, snapshot)
with service.repository.transaction() as connection:
rows = connection.execute(text("SELECT repair_id,data->>'status' AS recorded_status, "
"data->>'reason' AS reason FROM thoth_memory.archive_repairs "
"WHERE workspace_id=:w AND session_id=:s ORDER BY created_at"),
{"w": service.repository.workspace_id, "s": snapshot.manifest.id}).mappings().all()
return {"repairs": [dict(row) for row in rows]}
def apply(service, snapshot, cfg, repair_id, choice, *, activate=None):
_authorize(service, snapshot)
if choice != "reject":
service._admin()
with service.repository.operation() as repo:
value = _read(repo, snapshot, repair_id)
if value["choice"] is not None and value["choice"] != choice:
raise MemoryConflict("This repair already has a different recorded choice")
if value["choice"] is None and value["context"] != _context(snapshot):
raise MemoryConflict("Session decisions changed; reformulate the repair")
if choice == "reject":
value.update(choice=choice, status="rejected", actor=service.principal.subject)
_write(repo, snapshot, repair_id, value)
else:
option = next((o for o in value["options"] if o["id"] == choice), None)
if option is None:
raise ValueError("Select one of the reviewed repair choices")
if option["archive"] == "memory":
with repo.transaction():
if not value["saved"]:
current = repo.get(option["target_id"])
if current.revision != option["revision"]:
raise MemoryConflict("Memory changed; reformulate the repair")
repo.save(CardInput.model_validate(option["content"]),
card_id=option["target_id"])
value.update(choice=choice, saved=True, status="pending_activation",
actor=service.principal.subject,
saved_revision=repo.get(option["target_id"]).revision)
_write(repo, snapshot, repair_id, value)
elif repo.get(option["target_id"]).revision != value["saved_revision"]:
raise MemoryConflict("The repaired Memory was changed again; do not replay it")
result = service._propagate(repo, option["target_id"])
value.update(indexed=result["indexed"], status="active" if result["indexed"] else
"pending_activation")
_write(repo, snapshot, repair_id, value)
else:
_apply_evidence(service, repo, snapshot, cfg, repair_id, choice, value, option,
activate)
return show(service, snapshot, cfg, repair_id)
def _apply_evidence(service, repo, snapshot, cfg, repair_id, choice, receipt, option, activate):
archive = _archive(cfg, repo.workspace_id)
proposed = CuratedEvidence.model_validate(option["content"])
with archive.operation():
state = archive._state()
if state.get("import_writes") or (receipt["choice"] is None and state["pending"]):
raise MemoryConflict("Finish the pending Evidence consolidation before this repair")
files = archive._files(archive.evidence)
if {p: _digest(v) for p, v in files.items() if p != option["file"]} != option["other_files"]:
raise MemoryConflict("Other Evidence files changed; reconcile them before retrying")
existing = archive._units(files).get(option["target_id"])
if existing is None:
raise MemoryConflict("Evidence was removed after review")
_, current = existing
already_written = _content(current) == _content(proposed) and receipt["choice"] == choice
if not already_written and (receipt["saved"] or _content(current) != option["revision"]
or current.provenance != proposed.provenance):
raise MemoryConflict("Evidence changed; the approved correction cannot overwrite it")
# Commit approval before touching the filesystem; a restart can recover only this choice.
receipt.update(choice=choice, actor=service.principal.subject, status="applying")
_write(repo, snapshot, repair_id, receipt)
try:
if not already_written:
archive._save(proposed, expected_revision=option["revision"],
actor=service.principal.subject)
receipt.update(saved=True, status="pending_activation")
_write(repo, snapshot, repair_id, receipt)
archive._consolidate(service.principal.subject, activate)
except Exception: # noqa: BLE001 - durable approval covers file/index interruption.
receipt.update(status="pending_activation" if receipt["saved"] else "applying",
indexed=False)
_write(repo, snapshot, repair_id, receipt)
return
receipt.update(indexed=activate is not None,
status="active" if activate else "pending_activation")
_write(repo, snapshot, repair_id, receipt)