"""Thin command adapters for the authoritative Memory service.""" import json import re from contextlib import contextmanager from pathlib import Path import typer from pydantic import ValidationError 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.memory.models import CardInput, CardQuery, MemoryError from tht.memory.runtime import admin_service, memory_service from tht.ports.vector import VectorStoreError from tht.vectorstore.embeddings import EmbeddingsError memory_app = typer.Typer(help="Authoritative Memory cards and verified recall") DECISION_OPT = typer.Option(None, "--decision") def _output(value): typer.echo(json.dumps(value, ensure_ascii=False, default=str)) @contextmanager def _service(config): service = None try: cfg = _load_config_or_exit(config) service = memory_service(cfg) yield cfg, service except MemoryError as error: _output({"code": error.code, "message": str(error), "status": error.status}) raise typer.Exit(1) from None except (ValidationError, ValueError): _output({"code": "memory_invalid", "message": "Memory request is invalid", "status": 400}) raise typer.Exit(1) from None except (VectorStoreError, EmbeddingsError, OSError): _output({"code": "memory_unavailable", "message": "Memory operation is unavailable", "status": 503}) raise typer.Exit(1) from None finally: if service: service.close() @memory_app.command("admin") def admin_cmd(workspace: str = typer.Option(...), config: Path = CONFIG_OPT): """Backend-owned request snapshot. No DWH or session configuration is required.""" service = None try: if not re.fullmatch(r"[a-z][a-z0-9_-]{0,63}", workspace): raise ValueError("Invalid workspace") payload = json.loads(config.read_text()) service = admin_service(workspace, payload["runtime"]) request = payload.get("request", {}) action = payload["action"] if action == "cleanup": from tht.memory.cleanup import CleanupRequest, cleanup result = cleanup(service, CleanupRequest.model_validate(request)) elif action == "list": result = service.list(CardQuery.model_validate(request)) elif action == "show": result = service.get(request["id"]) elif action in {"create", "update"}: result = service.save(CardInput.model_validate(request["card"]), request.get("id") if action == "update" else None) elif action == "delete": result = service.delete(request["id"]) elif action == "pending": result = service.pending() elif action == "retry": result = service.retry(request["id"]) else: raise ValueError("Unknown Memory action") _output(result) except MemoryError as error: _output({"code": error.code, "message": str(error), "status": error.status}) raise typer.Exit(1) from None except (ValueError, KeyError, TypeError, OSError): _output({"code": "memory_invalid", "message": "Memory request is invalid", "status": 400}) raise typer.Exit(1) from None finally: if service: service.close() @memory_app.command("list") def list_cmd(filters: str = typer.Option("{}", "--filters"), json_out: bool = typer.Option(False, "--json"), config: Path = CONFIG_OPT): with _service(config) as (_, service): _output(service.list(CardQuery.model_validate_json(filters))) @memory_app.command("show") def show_cmd(mem_id: str, json_out: bool = typer.Option(False, "--json"), config: Path = CONFIG_OPT): with _service(config) as (_, service): _output(service.get(mem_id)) @memory_app.command("create") def create_cmd(data: Path = typer.Option(..., "--data"), json_out: bool = typer.Option(False, "--json"), config: Path = CONFIG_OPT): with _service(config) as (_, service): _output(service.save(CardInput.model_validate_json(data.read_text()))) @memory_app.command("update") def update_cmd(mem_id: str, data: Path = typer.Option(..., "--data"), json_out: bool = typer.Option(False, "--json"), config: Path = CONFIG_OPT): with _service(config) as (_, service): _output(service.save(CardInput.model_validate_json(data.read_text()), mem_id)) @memory_app.command("delete") def delete_cmd(mem_id: str, yes: bool = typer.Option(False, "--yes", "-y"), json_out: bool = typer.Option(False, "--json"), config: Path = CONFIG_OPT): with _service(config) as (_, service): service._admin() if not yes and not typer.confirm(f"Delete Memory card {mem_id} and its links?"): raise typer.Exit(1) _output(service.delete(mem_id)) @memory_app.command("pending") def pending_cmd(json_out: bool = typer.Option(False, "--json"), config: Path = CONFIG_OPT): with _service(config) as (_, service): _output(service.pending()) @memory_app.command("retry") def retry_cmd(mem_id: str, json_out: bool = typer.Option(False, "--json"), config: Path = CONFIG_OPT): with _service(config) as (_, service): _output(service.retry(mem_id)) @memory_app.command("index") def index_cmd(json_out: bool = typer.Option(False, "--json"), config: Path = CONFIG_OPT): with _service(config) as (_, service): _output(service.rebuild()) @memory_app.command("promote") def promote_cmd(session: str = typer.Option(..., "--session"), decision: list[int] = DECISION_OPT, preview: bool = typer.Option(False, "--preview"), json_out: bool = typer.Option(False, "--json"), config: Path = CONFIG_OPT): with _service(config) as (cfg, service): snapshot = load_snapshot_or_exit(cfg, session) if preview: _output(service.promotions(snapshot)[:5]) elif not decision: raise ValueError("Explicit decisions are required") else: results = service.promote(snapshot, decision) _output({"promoted": [r.get("card") for r in results], "indexed": all(r["indexed"] for r in results), "results": results}) @memory_app.command("propose") def propose_cmd(session: str = typer.Option(..., "--session"), data: Path = typer.Option(..., "--data"), config: Path = CONFIG_OPT): """Persist reviewer-grounded proposals without changing the Memory archive.""" from tht.cli.session_cmd import session_repository from tht.memory.review import validate_proposals with _service(config) as (cfg, service): snapshot = load_snapshot_or_exit(cfg, session) service._session(snapshot) proposals = validate_proposals(snapshot, json.loads(data.read_text())) session_repository(cfg).write_artifact(session, "memory_proposals", json.dumps([p.model_dump(mode="json") for p in proposals], ensure_ascii=False)) _output({"proposals": len(proposals), "saved_to_archive": False}) @memory_app.command("summary") def summary_cmd(session: str = typer.Option(..., "--session"), json_out: bool = typer.Option(False, "--json"), config: Path = CONFIG_OPT): from tht.memory.review import prepare with _service(config) as (cfg, service): _output(prepare(service, load_snapshot_or_exit(cfg, session))) @memory_app.command("repair-prepare") def repair_prepare_cmd(session: str = typer.Option(..., "--session"), proposal_json: str = typer.Option(..., "--proposal-json"), json_out: bool = typer.Option(False, "--json"), config: Path = CONFIG_OPT): from tht.archive_repair import RepairProposal, prepare if len(proposal_json.encode()) > 1_000_000: _output({"code": "memory_invalid", "message": "Repair proposal is too large", "status": 400}) raise typer.Exit(1) with _service(config) as (cfg, service): _output(prepare(service, load_snapshot_or_exit(cfg, session), cfg, RepairProposal.model_validate_json(proposal_json))) @memory_app.command("repair-target") def repair_target_cmd(session: str = typer.Option(..., "--session"), archive: str = typer.Option(..., "--archive"), target_id: str = typer.Option(..., "--target-id"), json_out: bool = typer.Option(False, "--json"), config: Path = CONFIG_OPT): from tht.archive_repair import target with _service(config) as (cfg, service): _output(target(service, load_snapshot_or_exit(cfg, session), cfg, archive, target_id)) @memory_app.command("repairs") def repairs_cmd(session: str = typer.Option(..., "--session"), json_out: bool = typer.Option(False, "--json"), config: Path = CONFIG_OPT): from tht.archive_repair import list_repairs with _service(config) as (cfg, service): _output(list_repairs(service, load_snapshot_or_exit(cfg, session))) @memory_app.command("repair-show") def repair_show_cmd(session: str = typer.Option(..., "--session"), repair_id: str = typer.Option(..., "--repair-id"), json_out: bool = typer.Option(False, "--json"), config: Path = CONFIG_OPT): from tht.archive_repair import show with _service(config) as (cfg, service): _output(show(service, load_snapshot_or_exit(cfg, session), cfg, repair_id)) @memory_app.command("repair-apply") def repair_apply_cmd(session: str = typer.Option(..., "--session"), repair_id: str = typer.Option(..., "--repair-id"), choice: str = typer.Option(..., "--choice"), json_out: bool = typer.Option(False, "--json"), config: Path = CONFIG_OPT): from tht.archive_repair import apply from tht.cli.preprocess_cmd import run_from_config def activate(snapshot): result = run_from_config(config, local_snapshot=snapshot) if result.status != "succeeded": raise RuntimeError("Evidence activation did not complete") with _service(config) as (cfg, service): _output(apply(service, load_snapshot_or_exit(cfg, session), cfg, repair_id, choice, activate=activate)) @memory_app.command("review-apply") def review_apply_cmd(session: str = typer.Option(..., "--session"), review_json: str = typer.Option(..., "--review-json"), json_out: bool = typer.Option(False, "--json"), config: Path = CONFIG_OPT): from tht.memory.review import ReviewResponse, apply with _service(config) as (cfg, service): _output(apply(service, load_snapshot_or_exit(cfg, session), ReviewResponse.model_validate_json(review_json))) @memory_app.command("save-one") def save_one_cmd(session: str = typer.Option(..., "--session"), decision: int = typer.Option(..., "--decision"), json_out: bool = typer.Option(False, "--json"), config: Path = CONFIG_OPT): with _service(config) as (cfg, service): results = service.promote(load_snapshot_or_exit(cfg, session), [decision]) _output({"upserted": sum(r["indexed"] for r in results), "decision_seq": decision, "indexed": all(r["indexed"] for r in results), "results": results}) @memory_app.command("search") def search_cmd(question: str, top: int = typer.Option(5, "--top"), session: str | None = typer.Option(None, "--session"), filters: str = typer.Option("{}", "--filters"), json_out: bool = typer.Option(False, "--json"), config: Path = CONFIG_OPT): from tht.cli.vector_cmd import make_embedder, open_searcher with _service(config) as (cfg, service): decisions = load_snapshot_or_exit(cfg, session).decisions if session else [] _output(service.recall(question, searcher=open_searcher(cfg), embedder=make_embedder(cfg.embeddings), top=top, decisions=decisions, scope=_recall_scope(cfg, filters))) def _recall_scope(cfg, filters): from tht.memory.retrieval import RecallScope context = {"database": cfg.database.database, "schema_name": cfg.database.db_schema} supplied = json.loads(filters) if not isinstance(supplied, dict) or any( key in supplied and supplied[key] != value for key, value in context.items() ): raise ValueError("Recall cannot override the configured database/schema context") return RecallScope.model_validate({**supplied, **context}) @memory_app.command("rules") def rules_cmd(question: str, session: str = typer.Option(..., "--session"), filters: str = typer.Option("{}", "--filters"), json_out: bool = typer.Option(False, "--json"), config: Path = CONFIG_OPT): """Consult SQL rules and explained errors in schema linking and SQL construction.""" from functools import cache from types import SimpleNamespace from tht.cli.vector_cmd import make_embedder, open_searcher from tht.phase import current_phase with _service(config) as (cfg, service): snapshot = load_snapshot_or_exit(cfg, session) service._session(snapshot) if current_phase(snapshot) not in {4, 6, 7}: raise ValueError("Memory rules are consulted in schema linking or SQL construction") scope = _recall_scope(cfg, filters) # Share one vector across both families, but only when the archive has cards. embedder = SimpleNamespace(embed_query=cache(make_embedder(cfg.embeddings).embed_query)) searcher = open_searcher(cfg) candidates = [] for family in ("sql_rule", "explained_error"): candidates.extend(service.retrieve(question, searcher=searcher, embedder=embedder, scope=scope, family=family, top=5)) _output([{**candidate.card.model_dump(mode="json"), "score": candidate.score, "retrieval_path": candidate.path, "consultative": True} for candidate in sorted(candidates, key=lambda c: (-c.score, c.card.id))[:10]]) def index_solved_session(cfg, session_id: str) -> int: """Recovery only: never recreate a deleted card from historical session artifacts.""" service = memory_service(cfg) try: result = service.retry_solved(load_snapshot_or_exit(cfg, session_id)) if not result["indexed"]: raise RuntimeError(result["error"]) return int(result["action"] == "upsert") finally: service.close() @memory_app.command("solved-index") def solved_index_cmd(session_id: str, json_out: bool = typer.Option(False, "--json"), config: Path = CONFIG_OPT): """Retry an existing authoritative exemplar's Qdrant projection; never import a session.""" with _service(config) as (cfg, service): _output(service.retry_solved(load_snapshot_or_exit(cfg, session_id))) @memory_app.command("solved-search") def solved_search_cmd(question: str, top: int = typer.Option(3, "--top"), filters: str = typer.Option("{}", "--filters"), json_out: bool = typer.Option(False, "--json"), config: Path = CONFIG_OPT): from tht.cli.vector_cmd import make_embedder, open_searcher with _service(config) as (cfg, service): try: result = service.recall(question, searcher=open_searcher(cfg), embedder=make_embedder(cfg.embeddings), top=top, solved=True, scope=_recall_scope(cfg, filters)) except (VectorStoreError, EmbeddingsError): typer.echo("Avviso: exemplar non disponibili; ricerca semantica non riuscita.", err=True) result = [] _output(result)