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

244 lines
9.9 KiB
Python

"""Phase machinery -- data-driven + effective_decisions (spec D15, F2, §4.8).
This is the single most important architectural fix vs the reference implementation: ALL helpers
consult effective_decisions() instead of raw list_decisions(), so the reopen-aware
view is consistent everywhere (fixes the bug where approved_ctes / advance_problems /
build_evidence conflated stale pre-reopen decisions with new ones).
Strada 2 (decisa in A5): effective_decisions + ladder if-phase-N che la consulta,
MAX_PHASE/PHASE_NAMES letti da workflow.yaml. L'evaluator generico dei prerequisites
di workflow.yaml (F2 pieno) entra in un secondo momento.
"""
from __future__ import annotations
import json
from pathlib import Path
from pydantic import ValidationError
from tht.decisions import DecisionRecord, list_decisions
from tht.session.models import SchemaLinking, SessionSnapshot
from tht.workflow import load_workflow
def _phase_num(subject: str) -> int | None:
"""subject nel formato 'phase:N' -> N, oppure None."""
if not subject.startswith("phase:"):
return None
try:
return int(subject.split(":", 1)[1])
except ValueError:
return None
def _decisions(source: Path | SessionSnapshot) -> list[DecisionRecord]:
return list(source.decisions) if isinstance(source, SessionSnapshot) else list_decisions(source)
def _artifact(source: Path | SessionSnapshot, key: str, filename: str) -> str | None:
if isinstance(source, SessionSnapshot):
return source.artifacts.get(key)
path = source / filename
return path.read_text() if path.exists() else None
def _audit_excluding_retracted(source: Path | SessionSnapshot) -> list[DecisionRecord]:
"""Tutto il ledger (append-only) tranne le decisioni ritirate e i marker di ritrazione.
Base per il fold di current_phase: il guard 'n == cur' del fold e' gia' reopen-aware
(una phase_approved:N dopo un reopen a M<N non fa avanzare perche' cur!=N)."""
all_d = _decisions(source)
retracted_seqs = {
d.retracts for d in all_d if d.type == "decision_retracted" and d.retracts is not None
}
return [
d
for d in all_d
if d.seq not in retracted_seqs and d.type != "decision_retracted"
]
def current_phase(source: Path | SessionSnapshot) -> int:
"""Fase corrente come fold cronologico sull'audit (con ritirate escluse).
cur parte da 1; ogni phase_approved/phase_auto_approved per la fase CORRENTE avanza
(guard 'n == cur' -- gia' reopen-aware: dopo un reopen a M, le vecchie approvazioni
di N>M non fanno avanzare finche' non si riapprova in ordine); phase_reopened torna
indietro. Terminale: max_phase + 1.
"""
wf = load_workflow()
max_plus_one = wf.max_phase + 1
cur = 1
for d in _audit_excluding_retracted(source):
n = _phase_num(d.subject)
if n is None:
continue
if d.type in ("phase_approved", "phase_auto_approved") and n == cur:
cur = min(cur + 1, max_plus_one)
elif d.type == "phase_reopened":
cur = max(1, min(cur, n))
return cur
def effective_decisions(source: Path | SessionSnapshot) -> list[DecisionRecord]:
"""La vista canonica 'effective as of pointer'. TUTTI gli helper non-fold devono usare questa.
Semantica: una decisione e' effective se appartiene a una fase <= current_phase.
Una decisione di fase 7 (es. sql_approved) e' stale quando current_phase=4 dopo un
rollback a F4, anche se fisicamente appare nel ledger.
Il fold di current_phase gestisce le *approvazioni* via guard 'n == cur'; qui
applichiamo la stessa nozione alle decisioni *sostanziali* (table_promoted, sql_approved,
cte_approved, ...): contano solo se la loro fase e' <= quella corrente.
Inoltre esclude le decisioni ritirate (decision_retracted) e i marker stessi.
"""
cur = current_phase(source)
out: list[DecisionRecord] = []
for d in _audit_excluding_retracted(source):
n = _phase_num(d.subject)
# Per le decisioni con subject "a nome" (es. cte_approved -> nome CTE, evidence_*
# -> id evidence) il subject non porta la fase: si usa la fase emittente registrata
# (d.phase, high-water-mark D15). Senza nessuno dei due (record storici) la
# decisione e' ammessa: non c'e' modo di datarla, e il subject non e' di fase.
if n is None:
n = d.phase
if n is None or n <= cur:
out.append(d)
return out
# --- AUTO-ADVANCE -----------------------------------------------------------
_AUTO_ADVANCE_PHASES = frozenset({2, 6}) # F2 Memorie, F6 CTE
_BOUNDARY_TYPES = frozenset({"phase_approved", "phase_auto_approved", "phase_reopened"})
_META_TYPES = frozenset(
{"phase_approved", "phase_auto_approved", "phase_reopened", "phase_skipped"}
)
def substantive_count_current_phase(source: Path | SessionSnapshot) -> int:
"""Numero di decisioni sostanziali dall'ultimo confine di fase (vista effective)."""
decs = effective_decisions(source)
start = 0
for i, d in enumerate(decs):
if d.type in _BOUNDARY_TYPES:
start = i + 1
return sum(1 for d in decs[start:] if d.type not in _META_TYPES)
def auto_advance_eligible(source: Path | SessionSnapshot) -> bool:
"""Vero sse la fase corrente puo' auto-avanzare (zero decisioni sostanziali + prereq ok)."""
cur = current_phase(source)
if cur not in _AUTO_ADVANCE_PHASES:
return False
if substantive_count_current_phase(source) > 0:
return False
return not advance_problems(source, cur)
# --- CTE helpers (consultano effective_decisions) ---------------------------
CTE_PLAN_FILE = "cte_plan.json"
def cte_plan(source: Path | SessionSnapshot) -> list[str]:
raw = _artifact(source, "cte_plan", CTE_PLAN_FILE)
if raw is None:
return []
return json.loads(raw)
def approved_ctes(source: Path | SessionSnapshot) -> set[str]:
"""Insieme dei CTE approvati, dalla vista effective (esclude stale post-reopen)."""
return {d.subject for d in effective_decisions(source) if d.type == "cte_approved"}
def next_cte(source: Path | SessionSnapshot) -> str | None:
approved = approved_ctes(source)
for name in cte_plan(source):
if name not in approved:
return name
return None
# --- advance_problems (ladder if-phase-N che consulta effective_decisions) ---
def _has_decision(source: Path | SessionSnapshot, type_: str) -> bool:
return any(d.type == type_ for d in effective_decisions(source))
def _has_decision_subject(source: Path | SessionSnapshot, type_: str, subject: str) -> bool:
return any(
d.type == type_ and d.subject == subject
for d in effective_decisions(source)
)
def advance_problems(source: Path | SessionSnapshot, phase: int) -> list[str]:
"""Prerequisiti minimi per chiudere `phase` (lista vuota = ok).
Ladder if-phase-N (Strada 2): la logica specifica resta, ma ogni lettura passa per
effective_decisions (fix D15). I prerequisiti sono anche documentati in workflow.yaml;
l'evaluator generico (F2 pieno) entra in un secondo momento.
"""
problems: list[str] = []
if phase == 4:
raw = _artifact(source, "schema_linking", "schema_linking.json")
if raw is None:
problems.append("schema_linking.json assente (Fase 4)")
else:
try:
linking = SchemaLinking.model_validate(json.loads(raw))
except (json.JSONDecodeError, ValidationError) as e:
problems.append(f"schema_linking.json non valido (Fase 4): {e}")
else:
promoted_tables = {
candidate.name
for candidate in linking.candidates
if candidate.kind == "table" and candidate.decision == "promoted"
}
if len(promoted_tables) > 1 and (
not linking.joins or not _has_decision(source, "join_modified")
):
problems.append(
"Fase 4: più tabelle promosse richiedono join strutturati in "
"schema_linking.json e una decisione join_modified del reviewer"
)
if phase == 3 and not _has_decision(source, "question_rewritten"):
problems.append("manca la decisione question_rewritten (Fase 3)")
if phase == 5:
raw = _artifact(source, "schema_linking", "schema_linking.json")
if raw is None:
problems.append("schema_linking.json assente (Fase 5)")
else:
try:
SchemaLinking.model_validate(json.loads(raw))
except (json.JSONDecodeError, ValidationError) as e:
problems.append(f"schema_linking.json non valido (Fase 5): {e}")
if phase == 6 and not _has_decision_subject(source, "phase_skipped", "phase:6"):
plan = cte_plan(source)
if not plan:
problems.append(
"Fase 6: nessun piano CTE (cte_plan.json) e nessun salto esplicito. "
"Approva un piano (reviewer_confirm kind:'cte_plan') oppure salta la Fase 6 "
"registrando una decisione phase_skipped subject phase:6."
)
else:
nc = next_cte(source)
if nc is not None:
problems.append(f"CTE non ancora approvato: {nc} (Fase 6)")
if phase == 7 and not _has_decision(source, "sql_approved"):
problems.append("manca la decisione sql_approved (Fase 7)")
if phase == 8 and not _has_decision(source, "memory_summary_reviewed"):
problems.append("Fase 8: il riepilogo Memory deve essere revisionato prima della chiusura")
if phase == 8 and not any(
d.type in ("datamart_requested", "datamart_declined")
for d in effective_decisions(source)
):
problems.append(
"Fase 8: nessuna risposta sulla generazione dbt del datamart "
"(manca una decisione datamart_requested o datamart_declined)."
)
return problems