refactor(memory): own solved-question lifecycle (#24)
This commit is contained in:
@@ -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
|
||||
@@ -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"])]
|
||||
@@ -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())
|
||||
|
||||
@@ -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)
|
||||
@@ -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
|
||||
Reference in New Issue
Block a user