Files
ThothII/harness/tht/cli/search_cmd.py

519 lines
20 KiB
Python

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
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/<id>/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)