"""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, }