From f1a9b567ba606e4de731289568b8f1c436984ea0 Mon Sep 17 00:00:00 2001 From: mptyl Date: Mon, 24 Aug 2026 21:59:47 +0200 Subject: [PATCH] feat(evidence): contribute to semantic stages (#42) --- harness/.pi/skills/tht-sessione/SKILL.md | 27 +++ .../modules/evidence/runtime-search.md | 26 +++ .../skills/tht-sessione/projection.md.tmpl | 2 + .../tests/fixtures/approved_cli_surface.json | 1 + harness/tests/test_cli_surface.py | 2 +- .../tests/test_evidence_facade_contract.py | 133 +++++++++++- .../tests/test_evidence_session_receipts.py | 38 ++++ harness/tests/test_pi_skill_projection.py | 3 +- harness/tests/test_search_pack.py | 30 ++- harness/tht/adapters/vector/qdrant.py | 27 ++- harness/tht/cli/search_cmd.py | 111 +++++++++- harness/tht/evidence/__init__.py | 16 +- harness/tht/evidence/search.py | 189 +++++++++++++++++- harness/tht/evidence/session.py | 34 +++- harness/tht/pi_skill_projection.py | 1 + harness/tht/session/filesystem_repository.py | 1 + harness/tht/session/postgres_repository.py | 1 + 17 files changed, 619 insertions(+), 23 deletions(-) create mode 100644 harness/.pi/skills/tht-sessione/modules/evidence/runtime-search.md create mode 100644 harness/tests/test_evidence_session_receipts.py diff --git a/harness/.pi/skills/tht-sessione/SKILL.md b/harness/.pi/skills/tht-sessione/SKILL.md index 25a40c04..6d82885a 100644 --- a/harness/.pi/skills/tht-sessione/SKILL.md +++ b/harness/.pi/skills/tht-sessione/SKILL.md @@ -167,6 +167,33 @@ persisted state is your only context. Bootstrap before doing anything else: If `status` is `finalized`, the session is read-only — do not resume; tell the reviewer it is complete. (The backend already refuses resume for finalized/archived sessions.) +## Evidence runtime contributor + +Evidence contributes to existing semantic stages; it is never a visible phase and does +not write decisions, canonical artifacts, or workflow state. Candidates are not truth: +show their provenance and let the reviewer decide. A formula is `kind=formula`, not a +separate store. + +Use the phase-appropriate Evidence purpose, and make every mapped stage search +independently: + +- `clarification` → `disambiguation`; +- `rewriting` → `rewriting`; +- `schema_linking` → `schema_linking`; +- `cte` and `final_sql` → `sql_generation`. + +After consuming the F1 retrieval pack, and before making a proposal in every other +mapped stage, call `tht search evidence "" +--stage --session --json` before making the stage proposal. Add +only the available approved context (`--concept`, `--table`, `--column`) and use +`--require-*` only for a mandatory constraint. In `final_sql`, include the approved CTE +plan in the current stage context. The command records only its minimal receipt. + +An `available` outcome with zero results is visible but does not block the stage. An +`unavailable` outcome blocks the calling stage: report the sanitized failure and retry +the same stage later. Never use a stale generation or retry with another purpose. +Never call Evidence from `memory` or `synthesis`. + ## Phase 1 — Clarification Prerequisite: you must already be in Phase 1. diff --git a/harness/.pi/skills/tht-sessione/modules/evidence/runtime-search.md b/harness/.pi/skills/tht-sessione/modules/evidence/runtime-search.md new file mode 100644 index 00000000..aa84e7ad --- /dev/null +++ b/harness/.pi/skills/tht-sessione/modules/evidence/runtime-search.md @@ -0,0 +1,26 @@ +## Evidence runtime contributor + +Evidence contributes to existing semantic stages; it is never a visible phase and does +not write decisions, canonical artifacts, or workflow state. Candidates are not truth: +show their provenance and let the reviewer decide. A formula is `kind=formula`, not a +separate store. + +Use the phase-appropriate Evidence purpose, and make every mapped stage search +independently: + +- `clarification` → `disambiguation`; +- `rewriting` → `rewriting`; +- `schema_linking` → `schema_linking`; +- `cte` and `final_sql` → `sql_generation`. + +After consuming the F1 retrieval pack, and before making a proposal in every other +mapped stage, call `tht search evidence "" +--stage --session --json` before making the stage proposal. Add +only the available approved context (`--concept`, `--table`, `--column`) and use +`--require-*` only for a mandatory constraint. In `final_sql`, include the approved CTE +plan in the current stage context. The command records only its minimal receipt. + +An `available` outcome with zero results is visible but does not block the stage. An +`unavailable` outcome blocks the calling stage: report the sanitized failure and retry +the same stage later. Never use a stale generation or retry with another purpose. +Never call Evidence from `memory` or `synthesis`. diff --git a/harness/.pi/skills/tht-sessione/projection.md.tmpl b/harness/.pi/skills/tht-sessione/projection.md.tmpl index e5a86e8e..59789b54 100644 --- a/harness/.pi/skills/tht-sessione/projection.md.tmpl +++ b/harness/.pi/skills/tht-sessione/projection.md.tmpl @@ -165,6 +165,8 @@ persisted state is your only context. Bootstrap before doing anything else: If `status` is `finalized`, the session is read-only — do not resume; tell the reviewer it is complete. (The backend already refuses resume for finalized/archived sessions.) +{{EVIDENCE_RUNTIME_SEARCH}} + {{DISAMBIGUATION_INSTRUCTIONS}} {{MEMORY_INSTRUCTIONS}} diff --git a/harness/tests/fixtures/approved_cli_surface.json b/harness/tests/fixtures/approved_cli_surface.json index c52211ee..bb0d44e6 100644 --- a/harness/tests/fixtures/approved_cli_surface.json +++ b/harness/tests/fixtures/approved_cli_surface.json @@ -31,6 +31,7 @@ "schema render", "schema suggest-fks", "search find", + "search evidence", "search pack", "session archive", "session check", diff --git a/harness/tests/test_cli_surface.py b/harness/tests/test_cli_surface.py index 91593cf5..cee69f84 100644 --- a/harness/tests/test_cli_surface.py +++ b/harness/tests/test_cli_surface.py @@ -32,7 +32,7 @@ def test_typer_tree_matches_the_approved_command_surface(): approved = _approved_surface() expected = set(approved["maintained"]) | set(approved["enhanced"]) - assert len(approved["maintained"]) == 58 + assert len(approved["maintained"]) == 59 assert len(approved["enhanced"]) == 8 assert len(approved["erased"]) == 14 assert not (expected & set(approved["erased"])) diff --git a/harness/tests/test_evidence_facade_contract.py b/harness/tests/test_evidence_facade_contract.py index 93b20330..9114f804 100644 --- a/harness/tests/test_evidence_facade_contract.py +++ b/harness/tests/test_evidence_facade_contract.py @@ -8,6 +8,7 @@ import pytest from tht.decisions import DecisionRecord from tht.evidence import ( + EvidenceSearchContext, acquire, active_searcher, build_preprocessing_pipeline, @@ -16,6 +17,7 @@ from tht.evidence import ( discover, project_session, resolve_citation, + search_evidence, ) from tht.evidence.contracts import ( AcquiredDocument, @@ -25,6 +27,7 @@ from tht.evidence.contracts import ( ) from tht.evidence.corpus.models import CanonicalDocument, CorpusManifest from tht.evidence.corpus.store import CorpusStore +from tht.ports.vector import VectorReadUnavailable from tht.session.models import Candidate, SchemaLinking @@ -186,13 +189,141 @@ def test_retrieval_entries_preserve_hit_order_and_existing_projection_shape(): SimpleNamespace(label="Second", status="reviewed", content="abcdefgh"), SimpleNamespace(label="First", status=None, content="12345678"), ] - assert build_retrieval_entries(hits, excerpt_chars=5) == [ {"title": "Second", "status": "reviewed", "excerpt": "abcde"}, {"title": "First", "status": None, "excerpt": "12345"}, ] +def test_typed_search_renders_one_stable_query_and_groups_fragments_by_evidence_unit(): + """Removing context rendering, hard filters, or grouping changes this public result.""" + class Searcher: + vector_generation = "gen:" + "a" * 32 + + def __init__(self): + self.calls = [] + + def search(self, embedding, **kwargs): + self.calls.append((embedding, kwargs)) + return [ + SimpleNamespace( + id="fragment:second", similarity=0.7, content="second excerpt", + title="Pediatric range", metadata={ + "evidence_id": "evidence:pediatric-range", "evidence_kind": "formula", + "document_id": "doc:range", "ordinal": 1, + "source_uri": "file:///curated/pediatric-range.md", + "provenance": {"source_file": "source/range.md"}, + }, + ), + SimpleNamespace( + id="fragment:first", similarity=0.9, content="first excerpt", + title="Pediatric range", metadata={ + "evidence_id": "evidence:pediatric-range", "evidence_kind": "formula", + "document_id": "doc:range", "ordinal": 0, + "source_uri": "file:///curated/pediatric-range.md", + "provenance": {"source_file": "source/range.md"}, + }, + ), + ] + + class Embedder: + def __init__(self): + self.queries = [] + + def embed_query(self, query): + self.queries.append(query) + return [0.25] + + searcher = Searcher() + embedder = Embedder() + outcome = search_evidence( + " Pazienti \"Età" + "\r\n" + " pediatrica ", + "schema_linking", + EvidenceSearchContext( + concepts=("pediatrica", "pediatrica", " Età "), + tables=("clinical.patient",), + columns=("clinical.patient.Age",), + required_kinds=("formula",), + required_concepts=("Età",), + required_tables=("clinical.patient",), + required_columns=("clinical.patient.Age",), + ), + searcher=searcher, + embedder=embedder, + ) + + rendered = ( + "Domanda: Pazienti \"Età\n pediatrica\n" + "Concetti: Età, pediatrica\n" + "Tabelle: clinical.patient\n" + "Colonne: clinical.patient.Age" + ) + assert embedder.queries == [rendered] + assert searcher.calls == [([0.25], { + "top_n": 10, + "kinds": ["evidence"], + "query_text": rendered, + "metadata_filter": { + "purpose": "schema_linking", + "required_kinds": ["formula"], + "required_concepts": ["Età"], + "required_tables": ["clinical.patient"], + "required_columns": ["clinical.patient.Age"], + }, + })] + assert outcome.status == "available" + assert outcome.vector_generation == "gen:" + "a" * 32 + assert [(item.evidence_id, item.excerpts, item.provenance, item.citation) for item in outcome.results] == [ + ("evidence:pediatric-range", ("first excerpt", "second excerpt"), + {"source_file": "source/range.md"}, "file:///curated/pediatric-range.md"), + ] + + +def test_typed_search_reports_vector_errors_as_unavailable_not_empty_results(): + class UnavailableSearcher: + vector_generation = "gen:" + "a" * 32 + + def search(self, _embedding, **_kwargs): + raise VectorReadUnavailable("reader unavailable") + + outcome = search_evidence( + "question", "rewriting", EvidenceSearchContext(), + searcher=UnavailableSearcher(), embedder=SimpleNamespace(embed_query=lambda _query: [0.25]), + ) + + assert outcome.status == "unavailable" + assert outcome.code == "vector_unavailable" + assert outcome.results == () + + +def test_typed_search_without_an_active_generation_is_unavailable_not_an_empty_search(): + outcome = search_evidence( + "question", "rewriting", EvidenceSearchContext(), + searcher=SimpleNamespace(), embedder=SimpleNamespace(embed_query=lambda _query: [0.25]), + ) + + assert outcome.status == "unavailable" + assert outcome.code == "active_corpus_unavailable" + + +def test_typed_search_reports_a_malformed_fragment_payload_as_unavailable(): + class Searcher: + vector_generation = "gen:" + "a" * 32 + + def search(self, _embedding, **_kwargs): + return [SimpleNamespace( + id="fragment:bad", similarity=0.5, title="Bad", content="bad", + metadata={"evidence_id": "evidence:bad"}, + )] + + outcome = search_evidence( + "question", "rewriting", EvidenceSearchContext(), + searcher=Searcher(), embedder=SimpleNamespace(embed_query=lambda _query: [0.25]), + ) + + assert outcome.status == "unavailable" + assert outcome.code == "evidence_search_unavailable" + def _canonical_store(root, evidence_id): content = f"# {evidence_id}\n" digest = hashlib.sha256(content.encode()).hexdigest() diff --git a/harness/tests/test_evidence_session_receipts.py b/harness/tests/test_evidence_session_receipts.py new file mode 100644 index 00000000..e8f47365 --- /dev/null +++ b/harness/tests/test_evidence_session_receipts.py @@ -0,0 +1,38 @@ +import json +from datetime import UTC, datetime + +from tht.evidence import EvidenceReceipt, replace_evidence_receipt +from tht.session.filesystem_repository import FilesystemSessionRepository +from tht.session.models import PrincipalContext, SessionManifest + + +def test_receipt_replaces_only_its_semantic_stage_without_copying_evidence_content(tmp_path): + repository = FilesystemSessionRepository( + tmp_path, "workspace-a", PrincipalContext(issuer="test", subject="operator"), + root=tmp_path / "sessions", + ) + session_id = "c2cdbbf5-ae30-432e-8298-837bf55abad9" + repository.create(SessionManifest( + id=session_id, question="q", database="d", schema="s", created_at=datetime(2026, 8, 24, tzinfo=UTC), + )) + + replace_evidence_receipt(repository, session_id, EvidenceReceipt( + stage="clarification", purpose="disambiguation", vector_generation="gen:" + "a" * 32, + evidence_ids=("evidence:age",), + )) + replace_evidence_receipt(repository, session_id, EvidenceReceipt( + stage="clarification", purpose="disambiguation", vector_generation="gen:" + "b" * 32, + evidence_ids=(), + )) + replace_evidence_receipt(repository, session_id, EvidenceReceipt( + stage="cte", purpose="sql_generation", vector_generation="gen:" + "b" * 32, + evidence_ids=("evidence:age", "evidence:procedure"), + )) + + stored = json.loads(repository.read_artifact(session_id, "evidence_receipts")) + assert stored == [ + {"stage": "clarification", "purpose": "disambiguation", "vector_generation": "gen:" + "b" * 32, + "evidence_ids": []}, + {"stage": "cte", "purpose": "sql_generation", "vector_generation": "gen:" + "b" * 32, + "evidence_ids": ["evidence:age", "evidence:procedure"]}, + ] diff --git a/harness/tests/test_pi_skill_projection.py b/harness/tests/test_pi_skill_projection.py index 4a9d1858..bd40e50f 100644 --- a/harness/tests/test_pi_skill_projection.py +++ b/harness/tests/test_pi_skill_projection.py @@ -9,13 +9,14 @@ from tht.pi_skill_projection import ( render_projection, ) -BASELINE_SHA256 = "626a794071c095a4f20fffabb3bab901f05c101590adbdc58e45adfae56f3219" +BASELINE_SHA256 = "bb6daa6fe83d22f5e701025c5334d72ec9eb35349e90d24cb9d7f6290d0fecfe" def test_modular_pi_skill_renders_the_byte_identical_approved_projection(): rendered = render_projection() assert FRAGMENT_ORDER == ( + ("{{EVIDENCE_RUNTIME_SEARCH}}", "evidence/runtime-search.md"), ("{{DISAMBIGUATION_OPEN_AMBIGUITY}}", "disambiguation/open-ambiguity.md"), ("{{DISAMBIGUATION_INSTRUCTIONS}}", "disambiguation/phase-1.md"), ("{{MEMORY_INSTRUCTIONS}}", "memory/phase-2.md"), diff --git a/harness/tests/test_search_pack.py b/harness/tests/test_search_pack.py index 35421d9e..f41451f3 100644 --- a/harness/tests/test_search_pack.py +++ b/harness/tests/test_search_pack.py @@ -5,7 +5,9 @@ from types import SimpleNamespace from typer.testing import CliRunner from tht.cli import app -from tht.config import load_config +from tht.config import load_config, workspace_id_for_config +from tht.evidence.corpus.models import CorpusManifest +from tht.evidence.corpus.store import CorpusStore from tht.jobs.dwh_pipeline import DwhPreprocessPipeline, config_dwh_binding from tht.mschema.models import ColumnPhysical, PhysicalSchema, TablePhysical from tht.ports.vector import VectorReadUnavailable @@ -75,6 +77,19 @@ def _workspace(tmp_path, with_session=None): ], lsh_filenames=("s_lsh.pkl", "s_minhashes.pkl", "s_meta.json"), ).run() + corpus = CorpusStore(tmp_path / "corpus") + generation = "gen:" + "a" * 32 + active = corpus.stage( + CorpusManifest( + vector_generation=generation, + embedding_model="test", + embedding_dimensions=1, + metadata={"workspace_id": workspace_id_for_config(load_config(cfg), cfg)}, + ), + {}, + generation=generation, + ) + corpus.publish(active) if with_session: sdir = tmp_path / "sessions" / with_session sdir.mkdir(parents=True) @@ -111,7 +126,7 @@ def test_pack_single_embed_and_sections(tmp_path, monkeypatch): _patch(monkeypatch, emb, searcher) res = CliRunner().invoke(app, ["search", "pack", "quanti pazienti", "-c", str(cfg)]) assert res.exit_code == 0, res.output - assert emb.calls == 1 # UN solo embedding per le tre ricerche + assert emb.calls == 2 # schema/memory share one; Evidence embeds its own deterministic text assert [call["kinds"] for call in searcher.calls] == [ ["schema_table", "schema_column"], ["solved_question"], @@ -167,3 +182,14 @@ def test_pack_degrades_gracefully(tmp_path, monkeypatch): data = json.loads(res.output[res.output.index("{"):]) assert data["tables"] == [] and data["evidence"] == [] and data["solved"] == [] assert any("retrieval non disponibile" in w for w in data["warnings"]) + + +def test_pack_refuses_to_treat_an_unavailable_evidence_corpus_as_no_matches(tmp_path, monkeypatch): + cfg = _workspace(tmp_path) + (tmp_path / "corpus" / "ACTIVE").unlink() + _patch(monkeypatch, _FakeEmbedder(), _FakeSearcher()) + + result = CliRunner().invoke(app, ["search", "pack", "q", "-c", str(cfg)]) + + assert result.exit_code == 1 + assert "Evidence non disponibile" in result.output diff --git a/harness/tht/adapters/vector/qdrant.py b/harness/tht/adapters/vector/qdrant.py index 5c7fbd87..b79cfc87 100644 --- a/harness/tht/adapters/vector/qdrant.py +++ b/harness/tht/adapters/vector/qdrant.py @@ -143,7 +143,13 @@ class QdrantVectorStore: filter_must.append(self._semantic_kind_filter(allowed_record_kinds)) filter_must.append({"key": "record_kind", "match": {"any": allowed_record_kinds}}) if metadata_filter is not None: - if set(metadata_filter) != {"vector_generation", "document_ids", "workspace_id"}: + allowed_filters = { + "vector_generation", "document_ids", "workspace_id", "purpose", + "required_kinds", "required_concepts", "required_tables", "required_columns", + } + if not {"vector_generation", "document_ids", "workspace_id"} <= set(metadata_filter) or ( + set(metadata_filter) - allowed_filters + ): raise VectorStoreError("Unsupported vector metadata filter") generation = metadata_filter["vector_generation"] document_ids = metadata_filter["document_ids"] @@ -160,6 +166,25 @@ class QdrantVectorStore: {"key": "vector_generation", "match": {"value": generation}}, {"key": "document_id", "match": {"any": document_ids}}, ]) + purpose = metadata_filter.get("purpose") + if purpose is not None: + if not isinstance(purpose, str): + raise VectorStoreError("Invalid vector metadata filter") + filter_must.append({"key": "purposes", "match": {"value": purpose}}) + required_kinds = metadata_filter.get("required_kinds", []) + if not isinstance(required_kinds, list) or not all(isinstance(item, str) for item in required_kinds): + raise VectorStoreError("Invalid vector metadata filter") + if required_kinds: + filter_must.append({"key": "evidence_kind", "match": {"any": required_kinds}}) + for filter_key, payload_key in ( + ("required_concepts", "scope.concepts"), + ("required_tables", "scope.tables"), + ("required_columns", "scope.columns"), + ): + values = metadata_filter.get(filter_key, []) + if not isinstance(values, list) or not all(isinstance(item, str) for item in values): + raise VectorStoreError("Invalid vector metadata filter") + filter_must.extend({"key": payload_key, "match": {"value": item}} for item in values) if query_text is None: if allowed_record_kinds == ["evidence"]: raise VectorStoreError("Evidence hybrid query text is required") diff --git a/harness/tht/cli/search_cmd.py b/harness/tht/cli/search_cmd.py index a17bcef9..90f9b0bc 100644 --- a/harness/tht/cli/search_cmd.py +++ b/harness/tht/cli/search_cmd.py @@ -14,6 +14,14 @@ KIND_MAP = { "formula": [], # solo formula store (D14b), niente LSH/vector } +_STAGE_PURPOSES = { + "clarification": "disambiguation", + "rewriting": "rewriting", + "schema_linking": "schema_linking", + "cte": "sql_generation", + "final_sql": "sql_generation", +} + # Default di `--top` per le famiglie diverse da `schema` (numero di risultati). Per `schema` # `--top` indica il numero di TABELLE candidate ed e' configurabile via `search.top_schema_tables` # (recupero ancorato alle tabelle: di ognuna si rendono tutte le colonne + FK). @@ -31,6 +39,79 @@ def _leased_dwh_snapshot(cfg, context: typer.Context): return snapshot +@search_app.command("evidence") +def evidence_search_cmd( + ctx: typer.Context, + query: str = typer.Argument(..., help="Domanda o contesto dello stage."), + stage: str = typer.Option(..., "--stage", help="Stage semantico chiamante."), + config: Path = CONFIG_OPT, + session: str | None = typer.Option(None, "--session", help="Sessione per la ricevuta minima."), + concept: list[str] = typer.Option([], "--concept"), + table: list[str] = typer.Option([], "--table"), + column: list[str] = typer.Option([], "--column"), + require_kind: list[str] = typer.Option([], "--require-kind"), + require_concept: list[str] = typer.Option([], "--require-concept"), + require_table: list[str] = typer.Option([], "--require-table"), + require_column: list[str] = typer.Option([], "--require-column"), + top: int = typer.Option(10, "--top"), + json_out: bool = typer.Option(False, "--json"), +) -> None: + """Run one typed, purpose-bound Evidence search for a semantic workflow stage.""" + from tht.cli.vector_cmd import make_embedder, open_searcher, require_vector_cfg + from tht.evidence import ( + EvidenceReceipt, + EvidenceSearchContext, + active_searcher, + replace_evidence_receipt, + search_evidence, + validate_corpus_workspace, + ) + + purpose = _STAGE_PURPOSES.get(stage) + if purpose is None: + raise typer.BadParameter("stage must be clarification, rewriting, schema_linking, cte, or final_sql") + cfg = _load_config_or_exit(config) + workspace_id = workspace_id_for_config(cfg, config) + validate_corpus_workspace(cfg, workspace_id) + require_vector_cfg(cfg) + outcome = search_evidence( + query, purpose, + EvidenceSearchContext( + concepts=tuple(concept), tables=tuple(table), columns=tuple(column), + required_kinds=tuple(require_kind), required_concepts=tuple(require_concept), + required_tables=tuple(require_table), required_columns=tuple(require_column), + ), + searcher=active_searcher(cfg, open_searcher(cfg), workspace_id=workspace_id), + embedder=make_embedder(cfg.embeddings), top_n=top, + ) + if outcome.status == "unavailable": + payload = {"status": outcome.status, "code": outcome.code, "message": outcome.message} + if json_out: + typer.echo(json.dumps(payload, ensure_ascii=False)) + else: + typer.secho(f"ERRORE: {outcome.message}", fg=typer.colors.RED, err=True) + raise typer.Exit(1) + if session: + from tht.cli.session_cmd import load_session_or_exit, session_repository + + load_session_or_exit(cfg, session) + replace_evidence_receipt(session_repository(cfg), session, EvidenceReceipt( + stage=stage, purpose=purpose, vector_generation=outcome.vector_generation or "", + evidence_ids=tuple(result.evidence_id for result in outcome.results), + )) + payload = { + "status": "available", "vector_generation": outcome.vector_generation, + "results": [ + {"evidence_id": result.evidence_id, "title": result.title, "kind": result.kind, + "excerpts": list(result.excerpts), "provenance": result.provenance, + "citation": result.citation, "document_id": result.document_id} + for result in outcome.results + ], + } + if json_out: + typer.echo(json.dumps(payload, ensure_ascii=False, indent=2)) + + @search_app.command("find") def search_cmd( ctx: typer.Context, @@ -254,8 +335,10 @@ def pack_cmd( from tht.cli.vector_cmd import make_embedder, open_searcher, require_vector_cfg from tht.evidence import ( + EvidenceSearchContext, active_searcher, build_retrieval_entries, + search_evidence, validate_corpus_workspace, ) from tht.memory import SOLVED_KIND @@ -274,6 +357,7 @@ def pack_cmd( evidence: list[dict] = [] solved: list[dict] = [] warnings: list[str] = [] + evidence_outcome = None degrade = (VectorStoreError, VectorReadUnavailable, EmbeddingsError, OperationalError) vec = None @@ -309,12 +393,19 @@ def pack_cmd( except degrade as e: warnings.append(f"ricerca schema fallita ({e})") try: - ev = combined_search( - keyword=question, lsh_hits=None, store=searcher, embedder=embedder, - top=PACK_EVIDENCE_TOP, rrf_k=cfg.search.rrf_k, - kinds=KIND_MAP["evidence"], query_vec=vec, + evidence_outcome = search_evidence( + question, "disambiguation", EvidenceSearchContext(), searcher=searcher, + embedder=embedder, top_n=PACK_EVIDENCE_TOP, ) - evidence = build_retrieval_entries(ev, excerpt_chars=PACK_EXCERPT_CHARS) + if evidence_outcome.status == "available": + evidence = build_retrieval_entries(evidence_outcome.results, excerpt_chars=PACK_EXCERPT_CHARS) + else: + typer.secho( + f"ERRORE: Evidence non disponibile ({evidence_outcome.code})", + fg=typer.colors.RED, + err=True, + ) + raise typer.Exit(code=1) except degrade as e: warnings.append(f"ricerca evidence fallita ({e})") try: @@ -367,9 +458,17 @@ def pack_cmd( if session: from tht.cli.session_cmd import load_session_or_exit, session_repository + from tht.evidence import EvidenceReceipt, replace_evidence_receipt load_session_or_exit(cfg, session) - session_repository(cfg).write_artifact(session, "retrieval_pack", md) + repository = session_repository(cfg) + repository.write_artifact(session, "retrieval_pack", md) + if evidence_outcome is not None and evidence_outcome.status == "available": + replace_evidence_receipt(repository, session, EvidenceReceipt( + stage="clarification", purpose="disambiguation", + vector_generation=evidence_outcome.vector_generation or "", + evidence_ids=tuple(result.evidence_id for result in evidence_outcome.results), + )) if not json_out: typer.secho( "OK: retrieval pack scritto " diff --git a/harness/tht/evidence/__init__.py b/harness/tht/evidence/__init__.py index 184bdb92..46205a3e 100644 --- a/harness/tht/evidence/__init__.py +++ b/harness/tht/evidence/__init__.py @@ -39,12 +39,18 @@ from tht.evidence.preprocessing import EvidenceEmbedder, build_preprocessing_pip from tht.evidence.search import ( ActiveEvidenceSearcher, CorpusWorkspaceMismatchError, + EvidenceQueryEmbedder, + EvidenceResult, + EvidenceSearchContext, + EvidenceSearchOutcome, active_searcher, build_retrieval_entries, + render_evidence_query, resolve_citation, + search_evidence, validate_corpus_workspace, ) -from tht.evidence.session import project_session +from tht.evidence.session import EvidenceReceipt, project_session, replace_evidence_receipt from tht.evidence.sources import build_sources __all__ = [ @@ -56,8 +62,13 @@ __all__ = [ "EvidenceManifest", "EvidencePreparationError", "EvidencePreparationReport", + "EvidenceQueryEmbedder", + "EvidenceReceipt", "EvidenceResolutionReport", "EvidenceRestructurer", + "EvidenceResult", + "EvidenceSearchContext", + "EvidenceSearchOutcome", "EvidenceSource", "EvidenceSourceError", "EvidenceSourceErrorCategory", @@ -82,8 +93,11 @@ __all__ = [ "parse_curated_markdown", "prepare_workspace_evidence", "project_session", + "render_evidence_query", + "replace_evidence_receipt", "resolve_citation", "resolve_workspace_evidence", + "search_evidence", "validate_corpus_workspace", "validate_namespaced_value", "validate_safe_metadata", diff --git a/harness/tht/evidence/search.py b/harness/tht/evidence/search.py index ec7210b9..ca42baf3 100644 --- a/harness/tht/evidence/search.py +++ b/harness/tht/evidence/search.py @@ -1,15 +1,175 @@ """Evidence-owned runtime lookup bound to the atomically active corpus generation.""" import re +import unicodedata +from dataclasses import dataclass +from typing import Literal, Protocol +from pydantic import BaseModel, ConfigDict + +from tht.evidence.canonical import EvidenceKind, EvidencePurpose from tht.evidence.corpus.store import CorpusStore -from tht.ports.vector import VectorStoreError +from tht.ports.vector import VectorReadUnavailable, VectorStoreError class CorpusWorkspaceMismatchError(RuntimeError): """The configured workspace does not own the persisted corpus.""" +class EvidenceQueryEmbedder(Protocol): + def embed_query(self, query: str) -> list[float]: ... + + +class EvidenceSearchContext(BaseModel): + """Optional query enrichment and explicit, server-enforced Evidence constraints.""" + + model_config = ConfigDict(frozen=True, extra="forbid") + + concepts: tuple[str, ...] = () + tables: tuple[str, ...] = () + columns: tuple[str, ...] = () + required_kinds: tuple[EvidenceKind, ...] = () + required_concepts: tuple[str, ...] = () + required_tables: tuple[str, ...] = () + required_columns: tuple[str, ...] = () + + +@dataclass(frozen=True) +class EvidenceResult: + evidence_id: str + title: str + kind: str + excerpts: tuple[str, ...] + provenance: dict + citation: str + document_id: str + score: float + + +@dataclass(frozen=True) +class EvidenceSearchOutcome: + status: Literal["available", "unavailable"] + vector_generation: str | None + results: tuple[EvidenceResult, ...] = () + code: str | None = None + message: str | None = None + + @classmethod + def unavailable(cls, code: str, message: str) -> "EvidenceSearchOutcome": + return cls("unavailable", None, (), code, message[:240]) + + +def _normalized_values(values: tuple[str, ...]) -> tuple[str, ...]: + normalized = { + unicodedata.normalize("NFC", value).strip() + for value in values + if isinstance(value, str) and unicodedata.normalize("NFC", value).strip() + } + return tuple(sorted(normalized)) + + +def render_evidence_query(query: str, context: EvidenceSearchContext) -> str: + """Build the one exact query text shared by dense and BM25 retrieval.""" + question = unicodedata.normalize("NFC", query).replace("\r\n", "\n").replace("\r", "\n").strip() + if not question: + raise ValueError("Evidence query must not be empty") + sections = [("Domanda", question)] + for label, values in ( + ("Concetti", _normalized_values(context.concepts)), + ("Tabelle", _normalized_values(context.tables)), + ("Colonne", _normalized_values(context.columns)), + ): + if values: + sections.append((label, ", ".join(values))) + return "\n".join(f"{label}: {value}" for label, value in sections) + + +def _required_metadata_filter(purpose: EvidencePurpose, context: EvidenceSearchContext) -> dict[str, object]: + return { + "purpose": purpose, + "required_kinds": list(_normalized_values(context.required_kinds)), + "required_concepts": list(_normalized_values(context.required_concepts)), + "required_tables": list(_normalized_values(context.required_tables)), + "required_columns": list(_normalized_values(context.required_columns)), + } + + +def _active_vector_generation(searcher) -> str | None: + supplied = getattr(searcher, "vector_generation", None) + if isinstance(supplied, str) and supplied: + return supplied + corpus = getattr(searcher, "corpus", None) + if corpus is None: + return None + with corpus.writer_lock(): + manifest = corpus.active_manifest() + return manifest.vector_generation if manifest is not None else None + + +def _group_evidence_fragments(hits) -> tuple[EvidenceResult, ...]: + grouped: dict[str, list] = {} + for hit in hits: + metadata = getattr(hit, "metadata", {}) + evidence_id = metadata.get("evidence_id") if isinstance(metadata, dict) else None + if not isinstance(evidence_id, str) or not evidence_id: + raise VectorStoreError("Evidence search returned malformed payload") + grouped.setdefault(evidence_id, []).append(hit) + results = [] + for evidence_id, fragments in grouped.items(): + ordered = sorted( + fragments, + key=lambda item: (-float(item.similarity), int(item.metadata.get("ordinal", 0)), item.id), + ) + first = ordered[0] + metadata = first.metadata + citation = metadata.get("source_uri") + document_id = metadata.get("document_id") + evidence_kind = metadata.get("evidence_kind") + if not all(isinstance(value, str) and value for value in (citation, document_id, evidence_kind)): + raise VectorStoreError("Evidence search returned malformed payload") + results.append(EvidenceResult( + evidence_id=evidence_id, + title=str(first.title), + kind=evidence_kind, + excerpts=tuple(str(item.content) for item in ordered), + provenance=dict(metadata.get("provenance", {})), + citation=citation, + document_id=document_id, + score=float(first.similarity), + )) + return tuple(sorted(results, key=lambda item: (-item.score, item.evidence_id))) + + +def search_evidence( + query: str, + purpose: EvidencePurpose, + context: EvidenceSearchContext, + *, + searcher, + embedder: EvidenceQueryEmbedder, + top_n: int = 10, +) -> EvidenceSearchOutcome: + """Search the active generation once, with no stale-generation or purpose fallback.""" + try: + rendered = render_evidence_query(query, context) + generation = _active_vector_generation(searcher) + if generation is None: + return EvidenceSearchOutcome.unavailable("active_corpus_unavailable", "Active Evidence corpus is unavailable") + query_embedding = embedder.embed_query(rendered) + hits = searcher.search( + query_embedding, + top_n=top_n, + kinds=["evidence"], + query_text=rendered, + metadata_filter=_required_metadata_filter(purpose, context), + ) + return EvidenceSearchOutcome("available", generation, _group_evidence_fragments(hits)) + except VectorReadUnavailable: + return EvidenceSearchOutcome.unavailable("vector_unavailable", "Evidence vector search is unavailable") + except (VectorStoreError, CorpusWorkspaceMismatchError): + return EvidenceSearchOutcome.unavailable("evidence_search_unavailable", "Evidence search is unavailable") + + class ActiveEvidenceSearcher: """Searcher facade that enforces ACTIVE generation predicates before LIMIT.""" @@ -78,15 +238,17 @@ class ActiveEvidenceSearcher: if by_generation and (not isinstance(query_text, str) or query_text.strip() == ""): raise VectorStoreError("Evidence hybrid query text is required") for generation, document_ids in sorted(by_generation.items()): + filters = dict(metadata_filter or {}) + filters.update({ + "vector_generation": generation, + "document_ids": sorted(document_ids), + "workspace_id": workspace_id, + }) hits.extend(self.delegate.search( embedding, top_n=top_n, kinds=["evidence"], query_text=query_text, query_language=query_language or self.evidence_language, - metadata_filter={ - "vector_generation": generation, - "document_ids": sorted(document_ids), - "workspace_id": workspace_id, - }, + metadata_filter=filters, )) return sorted(hits, key=lambda hit: (-hit.similarity, hit.id))[:top_n] @@ -144,9 +306,12 @@ def build_retrieval_entries(results, *, excerpt_chars: int) -> list[dict]: """Project ordered Evidence search hits into the retrieval-pack shape.""" return [ { - "title": result.label, - "status": result.status, - "excerpt": result.content[:excerpt_chars], + "title": getattr(result, "title", getattr(result, "label", "")), + "status": getattr(result, "status", None), + "excerpt": ( + result.excerpts[0] if isinstance(result, EvidenceResult) and result.excerpts + else getattr(result, "content", "") + )[:excerpt_chars], } for result in results ] @@ -155,8 +320,14 @@ def build_retrieval_entries(results, *, excerpt_chars: int) -> list[dict]: __all__ = [ "ActiveEvidenceSearcher", "CorpusWorkspaceMismatchError", + "EvidenceQueryEmbedder", + "EvidenceResult", + "EvidenceSearchContext", + "EvidenceSearchOutcome", "active_searcher", "build_retrieval_entries", + "render_evidence_query", "resolve_citation", + "search_evidence", "validate_corpus_workspace", ] diff --git a/harness/tht/evidence/session.py b/harness/tht/evidence/session.py index 37f27066..bce7fb43 100644 --- a/harness/tht/evidence/session.py +++ b/harness/tht/evidence/session.py @@ -1,5 +1,7 @@ """Evidence-specific projection into persisted session artifacts.""" +import json +from dataclasses import dataclass from pathlib import Path from typing import TYPE_CHECKING @@ -56,4 +58,34 @@ def project_session( return list(entries.values()) -__all__ = ["project_session"] +@dataclass(frozen=True) +class EvidenceReceipt: + stage: str + purpose: str + vector_generation: str + evidence_ids: tuple[str, ...] + + def payload(self) -> dict: + return { + "stage": self.stage, + "purpose": self.purpose, + "vector_generation": self.vector_generation, + "evidence_ids": list(self.evidence_ids), + } + + +def replace_evidence_receipt(repository, session_id: str, receipt: EvidenceReceipt) -> None: + """Replace the one minimal receipt for a semantic stage; preserve other stages.""" + current = repository.read_artifact(session_id, "evidence_receipts") + try: + receipts = json.loads(current) if current else [] + except json.JSONDecodeError as error: + raise ValueError("Evidence receipts artifact is malformed") from error + if not isinstance(receipts, list): + raise TypeError("Evidence receipts artifact is malformed") + replaced = [item for item in receipts if isinstance(item, dict) and item.get("stage") != receipt.stage] + replaced.append(receipt.payload()) + repository.write_artifact(session_id, "evidence_receipts", json.dumps(replaced, ensure_ascii=False, indent=2) + "\n") + + +__all__ = ["EvidenceReceipt", "project_session", "replace_evidence_receipt"] diff --git a/harness/tht/pi_skill_projection.py b/harness/tht/pi_skill_projection.py index ad3eb81b..d5e49f63 100644 --- a/harness/tht/pi_skill_projection.py +++ b/harness/tht/pi_skill_projection.py @@ -12,6 +12,7 @@ PROJECTION_PATH = SKILL_ROOT / "SKILL.md" # This tuple is the composition contract. Never derive it from directory order. FRAGMENT_ORDER = ( + ("{{EVIDENCE_RUNTIME_SEARCH}}", "evidence/runtime-search.md"), ("{{DISAMBIGUATION_OPEN_AMBIGUITY}}", "disambiguation/open-ambiguity.md"), ("{{DISAMBIGUATION_INSTRUCTIONS}}", "disambiguation/phase-1.md"), ("{{MEMORY_INSTRUCTIONS}}", "memory/phase-2.md"), diff --git a/harness/tht/session/filesystem_repository.py b/harness/tht/session/filesystem_repository.py index 6787f9bc..0f92a28e 100644 --- a/harness/tht/session/filesystem_repository.py +++ b/harness/tht/session/filesystem_repository.py @@ -23,6 +23,7 @@ _ARTIFACT_FILES = { "question": "question.md", "schema_linking": "schema_linking.json", "evidence": "evidence.json", + "evidence_receipts": "evidence_receipts.json", "sql_final": "sql_final.sql", "validation_report": "validation_report.md", "retrieval_pack": "retrieval_pack.md", diff --git a/harness/tht/session/postgres_repository.py b/harness/tht/session/postgres_repository.py index 771d8f2c..68eb4a92 100644 --- a/harness/tht/session/postgres_repository.py +++ b/harness/tht/session/postgres_repository.py @@ -29,6 +29,7 @@ _ARTIFACT_KEYS = { "question", "schema_linking", "evidence", + "evidence_receipts", "sql_final", "validation_report", "retrieval_pack",