Add catalog-owned logical relationships and runtime snapshots, extend the database-management UI and validation coverage, and document the updated operational workflow. Keep active sensitive-generation status in a tooltip and indicator, and update the layout E2E to follow the history action in its new database-scoped location.
526 lines
20 KiB
Python
526 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.mschema.context import SchemaContextError, load_schema_context
|
|
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)
|
|
try:
|
|
schema_context = load_schema_context(cfg, physical_file=phys_file)
|
|
except SchemaContextError as exc:
|
|
typer.secho(f"ERRORE: {exc}", fg=typer.colors.RED, err=True)
|
|
raise typer.Exit(code=1) from None
|
|
|
|
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
|
|
|
|
selected = [t for t, _ in ranked]
|
|
mschema = to_mschema_text(
|
|
schema_context.physical,
|
|
schema_context.annotations,
|
|
tables=selected,
|
|
effective_relationships=schema_context.effective_relationships,
|
|
)
|
|
|
|
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)
|