"""Phase machinery -- data-driven + effective_decisions (spec D15, F2, ยง4.8). This is the single most important architectural fix vs ChironeWp3: 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 nsp.decisions import DecisionRecord, list_decisions from nsp.session.models import SchemaLinking from nsp.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 _audit_excluding_retracted(session_dir: Path) -> 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 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(session_dir): 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(session_dir: Path) -> 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(session_dir) out: list[DecisionRecord] = [] for d in _audit_excluding_retracted(session_dir): n = _phase_num(d.subject) # decisioni senza subject di fase (es. concept_clarified puo' avere subject libero): # le ammettiamo (non sono legate a una fase specifica da filtrare). 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(session_dir: Path) -> int: """Numero di decisioni sostanziali dall'ultimo confine di fase (vista effective).""" decs = effective_decisions(session_dir) 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(session_dir: Path) -> bool: """Vero sse la fase corrente puo' auto-avanzare (zero decisioni sostanziali + prereq ok).""" cur = current_phase(session_dir) if cur not in _AUTO_ADVANCE_PHASES: return False if substantive_count_current_phase(session_dir) > 0: return False return not advance_problems(session_dir, cur) # --- CTE helpers (consultano effective_decisions) --------------------------- CTE_PLAN_FILE = "cte_plan.json" def cte_plan(session_dir: Path) -> list[str]: path = session_dir / CTE_PLAN_FILE if not path.exists(): return [] return json.loads(path.read_text()) def approved_ctes(session_dir: Path) -> set[str]: """Insieme dei CTE approvati, dalla vista effective (esclude stale post-reopen).""" return {d.subject for d in effective_decisions(session_dir) if d.type == "cte_approved"} def next_cte(session_dir: Path) -> str | None: approved = approved_ctes(session_dir) for name in cte_plan(session_dir): if name not in approved: return name return None # --- advance_problems (ladder if-phase-N che consulta effective_decisions) --- def _has_decision(session_dir: Path, type_: str) -> bool: return any(d.type == type_ for d in effective_decisions(session_dir)) def _has_decision_subject(session_dir: Path, type_: str, subject: str) -> bool: return any( d.type == type_ and d.subject == subject for d in effective_decisions(session_dir) ) def advance_problems(session_dir: Path, 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 == 3 and not _has_decision(session_dir, "question_rewritten"): problems.append("manca la decisione question_rewritten (Fase 3)") if phase == 5: path = session_dir / "schema_linking.json" if not path.exists(): problems.append("schema_linking.json assente (Fase 5)") else: try: SchemaLinking.model_validate(json.loads(path.read_text())) 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(session_dir, "phase_skipped", "phase:6"): plan = cte_plan(session_dir) 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(session_dir) if nc is not None: problems.append(f"CTE non ancora approvato: {nc} (Fase 6)") if phase == 7 and not _has_decision(session_dir, "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(session_dir) ): problems.append( "Fase 8: nessuna risposta sulla generazione dbt del datamart " "(manca una decisione datamart_requested o datamart_declined)." ) return problems