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

314 lines
12 KiB
Python

"""Reviewer-edited additions/updates, grounded in effective approved session decisions."""
import hashlib
import json
from uuid import NAMESPACE_URL, uuid5
from pydantic import BaseModel, ConfigDict, Field
from sqlalchemy import text
from tht.phase import current_phase, effective_decisions
from .models import CardInput, MemoryConflict
from .solved import _build_solved_snapshot
APPROVED_SOURCES = {
"concept_clarified",
"join_modified",
"column_corrected",
"concept_formula_approved",
"cte_corrected",
"cte_approved",
"sql_approved",
}
class Proposal(BaseModel):
model_config = ConfigDict(extra="forbid")
id: str = Field(pattern=r"^[a-zA-Z0-9_-]{1,80}$")
source_seqs: list[int] = Field(min_length=1, max_length=20)
card: CardInput
target_id: str | None = None
target_revision: str | None = None
reason: str = Field(min_length=1, max_length=10000)
class Selection(BaseModel):
model_config = ConfigDict(extra="forbid")
id: str
card: CardInput
class ReviewResponse(BaseModel):
model_config = ConfigDict(extra="forbid")
summary_id: str
items: list[Selection] = Field(max_length=20)
def digest(value):
return hashlib.sha256(
json.dumps(value, sort_keys=True, ensure_ascii=False).encode()
).hexdigest()
def context_hash(snapshot):
return digest(
{
"decisions": [
d.model_dump(mode="json")
for d in effective_decisions(snapshot)
if d.type != "memory_summary_reviewed"
],
"proposals": snapshot.artifacts.get("memory_proposals"),
"sql": snapshot.artifacts.get("sql_final"),
}
)
def solved_snapshot(snapshot):
linking = json.loads(snapshot.artifacts.get("schema_linking", "{}"))
tables = {
c["name"]
for c in linking.get("candidates", [])
if c.get("kind") == "table" and c.get("decision") == "promoted"
}
return _build_solved_snapshot(snapshot, tables)
def validate_proposals(snapshot, raw):
if not isinstance(raw, list) or len(raw) > 20:
raise ValueError("Memory proposals must be a list of at most 20 cards")
proposals = [Proposal.model_validate(p) for p in raw]
effective = {d.seq: d for d in effective_decisions(snapshot)}
if len({p.id for p in proposals}) != len(proposals):
raise ValueError("Memory proposal identities must be unique")
for p in proposals:
sources = [effective.get(seq) for seq in p.source_seqs]
if any(d is None or d.type not in APPROVED_SOURCES for d in sources):
raise MemoryConflict("Memory proposals require effective approved source decisions")
if p.card.family == "explained_error" and not any(d.rationale.strip() for d in sources):
raise MemoryConflict("Explained errors require an approved explanation")
if p.card.family == "solved_question":
solved = solved_snapshot(snapshot)
if p.card.sql != solved.metadata["sql"]:
raise MemoryConflict("Exemplar SQL must match the current approved solution")
if bool(p.target_id) != bool(p.target_revision):
raise ValueError("Updates require the identity and revision of the card being replaced")
return proposals
def prepare(service, snapshot):
service._session(snapshot)
if current_phase(snapshot) != 8 or snapshot.manifest.status in {"finalized", "archived"}:
raise MemoryConflict("The Memory summary is reviewed at the end of F8")
with service.repository.transaction() as connection:
receipt = (
connection.execute(
text(
"SELECT summary_id,result FROM thoth_memory.reviews "
"WHERE workspace_id=:w AND session_id=:s AND result->>'context_hash'=:h "
"ORDER BY created_at DESC LIMIT 1"
),
{
"w": service.repository.workspace_id,
"s": snapshot.manifest.id,
"h": context_hash(snapshot),
},
)
.mappings()
.first()
)
if receipt:
return {
"reviewed": True,
"summary_id": receipt["summary_id"],
"saved": len(receipt["result"]["saved"]),
}
raw = json.loads(snapshot.artifacts.get("memory_proposals", "[]"))
proposals = validate_proposals(snapshot, raw)
covered = {seq for p in proposals for seq in p.source_seqs}
for d in service.promotions(snapshot):
if d["decision_seq"] not in covered:
proposals.append(
Proposal(
id=f"decision-{d['decision_seq']}",
source_seqs=[d["decision_seq"]],
reason="Reusable domain clarification",
card=CardInput(
family="domain_clarification",
subject=d["subject"],
detail=d["detail"],
rationale=d["rationale"],
question=d["question_context"],
scope=service.repository.workspace_id,
concepts=[d["subject"]],
),
)
)
if not any(p.card.family == "solved_question" for p in proposals):
solved = solved_snapshot(snapshot)
approved = [d for d in effective_decisions(snapshot) if d.type == "sql_approved"][-1]
proposals.append(
Proposal(
id="solved-question",
source_seqs=[approved.seq],
reason="Approved solution, for consultation in future questions",
card=CardInput(
family="solved_question",
subject=solved.title,
question=solved.content,
sql=solved.metadata["sql"],
scope=service.repository.workspace_id,
dependencies=[
{
"database": snapshot.manifest.database,
"schema_name": snapshot.manifest.db_schema,
"table": table,
}
for table in solved.metadata["tables"]
],
),
)
)
if len(proposals) > 20:
raise MemoryConflict("Reduce the final Memory summary to at most 20 cards")
items = []
targets = set()
content = {}
aliases = {}
for p in proposals:
before = None
if p.target_id:
old = service.repository.get(p.target_id)
if old.revision != p.target_revision:
raise MemoryConflict(
"A Memory card changed; refresh the proposed update before review"
)
if p.target_id in targets:
raise MemoryConflict("Propose only one update to each Memory card")
targets.add(p.target_id)
before = old.model_dump(mode="json")
key = digest(p.card.model_dump(mode="json"))
duplicate = service.repository.exact_match(p.card)
if duplicate and (not p.target_id or duplicate.id == p.target_id):
aliases["proposal:" + p.id] = duplicate.id
continue
if not p.target_id and key in content:
aliases["proposal:" + p.id] = "proposal:" + content[key]
continue
content[key] = p.id
items.append({**p.model_dump(mode="json"), "before": before})
for item in items:
for link in item["card"]["links"]:
link["target_id"] = aliases.get(link["target_id"], link["target_id"])
return {
"summary_id": digest(
{
"items": items,
"decisions": [d.model_dump(mode="json") for d in effective_decisions(snapshot)],
}
),
"items": items,
}
def apply(service, snapshot, response: ReviewResponse):
service._session(snapshot)
request_hash = digest(response.model_dump(mode="json"))
session_id = snapshot.manifest.id
with service.repository.operation() as repo:
with repo.transaction() as connection:
params = {"w": repo.workspace_id, "s": session_id, "id": response.summary_id}
receipt = (
connection.execute(
text(
"SELECT request_hash,result FROM thoth_memory.reviews "
"WHERE workspace_id=:w AND session_id=:s AND summary_id=:id"
),
params,
)
.mappings()
.first()
)
if receipt:
if receipt["request_hash"] != request_hash:
raise MemoryConflict(
"This summary was already reviewed with different selections"
)
saved = receipt["result"]["saved"]
else:
# Resolve the same locked repository for preview and optimistic update checks.
original = service.repository
service.repository = repo
try:
summary = prepare(service, snapshot)
finally:
service.repository = original
if summary.get("reviewed") or summary["summary_id"] != response.summary_id:
raise MemoryConflict("The Memory summary changed; review it again")
choices = {choice.id: choice for choice in response.items}
candidates = {p["id"]: p for p in summary["items"]}
if len(choices) != len(response.items) or choices.keys() - candidates.keys():
raise ValueError("Memory review contains duplicate or unknown choices")
selected = validate_proposals(
snapshot,
[
{
**{k: v for k, v in candidates[key].items() if k != "before"},
"card": choice.card.model_dump(mode="json"),
}
for key, choice in choices.items()
],
)
identities = {
p.id: p.target_id
or "mem-"
+ str(uuid5(NAMESPACE_URL, f"thothii:{repo.workspace_id}:{session_id}:{p.id}"))
for p in selected
}
saved = []
for p in selected:
# Create/update every selected card before inserting links among new cards.
identity = repo.save(
p.card.model_copy(update={"links": []}),
card_id=p.target_id,
new_id=identities[p.id],
source_key=f"review:{session_id}:{p.id}",
session_id=session_id,
decision_seq=p.source_seqs[0],
)
saved.append({"id": identity, "proposal_id": p.id})
for p in selected:
value = p.card.model_dump(mode="json")
for link in value["links"]:
if link["target_id"].startswith("proposal:"):
target = link["target_id"].removeprefix("proposal:")
if target not in identities:
raise ValueError("Select the linked card or remove its link")
link["target_id"] = identities[target]
repo.save(CardInput.model_validate(value), card_id=identities[p.id])
connection.execute(
text(
"INSERT INTO thoth_memory.reviews "
"(workspace_id,session_id,summary_id,request_hash,result) "
"VALUES (:w,:s,:id,:hash,CAST(:result AS jsonb))"
),
{
**params,
"hash": request_hash,
"result": json.dumps(
{
"saved": saved,
"context_hash": context_hash(snapshot),
}
),
},
)
results = [service._propagate(repo, item["id"]) for item in saved]
return {
"saved": len(saved),
"declined": len(summary["items"]) - len(saved) if not receipt else None,
"indexed": all(r["indexed"] for r in results),
"results": results,
}