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<N don't advance). - effective_decisions = decisions whose phase <= current_phase. A sql_approved at phase 7 is stale when current_phase=4 after a rollback to F4, even if in the ledger. Rollback to F4 does NOT invalidate decisions of phases 1-3 (they stay effective). - decision_retracted markers excluded (audit-only). Also: session/models.py ported (SchemaLinking + Candidate with grounded_values D14a + concept_formulas D14b). MAX_PHASE/PHASE_NAMES read from workflow.yaml via load_workflow() (no duplication). Strada 2: ladder if-phase-N kept for now, generic prerequisites evaluator (F2 full) deferred. 7 phase tests + 20 total passing.
This commit is contained in:
@@ -0,0 +1,204 @@
|
||||
"""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<N non fa avanzare perche' cur!=N)."""
|
||||
all_d = list_decisions(session_dir)
|
||||
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(session_dir: Path) -> 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
|
||||
@@ -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"}
|
||||
@@ -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
|
||||
Reference in New Issue
Block a user