diff --git a/harness/tht/cli/__init__.py b/harness/tht/cli/__init__.py index 913bf9f2..f9c56709 100644 --- a/harness/tht/cli/__init__.py +++ b/harness/tht/cli/__init__.py @@ -37,8 +37,13 @@ def main( from tht.cli.config_cmd import config_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.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.vector_cmd import vector_app # noqa: E402 @@ -47,3 +52,8 @@ app.add_typer(config_app, name="config") app.add_typer(schema_app, name="schema") app.add_typer(session_app, name="session") app.add_typer(vector_app, name="vector") +app.add_typer(memory_app, name="memory") +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") diff --git a/harness/tht/cli/db_cmd.py b/harness/tht/cli/db_cmd.py new file mode 100644 index 00000000..549262d1 --- /dev/null +++ b/harness/tht/cli/db_cmd.py @@ -0,0 +1,175 @@ +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 can_create_in_schema, make_engine, ping, writable_tables +from tht.db.fetch_ca import CaFetchError, describe_pem, fetch_chain_pem, parse_host_port + +db_app = typer.Typer(help="Operazioni sul database target") + + +def _ping_rest(cfg, schema: str) -> None: + """Health check via REST. Il read-only è garantito strutturalmente dall'API + (ammette solo SELECT/WITH): non serve il controllo dei privilegi di scrittura.""" + from tht.rest.client import RestClient, RestError + + try: + info = RestClient(cfg.rest).ping() + except RestError as e: + typer.secho(f"ERRORE di connessione: {e}", fg=typer.colors.RED, err=True) + raise typer.Exit(code=1) + if not info.get("db_connected") or not info.get("schema_accessible"): + typer.secho( + f"ERRORE: DWH non accessibile via REST (risposta: {info}).", + fg=typer.colors.RED, err=True, + ) + raise typer.Exit(code=1) + typer.secho( + f"OK: connesso via REST a {cfg.rest.base_url} (schema {schema})", fg=typer.colors.GREEN + ) + typer.secho( + "OK: accesso read-only garantito dall'API (solo SELECT/WITH).", fg=typer.colors.GREEN + ) + + +@db_app.command("ping") +def ping_cmd(config: Path = CONFIG_OPT) -> None: + """Testa la connessione e verifica che l'utente sia effettivamente read-only.""" + 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) + schema = cfg.database.db_schema + if cfg.database.transport == "rest": + _ping_rest(cfg, schema) + return + engine = make_engine(cfg.database) + try: + ping(engine) + except OperationalError as e: + typer.secho(f"ERRORE di connessione: {e.orig}", fg=typer.colors.RED, err=True) + raise typer.Exit(code=1) + typer.secho(f"OK: connesso a {cfg.database.database} (schema {schema})", fg=typer.colors.GREEN) + + writable = writable_tables(engine, schema) + can_create = can_create_in_schema(engine, schema) + if writable or can_create: + typer.secho( + f"ERRORE: l'utente '{cfg.database.user}' NON e' read-only.", fg=typer.colors.RED, err=True + ) + if writable: + typer.echo(f" Tabelle scrivibili: {', '.join(writable[:10])}", err=True) + if can_create: + typer.echo(f" L'utente puo' creare oggetti nello schema {schema}.", err=True) + typer.echo(" Crea un ruolo read-only con scripts/create_readonly_role.sql.", err=True) + raise typer.Exit(code=2) + typer.secho("OK: l'utente e' read-only sullo schema target.", fg=typer.colors.GREEN) + + +@db_app.command("fetch-ca") +def fetch_ca_cmd( + config: Path = CONFIG_OPT, + out: Path = typer.Option( + Path("config/ca-chain.pem"), "--out", "-o", help="File PEM di destinazione." + ), +) -> None: + """Scarica la catena CA presentata dal server REST e la salva in un bundle PEM. + + Utile sulle postazioni *workstation* dietro una CA interna: il file va poi puntato + con `THT_SSL_CA` nel `.env` (lo consuma `requests`). NON tocca il trust store dell'OS. + """ + 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) + if cfg.rest is None: + typer.secho( + "ERRORE: comando pensato per postazioni con accesso REST. " + "Configura la sezione `rest` (base_url) nel workspace yaml.", + fg=typer.colors.RED, err=True, + ) + raise typer.Exit(code=1) + + # Host del DWH REST; se il vector REST è su host diverso, aggiungi anche quello. + targets: list[tuple[str, int]] = [parse_host_port(cfg.rest.base_url)] + if cfg.vector_rest is not None: + vec = parse_host_port(cfg.vector_rest.base_url) + if vec not in targets: + targets.append(vec) + + pems: list[str] = [] + seen: set[str] = set() + try: + for host, port in targets: + typer.echo(f"Recupero catena TLS da {host}:{port} ...") + for pem in fetch_chain_pem(host, port, cfg.rest.timeout): + key = pem.strip() + if key not in seen: + seen.add(key) + pems.append(pem if pem.endswith("\n") else pem + "\n") + except CaFetchError as e: + typer.secho(f"ERRORE: {e}", fg=typer.colors.RED, err=True) + raise typer.Exit(code=1) + + out.parent.mkdir(parents=True, exist_ok=True) + out.write_text("".join(pems)) + abs_out = out.resolve() + + typer.secho( + f"OK: salvati {len(pems)} certificati -> {abs_out}", fg=typer.colors.GREEN + ) + typer.echo("Controlla che corrispondano alla CA interna attesa:") + for i, pem in enumerate(pems, 1): + desc = describe_pem(pem) + if desc: + typer.echo(f" [{i}] {desc}") + typer.echo( + "NOTA: il recupero non verifica la fiducia; la verifica avviene al prossimo " + "`tht db ping`, quando THT_SSL_CA punta a questo bundle." + ) + + # Controlla se la config *effettiva* (THT_SSL_CA nel .env + ssl_ca nel workspace yaml) + # gia' punta al bundle appena salvato: in caso contrario, il certificato non e' + # ancora collegato e il `db ping` fallirebbe ancora con certificate verify failed. + endpoints = [("rest", cfg.rest)] + if cfg.vector_rest is not None: + endpoints.append(("vector_rest", cfg.vector_rest)) + not_wired = [name for name, ep in endpoints if not _points_to(ep.ssl_ca, abs_out)] + + if not not_wired: + typer.secho( + "OK: THT_SSL_CA punta gia' a questo bundle (rest" + + (", vector_rest" if cfg.vector_rest is not None else "") + + "). Lancia `tht db ping`.", + fg=typer.colors.GREEN, + ) + return + + typer.secho( + "\nATTENZIONE: il certificato non e' ancora collegato " + f"(ssl_ca di {', '.join(not_wired)} non punta a questo bundle). Per attivarlo:", + fg=typer.colors.YELLOW, + ) + typer.echo(f" 1. nel .env imposta: THT_SSL_CA={abs_out}") + typer.echo( + " 2. in il workspace yaml, sotto `rest:`" + + (" e `vector_rest:`" if cfg.vector_rest is not None else "") + + ", decommenta/aggiungi:" + ) + typer.echo(" ssl_ca: ${THT_SSL_CA}") + typer.echo("Poi rilancia `tht db ping`.") + + +def _points_to(ssl_ca: str | None, target: Path) -> bool: + """True se `ssl_ca` (path risolto) coincide col bundle appena salvato.""" + if not ssl_ca: + return False + try: + return Path(ssl_ca).resolve() == target + except (OSError, ValueError): + return False diff --git a/harness/tht/cli/decision_cmd.py b/harness/tht/cli/decision_cmd.py new file mode 100644 index 00000000..fa9cd33b --- /dev/null +++ b/harness/tht/cli/decision_cmd.py @@ -0,0 +1,89 @@ +from pathlib import Path +from typing import get_args + +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 + +decision_app = typer.Typer(help="Decisioni del reviewer (per sessione, append-only)") + + +@decision_app.command("add") +def add_cmd( + session: str = typer.Option(..., "--session", help="Id della sessione."), + type: str = typer.Option(..., "--type", help="Tipo di decisione."), + subject: str = typer.Option(..., "--subject", help="Oggetto (tabella, colonna, concetto)."), + detail: str = typer.Option("", "--detail"), + rationale: str = typer.Option("", "--rationale"), + config: Path = CONFIG_OPT, +) -> None: + """Registra una decisione del reviewer nella sessione.""" + from tht.decisions import DecisionType, append_decision + + cfg = _load_config_or_exit(config) + load_session_or_exit(cfg, session) + valid = get_args(DecisionType) + if type not in valid: + typer.secho( + f"ERRORE: tipo '{type}' non valido. Tipi: {', '.join(valid)}", + fg=typer.colors.RED, err=True, + ) + raise typer.Exit(code=1) + from tht.cli.phase_cmd import require_phase_or_exit + from tht.workflow import load_workflow + + require_phase_or_exit(cfg, session, load_workflow().decision_min_phase(type)) + sdir = session_dir(cfg, session) + if type == "cte_approved": + # Un CTE si approva solo se appartiene al piano persistito. Senza piano + # (o con subject fuori piano) l'approvazione e' priva di significato: + # rifiutala (exit 5) invece di sporcare il ledger. Regressione quo-8, + # dove fu registrato un cte_approved:cte_plan senza alcun cte_plan.json. + from tht.phase import cte_plan + + plan = cte_plan(sdir) + if not plan: + typer.secho( + "ERRORE: nessun piano CTE (cte_plan.json) in sessione. Persisti prima " + "il piano con reviewer_confirm kind:'cte_plan', poi approva i CTE.", + fg=typer.colors.RED, err=True, + ) + raise typer.Exit(code=5) + if subject not in plan: + typer.secho( + f"ERRORE: '{subject}' non e' un CTE del piano ({', '.join(plan)}). " + "Usa un nome di CTE del piano.", + fg=typer.colors.RED, err=True, + ) + raise typer.Exit(code=5) + record = append_decision( + sdir, type=type, subject=subject, + detail=detail, rationale=rationale, + ) + typer.secho(f"OK: decisione [{record.seq}] {record.type}: {record.subject}", + fg=typer.colors.GREEN) + + +@decision_app.command("list") +def list_cmd( + session: str = typer.Option(..., "--session"), + config: Path = CONFIG_OPT, +) -> None: + """Elenca le decisioni della sessione.""" + from tht.decisions import list_decisions + + cfg = _load_config_or_exit(config) + load_session_or_exit(cfg, session) + decisions = list_decisions(session_dir(cfg, session)) + if not decisions: + typer.echo("Nessuna decisione registrata.") + return + for d in decisions: + line = f"[{d.seq}] {d.ts:%Y-%m-%d %H:%M} {d.type}: {d.subject}" + if d.detail: + line += f" — {d.detail}" + if d.rationale: + line += f" ({d.rationale})" + typer.echo(line) diff --git a/harness/tht/cli/evidence_cmd.py b/harness/tht/cli/evidence_cmd.py new file mode 100644 index 00000000..fcd3165a --- /dev/null +++ b/harness/tht/cli/evidence_cmd.py @@ -0,0 +1,76 @@ +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._guards import require_vector_write_allowed + +evidence_app = typer.Typer(help="Generazione e gestione delle evidence") + + +def evidence_root(cfg) -> Path: + return cfg.paths.artifacts / "evidence" + + +def _require_evidence_cfg(cfg): + if cfg.evidence is None: + typer.secho( + "ERRORE: sezione `evidence` mancante nel workspace yaml (serve `source_root`).", + fg=typer.colors.RED, err=True, + ) + raise typer.Exit(code=1) + return cfg.evidence + + +@evidence_app.command("extract") +def extract_cmd(config: Path = CONFIG_OPT) -> None: + """Rispecchia la cartella curata dell'ETL in artifacts/evidence/. + + artifacts/evidence/ e' un artefatto derivato e rigenerabile: viene riallineato a + ogni estrazione (le evidence rimosse a monte spariscono anche qui). La gerarchia + per dominio della cartella sorgente viene preservata. + """ + import shutil + + from tht.evidence.extract import load_curated + + cfg = _load_config_or_exit(config) + sources = _require_evidence_cfg(cfg) + out_dir = evidence_root(cfg) + if out_dir.exists(): + shutil.rmtree(out_dir) + written = 0 + for rel, doc in load_curated(sources): + doc.save(out_dir / rel) + written += 1 + if written == 0: + typer.secho( + f"ATTENZIONE: nessuna evidence trovata in " + f"{sources.source_root / sources.evidence_dir}", + fg=typer.colors.YELLOW, err=True, + ) + raise typer.Exit(code=1) + typer.secho(f"OK: {written} evidence rispecchiate in {out_dir}", fg=typer.colors.GREEN) + + +@evidence_app.command("index") +def index_cmd(config: Path = CONFIG_OPT) -> None: + """Embedda e sincronizza su pgvector tutte le evidence presenti in artifacts/.""" + from tht.cli.vector_cmd import ( + _print_stats, + make_embedder, + open_store, + require_vector_cfg, + ) + from tht.evidence.model import load_evidence_dir + from tht.vectorstore.records import evidence_records + + cfg = _load_config_or_exit(config) + require_vector_write_allowed(cfg, "evidence index") + require_vector_cfg(cfg) + docs = load_evidence_dir(evidence_root(cfg)) + records = evidence_records(docs, cfg.vector.max_chunk_chars) + store = open_store(cfg, "evidence") + stats = store.sync(records, make_embedder(cfg.embeddings), kinds={"evidence"}) + _print_stats(stats) diff --git a/harness/tht/cli/memory_cmd.py b/harness/tht/cli/memory_cmd.py new file mode 100644 index 00000000..17de2345 --- /dev/null +++ b/harness/tht/cli/memory_cmd.py @@ -0,0 +1,385 @@ +# TODO (drop registry, decisione spec 5): questo modulo e' portato col modello +# registry intatto (load_registry/promote/update_record/delete_record). Le memory +# dovrebbero vivere SOLO nel vectordb (metadata arricchito con subject/detail/rationale +# in Onda 3.1). Riscrivere: promote -> upsert batch vectordb; list/show -> scan +# vectordb; delete -> metadata.status="superseded"; index/clear -> droppati. +# Task separato: la validazione richiede L2 (vectordb reale). +import json +from pathlib import Path + +import typer +from sqlalchemy.exc import OperationalError, ProgrammingError + +from tht.cli.config_cmd import CONFIG_OPT +from tht.cli.schema_cmd import _load_config_or_exit +from tht.cli._guards import require_server_profile, require_vector_write_allowed +from tht.cli.session_cmd import load_session_or_exit, session_dir +from tht.cli.vector_cmd import require_vector_cfg + +memory_app = typer.Typer(help="Review memory (registro canonico + indice pgvector)") + + +def registry_path(cfg) -> Path: + return cfg.paths.artifacts / "memory" / "registry.jsonl" + + +def _resync_memory(cfg): + """Risincronizza l'indice pgvector col registro corrente (incrementale).""" + from tht.cli.vector_cmd import make_embedder, open_store + from tht.memory import load_registry, memory_vector_records + + records = memory_vector_records(load_registry(registry_path(cfg))) + store = open_store(cfg, "memory") + return store.sync(records, make_embedder(cfg.embeddings), kinds={"memory"}) + + +@memory_app.command("promote") +def promote_cmd( + session: str = typer.Option(..., "--session"), + decision: list[int] = typer.Option(None, "--decision", help="Seq da promuovere (ripetibile)."), + preview: bool = typer.Option(False, "--preview", help="Mostra i candidati in JSON, non scrive."), + json_out: bool = typer.Option(False, "--json", help="Output JSON (per Pi)."), + config: Path = CONFIG_OPT, +) -> None: + """Promuove le decisioni SCELTE nel registro globale. Usa --preview per vedere i candidati.""" + import json as _json + + from tht.memory import promote + + cfg = _load_config_or_exit(config) + manifest = load_session_or_exit(cfg, session) + + if preview: + from tht.memory import ( + MAX_PROMOTION_CANDIDATES, preview_promotions, reusable_promotions, + ) + sdir = session_dir(cfg, session) + cand = preview_promotions(sdir, manifest, registry_path(cfg)) + extra = len(reusable_promotions(sdir, manifest, registry_path(cfg))) - len(cand) + payload = [ + {"decision_seq": c.decision_seq, "type": c.type, "subject": c.subject, + "detail": c.detail, "rationale": c.rationale, + "question_context": c.question_context, + "tables": c.tables, "concepts": c.concepts} + for c in cand + ] + if json_out: + typer.echo(_json.dumps(payload, ensure_ascii=False, indent=2)) + elif not payload: + typer.secho("Nessun candidato da promuovere.", fg=typer.colors.YELLOW) + else: + for c in payload: + typer.echo(f" [{c['decision_seq']}] {c['type']}: {c['subject']}") + if extra > 0: + typer.secho( + f"NOTA: mostrati {len(cand)} candidati su {len(cand) + extra} riusabili " + f"(cap {MAX_PROMOTION_CANDIDATES}); gli altri non sono proposti.", + fg=typer.colors.YELLOW, err=True, + ) + return + + if not decision: + typer.secho("ERRORE: indica le decisioni con --decision (vedi `--preview`).", + fg=typer.colors.RED, err=True) + raise typer.Exit(code=1) + + require_server_profile(cfg, "memory promote") + require_vector_cfg(cfg) + promoted = promote( + session_dir(cfg, session), manifest, + seqs=list(decision), registry_path=registry_path(cfg), + ) + if not promoted: + msg = "Nessuna nuova promozione (gia' presenti o seq inesistenti)." + if json_out: + typer.echo(_json.dumps( + {"promoted": [], "indexed": False, "message": msg}, ensure_ascii=False)) + else: + typer.secho(msg, fg=typer.colors.YELLOW) + return + + # Promozione nel registro: riuscita. L'indicizzazione su pgvector puo' fallire + # (tabella mancante o vectordb irraggiungibile da questa postazione): in quel + # caso le memorie restano nel registro ma NON sono trovate da `tht memory + # search` finche' non si reindicizza sul server. `indexed` rende lo stato + # leggibile da Pi, cosi' il reviewer lo vede invece di perderlo nello stderr. + ids = [{"id": r.id, "type": r.type, "subject": r.subject} for r in promoted] + indexed = True + warning = None + try: + _resync_memory(cfg) + except (ProgrammingError, OperationalError): + indexed = False + warning = ( + f"{len(promoted)} memorie promosse nel registro, ma l'indice vettoriale " + "NON e' stato sincronizzato (tabella pgvector mancante o irraggiungibile): " + "NON saranno trovate da `tht memory search` finche' non reindicizzi sul " + "server (`tht vector init`, poi `tht memory index`)." + ) + + if json_out: + typer.echo(_json.dumps( + {"promoted": ids, "indexed": indexed, + "message": warning or f"{len(promoted)} memorie promosse e indicizzate."}, + ensure_ascii=False, indent=2)) + return + for r in promoted: + typer.echo(f" {r.id}: {r.type} {r.subject}") + if indexed: + typer.secho(f"OK: {len(promoted)} memorie promosse e indicizzate.", + fg=typer.colors.GREEN) + else: + typer.secho(f"ATTENZIONE: {warning}", fg=typer.colors.YELLOW) + + +@memory_app.command("clear") +def clear_cmd( + yes: bool = typer.Option(False, "--yes", "-y", help="Salta la richiesta di conferma."), + config: Path = CONFIG_OPT, +) -> None: + """Cancella TUTTA la review memory: registro canonico + indice pgvector (kind=memory).""" + from tht.cli.vector_cmd import make_embedder, open_store, require_direct_vector_cfg + from tht.memory import load_registry + + cfg = _load_config_or_exit(config) + require_server_profile(cfg, "memory clear") + registry = registry_path(cfg) + count = len(load_registry(registry)) + if count == 0: + typer.secho("Nessuna memoria da cancellare.", fg=typer.colors.YELLOW) + return + + if not yes and not typer.confirm( + f"Cancellare definitivamente {count} memorie (registro + indice)?" + ): + typer.secho("Annullato.", fg=typer.colors.YELLOW) + raise typer.Exit(code=1) + + # Indice pgvector: rimuove i record kind=memory (sync con insieme vuoto). + require_direct_vector_cfg(cfg) + store = open_store(cfg, "memory") + store.sync([], make_embedder(cfg.embeddings), kinds={"memory"}) + + # Registro canonico. + registry.unlink() + typer.secho(f"OK: {count} memorie cancellate.", fg=typer.colors.GREEN) + + +@memory_app.command("index") +def index_cmd(config: Path = CONFIG_OPT) -> None: + """Sincronizza il registro memory su pgvector (full-resync).""" + from tht.cli.vector_cmd import _print_stats + + cfg = _load_config_or_exit(config) + require_vector_write_allowed(cfg, "memory index") + require_vector_cfg(cfg) + _print_stats(_resync_memory(cfg)) + + +@memory_app.command("list") +def list_cmd( + type_: str = typer.Option(None, "--type", help="Filtra per tipo decisione."), + session: str = typer.Option(None, "--session", help="Filtra per sessione."), + table: str = typer.Option(None, "--table", help="Filtra per tabella coinvolta."), + concept: str = typer.Option(None, "--concept", help="Filtra per concetto."), + json_out: bool = typer.Option(False, "--json"), + config: Path = CONFIG_OPT, +) -> None: + """Elenca le memorie del registro (filtri combinati in AND).""" + import json as _json + + from rich.console import Console + from rich.table import Table + + from tht.memory import load_registry + + cfg = _load_config_or_exit(config) + recs = load_registry(registry_path(cfg)) + if type_: + recs = [r for r in recs if r.type == type_] + if session: + recs = [r for r in recs if r.session_id == session] + if table: + recs = [r for r in recs if table in r.tables] + if concept: + recs = [r for r in recs if concept in r.concepts] + + if json_out: + typer.echo(_json.dumps([r.model_dump(mode="json") for r in recs], + ensure_ascii=False, indent=2)) + return + if not recs: + typer.secho("Nessuna memoria nel registro.", fg=typer.colors.YELLOW) + return + t = Table(title="Review memory") + for col in ("Id", "Tipo", "Soggetto", "Sessione"): + t.add_column(col) + for r in recs: + t.add_row(r.id, r.type, r.subject, r.session_id) + Console().print(t) + + +@memory_app.command("show") +def show_cmd( + mem_id: str = typer.Argument(..., help="Id memoria (es. mem-0001)."), + json_out: bool = typer.Option(False, "--json"), + config: Path = CONFIG_OPT, +) -> None: + """Mostra una singola memoria.""" + import json as _json + + from tht.memory import load_registry + + cfg = _load_config_or_exit(config) + rec = {r.id: r for r in load_registry(registry_path(cfg))}.get(mem_id) + if rec is None: + typer.secho(f"ERRORE: memoria '{mem_id}' non trovata. Usa `tht memory list`.", + fg=typer.colors.RED, err=True) + raise typer.Exit(code=6) + if json_out: + typer.echo(_json.dumps(rec.model_dump(mode="json"), ensure_ascii=False, indent=2)) + return + for k, v in rec.model_dump(mode="json").items(): + typer.echo(f"{k}: {v}") + + +@memory_app.command("update") +def update_cmd( + mem_id: str = typer.Argument(..., help="Id memoria (es. mem-0001)."), + subject: str = typer.Option(None, "--subject"), + type_: str = typer.Option(None, "--type"), + detail: str = typer.Option(None, "--detail"), + rationale: str = typer.Option(None, "--rationale"), + question_context: str = typer.Option(None, "--question-context"), + tables: str = typer.Option(None, "--tables", help="CSV; \"\" per azzerare."), + concepts: str = typer.Option(None, "--concepts", help="CSV; \"\" per azzerare."), + config: Path = CONFIG_OPT, +) -> None: + """Modifica i campi di merito di una memoria (provenienza immutabile).""" + from typing import get_args + + from tht.memory import MemoryNotFound, update_record + from tht.decisions import DecisionType + + cfg = _load_config_or_exit(config) + + fields: dict = {} + for name, val in (("subject", subject), ("type", type_), ("detail", detail), + ("rationale", rationale), ("question_context", question_context)): + if val is not None: + fields[name] = val + if tables is not None: + fields["tables"] = [t.strip() for t in tables.split(",") if t.strip()] + if concepts is not None: + fields["concepts"] = [c.strip() for c in concepts.split(",") if c.strip()] + + if not fields: + typer.secho("ERRORE: nessun campo da modificare indicato.", + fg=typer.colors.RED, err=True) + raise typer.Exit(code=1) + if "type" in fields and fields["type"] not in get_args(DecisionType): + typer.secho(f"ERRORE: tipo '{fields['type']}' non valido.", + fg=typer.colors.RED, err=True) + raise typer.Exit(code=1) + + require_server_profile(cfg, "memory update") + require_vector_cfg(cfg) + try: + rec = update_record(registry_path(cfg), mem_id, fields) + except MemoryNotFound: + typer.secho(f"ERRORE: memoria '{mem_id}' non trovata. Usa `tht memory list`.", + fg=typer.colors.RED, err=True) + raise typer.Exit(code=6) + _resync_memory(cfg) + typer.secho(f"OK: {rec.id} aggiornata e reindicizzata.", fg=typer.colors.GREEN) + + +@memory_app.command("delete") +def delete_cmd( + mem_id: str = typer.Argument(..., help="Id memoria (es. mem-0001)."), + yes: bool = typer.Option(False, "--yes", "-y", help="Salta la conferma."), + config: Path = CONFIG_OPT, +) -> None: + """Cancella una singola memoria (registro + indice).""" + from tht.memory import MemoryNotFound, delete_record + + cfg = _load_config_or_exit(config) + require_server_profile(cfg, "memory delete") + require_vector_cfg(cfg) + if not yes and not typer.confirm(f"Cancellare definitivamente la memoria '{mem_id}'?"): + typer.secho("Annullato.", fg=typer.colors.YELLOW) + raise typer.Exit(code=1) + try: + delete_record(registry_path(cfg), mem_id) + except MemoryNotFound: + typer.secho(f"ERRORE: memoria '{mem_id}' non trovata. Usa `tht memory list`.", + fg=typer.colors.RED, err=True) + raise typer.Exit(code=6) + _resync_memory(cfg) + typer.secho(f"OK: {mem_id} cancellata e deindicizzata.", fg=typer.colors.GREEN) + + +@memory_app.command("search") +def search_cmd( + question: str = typer.Argument(..., help="Domanda o termini di ricerca."), + top: int = typer.Option(5, "--top"), + session: str = typer.Option( + None, "--session", + help="Esclude le memorie gia' decise (applicate o rifiutate) in questa sessione.", + ), + json_out: bool = typer.Option(False, "--json", help="Output JSON per Pi."), + config: Path = CONFIG_OPT, +) -> None: + """Cerca memorie riapplicabili, ordinate per similarita'. Mai applicate in automatico. + + Con `--session` non ripropone le memorie gia' decise in quella sessione (fix: + memorie scartate riproposte): rifiutate via `memory_rejected` o gia' applicate.""" + from rich.console import Console + from rich.table import Table + + from tht.cli.session_cmd import session_dir + from tht.cli.vector_cmd import make_embedder, open_searcher, require_vector_cfg + from tht.memory import decided_memory_ids, load_registry + from tht.decisions import list_decisions + + cfg = _load_config_or_exit(config) + require_vector_cfg(cfg) + excluded: set[str] = set() + if session is not None: + excluded = decided_memory_ids(list_decisions(session_dir(cfg, session))) + searcher = open_searcher(cfg) + embedder = make_embedder(cfg.embeddings) + hits = searcher.search(embedder.embed_query(question), top_n=top, kinds=["memory"]) + by_id = {r.id: r for r in load_registry(registry_path(cfg))} + + results = [] + for h in hits: + rec = by_id.get(h.ref) + if rec is None: + continue # indice piu' avanti del registro: ignora + if rec.id in excluded: + continue # gia' decisa in questa sessione: non riproporla + results.append({ + "id": rec.id, "type": rec.type, "subject": rec.subject, + "detail": rec.detail, "rationale": rec.rationale, + "question_context": rec.question_context, + "tables": rec.tables, "concepts": rec.concepts, + "session_id": rec.session_id, "score": round(h.similarity, 4), + }) + + if json_out: + typer.echo(json.dumps(results, ensure_ascii=False, indent=2)) + return + if not results: + typer.secho("Nessuna memoria candidata.", fg=typer.colors.YELLOW) + return + table = Table(title=f"Memorie candidate per: {question}") + table.add_column("Id") + table.add_column("Tipo") + table.add_column("Soggetto") + table.add_column("Contesto originale") + table.add_column("Score", justify="right") + for r in results: + table.add_row(r["id"], r["type"], r["subject"], + r["question_context"][:60], f"{r['score']:.3f}") + Console().print(table) diff --git a/harness/tht/cli/search_cmd.py b/harness/tht/cli/search_cmd.py new file mode 100644 index 00000000..12e2fc3e --- /dev/null +++ b/harness/tht/cli/search_cmd.py @@ -0,0 +1,182 @@ +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 + +KIND_MAP = { + "evidence": ["evidence"], + "schema": ["schema_table", "schema_column"], + "values": [], # solo LSH +} + +# Default di `--top` per le famiglie diverse da `schema` (numero di risultati). Per `schema` +# `--top` indica il numero di TABELLE candidate ed e' configurabile via `search.top_schema_tables` +# (recupero ancorato alle tabelle: di ognuna si rendono tutte le colonne + FK). +DEFAULT_TOP_FALLBACK = 10 + +search_app = typer.Typer(help="Ricerca semantica (evidence/schema/values) nel vectorstore") + + +@search_app.command("find") +def search_cmd( + keyword: str = typer.Argument(..., help="Termine da cercare, es. 'ablazione'."), + config: Path = CONFIG_OPT, + top: int | None = typer.Option( + None, "--top", + help="Max risultati; con --kind schema indica il numero di tabelle " + "(default: 12 tabelle per schema, 10 altrimenti).", + ), + kind: str = typer.Option( + None, "--kind", help="Filtra per famiglia: evidence | schema | values." + ), + explain: bool = typer.Option(False, "--explain", help="Mostra anche il testo matchato."), + json_out: bool = typer.Option( + False, "--json", help="Output JSON machine-readable per Pi (sopprime le tabelle a video)." + ), +) -> None: + """Ricerca combinata LSH + pgvector con ranking RRF spiegabile.""" + from rich.console import Console + from rich.table import Table + + from tht.cli.vector_cmd import make_embedder, open_searcher, require_vector_cfg + from tht.lshindex import LshIndexError, load_index, query_index + from tht.search import combined_search + + cfg = _load_config_or_exit(config) + require_vector_cfg(cfg) + if kind is not None and kind not in KIND_MAP: + typer.secho( + f"ERRORE: --kind sconosciuto: {kind} (validi: {', '.join(KIND_MAP)})", + fg=typer.colors.RED, err=True, + ) + raise typer.Exit(code=1) + + if top is None: + top = cfg.search.top_schema_tables if kind == "schema" else DEFAULT_TOP_FALLBACK + + lsh_hits = None + try: + lsh, minhashes, meta = load_index( + cfg.paths.indexes / "lsh", name=cfg.database.db_schema + ) + hits = query_index(lsh, minhashes, keyword, meta, top_n=top * 3) + lsh_hits = [(h.table, h.column, h.value, h.score) for h in hits] + except LshIndexError: + if not json_out: # in JSON mode lo stdout resta puro: niente warning umano + typer.secho( + "ATTENZIONE: indice LSH assente, ricerca solo vettoriale " + "(esegui `tht lsh build`).", fg=typer.colors.YELLOW, + ) + + if kind == "schema": + from tht.cli.schema_cmd import annotations_path, physical_path + from tht.mschema.models import Annotations, PhysicalSchema + from tht.mschema.render import to_mschema_text + from tht.search import schema_tables + + 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) + + candidates = combined_search( + keyword=keyword, lsh_hits=lsh_hits, + store=open_searcher(cfg), embedder=make_embedder(cfg.embeddings), + top=cfg.search.schema_chunk_pool, rrf_k=cfg.search.rrf_k, + kinds=KIND_MAP["schema"], + ) + ranked = schema_tables(candidates, top_tables=top) + if not ranked: + if json_out: + typer.echo(json.dumps({"tables": [], "mschema": ""}, ensure_ascii=False)) + return + typer.secho(f"Nessuna tabella candidata per '{keyword}'.", fg=typer.colors.YELLOW) + return + + physical = PhysicalSchema.from_yaml(phys_file) + annotations = Annotations.from_yaml(annotations_path(cfg)) + selected = [t for t, _ in ranked] + mschema = to_mschema_text(physical, annotations, tables=selected) + + if json_out: + typer.echo(json.dumps( + { + "tables": [{"name": n, "rrf": round(s, 6)} for n, s in ranked], + "mschema": mschema, + }, + ensure_ascii=False, indent=2, + )) + return + + reviewer = Table(title=f"Tabelle candidate per '{keyword}' (top {top}, RRF)") + reviewer.add_column("#", justify="right") + reviewer.add_column("Tabella") + reviewer.add_column("RRF", justify="right") + for i, (name, score) in enumerate(ranked, start=1): + reviewer.add_row(str(i), name, f"{score:.4f}") + Console().print(reviewer) + Console().print(mschema) + return + + if kind == "values": + results = [] + kinds = None + else: + kinds = KIND_MAP.get(kind) if kind else None + results = combined_search( + keyword=keyword, lsh_hits=lsh_hits if kind != "evidence" else None, + store=open_searcher(cfg), embedder=make_embedder(cfg.embeddings), + top=top, rrf_k=cfg.search.rrf_k, kinds=kinds, + ) + + if kind == "values": + if json_out: + typer.echo(json.dumps( + [{"table": t, "column": c, "value": v, "score": round(s, 6)} + for t, c, v, s in (lsh_hits or [])[:top]], + ensure_ascii=False, indent=2, + )) + return + if not lsh_hits: + typer.secho("Nessun match LSH.", fg=typer.colors.YELLOW) + return + table = Table(title=f"Match LSH per '{keyword}'") + table.add_column("Tabella.Colonna") + table.add_column("Valore") + table.add_column("Score", justify="right") + for t, c, v, s in lsh_hits[:top]: + table.add_row(f"{t}.{c}", v, f"{s:.3f}") + Console().print(table) + return + + if json_out: + typer.echo(json.dumps( + [r.model_dump() for r in results], ensure_ascii=False, indent=2 + )) + return + if not results: + typer.secho(f"Nessun candidato per '{keyword}'.", fg=typer.colors.YELLOW) + return + table = Table(title=f"Candidati per '{keyword}' (RRF, k={cfg.search.rrf_k})") + table.add_column("Candidato") + table.add_column("Tipo") + table.add_column("Segnali") + table.add_column("RRF", justify="right") + table.add_column("Status") + if explain: + table.add_column("Testo") + for r in results: + signals = " · ".join( + f"{name} #{s['rank']} ({s['score']})" for name, s in r.signals.items() + ) + row = [r.label, r.kind, signals, f"{r.rrf:.4f}", r.status] + if explain: + row.append((r.content[:120] + "…") if len(r.content) > 120 else r.content) + table.add_row(*row) + Console().print(table)