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 _run_transport(cfg, sql: str, *, limit: int): """Dispatch all'esecutore controllato secondo il transport (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 ) def do_run(cfg, sql: str, *, limit: int, offset: int = 0): """Esecuzione controllata secondo il transport configurato (direct|rest). Per offset == 0: path invariato (LIMIT iniettato dall'esecutore via AST, +1 per rilevare il troncamento). Per offset > 0: la query viene wrappata in `SELECT * FROM (...) LIMIT (N+1) OFFSET M`. Il +1 e' essenziale: l'esecutore vede gia' un LIMIT esterno, quindi il suo `_inject_limit` non inietta nulla (bail perche' un LIMIT E' PRESENTE, non assente) e non rileverebbe mai il troncamento. Recuperando N+1 righe qui calcoliamo noi `truncated = len(rows) > N` e ritagliamo a N. NON rimuovere il LIMIT del wrapper pensando sia ridondante: e' l'unico cap effettivo per il path con offset. """ if offset == 0: return _run_transport(cfg, sql, limit=limit) from dataclasses import replace from tht.execute.limit import inject_limit_offset probe_limit = limit + 1 wrapped = inject_limit_offset(sql, limit=probe_limit, offset=offset) # limit=probe_limit cosi' l'esecutore ritaglia a N+1 (non a N) e ci lascia la riga sonda. result = _run_transport(cfg, wrapped, limit=probe_limit) truncated = len(result.rows) > limit return replace(result, rows=result.rows[:limit], truncated=truncated) @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."), offset: int = typer.Option(0, "--offset", help="Riga di partenza (0-based) per il paging."), session: str = typer.Option(None, "--session"), json_out: bool = typer.Option(False, "--json", help="Output JSON puro per il backend (sopprime tabella rich)."), config: Path = CONFIG_OPT, ) -> None: """Esecuzione controllata con LIMIT iniettato; aggregati mostrati per interi.""" from tht.execute import ExecutionError 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, offset=offset) except ExecutionError as e: typer.secho(f"ERRORE: {e}", fg=typer.colors.RED, err=True) raise typer.Exit(code=1) if json_out: typer.echo(json.dumps( { "columns": result.columns, "rows": [list(r) for r in result.rows], "execution_ms": result.execution_ms, "truncated": result.truncated, "limit": effective_limit, "offset": offset, }, ensure_ascii=False, )) return from rich.console import Console from rich.table import Table from tht.execute.warnings import runtime_warnings, static_warnings 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, )