feat(harness): port memory/search/evidence/db/decision cmd (Onda 3.2)
5 cmd foglia portati con rename + drift fix: - memory_cmd: portato col modello registry INTATTO (TODO marker per il drop registry decisione spec 5 — task separato, richiede L2 per validare il rewrite su vectordb) - search_cmd: creata search_app sub-app (era funzione standalone in ChironeWp3), registrata come 'tht search find' - evidence_cmd, db_cmd, decision_cmd: port verbatim Drift fix decision_cmd: DECISION_MIN_PHASE.get(type,1) -> load_workflow().decision_min_phase(type). Check grep-per-file: ~15 residui nsp/PSD_SSL_CA nei messaggi utente fixati (nsp <cmd> -> tht <cmd>, nsp.yaml -> workspace yaml, PSD_SSL_CA -> THT_SSL_CA). Suite: 165 passed. tht --help ora mostra 10 sottocomandi.
This commit is contained in:
@@ -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")
|
||||
|
||||
@@ -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
|
||||
@@ -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)
|
||||
@@ -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)
|
||||
@@ -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 <seq> (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)
|
||||
@@ -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)
|
||||
Reference in New Issue
Block a user