Files
ThothII/harness/tht/cli/memory_cmd.py
T

568 lines
22 KiB
Python

# 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._guards import (
has_vector_write_rest,
require_server_profile,
require_vector_write_allowed,
)
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_snapshot_or_exit
from tht.cli.vector_cmd import require_vector_cfg
memory_app = typer.Typer(help="Review memory (registro canonico + indice pgvector)")
DECISION_OPT = typer.Option(None, "--decision", help="Seq da promuovere (ripetibile).")
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.adapters.factory import build_vector_store
from tht.cli.vector_cmd import make_embedder, sync_canonical_records
from tht.memory import load_registry, memory_vector_records
records = memory_vector_records(load_registry(registry_path(cfg)))
return sync_canonical_records(
"memory",
records,
store=build_vector_store(cfg, require_write=True),
embedder=make_embedder(cfg.embeddings),
)
def clear_memory_index(cfg):
from tht.adapters.factory import build_vector_store
from tht.cli.vector_cmd import make_embedder, open_store, require_direct_vector_cfg
if cfg.vectors is not None and cfg.vectors.type == "qdrant":
return build_vector_store(cfg, require_write=True).delete_kinds("memory", ["memory"])
require_direct_vector_cfg(cfg)
legacy_store = open_store(cfg, "memory")
legacy_store.sync([], make_embedder(cfg.embeddings), kinds={"memory"})
return 0
@memory_app.command("promote")
def promote_cmd(
session: str = typer.Option(..., "--session"),
decision: list[int] = DECISION_OPT,
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_snapshot
cfg = _load_config_or_exit(config)
snapshot = load_snapshot_or_exit(cfg, session)
if preview:
from tht.memory import (
MAX_PROMOTION_CANDIDATES,
preview_promotions_snapshot,
reusable_promotions_snapshot,
)
cand = preview_promotions_snapshot(snapshot, registry_path(cfg))
extra = len(reusable_promotions_snapshot(snapshot, 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_snapshot(snapshot, 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 (D11).
Promuove la decisione nel registro locale (idempotente) e fa un singolo upsert
con dedup hash client-side -- niente full-resync. Il factory seleziona il writer
REST su workstation oppure il writer pgvector diretto sul profilo server.
"""
import json as _json
from tht.adapters.factory import build_vector_store
from tht.cli.vector_cmd import make_embedder
from tht.memory import load_registry, promote_snapshot, save_one_memory
cfg = _load_config_or_exit(config)
snapshot = load_snapshot_or_exit(cfg, session)
require_vector_write_allowed(cfg, "memory save-one")
store = build_vector_store(cfg, require_write=True)
# 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_snapshot(snapshot, seqs=[decision], registry_path=registry_path(cfg))
records = [r for r in load_registry(registry_path(cfg)) if r.session_id == snapshot.manifest.id]
embedder = make_embedder(cfg.embeddings)
count = save_one_memory(records, decision, store=store, 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.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)
clear_memory_index(cfg)
# 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.decisions import DecisionType
from tht.memory import MemoryNotFound, update_record
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.vector_cmd import make_embedder, open_searcher, require_vector_cfg
from tht.memory import REUSABLE_TYPES, decided_memory_ids, load_registry
cfg = _load_config_or_exit(config)
require_vector_cfg(cfg)
excluded: set[str] = set()
if session is not None:
excluded = decided_memory_ids(load_snapshot_or_exit(cfg, session).decisions)
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.type not in REUSABLE_TYPES:
continue # record legacy non-concept: non e' una memory riusabile
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.adapters.factory import build_vector_store
from tht.cli.sql_cmd import promoted_tables_for
from tht.cli.vector_cmd import make_embedder
from tht.solved import build_solved_snapshot, save_solved_question
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"
)
store = build_vector_store(cfg, require_write=True)
record = build_solved_snapshot(load_snapshot_or_exit(cfg, session_id), promoted_tables_for(cfg, session_id))
return save_solved_question(
record,
store=store,
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.ports.vector import VectorReadUnavailable
from tht.solved import SOLVED_KIND
from tht.vectorstore.embeddings import EmbeddingsError
from tht.vectorstore.rest_client import VectorRestError
cfg = _load_config_or_exit(config)
require_vector_cfg(cfg)
# Degrado gentile: SKILL.md prescrive solved-search in F4/F6/F7 di ogni sessione,
# quindi vectordb/Ollama irraggiungibili non devono produrre un traceback grezzo
# nel transcript: avviso di una riga su stderr, stdout puro ([] in --json), exit 0.
try:
searcher = open_searcher(cfg)
embedder = make_embedder(cfg.embeddings)
hits = searcher.search(embedder.embed_query(question), top_n=top, kinds=[SOLVED_KIND])
except (VectorRestError, VectorReadUnavailable, EmbeddingsError, OperationalError) as e:
typer.secho(
f"ATTENZIONE: exemplar non disponibili ({e}). Prosegui senza.",
fg=typer.colors.YELLOW, err=True,
)
if json_out:
typer.echo("[]")
return
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)