Publish documentation / publish (push) Successful in 1m27s
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.
314 lines
12 KiB
Python
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,
|
|
}
|