diff --git a/PROJECT_STATE.md b/PROJECT_STATE.md index e369408e..71c54c04 100644 --- a/PROJECT_STATE.md +++ b/PROJECT_STATE.md @@ -1015,7 +1015,7 @@ search` (Phase 2) reuse: (ledger detail `seq:`) so it is never re-proposed. `tht memory promote`/`save-one` were added to the gate's anti-bypass FORBIDDEN list (model must go through the gate tool). - **Solved-question exemplars.** New vector kind `solved_question` reusing the existing - `memory` pgvector table (no server-side DDL); `harness/tht/solved.py` does a one-row + `memory` pgvector table (no server-side DDL); `harness/tht/memory/solved.py` does a one-row upsert keyed by a hash of question+SQL. CLI: `tht memory solved-index` / `solved-search`. `tht session finalize` auto-indexes the pair (best-effort: green line on upsert, cyan "già aggiornata" on dedup no-op, yellow warning + the recovery command diff --git a/docs/gestione-memory.md b/docs/gestione-memory.md index e2e07d87..68f8533f 100644 --- a/docs/gestione-memory.md +++ b/docs/gestione-memory.md @@ -25,7 +25,10 @@ F8: il reviewer decide se promuoverla L'invariante principale è `REUSABLE_TYPES = {"concept_clarified"}`: le sole memory generabili, salvabili, ricercabili e proponibili sono i concetti chiariti. Le decisioni `table_promoted`, `table_excluded`, `column_promoted` e analoghe restano decisioni locali alla domanda. -Implementazione principale: [harness/tht/memory.py](../harness/tht/memory.py:14) e [harness/tht/cli/memory_cmd.py](../harness/tht/cli/memory_cmd.py:368). +Implementazione principale: la façade [harness/tht/memory/](../harness/tht/memory/), +con le policy riusabili in +[harness/tht/memory/core.py](../harness/tht/memory/core.py), e l'adapter +[harness/tht/cli/memory_cmd.py](../harness/tht/cli/memory_cmd.py). ## I tre livelli della gestione @@ -90,7 +93,10 @@ La preview: 5. deduplica contenuti equivalenti; 6. propone al massimo cinque candidati. -Il codice applica il filtro e la deduplica in [harness/tht/memory.py](../harness/tht/memory.py:199); il gate applica un ulteriore filtro difensivo in [harness/.pi/extensions/tht-gate.js](../harness/.pi/extensions/tht-gate.js:628). +Il codice applica il filtro e la deduplica in +[harness/tht/memory/core.py](../harness/tht/memory/core.py); il gate applica un +ulteriore filtro difensivo in +[harness/.pi/extensions/gate/memory/index.js](../harness/.pi/extensions/gate/memory/index.js). Il reviewer vede un'unica checklist, preselezionata. Per ogni candidato: @@ -126,7 +132,8 @@ Dopo la promozione, `save-one` costruisce un solo `VectorRecord` e lo invia all' Il record vettoriale usa l'id `memory:mem-XXXX`, mentre i metadati conservano `subject`, `detail`, `rationale`, `tables`, `concepts` e il discriminante `kind`. L'hash SHA-256 del contenuto impedisce di ricalcolare embedding e upsert quando il testo non è cambiato. -Il comportamento è implementato in [harness/tht/memory.py](../harness/tht/memory.py:253) e [harness/tht/memory.py](../harness/tht/memory.py:305). +Il comportamento è implementato in +[harness/tht/memory/core.py](../harness/tht/memory/core.py). ### Fonte canonica attuale diff --git a/harness/tests/memory/test_solved_finalization.py b/harness/tests/memory/test_solved_finalization.py new file mode 100644 index 00000000..038edff8 --- /dev/null +++ b/harness/tests/memory/test_solved_finalization.py @@ -0,0 +1,87 @@ +from datetime import UTC, datetime +from pathlib import Path +from types import SimpleNamespace +import uuid + +from tht.cli import session_cmd +from tht.decisions import DecisionInput +from tht.session.filesystem_repository import FilesystemSessionRepository +from tht.session.models import PrincipalContext, SessionManifest + + +def test_finalize_commits_session_before_best_effort_post_commit_read_failure( + tmp_path, monkeypatch, capsys +): + repository = FilesystemSessionRepository( + tmp_path / "home", + "demo", + PrincipalContext(issuer="local", subject="reviewer"), + root=tmp_path / "sessions", + ) + session_id = str(uuid.uuid4()) + repository.create(SessionManifest( + id=session_id, + created_at=datetime(2026, 8, 24, tzinfo=UTC), + question="active patients", + database="analytics", + schema="mart", + )) + repository.write_artifact(session_id, "sql_final", "SELECT 1\n") + repository.write_artifact( + session_id, + "schema_linking", + '{"question":"active patients","candidates":[],"excluded":[]}', + ) + repository.append_decisions(session_id, [ + DecisionInput(type="sql_approved", subject="phase:7"), + *[ + DecisionInput(type="phase_approved", subject=f"phase:{phase}") + for phase in range(1, 9) + ], + ]) + cfg = SimpleNamespace( + execution=SimpleNamespace(forbidden_functions=[], max_preview_rows=10), + paths=SimpleNamespace(artifacts=tmp_path / "artifacts"), + embeddings=object(), + ) + + class FailingPostCommitRead: + def get(self, sid): + snapshot = repository.get(sid) + if snapshot.manifest.status == "finalized": + raise RuntimeError("finalized snapshot unavailable") + return snapshot + + def finalize(self, manifest, artifacts): + return repository.finalize(manifest, artifacts) + + failing_read_repository = FailingPostCommitRead() + + monkeypatch.setattr(session_cmd, "_load_config_or_exit", lambda config: cfg) + monkeypatch.setattr( + session_cmd, + "session_repository", + lambda config: failing_read_repository, + ) + monkeypatch.setattr(session_cmd, "load_snapshot_or_exit", lambda config, sid: repository.get(sid)) + monkeypatch.setattr(session_cmd, "session_problems", lambda config, sid: []) + monkeypatch.setattr("tht.cli.sql_cmd._load_physical_or_exit", lambda config: object()) + monkeypatch.setattr("tht.cli.sql_cmd.promoted_tables_for", lambda config, sid: set()) + monkeypatch.setattr("tht.cli.sql_cmd.do_explain", lambda config, sql: object()) + monkeypatch.setattr("tht.cli.sql_cmd.do_run", lambda config, sql, limit: object()) + monkeypatch.setattr( + "tht.sqlcheck.validate_sql", + lambda *args, **kwargs: SimpleNamespace(ok=True, errors=[], ast=object()), + ) + monkeypatch.setattr("tht.execute.warnings.plan_warnings", lambda *args: []) + monkeypatch.setattr("tht.execute.warnings.runtime_warnings", lambda *args: []) + monkeypatch.setattr("tht.execute.warnings.static_warnings", lambda *args: []) + monkeypatch.setattr("tht.report.render_validation_report", lambda **kwargs: "verified\n") + monkeypatch.setattr("tht.session.artifacts.build_evidence_entries", lambda *args: []) + + session_cmd.finalize_cmd(session_id, config=Path("unused.yaml")) + + captured = capsys.readouterr() + assert repository.get(session_id).manifest.status == "finalized" + assert "finalized snapshot unavailable" in captured.err + assert f"OK: sessione {session_id} finalizzata" in captured.out diff --git a/harness/tests/memory/test_solved_lifecycle.py b/harness/tests/memory/test_solved_lifecycle.py new file mode 100644 index 00000000..a7982ea4 --- /dev/null +++ b/harness/tests/memory/test_solved_lifecycle.py @@ -0,0 +1,233 @@ +from datetime import UTC, datetime +from types import SimpleNamespace + +import pytest + +from tht.decisions import DecisionRecord +from tht.memory import ( + SolvedIndexError, + index_solved_question, + index_solved_question_best_effort, + search_solved_questions, +) +from tht.session.models import SessionManifest, SessionSnapshot + + +def _decision(seq: int, type_: str, subject: str, detail: str = "") -> DecisionRecord: + return DecisionRecord( + seq=seq, + ts=datetime(2026, 8, 24, tzinfo=UTC), + type=type_, + subject=subject, + detail=detail, + ) + + +def _snapshot(*, sql: str | None = "SELECT 1\n", approved: bool = True) -> SessionSnapshot: + decisions = [_decision(1, "question_rewritten", "question", "rewritten question")] + if approved: + decisions.append(_decision(2, "sql_approved", "phase:7")) + decisions.extend( + _decision(index + 2, "phase_approved", f"phase:{index}") + for index in range(1, 9) + ) + return SessionSnapshot( + manifest=SessionManifest( + id="s1", + created_at=datetime(2026, 8, 24, tzinfo=UTC), + status="finalized", + question="original question", + database="analytics", + schema="mart", + ), + artifacts={} if sql is None else {"sql_final": sql}, + decisions=decisions, + ) + + +class Store: + def __init__(self): + self.hashes = {} + self.upserts = [] + + def existing_hashes(self, collection, kinds): + assert (collection, kinds) == ("memory", ["solved_question"]) + return self.hashes + + def upsert(self, collection, records): + self.upserts.append((collection, records)) + return len(records) + + +class Embedder: + def __init__(self): + self.documents = [] + self.queries = [] + + def embed_documents(self, documents): + self.documents.append(documents) + return [[0.1, 0.2]] + + def embed_query(self, question): + self.queries.append(question) + return [0.3, 0.4] + + +def test_memory_facade_indexes_finalized_question_with_compatible_record_and_dedup(): + store = Store() + embedder = Embedder() + snapshot = _snapshot() + + assert index_solved_question( + snapshot, + {"fact_z", "dim_a"}, + store=store, + embedder=embedder, + ) == 1 + + collection, rows = store.upserts[0] + assert collection == "memory" + assert len(rows) == 1 + row = rows[0] + assert row.record.model_dump() == { + "id": "solved:s1", + "kind": "solved_question", + "ref": "s1", + "title": "rewritten question", + "content": "rewritten question", + "metadata": { + "question": "rewritten question", + "sql": "SELECT 1", + "tables": ["dim_a", "fact_z"], + "session_id": "s1", + }, + } + assert embedder.documents == [["rewritten question"]] + + store.hashes = {row.record.id: row.content_hash} + assert index_solved_question( + snapshot, + {"fact_z", "dim_a"}, + store=store, + embedder=embedder, + ) == 0 + assert len(store.upserts) == 1 + assert embedder.documents == [["rewritten question"]] + + changed = snapshot.model_copy( + update={"artifacts": {"sql_final": "SELECT 2\n"}}, + ) + assert index_solved_question( + changed, + {"fact_z", "dim_a"}, + store=store, + embedder=embedder, + ) == 1 + assert store.upserts[-1][1][0].record.metadata["sql"] == "SELECT 2" + assert embedder.documents == [["rewritten question"], ["rewritten question"]] + + +def test_memory_facade_falls_back_to_manifest_question(): + snapshot = _snapshot().model_copy( + update={"decisions": [ + decision + for decision in _snapshot().decisions + if decision.type != "question_rewritten" + ]}, + ) + store = Store() + + assert index_solved_question( + snapshot, + None, + store=store, + embedder=Embedder(), + ) == 1 + assert store.upserts[0][1][0].record.content == "original question" + assert store.upserts[0][1][0].record.metadata["tables"] == [] + + +@pytest.mark.parametrize( + ("snapshot", "message"), + [ + (_snapshot(sql=None), "sql_final.sql assente"), + (_snapshot(approved=False), "decisione sql_approved assente"), + ], +) +def test_memory_facade_rejects_incomplete_solved_question(snapshot, message): + with pytest.raises(SolvedIndexError, match=message): + index_solved_question(snapshot, set(), store=Store(), embedder=Embedder()) + + +def test_memory_facade_keeps_solved_indexing_best_effort_after_finalization(): + snapshot = _snapshot() + + outcome = index_solved_question_best_effort( + snapshot, + set(), + store_factory=lambda: (_ for _ in ()).throw(RuntimeError("vector unavailable")), + embedder_factory=Embedder, + ) + + assert outcome.upserted is None + assert outcome.error == "vector unavailable" + assert snapshot.manifest.status == "finalized" + + +def test_memory_facade_searches_solved_questions_in_rank_order_with_compatible_payload(): + hits = [ + SimpleNamespace( + ref="s2", + content="second fallback question", + metadata={ + "session_id": "s2", + "question": "second question", + "sql": "SELECT 2", + "tables": ["fact_two"], + }, + similarity=0.93456, + ), + SimpleNamespace( + ref="s1", + content="first fallback question", + metadata={}, + similarity=0.81234, + ), + ] + + class Searcher: + def __init__(self): + self.calls = [] + + def search(self, embedding, *, top_n, kinds): + self.calls.append((embedding, top_n, kinds)) + return hits + + searcher = Searcher() + embedder = Embedder() + + results = search_solved_questions( + "similar question", + searcher=searcher, + embedder=embedder, + top=3, + ) + + assert results == [ + { + "session_id": "s2", + "question": "second question", + "sql": "SELECT 2", + "tables": ["fact_two"], + "score": 0.9346, + }, + { + "session_id": "s1", + "question": "first fallback question", + "sql": "", + "tables": [], + "score": 0.8123, + }, + ] + assert embedder.queries == ["similar question"] + assert searcher.calls == [([0.3, 0.4], 3, ["solved_question"])] diff --git a/harness/tests/test_adapter_command_regressions.py b/harness/tests/test_adapter_command_regressions.py index b8892877..c47eadfb 100644 --- a/harness/tests/test_adapter_command_regressions.py +++ b/harness/tests/test_adapter_command_regressions.py @@ -93,20 +93,19 @@ def test_solved_index_writes_through_writer_only_factory_store(monkeypatch): ) cfg = SimpleNamespace(embeddings=object(), vector_write_rest=object()) manifest = SimpleNamespace(id="s1") - solved_record = object() + snapshot = SimpleNamespace(manifest=manifest, decisions=[], artifacts={}) calls = [] - monkeypatch.setattr(memory_cmd, "load_snapshot_or_exit", lambda cfg, session: SimpleNamespace(manifest=manifest, decisions=[], artifacts={})) + monkeypatch.setattr(memory_cmd, "load_snapshot_or_exit", lambda cfg, session: snapshot) monkeypatch.setattr( "tht.adapters.factory.build_vector_store", lambda cfg, require_write: calls.append(require_write) or writer_only_store, ) monkeypatch.setattr("tht.cli.sql_cmd.promoted_tables_for", lambda *args: []) - monkeypatch.setattr("tht.solved.build_solved_snapshot", lambda *args: solved_record) monkeypatch.setattr( - "tht.solved.save_solved_question", - lambda record, *, store, embedder: int( - record is solved_record and store is writer_only_store + "tht.memory.index_solved_question", + lambda loaded, tables, *, store, embedder: int( + loaded is snapshot and tables == [] and store is writer_only_store ), ) monkeypatch.setattr("tht.cli.vector_cmd.make_embedder", lambda cfg: object()) diff --git a/harness/tests/test_solved_build.py b/harness/tests/test_solved_build.py deleted file mode 100644 index 2925f297..00000000 --- a/harness/tests/test_solved_build.py +++ /dev/null @@ -1,55 +0,0 @@ -"""L1: build del record solved_question dagli artefatti persistiti della sessione. - -Il record si costruisce SOLO da cio' che il workflow ha approvato: sql_final.sql -presente + decisione sql_approved nella vista effective; la domanda e' l'ultima -question_rewritten (fallback: la domanda del manifest).""" -from datetime import datetime - -import pytest - -from tht.decisions import append_decision -from tht.session.models import SessionManifest -from tht.solved import SolvedIndexError, build_solved_record - - -def _manifest() -> SessionManifest: - return SessionManifest( - id="s1", created_at=datetime(2026, 1, 1), question="domanda originale", - database="db", schema="public", - ) - - -def test_build_uses_rewritten_question_sql_and_tables(tmp_path): - (tmp_path / "sql_final.sql").write_text("SELECT 1\n") - append_decision(tmp_path, type="question_rewritten", subject="domanda", - detail="domanda riscritta esplicita") - append_decision(tmp_path, type="sql_approved", subject="phase:7") - for n in range(1, 8): - append_decision(tmp_path, type="phase_approved", subject=f"phase:{n}") - rec = build_solved_record(tmp_path, _manifest(), {"fact_x", "dim_y"}) - assert rec.id == "solved:s1" - assert rec.content == "domanda riscritta esplicita" - assert rec.metadata["sql"] == "SELECT 1" - assert rec.metadata["tables"] == ["dim_y", "fact_x"] # ordinate - - -def test_build_falls_back_to_manifest_question(tmp_path): - (tmp_path / "sql_final.sql").write_text("SELECT 1") - append_decision(tmp_path, type="sql_approved", subject="phase:7") - for n in range(1, 8): - append_decision(tmp_path, type="phase_approved", subject=f"phase:{n}") - rec = build_solved_record(tmp_path, _manifest(), None) - assert rec.content == "domanda originale" - assert rec.metadata["tables"] == [] - - -def test_build_requires_sql_file(tmp_path): - append_decision(tmp_path, type="sql_approved", subject="phase:7") - with pytest.raises(SolvedIndexError, match="sql_final.sql"): - build_solved_record(tmp_path, _manifest(), None) - - -def test_build_requires_sql_approved(tmp_path): - (tmp_path / "sql_final.sql").write_text("SELECT 1") - with pytest.raises(SolvedIndexError, match="sql_approved"): - build_solved_record(tmp_path, _manifest(), None) diff --git a/harness/tests/test_solved_question.py b/harness/tests/test_solved_question.py deleted file mode 100644 index 91c9e5d0..00000000 --- a/harness/tests/test_solved_question.py +++ /dev/null @@ -1,78 +0,0 @@ -"""L1: coppia domanda->SQL risolta (kind solved_question) — memoria attiva parte B. - -Una sessione finalizzata produce UN record nel vectordb (tabella `memory`, kind -dedicato): embedding = domanda riscritta, metadata = {question, sql, tables, -session_id}. Upsert one-row stile D11 (mai sync: il suo delete-stale cancellerebbe -i record delle altre sessioni). L'hash di dedup copre domanda+SQL, cosi' un -re-finalize che cambia solo l'SQL aggiorna comunque la riga. -""" -from unittest.mock import MagicMock - -from tht.solved import ( - SOLVED_KIND, - _solved_hash, - save_solved_question, - solved_question_record, -) -from tht.vectorstore.reader import tables_for_kinds - - -def _rec(**kw): - base = dict( - session_id="s1", question="quante ablazioni nel 2023", - sql="SELECT count(*) FROM fact_seeablazione", tables=["fact_seeablazione"], - ) - base.update(kw) - return solved_question_record(**base) - - -def test_record_shape(): - r = _rec() - assert r.id == "solved:s1" - assert r.kind == SOLVED_KIND - assert r.content == "quante ablazioni nel 2023" # embedding = solo la domanda - assert r.metadata["sql"].startswith("SELECT") - assert r.metadata["tables"] == ["fact_seeablazione"] - assert r.metadata["session_id"] == "s1" - - -def test_solved_kind_maps_to_memory_table(): - assert tables_for_kinds([SOLVED_KIND]) == ["memory"] - - -def test_save_upserts_single_row_into_memory_table(): - writer = MagicMock() - writer.existing_hashes.return_value = {} - writer.upsert.return_value = 1 - embedder = MagicMock() - embedder.embed_documents.return_value = [[0.1] * 8] - - assert save_solved_question(_rec(), store=writer, embedder=embedder) == 1 - writer.sync.assert_not_called() - table, rows = writer.upsert.call_args[0] - assert table == "memory" - assert len(rows) == 1 - assert rows[0].record.id == "solved:s1" - assert rows[0].record.kind == SOLVED_KIND - assert rows[0].record.metadata["sql"].startswith("SELECT") - - -def test_save_skips_when_question_and_sql_unchanged(): - r = _rec() - writer = MagicMock() - writer.existing_hashes.return_value = {r.id: _solved_hash(r)} - embedder = MagicMock() - assert save_solved_question(r, store=writer, embedder=embedder) == 0 - embedder.embed_documents.assert_not_called() - writer.upsert.assert_not_called() - - -def test_sql_change_alone_triggers_reupsert(): - old = _rec() - new = _rec(sql="SELECT 1") # stessa domanda, SQL diverso - writer = MagicMock() - writer.existing_hashes.return_value = {old.id: _solved_hash(old)} - writer.upsert.return_value = 1 - embedder = MagicMock() - embedder.embed_documents.return_value = [[0.0] * 4] - assert save_solved_question(new, store=writer, embedder=embedder) == 1 diff --git a/harness/tht/cli/memory_cmd.py b/harness/tht/cli/memory_cmd.py index c9ca446c..213a4be8 100644 --- a/harness/tht/cli/memory_cmd.py +++ b/harness/tht/cli/memory_cmd.py @@ -403,12 +403,12 @@ def index_solved_session(cfg, session_id: str) -> int: 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 + from tht.memory import index_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, + return index_solved_question( + load_snapshot_or_exit(cfg, session_id), + promoted_tables_for(cfg, session_id), store=store, embedder=make_embedder(cfg.embeddings), ) @@ -423,7 +423,7 @@ def solved_index_cmd( """Indicizza la coppia domanda->SQL nel semantic store (backfill; il finalize lo fa da solo).""" import json as _json - from tht.solved import SolvedIndexError + from tht.memory import SolvedIndexError cfg = _load_config_or_exit(config) require_vector_write_allowed(cfg, "memory solved-index") @@ -460,7 +460,7 @@ def solved_search_cmd( 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.memory import search_solved_questions from tht.vectorstore.embeddings import EmbeddingsError cfg = _load_config_or_exit(config) @@ -471,7 +471,12 @@ def solved_search_cmd( try: searcher = open_searcher(cfg) embedder = make_embedder(cfg.embeddings) - hits = searcher.search(embedder.embed_query(question), top_n=top, kinds=[SOLVED_KIND]) + results = search_solved_questions( + question, + searcher=searcher, + embedder=embedder, + top=top, + ) except (VectorStoreError, VectorReadUnavailable, EmbeddingsError, OperationalError) as e: typer.secho( f"ATTENZIONE: exemplar non disponibili ({e}). Prosegui senza.", @@ -480,16 +485,6 @@ def solved_search_cmd( 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 diff --git a/harness/tht/cli/search_cmd.py b/harness/tht/cli/search_cmd.py index de220f28..ec27254b 100644 --- a/harness/tht/cli/search_cmd.py +++ b/harness/tht/cli/search_cmd.py @@ -257,7 +257,7 @@ def pack_cmd( from tht.cli.vector_cmd import make_embedder, open_searcher, require_vector_cfg from tht.ports.vector import VectorReadUnavailable, VectorStoreError from tht.search import combined_search, schema_tables - from tht.solved import SOLVED_KIND + from tht.memory import SOLVED_KIND from tht.vectorstore.embeddings import EmbeddingsError cfg = _load_config_or_exit(config) diff --git a/harness/tht/cli/session_cmd.py b/harness/tht/cli/session_cmd.py index 10d8819f..bf8acf27 100644 --- a/harness/tht/cli/session_cmd.py +++ b/harness/tht/cli/session_cmd.py @@ -578,10 +578,11 @@ def finalize_cmd(session_id: str = typer.Argument(...), config: Path = CONFIG_OP # --- batteria di validazione su sql_final.sql --- assert sql is not None + promoted_tables = promoted_tables_for(cfg, session_id) check = validate_sql( sql, physical=_load_physical_or_exit(cfg), - promoted_tables=promoted_tables_for(cfg, session_id), + promoted_tables=promoted_tables, forbidden_functions=set(cfg.execution.forbidden_functions), ) if not check.ok: @@ -623,14 +624,29 @@ def finalize_cmd(session_id: str = typer.Argument(...), config: Path = CONFIG_OP repository, session_id, validation_report=report, evidence=evidence ) # --- memoria attiva (parte B): indicizza la coppia domanda->SQL, best-effort --- - # Import lazy: memory_cmd importa da session_cmd (un import top-level qui sarebbe - # circolare). Qualunque errore (writer key assente, VPN giu', Ollama spento) NON - # deve bloccare il finalize: l'indice e' derivato e recuperabile con - # `tht memory solved-index `. + # Qualunque errore (writer key assente, VPN giu', Ollama spento) NON deve + # bloccare il finalize: l'indice e' derivato e recuperabile con + # `tht memory solved-index `. Memory owns this best-effort policy; core + # has already committed the authoritative finalized snapshot above. try: - from tht.cli.memory_cmd import index_solved_session + from tht.adapters.factory import build_vector_store + from tht.cli.vector_cmd import make_embedder + from tht.memory import index_solved_question_best_effort - if index_solved_session(cfg, session_id): + finalized_snapshot = repository.get(session_id) + outcome = index_solved_question_best_effort( + finalized_snapshot, + promoted_tables, + store_factory=lambda: build_vector_store(cfg, require_write=True), + embedder_factory=lambda: make_embedder(cfg.embeddings), + ) + if outcome.error is not None: + typer.secho( + f"ATTENZIONE: coppia domanda->SQL non indicizzata ({outcome.error}). " + f"Recupera con `tht memory solved-index {session_id}`.", + fg=typer.colors.YELLOW, err=True, + ) + elif outcome.upserted: typer.secho( "OK: coppia domanda->SQL indicizzata nel vectordb (solved_question).", fg=typer.colors.GREEN, diff --git a/harness/tht/memory/__init__.py b/harness/tht/memory/__init__.py new file mode 100644 index 00000000..889e16ea --- /dev/null +++ b/harness/tht/memory/__init__.py @@ -0,0 +1,63 @@ +"""Public Memory facade for reusable decisions and solved-question exemplars.""" + +from .core import ( + EDITABLE_FIELDS, + MAX_PROMOTION_CANDIDATES, + REUSABLE_TYPES, + MemoryNotFound, + MemoryRecord, + decided_memory_ids, + declined_promotion_seqs, + delete_record, + load_registry, + memory_vector_record_for_decision, + memory_vector_records, + preview_promotions, + preview_promotions_snapshot, + promote, + promote_snapshot, + recall_memories, + reusable_promotions, + reusable_promotions_snapshot, + save_one_memory, + save_registry, + update_record, +) +from .solved import ( + SOLVED_KIND, + SolvedIndexError, + SolvedIndexOutcome, + index_solved_question, + index_solved_question_best_effort, + search_solved_questions, +) + +__all__ = [ + "EDITABLE_FIELDS", + "MAX_PROMOTION_CANDIDATES", + "REUSABLE_TYPES", + "SOLVED_KIND", + "MemoryNotFound", + "MemoryRecord", + "SolvedIndexError", + "SolvedIndexOutcome", + "decided_memory_ids", + "declined_promotion_seqs", + "delete_record", + "index_solved_question", + "index_solved_question_best_effort", + "load_registry", + "memory_vector_record_for_decision", + "memory_vector_records", + "preview_promotions", + "preview_promotions_snapshot", + "promote", + "promote_snapshot", + "recall_memories", + "reusable_promotions", + "reusable_promotions_snapshot", + "save_one_memory", + "save_registry", + "search_solved_questions", + "update_record", +] diff --git a/harness/tht/memory.py b/harness/tht/memory/core.py similarity index 100% rename from harness/tht/memory.py rename to harness/tht/memory/core.py diff --git a/harness/tht/solved.py b/harness/tht/memory/solved.py similarity index 54% rename from harness/tht/solved.py rename to harness/tht/memory/solved.py index e0b33ce9..1e4dcfa0 100644 --- a/harness/tht/solved.py +++ b/harness/tht/memory/solved.py @@ -6,18 +6,26 @@ nel semantic store workspace-scoped, nel gruppo logico `memory` con kind dedicat consulta nelle fasi F4/F6/F7 con `tht memory solved-search` come materiale di riferimento (exemplar), NON come decisione da ri-applicare. -Scrittura: SOLO upsert one-row stile D11 (`save_solved_question`). Questi record +Scrittura: SOLO upsert one-row stile D11 (`index_solved_question`). Questi record non passano MAI da `VectorStore.sync`/`RestVectorWriter.sync`: il passo delete-stale del sync, ricevendo il solo record corrente, cancellerebbe le coppie delle altre sessioni. Per lo stesso motivo l'hash di dedup e' calcolato qui (domanda+SQL) e non dal solo content come fa il sync. """ +from dataclasses import dataclass + +from tht.phase import effective_decisions +from tht.ports.vector import VectorWriteRecord +from tht.session.models import SessionSnapshot from tht.vectorstore.records import VectorRecord +from tht.vectorstore.store import content_hash + +from .core import question_context SOLVED_KIND = "solved_question" -def solved_question_record( +def _solved_question_record( *, session_id: str, question: str, sql: str, tables: list[str] ) -> VectorRecord: return VectorRecord( @@ -38,18 +46,14 @@ def solved_question_record( def _solved_hash(record: VectorRecord) -> str: # La domanda e' l'embedding (content); l'SQL vive solo nel metadata. L'hash # copre entrambi: un re-finalize che cambia solo l'SQL aggiorna la riga. - from tht.vectorstore.store import content_hash - return content_hash(record.content + "\n" + str(record.metadata.get("sql", ""))) -def save_solved_question(record: VectorRecord, *, store, embedder) -> int: +def _save_solved_question(record: VectorRecord, *, store, embedder) -> int: """Upsert one-row della coppia domanda->SQL via writer key (stesso pattern di save_one_memory, spec D11): hash dedup client-side, embedding solo se domanda o SQL sono cambiati. Ritorna il numero di righe upsertate (0 = invariata).""" - from tht.ports.vector import VectorWriteRecord - new_hash = _solved_hash(record) existing = store.existing_hashes("memory", [SOLVED_KIND]) if existing.get(record.id) == new_hash: @@ -65,39 +69,78 @@ class SolvedIndexError(Exception): """La sessione non ha (ancora) gli artefatti per il record solved_question.""" -def build_solved_record(session_dir, manifest, promoted_tables) -> VectorRecord: - """Costruisce il record dagli artefatti persistiti (vista effective D15): - richiede sql_final.sql e la decisione sql_approved; la domanda e' l'ultima - question_rewritten, fallback la domanda del manifest.""" - from tht.memory import question_context - from tht.phase import effective_decisions - - sql_file = session_dir / "sql_final.sql" - if not sql_file.exists(): - raise SolvedIndexError("sql_final.sql assente") - decisions = effective_decisions(session_dir) - if not any(d.type == "sql_approved" for d in decisions): - raise SolvedIndexError("decisione sql_approved assente") - return solved_question_record( - session_id=manifest.id, - question=question_context(decisions, manifest), - sql=sql_file.read_text().strip(), - tables=sorted(promoted_tables or set()), - ) - - -def build_solved_snapshot(snapshot, promoted_tables) -> VectorRecord: - from tht.memory import question_context - from tht.phase import effective_decisions - +def _build_solved_snapshot( + snapshot: SessionSnapshot, + promoted_tables: set[str] | None, +) -> VectorRecord: sql = snapshot.artifacts.get("sql_final") if sql is None: raise SolvedIndexError("sql_final.sql assente") decisions = effective_decisions(snapshot) if not any(d.type == "sql_approved" for d in decisions): raise SolvedIndexError("decisione sql_approved assente") - return solved_question_record( + return _solved_question_record( session_id=snapshot.manifest.id, question=question_context(decisions, snapshot.manifest), sql=sql.strip(), tables=sorted(promoted_tables or set()), ) + + +def index_solved_question( + snapshot: SessionSnapshot, + promoted_tables: set[str] | None, + *, + store, + embedder, +) -> int: + """Index the finalized question through Memory's one-row, idempotent policy.""" + return _save_solved_question( + _build_solved_snapshot(snapshot, promoted_tables), + store=store, + embedder=embedder, + ) + + +@dataclass(frozen=True) +class SolvedIndexOutcome: + upserted: int | None + error: str | None = None + + +def index_solved_question_best_effort( + snapshot: SessionSnapshot, + promoted_tables: set[str] | None, + *, + store_factory, + embedder_factory, +) -> SolvedIndexOutcome: + """Attempt derived indexing without turning it into workflow success state.""" + try: + upserted = index_solved_question( + snapshot, + promoted_tables, + store=store_factory(), + embedder=embedder_factory(), + ) + except Exception as error: + return SolvedIndexOutcome(upserted=None, error=str(error)) + return SolvedIndexOutcome(upserted=upserted) + + +def search_solved_questions(question: str, *, searcher, embedder, top: int = 3) -> list[dict]: + """Return solved-question exemplars in semantic-search rank order.""" + hits = searcher.search( + embedder.embed_query(question), + top_n=top, + kinds=[SOLVED_KIND], + ) + return [ + { + "session_id": hit.metadata.get("session_id", hit.ref), + "question": hit.metadata.get("question", hit.content), + "sql": hit.metadata.get("sql", ""), + "tables": hit.metadata.get("tables", []), + "score": round(hit.similarity, 4), + } + for hit in hits + ]