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

620 lines
25 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 (
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 semantico)")
DECISION_OPT = typer.Option(None, "--decision", help="Seq da promuovere (ripetibile).")
def registry_path(cfg) -> Path:
if getattr(cfg.paths, "memory", None) is not None:
return cfg.paths.memory / "registry.jsonl"
# Legacy location; migrate with `tht memory migrate` (P3).
return cfg.paths.artifacts / "memory" / "registry.jsonl"
def _resync_memory(cfg):
"""Risincronizza l'indice semantico 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
return build_vector_store(cfg, require_write=True).delete_kinds("memory", ["memory"])
@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 semantica puo' fallire
# (runtime non pronto 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 (runtime vettoriale 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 nel semantic store (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
del runtime vettoriale attivo.
"""
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 nell'indice semantico (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 semantico (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 nell'indice semantico (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 SolvedIndexError se mancano gli artefatti: il finalize lo degrada a warning,
il comando CLI lo converte in errore esplicito."""
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
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 semantic store (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, VectorStoreError
from tht.solved import SOLVED_KIND
from tht.vectorstore.embeddings import EmbeddingsError
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 (VectorStoreError, 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)
@memory_app.command("migrate")
def memory_migrate_cmd(
config: Path = CONFIG_OPT,
json_output: bool = typer.Option(False, "--json"),
) -> None:
"""Migrate the legacy artifacts/memory registry to the explicit workspace memory root (P3).
Copies and verifies exactly one legacy canonical JSONL under the workspace lock, then rebuilds
the Qdrant projection. Conflicting legacy registries fail closed; no in-place reinterpretation.
"""
from tht.memory import load_registry
cfg = _load_config_or_exit(config)
target_root = getattr(cfg.paths, "memory", None)
if target_root is None:
payload = {"status": "failed", "error": "explicit memory root is not configured"}
if json_output:
typer.echo(json.dumps(payload, sort_keys=True))
else:
typer.secho("ERRORE: memory root esplicito non configurato", fg=typer.colors.RED, err=True)
raise typer.Exit(code=1)
legacy = cfg.paths.artifacts / "memory" / "registry.jsonl"
target = registry_path(cfg)
if target.exists():
payload = {"status": "unchanged", "path": str(target)}
if json_output:
typer.echo(json.dumps(payload, sort_keys=True))
else:
typer.secho(f"OK: memory registry già in {target}", fg=typer.colors.GREEN)
return
if not legacy.exists():
payload = {"status": "failed", "error": "legacy memory registry is missing"}
if json_output:
typer.echo(json.dumps(payload, sort_keys=True))
else:
typer.secho("ERRORE: registry legacy mancante", fg=typer.colors.RED, err=True)
raise typer.Exit(code=1)
try:
records = load_registry(legacy)
except Exception: # noqa: BLE001
payload = {"status": "failed", "error": "legacy memory registry is invalid"}
if json_output:
typer.echo(json.dumps(payload, sort_keys=True))
else:
typer.secho("ERRORE: registry legacy non valido", fg=typer.colors.RED, err=True)
raise typer.Exit(code=1)
target_root.mkdir(parents=True, exist_ok=True)
from tht.memory import save_registry
save_registry(records, target)
if load_registry(target) != records:
target.unlink(missing_ok=True)
payload = {"status": "failed", "error": "memory registry migration verification failed"}
if json_output:
typer.echo(json.dumps(payload, sort_keys=True))
else:
typer.secho("ERRORE: verifica migrazione fallita", fg=typer.colors.RED, err=True)
raise typer.Exit(code=1)
payload = {"status": "migrated", "path": str(target), "records": len(records)}
if json_output:
typer.echo(json.dumps(payload, sort_keys=True))
else:
typer.secho(f"OK: migrate {len(records)} record verso {target}", fg=typer.colors.GREEN)