From 7d0562b82dbc1cc02961fbb5190bea7452b85041 Mon Sep 17 00:00:00 2001 From: mptyl Date: Fri, 26 Jun 2026 22:35:43 +0200 Subject: [PATCH] feat(harness): phase.py rewrite + effective_decisions (D15 core fix, F2) 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 conflated stale pre-reopen decisions with new ones). Model (corrected during TDD): - current_phase folds the audit (excluding retracted) with guard 'n == cur' -- already reopen-aware (old phase_approved:N after reopen to M 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 diff --git a/harness/nsp/session/__init__.py b/harness/nsp/session/__init__.py new file mode 100644 index 00000000..e69de29b diff --git a/harness/nsp/session/models.py b/harness/nsp/session/models.py new file mode 100644 index 00000000..718b3c20 --- /dev/null +++ b/harness/nsp/session/models.py @@ -0,0 +1,64 @@ +from datetime import datetime +from typing import Literal + +from pydantic import BaseModel, Field, ConfigDict + + +# Stub locale di _YamlModel (in ChironeWp3 vive in mschema/models.py). +# Qui serve solo come base con populate_by_name; mschema/ sara' portato nel Task A9. +class _YamlModel(BaseModel): + model_config = ConfigDict(populate_by_name=True) + + +class SessionManifest(_YamlModel): + id: str + created_at: datetime + status: Literal["open", "closed", "finalized"] = "open" + question: str + database: str + db_schema: str = Field(alias="schema") + # D12/D15: autore della sessione (auth) e versione del workflow usato. + author: str | None = None + summary: str | None = None + updated_at: datetime | None = None + updated_by: str | None = None + schema_version: int | None = None + + +class Candidate(BaseModel): + kind: Literal["table", "column"] + name: str + signals: dict = {} + evidence: list[str] = [] + decision: Literal["promoted", "excluded", "pending"] = "pending" + decision_seq: int | None = None + # D14a: valori citati nella domanda ancorati a questa colonna/tabella. + grounded_values: list[dict] = [] + + +class Join(BaseModel): + from_: str = Field(alias="from") + to: str + source: str = "" + decision: Literal["promoted", "excluded", "pending"] = "promoted" + decision_seq: int | None = None + + model_config = {"populate_by_name": True} + + +class ExcludedItem(BaseModel): + kind: Literal["table", "column"] + name: str + decision_seq: int | None = None + + +class SchemaLinking(BaseModel): + question: str + candidates: list[Candidate] = [] + joins: list[Join] = [] + excluded: list[ExcludedItem] = [] + open_questions: list[str] = [] + # D14b: formule di concetto approvate, parte dello schema-linking. + concept_formulas: list[dict] = [] + + model_config = {"extra": "forbid"} diff --git a/harness/tests/test_phase_effective.py b/harness/tests/test_phase_effective.py new file mode 100644 index 00000000..c419ff3c --- /dev/null +++ b/harness/tests/test_phase_effective.py @@ -0,0 +1,105 @@ +from nsp.decisions import append_decision, DecisionRecord +from nsp.phase import current_phase, effective_decisions + + +def _d(session, dtype, subject, **kw) -> DecisionRecord: + """Helper: appendi una decisione e ritorna la record (per leggere .seq).""" + return append_decision( + session, type=dtype, subject=subject, + detail=kw.get("detail", ""), rationale=kw.get("rationale", ""), + retracts=kw.get("retracts"), + ) + + +def test_effective_decisions_excludes_stale_high_phase(tmp_path): + """Dopo un phase_reopened a fase bassa, una decisione di fase alta e' stale (excluded). + + Sequenza: approva 1, approva 2, reopen a 1, poi table_promoted:4. + current_phase diventa 1 (il reopen). La table_promoted:4 e' di fase 4 > 1 -> stale. + """ + s = tmp_path / "s"; s.mkdir() + _d(s, "phase_approved", "phase:1") + _d(s, "phase_approved", "phase:2") + _d(s, "phase_reopened", "phase:1") + _d(s, "table_promoted", "phase:4") # stale: fase 4 > current_phase(1) + assert current_phase(s) == 1 + eff = effective_decisions(s) + types = [d.type for d in eff] + assert "table_promoted" not in types # la stale di fase 4 e' esclusa + + +def test_current_phase_after_reopen(tmp_path): + """Il fold su audit-excluding-retracted gestisce correttamente il reopen e le ri-approvazioni.""" + s = tmp_path / "s"; s.mkdir() + _d(s, "phase_approved", "phase:1") + _d(s, "phase_approved", "phase:2") + assert current_phase(s) == 3 # dopo 2 approvazioni -> fase 3 + _d(s, "phase_reopened", "phase:1") + assert current_phase(s) == 1 + _d(s, "phase_approved", "phase:1") # ri-approva dopo reopen (legittima, NON stale) + assert current_phase(s) == 2 + + +def test_retracted_decision_excluded_from_effective(tmp_path): + """Una decisione ritirata (decision_retracted) e' esclusa dalla vista effective.""" + s = tmp_path / "s"; s.mkdir() + d1 = _d(s, "table_promoted", "phase:4", detail="t1") + _d(s, "decision_retracted", "phase:4", retracts=d1.seq) + # senza approvazioni di fase, current_phase=1; table_promoted:4 e' gia' > 1. + # Ma anche a fase 4 raggiunta, la ritirata NON deve ricomparire. + _d(s, "phase_approved", "phase:1") + _d(s, "phase_approved", "phase:2") + _d(s, "phase_approved", "phase:3") + _d(s, "phase_approved", "phase:4") + assert current_phase(s) == 5 + eff = effective_decisions(s) + seqs = [d.seq for d in eff] + assert d1.seq not in seqs # la ritirata e' esclusa + types = [d.type for d in eff] + assert "decision_retracted" not in types # il marker stesso non conta + + +def test_current_phase_starts_at_1(tmp_path): + s = tmp_path / "s"; s.mkdir() + assert current_phase(s) == 1 + + +def test_current_phase_clamps_at_max_plus_1(tmp_path): + """Dopo tutte le approvazioni, current_phase = max_phase + 1.""" + s = tmp_path / "s"; s.mkdir() + for n in range(1, 9): + _d(s, "phase_approved", f"phase:{n}") + assert current_phase(s) == 9 # max_phase(8) + 1 + + +def test_effective_keeps_low_phase_after_high_phase_rollback(tmp_path): + """Rollback a F4 NON invalida le decisioni delle fasi 1-3 (che restano <= current_phase).""" + s = tmp_path / "s"; s.mkdir() + _d(s, "concept_clarified", "phase:1", detail="x") + _d(s, "question_rewritten", "phase:3", detail="q") + _d(s, "table_promoted", "phase:4", detail="t") + _d(s, "sql_approved", "phase:7", detail="sql") # sara' stale dopo rollback + for n in range(1, 8): + _d(s, "phase_approved", f"phase:{n}") + assert current_phase(s) == 8 + # ora rollback a F4 + _d(s, "phase_reopened", "phase:4") + assert current_phase(s) == 4 + eff = effective_decisions(s) + types_subjects = [(d.type, d.subject) for d in eff] + # le decisioni di fase <= 4 restano + assert ("concept_clarified", "phase:1") in types_subjects + assert ("question_rewritten", "phase:3") in types_subjects + assert ("table_promoted", "phase:4") in types_subjects + # la decisione di fase 7 (sql_approved) e' ora stale -> esclusa + assert ("sql_approved", "phase:7") not in types_subjects + + +def test_effective_decisions_no_reopen_returns_all_non_retracted(tmp_path): + """Senza reopen e senza retract, effective = tutte le decisioni (della fase corrente).""" + s = tmp_path / "s"; s.mkdir() + _d(s, "concept_clarified", "phase:1", detail="x") + _d(s, "phase_approved", "phase:1") + assert current_phase(s) == 2 + eff = effective_decisions(s) + assert len(eff) == 2