feat(harness): port sql/cte/datamart/lsh cmd (Onda 4) — CLI completa F1→F8

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.
This commit is contained in:
2026-06-27 14:16:35 +02:00
parent 99b01b0407
commit 91a374492c
5 changed files with 588 additions and 0 deletions
+8
View File
@@ -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")
+181
View File
@@ -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/<id>/ctes/<nome>.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/<id>/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)
+42
View File
@@ -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)
+97
View File
@@ -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)
+260
View File
@@ -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,
)