diff --git a/harness/tht/cli/__init__.py b/harness/tht/cli/__init__.py index f444f0e2..c706643a 100644 --- a/harness/tht/cli/__init__.py +++ b/harness/tht/cli/__init__.py @@ -36,6 +36,12 @@ def main( pass +from tht.cli.config_cmd import config_app # noqa: E402 from tht.cli.phase_cmd import phase_app # noqa: E402 +from tht.cli.schema_cmd import schema_app # noqa: E402 +from tht.cli.session_cmd import session_app # noqa: E402 app.add_typer(phase_app, name="phase") +app.add_typer(config_app, name="config") +app.add_typer(schema_app, name="schema") +app.add_typer(session_app, name="session") diff --git a/harness/tht/cli/config_cmd.py b/harness/tht/cli/config_cmd.py new file mode 100644 index 00000000..b5f55d3d --- /dev/null +++ b/harness/tht/cli/config_cmd.py @@ -0,0 +1,26 @@ +from pathlib import Path + +import typer + +from tht.config import ConfigError, load_config + +config_app = typer.Typer(help="Gestione configurazione") + +CONFIG_OPT = typer.Option( + Path("config/tht.yaml"), "--config", "-c", help="Percorso del file di configurazione." +) + + +@config_app.command("check") +def check(config: Path = CONFIG_OPT) -> None: + """Valida configurazione e variabili d'ambiente risolte.""" + try: + cfg = load_config(config) + except ConfigError as e: + typer.secho(f"ERRORE: {e}", fg=typer.colors.RED, err=True) + raise typer.Exit(code=1) + db = cfg.database + typer.secho(f"OK: configurazione valida ({config})", fg=typer.colors.GREEN) + typer.echo(f" profilo : {cfg.profile}") + typer.echo(f" database : {db.user}@{db.host}:{db.port}/{db.database} schema={db.db_schema}") + typer.echo(f" artifacts: {cfg.paths.artifacts} indexes: {cfg.paths.indexes}") diff --git a/harness/tht/cli/schema_cmd.py b/harness/tht/cli/schema_cmd.py new file mode 100644 index 00000000..0770ca38 --- /dev/null +++ b/harness/tht/cli/schema_cmd.py @@ -0,0 +1,158 @@ +from pathlib import Path + +import typer +from sqlalchemy.exc import OperationalError + +from tht.cli.config_cmd import CONFIG_OPT +from tht.config import ConfigError, load_config +from tht.db.connection import make_engine +from tht.db.introspect import IntrospectionError, introspect +from tht.db.sampling import add_examples +from tht.mschema.eligibility import classify_all + +schema_app = typer.Typer(help="Gestione mschema (rappresentazione canonica dello schema)") + + +def _load_config_or_exit(config: Path): + try: + return load_config(config) + except ConfigError as e: + typer.secho(f"ERRORE: {e}", fg=typer.colors.RED, err=True) + raise typer.Exit(code=1) + + +def physical_path(cfg) -> Path: + return cfg.paths.artifacts / "mschema" / "physical.yaml" + + +def annotations_path(cfg) -> Path: + return cfg.paths.artifacts / "mschema" / "annotations.yaml" + + +@schema_app.command("introspect") +def introspect_cmd(config: Path = CONFIG_OPT) -> None: + """Introspeziona lo schema target e genera artifacts/mschema/physical.yaml.""" + cfg = _load_config_or_exit(config) + if cfg.database.transport == "rest": + from tht.db.introspect import introspect_rest + from tht.db.sampling import add_examples_rest + from tht.rest.client import RestClient, RestError + + client = RestClient(cfg.rest) + try: + phys = introspect_rest( + client, database=cfg.database.database, schema=cfg.database.db_schema + ) + add_examples_rest(client, phys, cfg.examples) + classify_all(phys, cfg.eligibility) + except RestError as e: + typer.secho(f"ERRORE: {e}", fg=typer.colors.RED, err=True) + raise typer.Exit(code=1) + else: + engine = make_engine(cfg.database) + try: + phys = introspect( + engine, database=cfg.database.database, schema=cfg.database.db_schema + ) + add_examples(engine, phys, cfg.examples) + classify_all(phys, cfg.eligibility) + except (OperationalError, IntrospectionError) as e: + typer.secho(f"ERRORE: {e}", fg=typer.colors.RED, err=True) + raise typer.Exit(code=1) + out = physical_path(cfg) + phys.to_yaml(out) + n_cols = sum(len(t.columns) for t in phys.tables.values()) + n_ignored = sum( + 1 for t in phys.tables.values() for c in t.columns.values() if not c.eligible + ) + typer.secho( + f"OK: {len(phys.tables)} tabelle, {n_cols} colonne " + f"({n_ignored} ignorate: testo ampio) -> {out}", + fg=typer.colors.GREEN, + ) + + +@schema_app.command("check") +def check_cmd(config: Path = CONFIG_OPT) -> None: + """Confronta physical.yaml e annotations.yaml; segnala annotazioni orfane.""" + from tht.mschema.merge import find_orphans + from tht.mschema.models import Annotations, PhysicalSchema + + cfg = _load_config_or_exit(config) + phys_file = physical_path(cfg) + if not phys_file.exists(): + typer.secho( + f"ERRORE: {phys_file} non trovato. Esegui prima `tht schema introspect`.", + fg=typer.colors.RED, err=True, + ) + raise typer.Exit(code=1) + physical = PhysicalSchema.from_yaml(phys_file) + annotations = Annotations.from_yaml(annotations_path(cfg)) + + ignored = [ + f"{t}.{c} ({col.eligibility_reason})" + for t, table in physical.tables.items() + for c, col in table.columns.items() + if not col.eligible + ] + if ignored: + typer.secho( + f"Colonne ignorate (testo ampio, {len(ignored)}):", fg=typer.colors.YELLOW + ) + for line in ignored: + typer.echo(f" - {line}") + + orphans = find_orphans(physical, annotations) + if orphans: + typer.secho(f"ATTENZIONE: {len(orphans)} annotazioni orfane:", fg=typer.colors.YELLOW) + for o in orphans: + typer.echo(f" - {o}") + raise typer.Exit(code=3) + typer.secho("OK: nessuna annotazione orfana.", fg=typer.colors.GREEN) + + +@schema_app.command("render") +def render_cmd( + config: Path = CONFIG_OPT, + format: str = typer.Option( + "markdown", "--format", "-f", help="Formato: markdown | mschema-text | schema-dict" + ), + tables: list[str] = typer.Option( + None, "--table", "-t", help="Limita alle tabelle indicate (ripetibile)." + ), + output: Path = typer.Option(None, "--output", "-o", help="File di output (default stdout)."), +) -> None: + """Serializza mschema (physical + annotations) nel formato richiesto.""" + import json + + from tht.mschema.models import Annotations, PhysicalSchema + from tht.mschema.render import to_markdown, to_mschema_text, to_schema_dict + + cfg = _load_config_or_exit(config) + phys_file = physical_path(cfg) + if not phys_file.exists(): + typer.secho( + f"ERRORE: {phys_file} non trovato. Esegui prima `tht schema introspect`.", + fg=typer.colors.RED, err=True, + ) + raise typer.Exit(code=1) + physical = PhysicalSchema.from_yaml(phys_file) + annotations = Annotations.from_yaml(annotations_path(cfg)) + table_filter = list(tables) if tables else None + + if format == "markdown": + out = to_markdown(physical, annotations) + elif format == "mschema-text": + out = to_mschema_text(physical, annotations, tables=table_filter) + elif format == "schema-dict": + out = json.dumps(to_schema_dict(physical, annotations), ensure_ascii=False, indent=2) + else: + typer.secho(f"ERRORE: formato sconosciuto: {format}", fg=typer.colors.RED, err=True) + raise typer.Exit(code=1) + + if output: + output.parent.mkdir(parents=True, exist_ok=True) + output.write_text(out) + typer.secho(f"OK: scritto {output}", fg=typer.colors.GREEN) + else: + typer.echo(out) diff --git a/harness/tht/cli/session_cmd.py b/harness/tht/cli/session_cmd.py new file mode 100644 index 00000000..a3e2d7ab --- /dev/null +++ b/harness/tht/cli/session_cmd.py @@ -0,0 +1,268 @@ +from pathlib import Path + +import typer + +from tht.cli.config_cmd import CONFIG_OPT +from tht.cli.schema_cmd import _load_config_or_exit + +session_app = typer.Typer(help="Sessioni (directory artefatti)") + + +def session_dir(cfg, session_id: str) -> Path: + return cfg.paths.sessions / session_id + + +def load_session_or_exit(cfg, session_id: str): + from tht.session.store import SessionError, load_session + + try: + return load_session(session_id, cfg.paths.sessions) + except SessionError as e: + typer.secho(f"ERRORE: {e}", fg=typer.colors.RED, err=True) + raise typer.Exit(code=1) + + +@session_app.command("new") +def new_cmd( + question: str = typer.Argument(..., help="La domanda in linguaggio naturale."), + config: Path = CONFIG_OPT, +) -> None: + """Crea una sessione e stampa il suo id (ultima riga dell'output).""" + from tht.session.store import create_session + + cfg = _load_config_or_exit(config) + manifest = create_session(question, cfg.database, cfg.paths.sessions) + typer.secho(f"OK: sessione creata in {session_dir(cfg, manifest.id)}", fg=typer.colors.GREEN) + typer.echo(manifest.id) + + +@session_app.command("set-question") +def set_question_cmd( + session_id: str = typer.Argument(...), + question: str = typer.Option(..., "--question", "-q", + help="La domanda riscritta (chiara)."), + assumption: list[str] = typer.Option( + None, "--assumption", "-a", + help="Una assunzione (ripetibile). Testo libero, già formattato."), + config: Path = CONFIG_OPT, +) -> None: + """Scrive question.md (domanda riscritta + assunzioni) in modo deterministico.""" + from tht.session.store import set_question + + cfg = _load_config_or_exit(config) + load_session_or_exit(cfg, session_id) + path = set_question(session_id, question, assumption or [], cfg.paths.sessions) + typer.secho(f"OK: question.md aggiornato ({path}).", fg=typer.colors.GREEN) + + +@session_app.command("show") +def show_cmd(session_id: str = typer.Argument(...), config: Path = CONFIG_OPT) -> None: + """Stato della sessione: manifest + decisioni registrate (per la ripresa).""" + from tht.decisions import list_decisions + + cfg = _load_config_or_exit(config) + manifest = load_session_or_exit(cfg, session_id) + decisions = list_decisions(session_dir(cfg, session_id)) + typer.echo(f"id : {manifest.id}") + typer.echo(f"stato : {manifest.status}") + typer.echo(f"domanda : {manifest.question}") + typer.echo(f"target : {manifest.database} / {manifest.db_schema}") + typer.echo(f"decisioni: {len(decisions)}") + for d in decisions[-10:]: + typer.echo(f" [{d.seq}] {d.type}: {d.subject}" + (f" — {d.detail}" if d.detail else "")) + linking = session_dir(cfg, session_id) / "schema_linking.json" + typer.echo(f"schema_linking.json: {'presente' if linking.exists() else 'assente'}") + + +@session_app.command("close") +def close_cmd(session_id: str = typer.Argument(...), config: Path = CONFIG_OPT) -> None: + """Chiude la sessione (status=closed).""" + from tht.session.store import close_session + + cfg = _load_config_or_exit(config) + load_session_or_exit(cfg, session_id) + close_session(session_id, cfg.paths.sessions) + typer.secho(f"OK: sessione {session_id} chiusa.", fg=typer.colors.GREEN) + + +def session_problems(cfg, session_id: str) -> list[str]: + """Problemi del Blocco 3 (decisioni + schema_linking). Riusato da check e finalize.""" + import json + + from pydantic import ValidationError + + from tht.decisions import list_decisions + from tht.session.models import SchemaLinking + + sdir = session_dir(cfg, session_id) + problems: list[str] = [] + if not list_decisions(sdir): + problems.append("nessuna decisione registrata (review_decisions.jsonl vuoto o assente)") + linking_path = sdir / "schema_linking.json" + if not linking_path.exists(): + problems.append("schema_linking.json assente") + else: + try: + SchemaLinking.model_validate(json.loads(linking_path.read_text())) + except (json.JSONDecodeError, ValidationError) as e: + problems.append(f"schema_linking.json non valido: {e}") + return problems + + +@session_app.command("check") +def check_cmd(session_id: str = typer.Argument(...), config: Path = CONFIG_OPT) -> None: + """Gate oggettivo: decisioni presenti e schema_linking.json valido (dalla Fase 4).""" + from tht.phase import current_phase + from tht.workflow import load_workflow + + cfg = _load_config_or_exit(config) + load_session_or_exit(cfg, session_id) + cur = current_phase(session_dir(cfg, session_id)) + wf = load_workflow() + schema_linking_phase = wf.schema_linking_phase() + if cur < schema_linking_phase: + # Prima della fase schema-linking non e' ancora atteso: il gate Blocco 3 + # non e' applicabile. Riportarlo come errore manderebbe il workflow in loop. + nome = wf.phase_name(cur) + typer.secho( + f"Sessione {session_id} in Fase {cur} ({nome}): il gate Blocco 3 " + f"(schema_linking.json) si applica dalla Fase {schema_linking_phase} " + f"({wf.phase_name(schema_linking_phase)}). Niente da verificare ora: " + "prosegui con lo schema linking.", + fg=typer.colors.CYAN, + ) + return + problems = session_problems(cfg, session_id) + if problems: + typer.secho(f"Sessione {session_id} incompleta:", fg=typer.colors.YELLOW) + for p in problems: + typer.echo(f" - {p}") + raise typer.Exit(code=3) + typer.secho(f"OK: sessione {session_id} completa per il Blocco 3.", fg=typer.colors.GREEN) + + +ARTIFACT_FILES = [ + "session_manifest.yaml", "question.md", "schema_linking.json", "evidence.json", + "cte_tests.json", "sql_final.sql", "validation_report.md", "review_decisions.jsonl", +] + + +@session_app.command("finalize") +def finalize_cmd(session_id: str = typer.Argument(...), config: Path = CONFIG_OPT) -> None: + """Gate finale + batteria di validazione + artefatto di sessione completo.""" + import json + + from tht.cli.sql_cmd import ( + _load_physical_or_exit, + do_explain, + do_run, + promoted_tables_for, + ) + from tht.ctetest import CteError, load_cte_tests + from tht.execute import ExecutionError + from tht.execute.warnings import plan_warnings, runtime_warnings, static_warnings + from tht.report import extract_reviewer_notes, render_validation_report + from tht.session.artifacts import build_evidence_entries + from tht.decisions import list_decisions + from tht.session.models import SchemaLinking + from tht.sqlcheck import validate_sql + + cfg = _load_config_or_exit(config) + manifest = load_session_or_exit(cfg, session_id) + sdir = session_dir(cfg, session_id) + + from tht.phase import current_phase + from tht.workflow import load_workflow + + cur = current_phase(sdir) + wf = load_workflow() + if cur <= wf.max_phase: + typer.secho( + f"Finalize rifiutato: workflow non completo, sei in Fase {cur} " + f"({wf.phase_name(min(cur, wf.max_phase))}). Tutte le {wf.max_phase} fasi " + "devono essere approvate (anche dopo eventuali reopen).", + fg=typer.colors.YELLOW, err=True, + ) + raise typer.Exit(code=5) + + # --- gate di ingresso --- + problems = session_problems(cfg, session_id) + decisions = list_decisions(sdir) + sql_file = sdir / "sql_final.sql" + if not sql_file.exists(): + problems.append("sql_final.sql assente") + if not any(d.type == "sql_approved" for d in decisions): + problems.append( + "decisione sql_approved assente: la validazione semantica del reviewer " + "e' obbligatoria prima del finalize" + ) + cte_dir = sdir / "ctes" + if cte_dir.is_dir(): + try: + tested = {r.name for r in load_cte_tests(sdir)} + except CteError as e: + typer.secho( + f"Finalize rifiutato: impossibile leggere cte_tests.json: {e}", + fg=typer.colors.RED, err=True, + ) + raise typer.Exit(code=3) + for cte_file in sorted(cte_dir.glob("*.sql")): + if cte_file.stem not in tested: + problems.append(f"CTE mai testato: {cte_file.stem}") + if problems: + typer.secho(f"Finalize rifiutato per {session_id}:", fg=typer.colors.YELLOW) + for p in problems: + typer.echo(f" - {p}") + raise typer.Exit(code=3) + + # --- batteria di validazione su sql_final.sql --- + sql = sql_file.read_text() + check = validate_sql( + sql, + physical=_load_physical_or_exit(cfg), + promoted_tables=promoted_tables_for(cfg, session_id), + forbidden_functions=set(cfg.execution.forbidden_functions), + ) + if not check.ok: + typer.secho("Finalize rifiutato: sql_final.sql non passa la validazione statica:", + fg=typer.colors.RED, err=True) + for e in check.errors: + typer.echo(f" - {e}", err=True) + raise typer.Exit(code=3) + limit = cfg.execution.max_preview_rows + try: + plan = do_explain(cfg, sql) + result = do_run(cfg, sql, limit=limit) + except ExecutionError as e: + typer.secho(f"Finalize rifiutato: {e}", fg=typer.colors.RED, err=True) + raise typer.Exit(code=3) + + # --- validation_report.md (preservando le note del reviewer) --- + report_path = sdir / "validation_report.md" + existing_notes = ( + extract_reviewer_notes(report_path.read_text()) if report_path.exists() else "" + ) + report_path.write_text(render_validation_report( + session_id=session_id, check=check, plan=plan, + plan_warnings=plan_warnings(plan, cfg.execution), + result=result, runtime_warnings=runtime_warnings(result, cfg.execution), + static_warnings=static_warnings(check.ast), limit=limit, + reviewer_notes=existing_notes, + )) + + # --- evidence.json --- + linking = SchemaLinking.model_validate( + json.loads((sdir / "schema_linking.json").read_text()) + ) + entries = build_evidence_entries(decisions, linking, cfg.paths.artifacts / "evidence") + (sdir / "evidence.json").write_text(json.dumps(entries, ensure_ascii=False, indent=2)) + + # --- manifest + riepilogo --- + from tht.session.store import MANIFEST + + manifest.status = "finalized" + manifest.to_yaml(sdir / MANIFEST) + typer.secho(f"OK: sessione {session_id} finalizzata. Artefatti:", fg=typer.colors.GREEN) + for name in ARTIFACT_FILES: + state = "presente" if (sdir / name).exists() else "assente" + typer.echo(f" - {name}: {state}")