# 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 ( has_vector_write_rest, 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("save-one") def save_one_cmd( session: str = typer.Option(..., "--session"), decision: int = typer.Option( ..., "--decision", help="decision_seq della decisione da salvare come memoria." ), json_out: bool = typer.Option(False, "--json", help="Output JSON (per Pi)."), config: Path = CONFIG_OPT, ) -> None: """Upsert mirato (una riga) della memoria di una decisione su pgvector via writer key (D11). Promuove la decisione nel registro locale (idempotente) e fa un singolo upsert remoto con dedup hash client-side -- niente full-resync. Abilita il salvataggio di una memoria da postazione remota (workstation) con la sola writer key. """ import json as _json from tht.cli.vector_cmd import make_embedder from tht.memory import load_registry, promote, save_one_memory from tht.vectorstore.rest_client import VectorRestClient cfg = _load_config_or_exit(config) manifest = load_session_or_exit(cfg, session) require_vector_write_allowed(cfg, "memory save-one") if not has_vector_write_rest(cfg): typer.secho( "ERRORE: `memory save-one` richiede la sezione `vector_write_rest` con una " "API key di upsert nel workspace yaml (upsert remoto via writer key).", fg=typer.colors.RED, err=True, ) raise typer.Exit(code=4) sdir = session_dir(cfg, session) # Promuove la decisione scelta nel registro locale (idempotente: salta se gia' presente # o se stale post-rollback, perche' _compute_promotions usa la vista effective). promote(sdir, manifest, seqs=[decision], registry_path=registry_path(cfg)) records = [r for r in load_registry(registry_path(cfg)) if r.session_id == manifest.id] writer = VectorRestClient(cfg.vector_write_rest) embedder = make_embedder(cfg.embeddings) count = save_one_memory(records, decision, writer=writer, embedder=embedder) msg = ( f"{count} memoria salvata su pgvector (decision_seq {decision})." if count else f"Nessun upsert (decisione {decision} assente/stale o memoria gia' aggiornata)." ) if json_out: typer.echo(_json.dumps( {"upserted": count, "decision_seq": decision, "message": msg}, ensure_ascii=False)) return typer.secho(f"OK: {msg}", fg=typer.colors.GREEN if count else 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) def index_solved_session(cfg, session_id: str) -> int: """Indicizza la coppia domanda->SQL della sessione (kind solved_question). Solleva RuntimeError se manca la writer key e SolvedIndexError se mancano gli artefatti: il finalize li degrada a warning, il comando CLI li converte in errori espliciti.""" from tht.cli.sql_cmd import promoted_tables_for from tht.cli.vector_cmd import make_embedder from tht.solved import build_solved_record, save_solved_question from tht.vectorstore.rest_client import VectorRestClient if not has_vector_write_rest(cfg): raise RuntimeError( "vector_write_rest assente: la coppia domanda->SQL si indicizza solo con la " "writer key configurata nel workspace yaml" ) manifest = load_session_or_exit(cfg, session_id) record = build_solved_record( session_dir(cfg, session_id), manifest, promoted_tables_for(cfg, session_id) ) return save_solved_question( record, writer=VectorRestClient(cfg.vector_write_rest), embedder=make_embedder(cfg.embeddings), ) @memory_app.command("solved-index") def solved_index_cmd( session_id: str = typer.Argument(..., help="Id sessione con sql_final.sql approvato."), json_out: bool = typer.Option(False, "--json", help="Output JSON (per Pi)."), config: Path = CONFIG_OPT, ) -> None: """Indicizza la coppia domanda->SQL nel vectordb (backfill; il finalize lo fa da solo).""" import json as _json from tht.solved import SolvedIndexError cfg = _load_config_or_exit(config) require_vector_write_allowed(cfg, "memory solved-index") try: count = index_solved_session(cfg, session_id) except RuntimeError as e: typer.secho(f"ERRORE: {e}", fg=typer.colors.RED, err=True) raise typer.Exit(code=4) except SolvedIndexError as e: typer.secho(f"ERRORE: sessione {session_id} non indicizzabile: {e}", fg=typer.colors.RED, err=True) raise typer.Exit(code=3) msg = ( f"1 coppia domanda->SQL indicizzata (solved:{session_id})." if count else "Nessun upsert: coppia gia' aggiornata." ) if json_out: typer.echo(_json.dumps({"upserted": count, "id": f"solved:{session_id}"}, ensure_ascii=False)) return typer.secho(f"OK: {msg}", fg=typer.colors.GREEN) @memory_app.command("solved-search") def solved_search_cmd( question: str = typer.Argument(..., help="Domanda da confrontare con quelle risolte."), top: int = typer.Option(3, "--top"), json_out: bool = typer.Option(False, "--json", help="Output JSON (per Pi)."), config: Path = CONFIG_OPT, ) -> None: """Domande gia' risolte simili (kind solved_question): domanda, SQL e tabelle.""" from rich.console import Console from rich.table import Table from tht.cli.vector_cmd import make_embedder, open_searcher from tht.solved import SOLVED_KIND cfg = _load_config_or_exit(config) require_vector_cfg(cfg) searcher = open_searcher(cfg) embedder = make_embedder(cfg.embeddings) hits = searcher.search(embedder.embed_query(question), top_n=top, kinds=[SOLVED_KIND]) results = [ { "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 ] if json_out: typer.echo(json.dumps(results, ensure_ascii=False, indent=2)) return if not results: typer.secho("Nessuna domanda risolta simile.", fg=typer.colors.YELLOW) return table = Table(title=f"Domande risolte simili a: {question}") table.add_column("Sessione") table.add_column("Domanda") table.add_column("Tabelle") table.add_column("Score", justify="right") for r in results: table.add_row(r["session_id"], r["question"][:60], ", ".join(r["tables"]), f"{r['score']:.3f}") Console().print(table)