Files
ThothII/harness/tht/phase.py

242 lines
9.7 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 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