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 from tht.config import workspace_id_for_config KIND_MAP = { "evidence": ["evidence"], "schema": ["schema_table", "schema_column"], "values": [], # solo LSH "formula": ["evidence"], } _STAGE_PURPOSES = { "clarification": "disambiguation", "rewriting": "rewriting", "schema_linking": "schema_linking", "cte": "sql_generation", "final_sql": "sql_generation", } # 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") def search_formula_evidence(keyword: str, *, searcher, embedder, top: int): """Search only published Formula Evidence through the shared typed facade.""" from tht.evidence import EvidenceSearchContext, search_evidence return search_evidence( keyword, "sql_generation", EvidenceSearchContext(required_kinds=("formula",)), searcher=searcher, embedder=embedder, top_n=top, ) def _leased_dwh_snapshot(cfg, context: typer.Context): from tht.jobs.dwh_pipeline import lease_dwh_snapshot lease = lease_dwh_snapshot(cfg) snapshot = lease.__enter__() context.call_on_close(lambda: lease.__exit__(None, None, None)) return snapshot @search_app.command("evidence") def evidence_search_cmd( ctx: typer.Context, query: str = typer.Argument(..., help="Domanda o contesto dello stage."), stage: str = typer.Option(..., "--stage", help="Stage semantico chiamante."), config: Path = CONFIG_OPT, session: str | None = typer.Option(None, "--session", help="Sessione per la ricevuta minima."), concept: list[str] = typer.Option([], "--concept"), table: list[str] = typer.Option([], "--table"), column: list[str] = typer.Option([], "--column"), require_kind: list[str] = typer.Option([], "--require-kind"), require_concept: list[str] = typer.Option([], "--require-concept"), require_table: list[str] = typer.Option([], "--require-table"), require_column: list[str] = typer.Option([], "--require-column"), top: int = typer.Option(10, "--top"), json_out: bool = typer.Option(False, "--json"), ) -> None: """Run one typed, purpose-bound Evidence search for a semantic workflow stage.""" from tht.cli.vector_cmd import make_embedder, open_searcher, require_vector_cfg from tht.evidence import ( EvidenceReceipt, EvidenceSearchContext, active_searcher, replace_evidence_receipt, search_evidence, validate_corpus_workspace, ) purpose = _STAGE_PURPOSES.get(stage) if purpose is None: raise typer.BadParameter("stage must be clarification, rewriting, schema_linking, cte, or final_sql") cfg = _load_config_or_exit(config) workspace_id = workspace_id_for_config(cfg, config) validate_corpus_workspace(cfg, workspace_id) require_vector_cfg(cfg) outcome = search_evidence( query, purpose, EvidenceSearchContext( concepts=tuple(concept), tables=tuple(table), columns=tuple(column), required_kinds=tuple(require_kind), required_concepts=tuple(require_concept), required_tables=tuple(require_table), required_columns=tuple(require_column), ), searcher=active_searcher(cfg, open_searcher(cfg), workspace_id=workspace_id), embedder=make_embedder(cfg.embeddings), top_n=top, ) if outcome.status == "unavailable": payload = {"status": outcome.status, "code": outcome.code, "message": outcome.message} if json_out: typer.echo(json.dumps(payload, ensure_ascii=False)) else: typer.secho(f"ERRORE: {outcome.message}", fg=typer.colors.RED, err=True) raise typer.Exit(1) if session: from tht.cli.session_cmd import load_session_or_exit, session_repository load_session_or_exit(cfg, session) replace_evidence_receipt(session_repository(cfg), session, EvidenceReceipt( stage=stage, purpose=purpose, vector_generation=outcome.vector_generation or "", evidence_ids=tuple(result.evidence_id for result in outcome.results), )) payload = { "status": "available", "vector_generation": outcome.vector_generation, "results": [ {"evidence_id": result.evidence_id, "title": result.title, "kind": result.kind, "excerpts": list(result.excerpts), "provenance": result.provenance, "citation": result.citation, "document_id": result.document_id} for result in outcome.results ], } if json_out: typer.echo(json.dumps(payload, ensure_ascii=False, indent=2)) @search_app.command("find") def search_cmd( ctx: typer.Context, 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 | formula." ), 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 + semantic search 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.evidence import active_searcher, validate_corpus_workspace from tht.lshindex import LshIndexError, load_index, query_index from tht.search import combined_search cfg = _load_config_or_exit(config) workspace_id = workspace_id_for_config(cfg, config) validate_corpus_workspace(cfg, workspace_id) 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 if kind == "formula": require_vector_cfg(cfg) outcome = search_formula_evidence( keyword, searcher=active_searcher(cfg, open_searcher(cfg), workspace_id=workspace_id), embedder=make_embedder(cfg.embeddings), top=top, ) if outcome.status == "unavailable": payload = {"status": outcome.status, "code": outcome.code, "message": outcome.message} if json_out: typer.echo(json.dumps(payload, ensure_ascii=False)) else: typer.secho(f"ERRORE: {outcome.message}", fg=typer.colors.RED, err=True) raise typer.Exit(1) if json_out: typer.echo(json.dumps([ { "evidence_id": result.evidence_id, "title": result.title, "kind": result.kind, "excerpts": list(result.excerpts), "provenance": result.provenance, "citation": result.citation, "document_id": result.document_id, } for result in outcome.results ], ensure_ascii=False, indent=2)) return typer.secho( "ATTENZIONE: lo store formule legacy non viene più consultato; " "sono disponibili solo Formula Evidence pubblicate.", fg=typer.colors.YELLOW, err=True, ) if not outcome.results: typer.secho(f"Nessuna formula per '{keyword}'.", fg=typer.colors.YELLOW) return table = Table(title=f"Formula Evidence per '{keyword}'") table.add_column("Formula") table.add_column("Provenienza") table.add_column("Estratto") for result in outcome.results: excerpt = result.excerpts[0] if result.excerpts else "" table.add_row(result.title, result.citation, excerpt[:120]) Console().print(table) return dwh_snapshot = _leased_dwh_snapshot(cfg, ctx) require_vector_cfg(cfg) runtime_searcher = active_searcher( cfg, open_searcher(cfg), workspace_id=workspace_id, ) lsh_hits = None try: lsh, minhashes, meta = load_index( dwh_snapshot.lsh_dir, 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 preprocess dwh --steps lsh`).", fg=typer.colors.YELLOW, ) if kind == "schema": from tht.cli.schema_cmd import annotations_path from tht.mschema.models import Annotations, PhysicalSchema from tht.mschema.render import to_mschema_text from tht.search import schema_tables phys_file = dwh_snapshot.physical 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=runtime_searcher, 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=runtime_searcher, 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) # Dimensioni fisse del pack (niente config: il pack deve restare piccolo perche' # entra nel contesto del modello in un turno solo). PACK_EVIDENCE_TOP = 5 PACK_SOLVED_TOP = 3 PACK_EXCERPT_CHARS = 400 @search_app.command("pack") def pack_cmd( ctx: typer.Context, question: str = typer.Argument(..., help="La domanda in linguaggio naturale."), config: Path = CONFIG_OPT, session: str = typer.Option( None, "--session", help="Scrive il pack in sessions//retrieval_pack.md." ), json_out: bool = typer.Option(False, "--json", help="Output JSON (per Pi)."), ) -> None: """Context-pack F1: tabelle candidate + evidence + domande risolte in UNA chiamata. Un solo embedding della domanda, riusato per le tre ricerche vettoriali. Degrado gentile: se Ollama/vectordb non rispondono, le sezioni restano vuote con un'avvertenza (exit 0) — la sessione prosegue con le ricerche live. """ from sqlalchemy.exc import OperationalError from tht.cli.vector_cmd import make_embedder, open_searcher, require_vector_cfg from tht.evidence import ( EvidenceSearchContext, active_searcher, build_retrieval_entries, search_evidence, validate_corpus_workspace, ) from tht.memory import SOLVED_KIND from tht.ports.vector import VectorReadUnavailable, VectorStoreError from tht.search import combined_search, schema_tables from tht.vectorstore.embeddings import EmbeddingsError cfg = _load_config_or_exit(config) workspace_id = workspace_id_for_config(cfg, config) validate_corpus_workspace(cfg, workspace_id) dwh_snapshot = _leased_dwh_snapshot(cfg, ctx) require_vector_cfg(cfg) tables: list[dict] = [] evidence: list[dict] = [] solved: list[dict] = [] warnings: list[str] = [] evidence_outcome = None degrade = (VectorStoreError, VectorReadUnavailable, EmbeddingsError, OperationalError) vec = None searcher = embedder = None try: searcher = active_searcher( cfg, open_searcher(cfg), workspace_id=workspace_id, ) embedder = make_embedder(cfg.embeddings) vec = embedder.embed_query(question) except degrade as e: warnings.append(f"retrieval non disponibile ({e}): prosegui con le ricerche live") if vec is not None: descriptions: dict[str, str] = {} phys_file = dwh_snapshot.physical if phys_file.exists(): from tht.mschema.models import PhysicalSchema phys = PhysicalSchema.from_yaml(phys_file) descriptions = {t: tab.comment for t, tab in phys.tables.items()} try: cand = combined_search( keyword=question, lsh_hits=None, store=searcher, embedder=embedder, top=cfg.search.schema_chunk_pool, rrf_k=cfg.search.rrf_k, kinds=KIND_MAP["schema"], query_vec=vec, ) tables = [ {"name": n, "rrf": round(s, 6), "description": descriptions.get(n, "")} for n, s in schema_tables(cand, top_tables=cfg.search.top_schema_tables) ] except degrade as e: warnings.append(f"ricerca schema fallita ({e})") try: evidence_outcome = search_evidence( question, "disambiguation", EvidenceSearchContext(), searcher=searcher, embedder=embedder, top_n=PACK_EVIDENCE_TOP, ) if evidence_outcome.status == "available": evidence = build_retrieval_entries(evidence_outcome.results, excerpt_chars=PACK_EXCERPT_CHARS) else: typer.secho( f"ERRORE: Evidence non disponibile ({evidence_outcome.code})", fg=typer.colors.RED, err=True, ) raise typer.Exit(code=1) except degrade as e: warnings.append(f"ricerca evidence fallita ({e})") try: hits = searcher.search(vec, top_n=PACK_SOLVED_TOP, kinds=[SOLVED_KIND]) solved = [ { "session_id": h.metadata.get("session_id", h.ref), "question": h.metadata.get("question", h.content), "sql": h.metadata.get("sql", ""), "tables": h.metadata.get("tables", []), "score": round(h.similarity, 4), } for h in hits ] except degrade as e: warnings.append(f"solved-search fallita ({e})") for w in warnings: typer.secho(f"ATTENZIONE: {w}", fg=typer.colors.YELLOW, err=True) md_lines = ["# Retrieval pack", "", f"Domanda: {question}", ""] md_lines += [f"## Tabelle candidate (top {len(tables)}, vettoriale sull'intera domanda)", ""] if tables: for i, t in enumerate(tables, 1): desc = f" — {t['description']}" if t["description"] else "" md_lines.append(f"{i}. **{t['name']}**{desc} (rrf {t['rrf']})") else: md_lines.append("_nessuna (retrieval non disponibile o nessun match)_") md_lines += ["", "## Evidence rilevanti", ""] if evidence: for e in evidence: status = f" [{e['status']}]" if e["status"] else "" md_lines.append(f"- **{e['title']}**{status}: {e['excerpt']}") else: md_lines.append("_nessuna_") md_lines += ["", "## Domande risolte simili (exemplar di riferimento, NON decisioni)", ""] if solved: for s in solved: md_lines.append( f"### {s['question']} \n(sessione `{s['session_id']}`; " f"tabelle: {', '.join(s['tables']) or '-'})" ) if s["sql"]: md_lines += ["", "```sql", s["sql"], "```", ""] else: md_lines.append("_nessuna_") if warnings: md_lines += ["", "## Avvertenze", ""] + [f"- {w}" for w in warnings] md = "\n".join(md_lines) + "\n" if session: from tht.cli.session_cmd import load_session_or_exit, session_repository from tht.evidence import EvidenceReceipt, replace_evidence_receipt load_session_or_exit(cfg, session) repository = session_repository(cfg) repository.write_artifact(session, "retrieval_pack", md) if evidence_outcome is not None and evidence_outcome.status == "available": replace_evidence_receipt(repository, session, EvidenceReceipt( stage="clarification", purpose="disambiguation", vector_generation=evidence_outcome.vector_generation or "", evidence_ids=tuple(result.evidence_id for result in evidence_outcome.results), )) if not json_out: typer.secho( "OK: retrieval pack scritto " f"({len(tables)} tabelle, {len(evidence)} evidence, {len(solved)} solved).", fg=typer.colors.GREEN, ) if json_out: typer.echo(json.dumps( {"question": question, "tables": tables, "evidence": evidence, "solved": solved, "warnings": warnings}, ensure_ascii=False, indent=2, )) elif not session: typer.echo(md)