diff --git a/harness/tests/fixtures/approved_cli_surface.json b/harness/tests/fixtures/approved_cli_surface.json index bb0d44e6..bb1d960f 100644 --- a/harness/tests/fixtures/approved_cli_surface.json +++ b/harness/tests/fixtures/approved_cli_surface.json @@ -10,6 +10,7 @@ "decision add", "decision add-batch", "decision add-join-set", + "evidence evaluate", "evidence prepare", "evidence resolve", "evidence validate", diff --git a/harness/tests/test_cli_surface.py b/harness/tests/test_cli_surface.py index cee69f84..23331d3c 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"]) == 59 + assert len(approved["maintained"]) == 60 assert len(approved["enhanced"]) == 8 assert len(approved["erased"]) == 14 assert not (expected & set(approved["erased"])) diff --git a/harness/tests/test_corpus_pipeline.py b/harness/tests/test_corpus_pipeline.py index 867b1407..13543ab3 100644 --- a/harness/tests/test_corpus_pipeline.py +++ b/harness/tests/test_corpus_pipeline.py @@ -106,7 +106,7 @@ def item(name, fingerprint): def pipeline(tmp_path, source, *, embedder=None, vectors=None, model="model-a", policy=None, - retain=3): + retain=3, candidate_evaluator=None): return CorpusPipeline( store=CorpusStore(tmp_path / "corpus"), sources=[source], embedder=embedder or Embedder(), vector_store=vectors or Vectors(), @@ -114,6 +114,7 @@ def pipeline(tmp_path, source, *, embedder=None, vectors=None, model="model-a", chunk_policy=policy or ChunkPolicy(version="chunk-v1", max_chars=100), pipeline_version="evidence-v1", retain_published_generations=retain, + candidate_evaluator=candidate_evaluator, ) @@ -853,6 +854,55 @@ def test_dimension_mismatch_fails_before_vector_write_and_publish(tmp_path): assert candidate.store.active_generation() is None +def test_failed_candidate_evaluation_never_switches_the_active_generation(tmp_path): + from types import SimpleNamespace + + vectors = Vectors() + active = pipeline(tmp_path, Source([(item("one", "a"), "old")]), vectors=vectors).run().generation + candidate = pipeline( + tmp_path, + Source([(item("one", "b"), "new")]), + vectors=vectors, + candidate_evaluator=lambda manifest: SimpleNamespace(passed=False), + ) + + with pytest.raises(PipelineError, match="candidate retrieval evaluation failed"): + candidate.run() + + assert candidate.store.active_generation() == active + assert {record.record.metadata["vector_generation"] for record in vectors.records} == {active} + + +def test_job_failed_candidate_evaluation_never_switches_the_active_generation(tmp_path): + from types import SimpleNamespace + + vectors = Vectors() + active = pipeline(tmp_path, Source([(item("one", "a"), "old")]), vectors=vectors).run_as_job( + workspace_id="demo", + workspace_root=tmp_path, + config_fingerprint="sha256:" + "1" * 64, + input_fingerprint="sha256:" + "2" * 64, + ).generation + candidate = pipeline( + tmp_path, + Source([(item("one", "b"), "new")]), + vectors=vectors, + candidate_evaluator=lambda manifest: SimpleNamespace(passed=False), + ) + + result = candidate.run_as_job( + workspace_id="demo", + workspace_root=tmp_path, + config_fingerprint="sha256:" + "1" * 64, + input_fingerprint="sha256:" + "3" * 64, + ) + + assert result.status == "failed" + assert result.published is False + assert candidate.store.active_generation() == active + assert {record.record.metadata["vector_generation"] for record in vectors.records} == {active} + + def test_pipeline_marks_each_evidence_fragment_for_server_side_italian_bm25(tmp_path): vectors = Vectors() diff --git a/harness/tests/test_evidence_cli.py b/harness/tests/test_evidence_cli.py index ea182c50..f11f2c4f 100644 --- a/harness/tests/test_evidence_cli.py +++ b/harness/tests/test_evidence_cli.py @@ -115,3 +115,35 @@ def test_evidence_validate_json_is_pristine_and_reports_review_required(monkeypa "schemaVersion": 1, "status": "review_required", } + + +def test_evidence_evaluate_json_reports_the_read_only_generation(monkeypatch, tmp_path): + from tht.cli import evidence_cmd + + monkeypatch.setattr(evidence_cmd, "_canonical_worktree", lambda root: root) + monkeypatch.setattr(evidence_cmd, "evaluate_from_config", lambda *args, **kwargs: { + "workspaceRevision": "a" * 40, + "vectorGeneration": "gen:" + "1" * 32, + "rrf": {"algorithm": "rrf", "k": 60, "prefetch_limit_multiplier": 2}, + "passed": True, + "countsByExpectedKind": {"formula": 1}, + "queries": [], + }) + + result = CliRunner().invoke(app, [ + "evidence", "evaluate", str(tmp_path), "--json", "-c", str(tmp_path / "runtime.yaml"), + ]) + + assert result.exit_code == 0 + assert result.stderr == "" + assert json.loads(result.stdout) == { + "countsByExpectedKind": {"formula": 1}, + "operation": "evidence_evaluate", + "passed": True, + "queries": [], + "rrf": {"algorithm": "rrf", "k": 60, "prefetch_limit_multiplier": 2}, + "schemaVersion": 1, + "status": "passed", + "vectorGeneration": "gen:" + "1" * 32, + "workspaceRevision": "a" * 40, + } diff --git a/harness/tests/test_evidence_evaluation.py b/harness/tests/test_evidence_evaluation.py new file mode 100644 index 00000000..6bcc9c63 --- /dev/null +++ b/harness/tests/test_evidence_evaluation.py @@ -0,0 +1,185 @@ +import pytest + + +def _fixture(path): + path.write_text( + """ +schema_version: 1 +queries: + - id: lexical-code + query: Qual è il codice ICD-10? + profile: lexical + purpose: schema_linking + expected: [evidence:icd] + - id: semantic-age + query: Come distinguo i pazienti pediatrici? + profile: semantic + purpose: sql_generation + expected: [evidence:fascia-pediatrica] + - id: mixed-formula + query: Formula per patient.birth_date? + profile: mixed + purpose: sql_generation + expected: [evidence:fascia-pediatrica] +""".strip(), + encoding="utf-8", + ) + + +class _Embedder: + def embed_query(self, query): + return [float(len(query))] + + +class _Hit: + def __init__(self, evidence_id, kind, score): + self.id = evidence_id + ":fragment" + self.similarity = score + self.metadata = {"evidence_id": evidence_id, "evidence_kind": kind} + + +class _Searcher: + def __init__(self): + self.calls = [] + + def search(self, collections, embedding, **kwargs): + self.calls.append((collections, embedding, kwargs)) + mode = kwargs["retrieval_mode"] + if mode == "dense": + return [_Hit("evidence:fascia-pediatrica", "formula", 0.9)] + if mode == "bm25": + return [_Hit("evidence:icd", "enum", 0.8)] + return [ + _Hit("evidence:fascia-pediatrica", "formula", 0.95), + _Hit("evidence:icd", "enum", 0.8), + ] + + +def test_evaluation_reports_branch_and_fused_ranks_without_turning_diagnostics_into_gates(tmp_path): + from tht.evidence.evaluation import evaluate_retrieval, load_evaluation_fixture + + fixture_path = tmp_path / "evaluation.yaml" + _fixture(fixture_path) + searcher = _Searcher() + + report = evaluate_retrieval( + load_evaluation_fixture(fixture_path), + workspace_revision="a" * 40, + document_generations={"doc:one": "gen:" + "1" * 32}, + workspace_id="psd-clinical", + language="italian", + searcher=searcher, + embedder=_Embedder(), + ) + + assert report.passed is True + assert report.workspace_revision == "a" * 40 + assert report.vector_generation == "gen:" + "1" * 32 + assert report.rrf == {"algorithm": "rrf", "k": 60, "prefetch_limit_multiplier": 2} + assert report.queries[0].hit_at_5 is True + assert report.queries[0].hit_at_10 is True + assert report.queries[0].missing_expected == () + assert report.queries[0].expected[0].dense_rank is None + assert report.queries[0].expected[0].bm25_rank == 1 + assert report.queries[0].expected[0].fused_rank == 2 + assert report.counts_by_expected_kind == {"enum": 1, "formula": 2} + assert len(searcher.calls) == 9 + assert {call[2]["retrieval_mode"] for call in searcher.calls} == {"dense", "bm25", "fused"} + assert all(call[2]["metadata_filter"] == { + "workspace_id": "psd-clinical", + "vector_generation": "gen:" + "1" * 32, + "document_ids": ["doc:one"], + "purpose": call[2]["metadata_filter"]["purpose"], + "required_kinds": [], + "required_concepts": [], + "required_tables": [], + "required_columns": [], + } for call in searcher.calls) + + +def test_evaluation_fails_only_when_a_query_has_no_expected_fused_hit_in_top_ten(tmp_path): + from tht.evidence.evaluation import evaluate_retrieval, load_evaluation_fixture + + fixture_path = tmp_path / "evaluation.yaml" + _fixture(fixture_path) + + class MissingExpectedSearcher(_Searcher): + def search(self, collections, embedding, **kwargs): + self.calls.append((collections, embedding, kwargs)) + return [_Hit("evidence:other", "domain", 1.0)] + + report = evaluate_retrieval( + load_evaluation_fixture(fixture_path), + workspace_revision="a" * 40, + document_generations={"doc:one": "gen:" + "1" * 32}, + workspace_id="psd-clinical", + language="italian", + searcher=MissingExpectedSearcher(), + embedder=_Embedder(), + expected_kinds={"evidence:icd": "enum", "evidence:fascia-pediatrica": "formula"}, + ) + + assert report.passed is False + assert all(query.hit_at_5 is False and query.hit_at_10 is False for query in report.queries) + assert all(query.empty_result is False for query in report.queries) + assert report.queries[0].missing_expected == ("evidence:icd",) + assert report.counts_by_expected_kind == {"enum": 1, "formula": 2} + + +def test_evaluation_fixture_requires_all_retrieval_profiles(tmp_path): + from tht.evidence.evaluation import EvaluationFixtureError, load_evaluation_fixture + + path = tmp_path / "evaluation.yaml" + path.write_text( + """ +schema_version: 1 +queries: + - id: lexical-code + query: Qual è il codice ICD-10? + profile: lexical + purpose: schema_linking + expected: [evidence:icd] + - id: semantic-age + query: Come distinguo i pazienti pediatrici? + profile: semantic + purpose: sql_generation + expected: [evidence:fascia-pediatrica] +""".strip(), + encoding="utf-8", + ) + + with pytest.raises(EvaluationFixtureError, match="mixed"): + load_evaluation_fixture(path) + + +def test_evaluation_fixture_rejects_duplicate_ids_empty_expectations_and_private_purposes(tmp_path): + from tht.evidence.evaluation import EvaluationFixtureError, load_evaluation_fixture + + path = tmp_path / "evaluation.yaml" + path.write_text( + """ +schema_version: 1 +queries: + - id: duplicate + query: a + profile: lexical + purpose: private + expected: [] + - id: duplicate + query: b + profile: semantic + purpose: sql_generation + expected: [evidence:b] + - id: mixed + query: c + profile: mixed + purpose: rewriting + expected: [evidence:c] +""".strip(), + encoding="utf-8", + ) + + with pytest.raises(EvaluationFixtureError) as failure: + load_evaluation_fixture(path) + + assert {"duplicate", "expected", "purpose"} <= set(str(failure.value).split()) diff --git a/harness/tests/test_evidence_facade_contract.py b/harness/tests/test_evidence_facade_contract.py index 1cb56824..a4cd5cce 100644 --- a/harness/tests/test_evidence_facade_contract.py +++ b/harness/tests/test_evidence_facade_contract.py @@ -156,6 +156,7 @@ def test_preprocessing_factory_forwards_only_evidence_pipeline_dependencies(monk "retain_published_generations": 2, "workspace_id": None, "sparse_language": "italian", + "candidate_evaluator": None, } pipeline = build_preprocessing_pipeline(**dependencies) diff --git a/harness/tests/test_preprocess_cli.py b/harness/tests/test_preprocess_cli.py index ea75b0d2..473eeb87 100644 --- a/harness/tests/test_preprocess_cli.py +++ b/harness/tests/test_preprocess_cli.py @@ -237,12 +237,14 @@ def test_run_from_config_uses_runtime_identity_workspace_id(monkeypatch, tmp_pat pipeline_version, retain_published_generations, sparse_language, + candidate_evaluator, ): calls["init"] = { "embedding_model": embedding_model, "embedding_dimensions": embedding_dimensions, "pipeline_version": pipeline_version, "sparse_language": sparse_language, + "candidate_evaluator": candidate_evaluator, } def run_as_job(self, **kwargs): @@ -260,6 +262,7 @@ def test_run_from_config_uses_runtime_identity_workspace_id(monkeypatch, tmp_pat command.run_from_config(config) assert calls["init"]["sparse_language"] == "english" + assert callable(calls["init"]["candidate_evaluator"]) assert calls["run_as_job"]["workspace_id"] == "psd-clinical" assert calls["run_as_job"]["input_fingerprint"] != calls["run_as_job"]["config_fingerprint"] diff --git a/harness/tests/test_qdrant_vector_store.py b/harness/tests/test_qdrant_vector_store.py index ec17bf2e..e8546cf2 100644 --- a/harness/tests/test_qdrant_vector_store.py +++ b/harness/tests/test_qdrant_vector_store.py @@ -475,6 +475,35 @@ def test_evidence_search_uses_filtered_dense_and_bm25_prefetches_with_default_rr ] +@pytest.mark.parametrize("retrieval_mode, expected_query", [ + ("dense", None), + ("bm25", {"text": "cardiomiopatia", "model": "qdrant/bm25", "options": {"language": "italian"}}), +]) +def test_evidence_diagnostic_branch_searches_use_the_runtime_filter(retrieval_mode, expected_query): + fake = FakeQdrantHttp() + _ready_collection_with_bm25(fake) + store = _store(fake) + generation = "gen:" + "1" * 32 + + store.search( + ["evidence"], [0.2] * 1024, limit=10, kinds=["evidence"], + query_text="cardiomiopatia", query_language="italian", retrieval_mode=retrieval_mode, + metadata_filter={"workspace_id": "demo", "vector_generation": generation, "document_ids": ["doc:abc"]}, + ) + + query = next(call[2] for call in reversed(fake.calls) if call[1].endswith("/points/query")) + assert query["limit"] == 10 + assert query["filter"]["must"][-2:] == [ + {"key": "vector_generation", "match": {"value": generation}}, + {"key": "document_id", "match": {"any": ["doc:abc"]}}, + ] + if retrieval_mode == "dense": + assert query["vector"] == [0.2] * 1024 + else: + assert query["query"] == expected_query + assert query["using"] == "bm25" + + def test_search_filters_by_workspace_and_allowed_record_kinds(): fake = FakeQdrantHttp() store = _store(fake) diff --git a/harness/tht/adapters/vector/qdrant.py b/harness/tht/adapters/vector/qdrant.py index b79cfc87..ca0953a2 100644 --- a/harness/tht/adapters/vector/qdrant.py +++ b/harness/tht/adapters/vector/qdrant.py @@ -132,6 +132,7 @@ class QdrantVectorStore: metadata_filter: dict[str, object] | None = None, query_text: str | None = None, query_language: str | None = None, + retrieval_mode: str = "fused", ) -> list[VectorHit]: require_positive_limit(limit) self._validate_embedding(embedding, query=True) @@ -185,7 +186,40 @@ class QdrantVectorStore: 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 retrieval_mode not in {"fused", "dense", "bm25"}: + raise VectorStoreError("Evidence retrieval mode is invalid") + if retrieval_mode == "dense": + if allowed_record_kinds != ["evidence"]: + raise VectorStoreError("Evidence branch diagnostics are only available for Evidence") + response = self._call( + "POST", + f"/collections/{self._collection}/points/query", + { + "vector": embedding, + "limit": limit, + "with_payload": True, + "filter": {"must": filter_must}, + }, + ) + elif retrieval_mode == "bm25": + if allowed_record_kinds != ["evidence"]: + raise VectorStoreError("Evidence branch diagnostics are only available for Evidence") + if query_text is None or query_text.strip() == "" or query_language not in _BM25_LANGUAGES: + raise VectorStoreError("Evidence BM25 query is invalid") + self._ensure_collection(strict=False, require_bm25=True) + shared_filter = {"must": filter_must} + response = self._call( + "POST", + f"/collections/{self._collection}/points/query", + { + "query": self._bm25_document(query_text, query_language), + "using": "bm25", + "limit": limit, + "with_payload": True, + "filter": shared_filter, + }, + ) + elif query_text is None: if allowed_record_kinds == ["evidence"]: raise VectorStoreError("Evidence hybrid query text is required") response = self._call( diff --git a/harness/tht/cli/evidence_cmd.py b/harness/tht/cli/evidence_cmd.py index 2994a141..cc4b6184 100644 --- a/harness/tht/cli/evidence_cmd.py +++ b/harness/tht/cli/evidence_cmd.py @@ -10,6 +10,7 @@ from typing import Annotated import typer +from tht.cli.config_cmd import CONFIG_OPT from tht.evidence import ( EvidencePreparationError, PiEvidenceRestructurer, @@ -63,6 +64,48 @@ def _findings_payload(findings) -> list[dict[str, object]]: return [finding.__dict__ for finding in findings] +def evaluate_from_config( + workspace_root: Path, + config: Path, + *, + generation: str | None = None, +) -> dict[str, object]: + """Evaluate one stored Evidence generation without changing corpus or vectors.""" + from tht.adapters.factory import build_vector_store + from tht.cli.schema_cmd import _load_config_or_exit + from tht.cli.vector_cmd import make_embedder + from tht.evidence.canonical import load_curated_tree + from tht.evidence.corpus.store import CorpusStore + from tht.evidence.evaluation import evaluate_retrieval, load_evaluation_fixture + + cfg = _load_config_or_exit(config) + store = CorpusStore(cfg.paths.artifacts.parent / "corpus") + manifest = store.manifest(generation) if generation is not None else store.active_manifest() + if manifest is None: + raise RuntimeError("active Evidence generation is unavailable") + document_generations = manifest.metadata.get("document_generations") + if not isinstance(document_generations, dict): + raise TypeError("Evidence generation is invalid") + language = {"en": "english", "it": "italian"}.get(cfg.language) + if language is None: + raise RuntimeError("workspace language is unsupported for Qdrant BM25") + report = evaluate_retrieval( + load_evaluation_fixture(workspace_root / "evidence" / "evaluation.yaml"), + workspace_revision=cfg._workspace_revision, + vector_generation=manifest.vector_generation, + document_generations=document_generations, + workspace_id=cfg._workspace_id, + language=language, + searcher=build_vector_store(cfg), + embedder=make_embedder(cfg.embeddings), + expected_kinds={ + evidence.id: evidence.kind + for evidence in load_curated_tree(workspace_root / "evidence" / "curated") + }, + ) + return report.model_dump() + + @evidence_app.command("prepare") def prepare_cmd( workspace_root: Path, @@ -106,6 +149,36 @@ def validate_cmd( raise typer.Exit(code=1) +@evidence_app.command("evaluate") +def evaluate_cmd( + workspace_root: Path, + config: Path = CONFIG_OPT, + generation: str | None = typer.Option(None, "--generation"), + json_output: bool = typer.Option(False, "--json"), +) -> None: + """Evaluate active or selected Evidence retrieval generation without publishing it.""" + root = _canonical_worktree(workspace_root) + try: + payload = evaluate_from_config(root, config, generation=generation) + except Exception: # noqa: BLE001 - CLI reports a safe operational failure. + _emit({ + "schemaVersion": 1, + "operation": "evidence_evaluate", + "status": "failed", + "code": "evaluation_failed", + }, json_output) + raise typer.Exit(code=1) from None + payload = { + "schemaVersion": 1, + "operation": "evidence_evaluate", + "status": "passed" if payload["passed"] else "failed", + **payload, + } + _emit(payload, json_output) + if not payload["passed"]: + raise typer.Exit(code=1) + + @evidence_app.command("resolve") def resolve_cmd( workspace_root: Path, diff --git a/harness/tht/cli/preprocess_cmd.py b/harness/tht/cli/preprocess_cmd.py index ab3a64b4..a92c8aea 100644 --- a/harness/tht/cli/preprocess_cmd.py +++ b/harness/tht/cli/preprocess_cmd.py @@ -28,6 +28,53 @@ def _evidence_json_context(config: Path): return _load_config_or_exit(config) +def _evaluation_workspace_root(cfg) -> Path: + evidence = cfg.evidence + if evidence is None: + raise RuntimeError("Evidence evaluation fixture is unavailable") + if evidence.source_root is not None: + return evidence.source_root + filesystem_roots = [source.root for source in evidence.sources if source.type == "filesystem"] + if len(filesystem_roots) != 1: + raise RuntimeError("Evidence evaluation requires one filesystem workspace source") + root = filesystem_roots[0] + if root.name == "curated": + root = root.parent + if root.name == "evidence": + return root.parent + return root + + +def _candidate_evaluator(cfg, *, vector_store, embedder): + """Bind candidate publication to the same read-only retrieval evaluator as the CLI.""" + from tht.evidence.canonical import load_curated_tree + from tht.evidence.evaluation import evaluate_retrieval, load_evaluation_fixture + + workspace_root = _evaluation_workspace_root(cfg) + language = _bm25_language(cfg.language) + + def evaluate(manifest): + document_generations = manifest.metadata.get("document_generations") + if not isinstance(document_generations, dict): + raise TypeError("candidate Evidence generation is invalid") + return evaluate_retrieval( + load_evaluation_fixture(workspace_root / "evidence" / "evaluation.yaml"), + workspace_revision=cfg._workspace_revision, + vector_generation=manifest.vector_generation, + document_generations=document_generations, + workspace_id=cfg._workspace_id, + language=language, + searcher=vector_store, + embedder=embedder, + expected_kinds={ + evidence.id: evidence.kind + for evidence in load_curated_tree(workspace_root / "evidence" / "curated") + }, + ) + + return evaluate + + def _evidence_json_payload(cfg, payload: dict, *, code: str, error: str | None = None) -> dict: value = { **payload, @@ -112,15 +159,18 @@ def run_from_config(config: Path, *, dry_run: bool = False, resume: str | None = if cfg.embeddings is None: raise RuntimeError("embeddings are not configured") corpus_root = cfg.paths.artifacts.parent / "corpus" + vector_store = build_vector_store(cfg, require_write=True) + embedder = make_embedder(cfg.embeddings) pipeline = build_preprocessing_pipeline( store=CorpusStore(corpus_root), sources=build_sources(cfg.evidence), - embedder=make_embedder(cfg.embeddings), - vector_store=build_vector_store(cfg, require_write=True), + embedder=embedder, + vector_store=vector_store, embedding_model=cfg.embeddings.model, embedding_dimensions=cfg.embeddings.dim, chunk_policy=ChunkPolicy(version="chunk-v1", max_chars=cfg.vector.max_chunk_chars), pipeline_version="evidence-v1", retain_published_generations=cfg.vector.retain_published_generations, sparse_language=_bm25_language(cfg.language), + candidate_evaluator=_candidate_evaluator(cfg, vector_store=vector_store, embedder=embedder), ) def fingerprint(value: str) -> str: return "sha256:" + hashlib.sha256(value.encode()).hexdigest() diff --git a/harness/tht/evidence/corpus/pipeline.py b/harness/tht/evidence/corpus/pipeline.py index 319d4020..0b5973f2 100644 --- a/harness/tht/evidence/corpus/pipeline.py +++ b/harness/tht/evidence/corpus/pipeline.py @@ -7,7 +7,7 @@ import json import logging import re import uuid -from collections.abc import Mapping, Sequence +from collections.abc import Callable, Mapping, Sequence from dataclasses import asdict, dataclass, field from datetime import UTC from pathlib import Path @@ -127,6 +127,7 @@ class CorpusPipeline: vector_store: VectorStore, embedding_model: str, embedding_dimensions: int, chunk_policy: ChunkPolicy, pipeline_version: str, retain_published_generations: int = 3, workspace_id: str | None = None, sparse_language: str = "italian", + candidate_evaluator: Callable[[CorpusManifest], object] | None = None, ) -> None: self.store = store self.sources = sources @@ -143,6 +144,14 @@ class CorpusPipeline: if sparse_language not in {"english", "italian"}: raise ValueError("unsupported Qdrant BM25 language") self.sparse_language = sparse_language + self.candidate_evaluator = candidate_evaluator + + def _evaluate_candidate(self, manifest: CorpusManifest) -> None: + if self.candidate_evaluator is None: + return + report = self.candidate_evaluator(manifest) + if getattr(report, "passed", False) is not True: + raise PipelineError("candidate retrieval evaluation failed") def _assert_workspace_binding(self) -> None: manifest = self.store.active_manifest() @@ -658,6 +667,7 @@ class CorpusPipeline: raise generation = read(context, "plan.json")["generation"] try: + self._evaluate_candidate(CorpusManifest.model_validate(read(context, "manifest.json"))) self.store.publish(generation) except Exception: compensate(context) @@ -791,6 +801,7 @@ class CorpusPipeline: manifest, {document.document_id: document.content for document in documents}, generation=generation, ) + self._evaluate_candidate(manifest) self.store.publish(staged) self.gc(workspace_root=self.store.root.parent) except AtomicContentTooLargeError as error: diff --git a/harness/tht/evidence/evaluation.py b/harness/tht/evidence/evaluation.py new file mode 100644 index 00000000..749309bf --- /dev/null +++ b/harness/tht/evidence/evaluation.py @@ -0,0 +1,303 @@ +"""Read-only retrieval evaluation for published and candidate Evidence generations.""" + +from __future__ import annotations + +from dataclasses import dataclass +from pathlib import Path + +import yaml + +from tht.evidence.canonical import EVIDENCE_PURPOSES +from tht.evidence.search import EvidenceSearchContext, render_evidence_query + +_PROFILES = frozenset({"lexical", "semantic", "mixed"}) + + +class EvaluationFixtureError(ValueError): + """The versioned retrieval fixture is not safe to use as a publication gate.""" + + +class EvaluationError(RuntimeError): + """The configured Evidence generation could not be evaluated safely.""" + + +@dataclass(frozen=True) +class EvaluationQuery: + query_id: str + query: str + profile: str + purpose: str + expected: tuple[str, ...] + + +@dataclass(frozen=True) +class EvaluationFixture: + queries: tuple[EvaluationQuery, ...] + + +@dataclass(frozen=True) +class ExpectedEvidenceReport: + evidence_id: str + kind: str | None + dense_rank: int | None + bm25_rank: int | None + fused_rank: int | None + + +@dataclass(frozen=True) +class EvaluationQueryReport: + query_id: str + profile: str + purpose: str + hit_at_5: bool + hit_at_10: bool + missing_expected: tuple[str, ...] + empty_result: bool + expected: tuple[ExpectedEvidenceReport, ...] + + +@dataclass(frozen=True) +class EvaluationReport: + workspace_revision: str + vector_generation: str | None + rrf: dict[str, object] + passed: bool + queries: tuple[EvaluationQueryReport, ...] + counts_by_expected_kind: dict[str, int] + + def model_dump(self) -> dict[str, object]: + return { + "workspaceRevision": self.workspace_revision, + "vectorGeneration": self.vector_generation, + "rrf": self.rrf, + "passed": self.passed, + "countsByExpectedKind": self.counts_by_expected_kind, + "queries": [ + { + "id": query.query_id, + "profile": query.profile, + "purpose": query.purpose, + "hitAt5": query.hit_at_5, + "hitAt10": query.hit_at_10, + "missingExpected": list(query.missing_expected), + "emptyResult": query.empty_result, + "expected": [ + { + "evidenceId": expected.evidence_id, + "kind": expected.kind, + "denseRank": expected.dense_rank, + "bm25Rank": expected.bm25_rank, + "fusedRank": expected.fused_rank, + } + for expected in query.expected + ], + } + for query in self.queries + ], + } + + +def load_evaluation_fixture(path: Path) -> EvaluationFixture: + """Load the deliberately small, complete v1 evaluation fixture.""" + try: + raw = yaml.safe_load(path.read_text(encoding="utf-8")) + except (OSError, UnicodeError, yaml.YAMLError) as error: + raise EvaluationFixtureError("evaluation fixture is unreadable") from error + if not isinstance(raw, dict) or set(raw) != {"schema_version", "queries"}: + raise EvaluationFixtureError("evaluation fixture schema is invalid") + if raw.get("schema_version") != 1 or not isinstance(raw.get("queries"), list): + raise EvaluationFixtureError("evaluation fixture schema is invalid") + + queries: list[EvaluationQuery] = [] + errors: set[str] = set() + ids: set[str] = set() + profiles: set[str] = set() + for entry in raw["queries"]: + entry_errors: set[str] = set() + if not isinstance(entry, dict) or set(entry) != {"id", "query", "profile", "purpose", "expected"}: + errors.add("schema") + continue + query_id = entry["id"] + query = entry["query"] + profile = entry["profile"] + purpose = entry["purpose"] + expected = entry["expected"] + if not isinstance(query_id, str) or not query_id.strip() or query_id in ids: + entry_errors.add("duplicate") + else: + ids.add(query_id) + if not isinstance(query, str) or not query.strip(): + entry_errors.add("query") + if profile not in _PROFILES: + entry_errors.add("profile") + else: + profiles.add(profile) + if purpose not in EVIDENCE_PURPOSES: + entry_errors.add("purpose") + if ( + not isinstance(expected, list) + or not expected + or any(not isinstance(value, str) or not value.strip() for value in expected) + ): + entry_errors.add("expected") + errors.update(entry_errors) + if not entry_errors: + queries.append(EvaluationQuery(query_id, query, profile, purpose, tuple(expected))) + missing_profiles = _PROFILES - profiles + if missing_profiles: + errors.update(missing_profiles) + if errors: + raise EvaluationFixtureError(" ".join(sorted(errors))) + return EvaluationFixture(tuple(queries)) + + +def _ranked_evidence(hits) -> tuple[dict[str, int], dict[str, str]]: + ranks: dict[str, int] = {} + kinds: dict[str, str] = {} + for hit in hits: + metadata = getattr(hit, "metadata", None) + if not isinstance(metadata, dict): + raise EvaluationError("evaluation search returned malformed Evidence payload") + evidence_id = metadata.get("evidence_id") + kind = metadata.get("evidence_kind") + if not isinstance(evidence_id, str) or not evidence_id or not isinstance(kind, str) or not kind: + raise EvaluationError("evaluation search returned malformed Evidence payload") + if evidence_id not in ranks: + ranks[evidence_id] = len(ranks) + 1 + kinds[evidence_id] = kind + return ranks, kinds + + +def _search_generation( + searcher, + embedding: list[float], + *, + rendered_query: str, + purpose: str, + workspace_id: str, + generation: str, + document_ids: list[str], + language: str, + retrieval_mode: str, +): + return searcher.search( + ["evidence"], + embedding, + limit=10, + kinds=["evidence"], + query_text=rendered_query, + query_language=language, + retrieval_mode=retrieval_mode, + metadata_filter={ + "workspace_id": workspace_id, + "vector_generation": generation, + "document_ids": document_ids, + "purpose": purpose, + "required_kinds": [], + "required_concepts": [], + "required_tables": [], + "required_columns": [], + }, + ) + + +def evaluate_retrieval( + fixture: EvaluationFixture, + *, + workspace_revision: str, + document_generations: dict[str, str], + workspace_id: str, + language: str, + searcher, + embedder, + vector_generation: str | None = None, + expected_kinds: dict[str, str] | None = None, +) -> EvaluationReport: + """Evaluate a generation with the runtime hybrid request plus branch diagnostics.""" + if not document_generations: + raise EvaluationError("evaluation requires indexed Evidence documents") + by_generation: dict[str, list[str]] = {} + for document_id, generation in document_generations.items(): + if not isinstance(document_id, str) or not isinstance(generation, str) or not generation: + raise EvaluationError("evaluation document generations are invalid") + by_generation.setdefault(generation, []).append(document_id) + for document_ids in by_generation.values(): + document_ids.sort() + + reports: list[EvaluationQueryReport] = [] + expected_kinds = expected_kinds or {} + for query in fixture.queries: + rendered = render_evidence_query(query.query, EvidenceSearchContext()) + embedding = embedder.embed_query(rendered) + branch_hits = {"dense": [], "bm25": [], "fused": []} + for generation, document_ids in sorted(by_generation.items()): + for mode, hits in branch_hits.items(): + hits.extend(_search_generation( + searcher, + embedding, + rendered_query=rendered, + purpose=query.purpose, + workspace_id=workspace_id, + generation=generation, + document_ids=document_ids, + language=language, + retrieval_mode=mode, + )) + ranks_by_branch: dict[str, dict[str, int]] = {} + kinds_by_branch: dict[str, dict[str, str]] = {} + for mode, hits in branch_hits.items(): + ordered = sorted(hits, key=lambda hit: (-float(hit.similarity), str(hit.id))) + ranks_by_branch[mode], kinds_by_branch[mode] = _ranked_evidence(ordered) + expected = [] + for evidence_id in query.expected: + kind = expected_kinds.get(evidence_id) or next(( + kinds_by_branch[mode][evidence_id] + for mode in ("fused", "dense", "bm25") + if evidence_id in kinds_by_branch[mode] + ), None) + expected.append(ExpectedEvidenceReport( + evidence_id=evidence_id, + kind=kind, + dense_rank=ranks_by_branch["dense"].get(evidence_id), + bm25_rank=ranks_by_branch["bm25"].get(evidence_id), + fused_rank=ranks_by_branch["fused"].get(evidence_id), + )) + fused_ranks = ranks_by_branch["fused"] + missing = tuple(item.evidence_id for item in expected if item.fused_rank is None) + reports.append(EvaluationQueryReport( + query_id=query.query_id, + profile=query.profile, + purpose=query.purpose, + hit_at_5=any(item.fused_rank is not None and item.fused_rank <= 5 for item in expected), + hit_at_10=any(item.fused_rank is not None and item.fused_rank <= 10 for item in expected), + missing_expected=missing, + empty_result=not fused_ranks, + expected=tuple(expected), + )) + counts: dict[str, int] = {} + for query in reports: + for expected in query.expected: + if expected.kind is not None: + counts[expected.kind] = counts.get(expected.kind, 0) + 1 + evaluated = vector_generation or (next(iter(by_generation)) if len(by_generation) == 1 else None) + return EvaluationReport( + workspace_revision=workspace_revision, + vector_generation=evaluated, + rrf={"algorithm": "rrf", "k": 60, "prefetch_limit_multiplier": 2}, + passed=all(query.hit_at_10 for query in reports), + queries=tuple(reports), + counts_by_expected_kind=dict(sorted(counts.items())), + ) + + +__all__ = [ + "EvaluationError", + "EvaluationFixture", + "EvaluationFixtureError", + "EvaluationQuery", + "EvaluationQueryReport", + "EvaluationReport", + "ExpectedEvidenceReport", + "evaluate_retrieval", + "load_evaluation_fixture", +] diff --git a/harness/tht/evidence/preprocessing.py b/harness/tht/evidence/preprocessing.py index be2f269b..aff85696 100644 --- a/harness/tht/evidence/preprocessing.py +++ b/harness/tht/evidence/preprocessing.py @@ -1,5 +1,6 @@ """Explicit construction boundary for Evidence preprocessing.""" +from collections.abc import Callable from typing import Protocol from tht.evidence.contracts import EvidenceSource @@ -26,6 +27,7 @@ def build_preprocessing_pipeline( retain_published_generations: int = 3, workspace_id: str | None = None, sparse_language: str = "italian", + candidate_evaluator: Callable[[object], object] | None = None, ) -> CorpusPipeline: """Construct preprocessing from the bounded infrastructure supplied by core.""" return CorpusPipeline( @@ -40,6 +42,7 @@ def build_preprocessing_pipeline( retain_published_generations=retain_published_generations, workspace_id=workspace_id, sparse_language=sparse_language, + candidate_evaluator=candidate_evaluator, ) diff --git a/harness/tht/ports/vector.py b/harness/tht/ports/vector.py index 0c4bfa30..6bcae8dc 100644 --- a/harness/tht/ports/vector.py +++ b/harness/tht/ports/vector.py @@ -79,6 +79,7 @@ class VectorStore(Protocol): metadata_filter: dict[str, object] | None = None, query_text: str | None = None, query_language: str | None = None, + retrieval_mode: str = "fused", ) -> list[VectorHit]: ... def existing_hashes(self, collection: str, kinds: list[str]) -> dict[str, str]: ...