From 91a374492cde69be65deccc786df9888d90bf6e9 Mon Sep 17 00:00:00 2001 From: mptyl Date: Sat, 27 Jun 2026 14:16:35 +0200 Subject: [PATCH] =?UTF-8?q?feat(harness):=20port=20sql/cte/datamart/lsh=20?= =?UTF-8?q?cmd=20(Onda=204)=20=E2=80=94=20CLI=20completa=20F1=E2=86=92F8?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Ultima onda CLI. 4 cmd portati con rename + grep-per-file (3 residui nsp nei messaggi fixati). Nessun drift costanti phase in questi cmd. La CLI tht e' ora COMPLETA: 14 gruppi di comandi (phase config schema session vector memory search evidence db decision sql cte datamart lsh). tht --help li list tutti. Suite: 165 passed. Il loop skill->LLM->gate ora ha tutti i comandi che la skill chiamera'. Resta: skill riscritta (S), setup pre-sessione (0b), sessione L2 manuale. --- harness/tht/cli/__init__.py | 8 + harness/tht/cli/cte_cmd.py | 181 ++++++++++++++++++++++ harness/tht/cli/datamart_cmd.py | 42 ++++++ harness/tht/cli/lsh_cmd.py | 97 ++++++++++++ harness/tht/cli/sql_cmd.py | 260 ++++++++++++++++++++++++++++++++ 5 files changed, 588 insertions(+) create mode 100644 harness/tht/cli/cte_cmd.py create mode 100644 harness/tht/cli/datamart_cmd.py create mode 100644 harness/tht/cli/lsh_cmd.py create mode 100644 harness/tht/cli/sql_cmd.py diff --git a/harness/tht/cli/__init__.py b/harness/tht/cli/__init__.py index f9c56709..edea3b7c 100644 --- a/harness/tht/cli/__init__.py +++ b/harness/tht/cli/__init__.py @@ -37,14 +37,18 @@ def main( from tht.cli.config_cmd import config_app # noqa: E402 +from tht.cli.cte_cmd import cte_app # noqa: E402 +from tht.cli.datamart_cmd import datamart_app # noqa: E402 from tht.cli.db_cmd import db_app # noqa: E402 from tht.cli.decision_cmd import decision_app # noqa: E402 from tht.cli.evidence_cmd import evidence_app # noqa: E402 +from tht.cli.lsh_cmd import lsh_app # noqa: E402 from tht.cli.memory_cmd import memory_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.search_cmd import search_app # noqa: E402 from tht.cli.session_cmd import session_app # noqa: E402 +from tht.cli.sql_cmd import sql_app # noqa: E402 from tht.cli.vector_cmd import vector_app # noqa: E402 app.add_typer(phase_app, name="phase") @@ -57,3 +61,7 @@ app.add_typer(search_app, name="search") app.add_typer(evidence_app, name="evidence") app.add_typer(db_app, name="db") app.add_typer(decision_app, name="decision") +app.add_typer(sql_app, name="sql") +app.add_typer(cte_app, name="cte") +app.add_typer(datamart_app, name="datamart") +app.add_typer(lsh_app, name="lsh") diff --git a/harness/tht/cli/cte_cmd.py b/harness/tht/cli/cte_cmd.py new file mode 100644 index 00000000..1962bf70 --- /dev/null +++ b/harness/tht/cli/cte_cmd.py @@ -0,0 +1,181 @@ +import hashlib +from datetime import UTC, datetime +from pathlib import Path + +import typer + +from tht.cli.config_cmd import CONFIG_OPT +from tht.cli.schema_cmd import _load_config_or_exit +from tht.cli.session_cmd import load_session_or_exit, session_dir +from tht.cli.sql_cmd import promoted_tables_for, require_action + +_RULE6_HINT = ( + " — il file del CTE deve contenere SOLO il blocco WITH ... AS (...), " + "senza SELECT finale: la SELECT la aggiunge tht in fase di test." +) + +cte_app = typer.Typer(help="Test controllato dei CTE proposti (Agent View Generation)") + + +@cte_app.command("test") +def test_cmd( + name: str = typer.Argument(..., help="Nome del CTE (file sessions//ctes/.sql)."), + session: str = typer.Option(..., "--session"), + json_out: bool = typer.Option(False, "--json", help="Emette l'esito come JSON (per Pi)."), + config: Path = CONFIG_OPT, +) -> None: + """Valida e testa un CTE: wrap su ultimo CTE + esecuzione nelle 4 reti.""" + from rich.console import Console + from rich.table import Table + + from tht.cli.sql_cmd import do_run + from tht.ctetest import CteError, CteTestRecord, append_cte_test, build_test_sql, has_trailing_select + from tht.execute import ExecutionError + from tht.execute.warnings import runtime_warnings + + cfg = _load_config_or_exit(config) + require_action(cfg, "cte_test") + load_session_or_exit(cfg, session) + sdir = session_dir(cfg, session) + + from tht.cli.phase_cmd import require_phase_or_exit + from tht.phase import next_cte + + require_phase_or_exit(cfg, session, 6) + nxt = next_cte(sdir) + if nxt is not None and name != nxt: + typer.secho( + f"ERRORE: ordine CTE. Ora tocca a '{nxt}' (il primo CTE del piano non " + f"ancora approvato), non a '{name}'. Testa '{nxt}', poi falla approvare " + f"dal reviewer (decisione cte_approved); solo allora potrai testare il " + f"CTE successivo.", + fg=typer.colors.RED, err=True, + ) + raise typer.Exit(code=5) + + cte_file = sdir / "ctes" / f"{name}.sql" + if not cte_file.exists(): + typer.secho(f"ERRORE: file CTE non trovato: {cte_file}", fg=typer.colors.RED, err=True) + raise typer.Exit(code=1) + cte_sql = cte_file.read_text() + sql_hash = hashlib.sha256(cte_sql.encode()).hexdigest() + + def _record_error(message: str) -> None: + record = CteTestRecord( + name=name, ts=datetime.now(UTC), sql_hash=sql_hash, + status="error", error=message, + ) + append_cte_test(sdir, record) + if json_out: + typer.echo(record.model_dump_json()) + else: + typer.secho(f"ERRORE: {message}", fg=typer.colors.RED, err=True) + raise typer.Exit(code=1) + + try: + test_sql = build_test_sql(cte_sql) + except CteError as e: + msg = str(e) + if "non parsabile" in msg and has_trailing_select(cte_sql): + msg += _RULE6_HINT + _record_error(msg) + + # validazione statica (esce con messaggi propri se invalida; qui la + # intercettiamo per registrare comunque l'esito nel cte_tests.json) + from tht.cli.sql_cmd import _load_physical_or_exit + from tht.sqlcheck import validate_sql + + check = validate_sql( + test_sql, + physical=_load_physical_or_exit(cfg), + promoted_tables=promoted_tables_for(cfg, session), + forbidden_functions=set(cfg.execution.forbidden_functions), + ) + for w in check.warnings: + typer.secho(f" warning: {w}", fg=typer.colors.YELLOW) + if not check.ok: + msg = "; ".join(check.errors) + if any("statement" in e for e in check.errors): + msg += _RULE6_HINT + _record_error(msg) + + try: + result = do_run(cfg, test_sql, limit=cfg.execution.max_preview_rows) + except ExecutionError as e: + _record_error(str(e)) + + warnings = runtime_warnings(result, cfg.execution) + record = CteTestRecord( + name=name, ts=datetime.now(UTC), sql_hash=sql_hash, status="ok", + columns=result.columns, row_sample=len(result.rows), + execution_ms=result.execution_ms, warnings=warnings, + ) + append_cte_test(sdir, record) + + if json_out: + typer.echo(record.model_dump_json()) + return + + typer.echo(f"CTE {name} — {result.execution_ms} ms") + table = Table() + for col in result.columns: + table.add_column(col) + for row in result.rows: + table.add_row(*[str(v) for v in row]) + Console().print(table) + for w in warnings: + typer.secho(f" warning: {w}", fg=typer.colors.YELLOW) + typer.secho(f"OK: esito registrato in {sdir / 'cte_tests.json'}", fg=typer.colors.GREEN) + + +@cte_app.command("plan") +def plan_cmd( + session: str = typer.Option(..., "--session"), + name: list[str] = typer.Option(..., "--name", help="Nome CTE (ripetibile, in ordine)."), + config: Path = CONFIG_OPT, +) -> None: + """Persiste il piano CTE ordinato (sessions//cte_plan.json).""" + import json + + from tht.phase import CTE_PLAN_FILE + + cfg = _load_config_or_exit(config) + load_session_or_exit(cfg, session) + from tht.cli.phase_cmd import require_phase_or_exit + require_phase_or_exit(cfg, session, 6) + sdir = session_dir(cfg, session) + sdir.mkdir(parents=True, exist_ok=True) + (sdir / CTE_PLAN_FILE).write_text(json.dumps(name, ensure_ascii=False)) + typer.secho(f"OK: piano CTE salvato ({len(name)} CTE) in {sdir / CTE_PLAN_FILE}.", + fg=typer.colors.GREEN) + + +@cte_app.command("list") +def list_cmd( + session: str = typer.Option(..., "--session"), + config: Path = CONFIG_OPT, +) -> None: + """Ultimo esito registrato per ogni CTE della sessione.""" + from tht.ctetest import CteError, load_cte_tests + + cfg = _load_config_or_exit(config) + load_session_or_exit(cfg, session) + try: + records = load_cte_tests(session_dir(cfg, session)) + except CteError as e: + typer.secho(f"ERRORE: impossibile leggere cte_tests.json: {e}", + fg=typer.colors.RED, err=True) + raise typer.Exit(code=1) + if not records: + typer.echo("Nessun test CTE registrato.") + return + latest = {} + for r in records: + latest[r.name] = r + for name, r in sorted(latest.items()): + line = f"{name}: {r.status} ({r.ts:%Y-%m-%d %H:%M}, {r.execution_ms} ms)" + if r.error: + line += f" — {r.error}" + if r.warnings: + line += f" [{len(r.warnings)} warning]" + typer.echo(line) diff --git a/harness/tht/cli/datamart_cmd.py b/harness/tht/cli/datamart_cmd.py new file mode 100644 index 00000000..b7aeea03 --- /dev/null +++ b/harness/tht/cli/datamart_cmd.py @@ -0,0 +1,42 @@ +from pathlib import Path + +import typer + +from tht.cli.config_cmd import CONFIG_OPT +from tht.cli.schema_cmd import _load_config_or_exit +from tht.cli.session_cmd import load_session_or_exit, session_dir + +datamart_app = typer.Typer(help="Fase 8 — generazione dbt del datamart (hook)") + + +@datamart_app.command("generate") +def generate_cmd( + session: str = typer.Option(..., "--session", help="Id della sessione."), + mode: str = typer.Option( + ..., "--mode", + help="Variante paziente: 'clear' (in chiaro) o 'pseudonymized' (mascherato).", + ), + config: Path = CONFIG_OPT, +) -> None: + """Innesca la generazione dbt del datamart (Fase 8). Hook: stub non operativo.""" + from tht.cli.phase_cmd import require_phase_or_exit + from tht.datamart import DATAMART_MODES, generate_dbt_datamart + + cfg = _load_config_or_exit(config) + load_session_or_exit(cfg, session) + if mode not in DATAMART_MODES: + typer.secho( + f"ERRORE: mode '{mode}' non valido. Valori: {', '.join(DATAMART_MODES)}", + fg=typer.colors.RED, err=True, + ) + raise typer.Exit(code=1) + require_phase_or_exit(cfg, session, 8) + try: + generate_dbt_datamart(session_dir(cfg, session), mode) + except NotImplementedError as e: + typer.secho( + f"NOTA: {e} (hook, mode={mode}). Nessuna azione eseguita.", + fg=typer.colors.YELLOW, + ) + raise typer.Exit(code=0) + typer.secho(f"OK: dbt del datamart generato (mode={mode}).", fg=typer.colors.GREEN) diff --git a/harness/tht/cli/lsh_cmd.py b/harness/tht/cli/lsh_cmd.py new file mode 100644 index 00000000..c55c90af --- /dev/null +++ b/harness/tht/cli/lsh_cmd.py @@ -0,0 +1,97 @@ +from pathlib import Path + +import typer + +from tht.cli.config_cmd import CONFIG_OPT +from tht.cli.schema_cmd import _load_config_or_exit, physical_path + +lsh_app = typer.Typer(help="Indice LSH su valori dei campi (derivato, rigenerabile)") + + +def _lsh_dir(cfg) -> Path: + return cfg.paths.indexes / "lsh" + + +@lsh_app.command("build") +def build_cmd(config: Path = CONFIG_OPT) -> None: + """Costruisce l'indice LSH dai valori del database e lo salva su pickle.""" + from tht.lshindex import build_index, save_index + 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) + from tht.cli.schema_cmd import annotations_path + + annotations = Annotations.from_yaml(annotations_path(cfg)) + + typer.echo("Estrazione valori (i più frequenti) dalle colonne testuali eligible...") + if cfg.database.transport == "rest": + from tht.db.sampling import unique_values_for_lsh_rest + from tht.rest.client import RestClient + + values, skipped, truncated = unique_values_for_lsh_rest( + RestClient(cfg.rest), physical, cfg.lsh, annotations + ) + else: + from tht.db.connection import make_engine + from tht.db.sampling import unique_values_for_lsh + + values, skipped, truncated = unique_values_for_lsh( + make_engine(cfg.database), physical, cfg.lsh, annotations + ) + n_values = sum(len(v) for t in values.values() for v in t.values()) + typer.echo(f" {n_values} valori da {sum(len(t) for t in values.values())} colonne") + for s in skipped: + typer.secho(f" saltata {s.table}.{s.column}: {s.reason}", fg=typer.colors.YELLOW) + for t in truncated: + typer.secho( + f" troncata {t.table}.{t.column}: indicizzati i {t.indexed} valori più frequenti " + f"(limite max_values_per_column raggiunto; altri valori distinti NON indicizzati)", + fg=typer.colors.YELLOW, + ) + + lsh, minhashes = build_index(values, cfg.lsh, verbose=True) + save_index(lsh, minhashes, cfg.lsh, _lsh_dir(cfg), name=cfg.database.db_schema) + typer.secho( + f"OK: indice LSH ({len(minhashes)} entry) -> {_lsh_dir(cfg)}", fg=typer.colors.GREEN + ) + + +@lsh_app.command("query") +def query_cmd( + keyword: str = typer.Argument(..., help="Termine da cercare, es. 'ablazione'."), + config: Path = CONFIG_OPT, + top: int = typer.Option(10, "--top", help="Numero massimo di risultati."), +) -> None: + """Probe visuale: mostra i candidati LSH per un termine, con score Jaccard.""" + from rich.console import Console + from rich.table import Table + + from tht.lshindex import LshIndexError, load_index, query_index + + cfg = _load_config_or_exit(config) + try: + lsh, minhashes, meta = load_index(_lsh_dir(cfg), name=cfg.database.db_schema) + except LshIndexError as e: + typer.secho(f"ERRORE: {e}", fg=typer.colors.RED) + raise typer.Exit(code=1) + + hits = query_index(lsh, minhashes, keyword, meta, top_n=top) + if not hits: + typer.secho(f"Nessun candidato LSH per '{keyword}'.", fg=typer.colors.YELLOW) + return + table = Table(title=f"Candidati LSH per '{keyword}' ({len(hits)})") + table.add_column("Tabella") + table.add_column("Colonna") + table.add_column("Valore") + table.add_column("Score", justify="right") + for h in hits: + table.add_row(h.table, h.column, h.value, f"{h.score:.3f}") + Console().print(table) diff --git a/harness/tht/cli/sql_cmd.py b/harness/tht/cli/sql_cmd.py new file mode 100644 index 00000000..336144d4 --- /dev/null +++ b/harness/tht/cli/sql_cmd.py @@ -0,0 +1,260 @@ +import json +from pathlib import Path + +import typer + +from tht.cli.config_cmd import CONFIG_OPT +from tht.cli.schema_cmd import _load_config_or_exit, physical_path + +sql_app = typer.Typer(help="Validazione ed esecuzione controllata di SQL (read-only)") + + +def _read_sql(file: Path) -> str: + if not file.exists(): + typer.secho(f"ERRORE: file non trovato: {file}", fg=typer.colors.RED, err=True) + raise typer.Exit(code=1) + return file.read_text() + + +def _load_physical_or_exit(cfg): + from tht.mschema.models import PhysicalSchema + + 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) + return PhysicalSchema.from_yaml(phys_file) + + +def require_action(cfg, action: str) -> None: + if action not in cfg.execution.allow: + typer.secho( + f"ERRORE: azione '{action}' non consentita dalla policy " + f"(execution.allow = {cfg.execution.allow}).", + fg=typer.colors.RED, err=True, + ) + raise typer.Exit(code=1) + + +def promoted_tables_for(cfg, session_id: str | None) -> set[str] | None: + if session_id is None: + return None + linking_path = cfg.paths.sessions / session_id / "schema_linking.json" + if not linking_path.exists(): + return None + from tht.session.models import SchemaLinking + + linking = SchemaLinking.model_validate(json.loads(linking_path.read_text())) + return { + c.name for c in linking.candidates + if c.kind == "table" and c.decision == "promoted" + } + + +def validate_or_exit(cfg, sql: str, session_id: str | None): + """Validazione statica; stampa errori/warning. Exit 1 sugli errori.""" + from tht.sqlcheck import validate_sql + + result = validate_sql( + sql, + physical=_load_physical_or_exit(cfg), + promoted_tables=promoted_tables_for(cfg, session_id), + forbidden_functions=set(cfg.execution.forbidden_functions), + ) + for w in result.warnings: + typer.secho(f" warning: {w}", fg=typer.colors.YELLOW) + if not result.ok: + for e in result.errors: + typer.secho(f" ERRORE: {e}", fg=typer.colors.RED, err=True) + raise typer.Exit(code=1) + return result + + +def _ro_engine(cfg): + """Engine sul target con search_path impostato allo schema (nomi non qualificati).""" + from sqlalchemy import create_engine + + db = cfg.database + url = f"postgresql+psycopg2://{db.user}:{db.password}@{db.host}:{db.port}/{db.database}" + return create_engine( + url, echo=False, + connect_args={"options": f"-csearch_path={db.db_schema}"}, + ) + + +def _rest_client(cfg): + from tht.rest.client import RestClient + + return RestClient(cfg.rest) + + +def do_explain(cfg, sql: str): + """EXPLAIN secondo il transport configurato (direct|rest).""" + if cfg.database.transport == "rest": + from tht.rest.execute import explain_rest + + return explain_rest(_rest_client(cfg), sql) + from tht.execute import explain + + return explain(_ro_engine(cfg), sql, timeout_ms=cfg.execution.statement_timeout_ms) + + +def do_run(cfg, sql: str, *, limit: int): + """Esecuzione controllata secondo il transport configurato (direct|rest).""" + if cfg.database.transport == "rest": + from tht.rest.execute import run_controlled_rest + + return run_controlled_rest(_rest_client(cfg), sql, limit=limit) + from tht.execute import run_controlled + + return run_controlled( + _ro_engine(cfg), sql, limit=limit, timeout_ms=cfg.execution.statement_timeout_ms + ) + + +@sql_app.command("validate") +def validate_cmd( + file: Path = typer.Argument(..., help="File SQL da validare."), + session: str = typer.Option(None, "--session", help="Verifica anche il perimetro promosso."), + config: Path = CONFIG_OPT, +) -> None: + """Parse, read-only strutturale, blacklist funzioni, oggetti vs mschema.""" + cfg = _load_config_or_exit(config) + validate_or_exit(cfg, _read_sql(file), session) + typer.secho("OK: SQL valido (statico).", fg=typer.colors.GREEN) + + +@sql_app.command("explain") +def explain_cmd( + file: Path = typer.Argument(...), + session: str = typer.Option(None, "--session"), + config: Path = CONFIG_OPT, +) -> None: + """EXPLAIN (FORMAT JSON) con sintesi e warning dal piano. Mai ANALYZE.""" + from tht.execute import ExecutionError + from tht.execute.warnings import plan_warnings + + cfg = _load_config_or_exit(config) + require_action(cfg, "explain") + sql = _read_sql(file) + validate_or_exit(cfg, sql, session) + try: + plan = do_explain(cfg, sql) + except ExecutionError as e: + typer.secho(f"ERRORE: {e}", fg=typer.colors.RED, err=True) + raise typer.Exit(code=1) + typer.echo(f"costo totale stimato: {plan.total_cost}") + typer.echo(f"righe stimate: {plan.plan_rows}") + typer.echo(f"nodi del piano: {', '.join(plan.node_types)}") + for w in plan_warnings(plan, cfg.execution): + typer.secho(f" warning: {w}", fg=typer.colors.YELLOW) + + +@sql_app.command("preview") +def preview_cmd( + file: Path = typer.Argument(...), + limit: int = typer.Option(None, "--limit", help="Default: execution.max_preview_rows."), + session: str = typer.Option(None, "--session"), + config: Path = CONFIG_OPT, +) -> None: + """Esecuzione controllata con LIMIT iniettato; aggregati mostrati per interi.""" + from rich.console import Console + from rich.table import Table + + from tht.execute import ExecutionError + from tht.execute.warnings import runtime_warnings, static_warnings + + cfg = _load_config_or_exit(config) + require_action(cfg, "preview") + sql = _read_sql(file) + check = validate_or_exit(cfg, sql, session) + effective_limit = limit if limit is not None else cfg.execution.max_preview_rows + try: + result = do_run(cfg, sql, limit=effective_limit) + except ExecutionError as e: + typer.secho(f"ERRORE: {e}", fg=typer.colors.RED, err=True) + raise typer.Exit(code=1) + + cells = len(result.rows) * len(result.columns) + is_aggregate = ( + "aggregate" in cfg.execution.allow + and not result.truncated + and cells <= cfg.execution.max_aggregate_cells + ) + title = "Risultato aggregato" if is_aggregate else f"Preview (limit {effective_limit})" + # titolo come riga di testo (non come title della tabella rich, che verrebbe + # spezzato sulla larghezza ridotta della tabella per query strette) + typer.echo(f"{title} — {result.execution_ms} ms") + table = Table() + for col in result.columns: + table.add_column(col) + for row in result.rows: + table.add_row(*[str(v) for v in row]) + Console().print(table) + + for w in static_warnings(check.ast) + runtime_warnings(result, cfg.execution): + typer.secho(f" warning: {w}", fg=typer.colors.YELLOW) + + +def _session_sql_file(cfg, session_id: str) -> Path: + from tht.cli.session_cmd import load_session_or_exit, session_dir + + load_session_or_exit(cfg, session_id) + sql_file = session_dir(cfg, session_id) / "sql_final.sql" + if not sql_file.exists(): + typer.secho(f"ERRORE: {sql_file} non trovato.", fg=typer.colors.RED, err=True) + raise typer.Exit(code=1) + return sql_file + + +@sql_app.command("save") +def save_cmd( + dest: Path = typer.Argument(..., help="Percorso di destinazione del file SQL."), + session: str = typer.Option(..., "--session"), + config: Path = CONFIG_OPT, +) -> None: + """Salva una copia di sql_final.sql nel percorso indicato (su richiesta esplicita).""" + cfg = _load_config_or_exit(config) + sql_file = _session_sql_file(cfg, session) + dest.parent.mkdir(parents=True, exist_ok=True) + dest.write_text(sql_file.read_text()) + typer.secho(f"OK: SQL salvato in {dest}", fg=typer.colors.GREEN) + + +@sql_app.command("export") +def export_cmd( + dest: Path = typer.Argument(..., help="Percorso del CSV di destinazione."), + session: str = typer.Option(..., "--session"), + config: Path = CONFIG_OPT, +) -> None: + """Esegue sql_final.sql nelle 4 reti e scrive i risultati in CSV + (cap: execution.max_export_rows).""" + import csv as csv_mod + + from tht.execute import ExecutionError + + cfg = _load_config_or_exit(config) + require_action(cfg, "export") + sql_file = _session_sql_file(cfg, session) + sql = sql_file.read_text() + validate_or_exit(cfg, sql, session) + try: + result = do_run(cfg, sql, limit=cfg.execution.max_export_rows) + except ExecutionError as e: + typer.secho(f"ERRORE: {e}", fg=typer.colors.RED, err=True) + raise typer.Exit(code=1) + + dest.parent.mkdir(parents=True, exist_ok=True) + with dest.open("w", newline="") as f: + writer = csv_mod.writer(f) + writer.writerow(result.columns) + writer.writerows(result.rows) + typer.secho(f"OK: {len(result.rows)} righe esportate in {dest}", fg=typer.colors.GREEN) + if result.truncated: + typer.secho( + f" warning: risultato troncato al cap di {cfg.execution.max_export_rows} righe " + f"(execution.max_export_rows)", fg=typer.colors.YELLOW, + )