From 3420c57c8b548e589dfb959123febc53575593ce Mon Sep 17 00:00:00 2001 From: mptyl Date: Mon, 24 Aug 2026 21:33:48 +0200 Subject: [PATCH] feat(evidence): build semantic fragments from typed units --- harness/tests/test_corpus_chunk.py | 170 +++++++++++++++++++++++ harness/tests/test_corpus_models.py | 28 ++++ harness/tests/test_corpus_normalize.py | 53 +++++++ harness/tests/test_corpus_pipeline.py | 126 +++++++++++++++++ harness/tht/evidence/canonical.py | 10 +- harness/tht/evidence/corpus/chunk.py | 163 +++++++++++++++++++++- harness/tht/evidence/corpus/models.py | 37 +++++ harness/tht/evidence/corpus/normalize.py | 39 +++++- harness/tht/evidence/corpus/pipeline.py | 30 +++- 9 files changed, 642 insertions(+), 14 deletions(-) diff --git a/harness/tests/test_corpus_chunk.py b/harness/tests/test_corpus_chunk.py index 804b95d2..f460d87b 100644 --- a/harness/tests/test_corpus_chunk.py +++ b/harness/tests/test_corpus_chunk.py @@ -2,6 +2,7 @@ import hashlib import pytest +from tht.evidence.canonical import CuratedEvidence from tht.evidence.corpus.chunk import ChunkPolicy, chunk from tht.evidence.corpus.models import CanonicalDocument, CorpusManifest @@ -33,6 +34,50 @@ def other_document(content: str) -> CanonicalDocument: ) +def curated_formula_document(*, sql: str = "CASE WHEN age < 18 THEN 'pediatric' END") -> CanonicalDocument: + evidence = CuratedEvidence.model_validate( + { + "schema_version": 1, + "id": "evidence:fascia-pediatrica", + "title": "Fascia pediatrica", + "kind": "formula", + "purposes": ["sql_generation", "schema_linking"], + "applies_to": { + "concepts": ["fascia pediatrica"], + "tables": ["clinical.patient"], + "columns": ["clinical.patient.birth_date"], + }, + "language": "it", + "provenance": { + "source_file": "source/paziente.md", + "source_sha256": "sha256:" + "a" * 64, + "supporting_excerpts": ["I pazienti pediatrici hanno età inferiore a 18 anni."], + }, + "review_items": [], + "payload": { + "concept": "fascia pediatrica", + "columns": ["clinical.patient.birth_date"], + "sql": sql, + }, + } + ) + return curated_document(evidence) + + +def curated_document(evidence: CuratedEvidence) -> CanonicalDocument: + content = "canonical curated Evidence" + return CanonicalDocument( + document_id="doc:" + evidence.id.removeprefix("evidence:"), + source_id="curated:" + evidence.id.removeprefix("evidence:"), + source_uri=f"file:///safe/curated/{evidence.kind}/{evidence.id.removeprefix('evidence:')}.md", + source_fingerprint="sha256:" + "b" * 64, + content_hash="sha256:" + hashlib.sha256(content.encode()).hexdigest(), + title=evidence.title, + content=content, + media_type="text/markdown", + pipeline_version="pipe:v1", + metadata={"curated_evidence": evidence.model_dump(mode="json")}, + ) def test_chunk_ids_are_stable_for_same_content_and_repeat_runs(): policy = ChunkPolicy(version="paragraph:v1", max_chars=8) first = chunk(document("A\n\nB"), policy) @@ -116,3 +161,128 @@ def test_empty_document_has_no_chunks_and_invalid_policy_is_rejected(): assert chunk(document(""), ChunkPolicy(version="v1", max_chars=4)) == [] with pytest.raises(ValueError): ChunkPolicy(version="v1", max_chars=0) + + +def test_typed_formula_is_rendered_as_one_traceable_semantic_fragment(): + fragments = chunk(curated_formula_document(), ChunkPolicy(version="semantic:v1", max_chars=4000)) + + assert len(fragments) == 1 + fragment = fragments[0] + assert "Formula: Fascia pediatrica" in fragment.content + assert "Scopi: sql_generation, schema_linking" in fragment.content + assert "Concetto: fascia pediatrica" in fragment.content + assert "Colonne: clinical.patient.birth_date" in fragment.content + assert "SQL: CASE WHEN age < 18 THEN 'pediatric' END" in fragment.content + assert "Provenienza: source/paziente.md" in fragment.content + assert fragment.metadata["evidence_id"] == "evidence:fascia-pediatrica" + assert fragment.metadata["evidence_kind"] == "formula" + assert fragment.metadata["purposes"] == ["sql_generation", "schema_linking"] + assert fragment.metadata["scope"] == { + "concepts": ["fascia pediatrica"], + "tables": ["clinical.patient"], + "columns": ["clinical.patient.birth_date"], + } + assert fragment.metadata["language"] == "it" + assert fragment.metadata["provenance"]["source_file"] == "source/paziente.md" + + +def test_oversized_typed_atomic_content_fails_instead_of_being_split(): + document = curated_formula_document(sql="CASE WHEN age < 18 THEN " + "x" * 200 + " END") + + with pytest.raises(ValueError, match="atomic_content_too_large") as caught: + chunk(document, ChunkPolicy(version="semantic:v1", max_chars=120)) + + assert caught.value.review_item.code == "atomic_content_too_large" + assert caught.value.review_item.field == "formula.sql" + + +def test_enum_value_meaning_pairs_are_atomic_and_fragment_ids_are_deterministic(): + payload = curated_formula_document().metadata["curated_evidence"] + evidence = CuratedEvidence.model_validate({ + **payload, + "id": "evidence:stato-ricovero", + "title": "Stato ricovero", + "kind": "enum", + "payload": { + "column": "clinical.admission.status", + "values": {"A": "Attivo", "D": "Dimesso"}, + }, + }) + document = curated_document(evidence) + + first = chunk(document, ChunkPolicy(version="semantic:v1", max_chars=4000)) + second = chunk(document, ChunkPolicy(version="semantic:v1", max_chars=4000)) + + assert first == second + assert [fragment.ordinal for fragment in first] == [0, 1] + assert "Valore: A\nSignificato: Attivo" in first[0].content + assert "Valore: D\nSignificato: Dimesso" in first[1].content + + +def test_italian_evidence_uses_an_italian_kind_heading(): + base = curated_formula_document().metadata["curated_evidence"] + evidence = CuratedEvidence.model_validate({ + **base, + "id": "evidence:regola-ricovero", + "title": "Regola ricovero", + "kind": "domain", + "payload": {"rule": "Il ricovero richiede una data di ammissione."}, + }) + + fragment = chunk(curated_document(evidence), ChunkPolicy(version="semantic:v1", max_chars=4000))[0] + + assert fragment.content.startswith("Dominio: Regola ricovero\n") + + +def test_domain_rule_remains_atomic_across_blank_paragraphs(): + base = curated_formula_document().metadata["curated_evidence"] + evidence = CuratedEvidence.model_validate({ + **base, + "id": "evidence:regole-ricovero", + "title": "Regole ricovero", + "kind": "domain", + "payload": {"rule": "La data di ammissione è obbligatoria.\n\nLa data di dimissione segue l'ammissione."}, + }) + + fragments = chunk(curated_document(evidence), ChunkPolicy(version="semantic:v1", max_chars=4000)) + + assert len(fragments) == 1 + assert "Regola: La data di ammissione è obbligatoria.\n\nLa data di dimissione segue l'ammissione." in fragments[0].content + + +@pytest.mark.parametrize( + ("kind", "payload", "field"), + [ + ("domain", {"rule": "x" * 200}, "domain.rule"), + ( + "mapping", + { + "concept": "ricovero", + "tables": ["clinical.admission"], + "columns": ["clinical.admission.status"], + }, + "mapping", + ), + ( + "reference", + {"url": "https://example.test/guide", "label": "Guida", "description": "x" * 200}, + "reference.url", + ), + ("normalization", {"input": "a", "output": "b", "rule": "x" * 200}, "normalization.rule"), + ("glossary", {"definition": "x" * 200}, "glossary.definition"), + ("example", {"question": "x" * 200, "interpretation": "attesa"}, "example"), + ], +) +def test_other_typed_atomic_content_blocks_with_a_stable_review_item(kind, payload, field): + base = curated_formula_document().metadata["curated_evidence"] + evidence = CuratedEvidence.model_validate({ + **base, + "id": f"evidence:{kind}-test", + "kind": kind, + "payload": payload, + }) + + with pytest.raises(ValueError, match="atomic_content_too_large") as caught: + chunk(curated_document(evidence), ChunkPolicy(version="semantic:v1", max_chars=120)) + + assert caught.value.review_item.field == field diff --git a/harness/tests/test_corpus_models.py b/harness/tests/test_corpus_models.py index 3253f882..97b38317 100644 --- a/harness/tests/test_corpus_models.py +++ b/harness/tests/test_corpus_models.py @@ -148,6 +148,34 @@ def test_manifest_rejects_inconsistent_pipeline_versions(): CorpusManifest(pipeline_version="evidence-v1", documents=[wrong]) +def test_typed_evidence_fragment_metadata_must_be_complete_and_allowlisted(): + metadata = { + "evidence_id": "evidence:fascia-pediatrica", + "evidence_kind": "formula", + "purposes": ["sql_generation"], + "scope": {"concepts": ["fascia pediatrica"], "tables": [], "columns": []}, + "language": "it", + "provenance": { + "source_file": "source/paziente.md", + "source_sha256": "sha256:" + "a" * 64, + "supporting_excerpts": ["Pazienti con età inferiore a 18 anni."], + }, + } + assert CanonicalChunk.model_validate({**chunk().model_dump(), "metadata": metadata}).metadata == metadata + + with pytest.raises(ValidationError, match="typed Evidence metadata"): + CanonicalChunk.model_validate({ + **chunk().model_dump(), + "metadata": {**metadata, "unreviewed_evidence_note": "not allowed"}, + }) + + with pytest.raises(ValidationError, match="typed Evidence metadata"): + CanonicalChunk.model_validate({ + **chunk().model_dump(), + "metadata": {"evidence_kind": "formula"}, + }) + + def test_vector_generation_requires_embedding_compatibility(): with pytest.raises(ValidationError, match="vector_generation"): CorpusManifest(pipeline_version="evidence-v1", vector_generation="generation:one") diff --git a/harness/tests/test_corpus_normalize.py b/harness/tests/test_corpus_normalize.py index 96d27c53..9262e896 100644 --- a/harness/tests/test_corpus_normalize.py +++ b/harness/tests/test_corpus_normalize.py @@ -52,6 +52,59 @@ def test_frontmatter_can_end_at_eof_without_inventing_content(): assert document.content == "" +def test_normalize_preserves_validated_curated_evidence_for_semantic_projection(): + source = SourceObject( + source_id="filesystem:fascia-pediatrica", + uri="file:///safe/evidence/curated/formula/fascia-pediatrica.md", + fingerprint="sha256:" + "a" * 64, + metadata={"relative_path": "curated/formula/fascia-pediatrica.md"}, + ) + raw = ( + "---\n" + "schema_version: 1\n" + "id: evidence:fascia-pediatrica\n" + "title: Fascia pediatrica\n" + "kind: formula\n" + "purposes: [sql_generation]\n" + "applies_to:\n" + " columns: [clinical.patient.birth_date]\n" + "language: it\n" + "provenance:\n" + " source_file: source/paziente.md\n" + " source_sha256: sha256:" + "b" * 64 + "\n" + " supporting_excerpts: [Pazienti con età inferiore a 18 anni.]\n" + "review_items: []\n" + "formula:\n" + " concept: fascia pediatrica\n" + " columns: [clinical.patient.birth_date]\n" + " sql: CASE WHEN age < 18 THEN 'pediatric' END\n" + "---\n" + ).encode() + + document = normalize(AcquiredDocument(source=source, content=raw, media_type="text/markdown"), "pipe:v1") + + assert document.content == raw.decode() + assert document.title == "Fascia pediatrica" + assert document.metadata["curated_evidence"]["id"] == "evidence:fascia-pediatrica" + assert document.metadata["curated_evidence"]["payload"]["concept"] == "fascia pediatrica" + + +def test_only_curated_markdown_is_treated_as_canonical_evidence_without_relative_metadata(): + source = SourceObject( + source_id="filesystem:notes", + uri="file:///safe/evidence/curated/formula/notes.txt", + fingerprint="sha256:" + "a" * 64, + ) + + document = normalize( + AcquiredDocument(source=source, content=b"ordinary note", media_type="text/plain"), + "pipe:v1", + ) + + assert document.content == "ordinary note" + assert "curated_evidence" not in document.metadata + + @pytest.mark.parametrize( "frontmatter", [ diff --git a/harness/tests/test_corpus_pipeline.py b/harness/tests/test_corpus_pipeline.py index 101343a4..867b1407 100644 --- a/harness/tests/test_corpus_pipeline.py +++ b/harness/tests/test_corpus_pipeline.py @@ -2,6 +2,7 @@ from datetime import UTC, datetime, timedelta import pytest +from tht.evidence.canonical import CuratedEvidence, dump_curated_markdown from tht.evidence.contracts import AcquiredDocument, SourceObject from tht.evidence.corpus.chunk import ChunkPolicy from tht.evidence.corpus.models import CanonicalChunk, CanonicalDocument, CorpusManifest @@ -116,6 +117,131 @@ def pipeline(tmp_path, source, *, embedder=None, vectors=None, model="model-a", ) +def test_pipeline_embeds_validated_curated_evidence_as_semantic_fragments(tmp_path): + evidence = CuratedEvidence.model_validate( + { + "schema_version": 1, + "id": "evidence:fascia-pediatrica", + "title": "Fascia pediatrica", + "kind": "formula", + "purposes": ["sql_generation"], + "applies_to": {"columns": ["clinical.patient.birth_date"]}, + "language": "it", + "provenance": { + "source_file": "source/paziente.md", + "source_sha256": "sha256:" + "a" * 64, + "supporting_excerpts": ["Pazienti con età inferiore a 18 anni."], + }, + "review_items": [], + "payload": { + "concept": "fascia pediatrica", + "columns": ["clinical.patient.birth_date"], + "sql": "CASE WHEN age < 18 THEN 'pediatric' END", + }, + } + ) + source_item = SourceObject( + source_id="fs:curated-formula", + uri="file:///safe/curated/formula/fascia-pediatrica.md", + fingerprint="sha256:" + "b" * 64, + metadata={"relative_path": "curated/formula/fascia-pediatrica.md"}, + ) + embedder = Embedder() + vectors = Vectors() + + result = pipeline( + tmp_path, + Source([(source_item, dump_curated_markdown(evidence))]), + embedder=embedder, + vectors=vectors, + policy=ChunkPolicy(version="chunk-v1", max_chars=4000), + ).run() + + assert result.status == "succeeded" + assert len(result.manifest.chunks) == 1 + assert "Formula: Fascia pediatrica" in embedder.calls[0] + assert vectors.records[0].record.metadata["evidence_id"] == evidence.id + assert vectors.records[0].record.metadata["provenance"]["source_file"] == "source/paziente.md" + + +def test_pipeline_exposes_atomic_content_review_code_when_candidate_is_blocked(tmp_path): + evidence = CuratedEvidence.model_validate( + { + "schema_version": 1, + "id": "evidence:formula-lunga", + "title": "Formula lunga", + "kind": "formula", + "purposes": ["sql_generation"], + "language": "it", + "provenance": { + "source_file": "source/paziente.md", + "source_sha256": "sha256:" + "a" * 64, + "supporting_excerpts": ["Una formula molto lunga."], + }, + "review_items": [], + "payload": {"concept": "formula lunga", "columns": [], "sql": "x" * 200}, + } + ) + source_item = SourceObject( + source_id="fs:formula-lunga", + uri="file:///safe/curated/formula/formula-lunga.md", + fingerprint="sha256:" + "b" * 64, + metadata={"relative_path": "curated/formula/formula-lunga.md"}, + ) + + result = pipeline( + tmp_path, + Source([(source_item, dump_curated_markdown(evidence))]), + policy=ChunkPolicy(version="chunk-v1", max_chars=120), + ).run() + + assert result.status == "blocked" + assert result.published is False + assert [item.code for item in result.review_items] == ["atomic_content_too_large"] + assert result.review_items[0].field == "formula.sql" + + +def test_job_pipeline_persists_atomic_content_review_item_when_candidate_is_blocked(tmp_path): + evidence = CuratedEvidence.model_validate( + { + "schema_version": 1, + "id": "evidence:formula-lunga-job", + "title": "Formula lunga", + "kind": "formula", + "purposes": ["sql_generation"], + "language": "it", + "provenance": { + "source_file": "source/paziente.md", + "source_sha256": "sha256:" + "a" * 64, + "supporting_excerpts": ["Una formula molto lunga."], + }, + "review_items": [], + "payload": {"concept": "formula lunga", "columns": [], "sql": "x" * 200}, + } + ) + source_item = SourceObject( + source_id="fs:formula-lunga-job", + uri="file:///safe/curated/formula/formula-lunga-job.md", + fingerprint="sha256:" + "b" * 64, + metadata={"relative_path": "curated/formula/formula-lunga-job.md"}, + ) + + result = pipeline( + tmp_path, + Source([(source_item, dump_curated_markdown(evidence))]), + policy=ChunkPolicy(version="chunk-v1", max_chars=120), + ).run_as_job( + workspace_id="demo", + workspace_root=tmp_path, + config_fingerprint="sha256:" + "1" * 64, + input_fingerprint="sha256:" + "2" * 64, + ) + + assert result.status == "blocked" + assert result.published is False + assert [item.code for item in result.review_items] == ["atomic_content_too_large"] + + def test_pipeline_routes_source_io_through_evidence_facade(tmp_path, monkeypatch): import tht.evidence.acquisition as evidence_acquisition diff --git a/harness/tht/evidence/canonical.py b/harness/tht/evidence/canonical.py index 289f1e36..596f77a1 100644 --- a/harness/tht/evidence/canonical.py +++ b/harness/tht/evidence/canonical.py @@ -20,7 +20,7 @@ class StrictModel(BaseModel): model_config = ConfigDict(extra="forbid") -EvidenceKind = Literal[ +EVIDENCE_KINDS = ( "glossary", "domain", "enum", @@ -29,13 +29,15 @@ EvidenceKind = Literal[ "normalization", "formula", "reference", -] -EvidencePurpose = Literal[ +) +EVIDENCE_PURPOSES = ( "disambiguation", "rewriting", "schema_linking", "sql_generation", -] +) +EvidenceKind = Literal[*EVIDENCE_KINDS] +EvidencePurpose = Literal[*EVIDENCE_PURPOSES] _IDENTIFIER = r"[A-Za-z_][A-Za-z0-9_$]*" _TABLE_IDENTIFIER = re.compile(rf"^{_IDENTIFIER}\.{_IDENTIFIER}$") _COLUMN_IDENTIFIER = re.compile(rf"^{_IDENTIFIER}\.{_IDENTIFIER}\.{_IDENTIFIER}$") diff --git a/harness/tht/evidence/corpus/chunk.py b/harness/tht/evidence/corpus/chunk.py index 5c868295..ae486add 100644 --- a/harness/tht/evidence/corpus/chunk.py +++ b/harness/tht/evidence/corpus/chunk.py @@ -3,8 +3,10 @@ import hashlib import json import re +from collections.abc import Mapping from dataclasses import asdict, dataclass +from tht.evidence.canonical import CuratedEvidence, ReviewItem from tht.evidence.corpus.models import CanonicalChunk, CanonicalDocument @@ -20,6 +22,32 @@ class ChunkPolicy: raise ValueError("max_chars must be greater than zero") +class AtomicContentTooLargeError(ValueError): + """A semantic Evidence element exceeds the configured embedding boundary.""" + + code = "atomic_content_too_large" + + def __init__(self, evidence: CuratedEvidence) -> None: + super().__init__(self.code) + field = { + "formula": "formula.sql", + "enum": "enum.values", + "mapping": "mapping", + "normalization": "normalization.rule", + "glossary": "glossary.definition", + "domain": "domain.rule", + "example": "example", + "reference": "reference.url", + }[evidence.kind] + self.review_item = ReviewItem( + code=self.code, + message=( + f"Il contenuto atomico di {evidence.id} supera max_chunk_chars e non può essere diviso." + ), + field=field, + ) + + def _hash(text: str) -> str: return hashlib.sha256(text.encode("utf-8")).hexdigest() @@ -43,11 +71,143 @@ def _policy_fingerprint(policy: ChunkPolicy) -> str: return f"sha256:{_hash(serialized)}" +def _curated_evidence(document: CanonicalDocument) -> CuratedEvidence | None: + raw = document.metadata.get("curated_evidence") + if raw is None: + return None + if not isinstance(raw, Mapping): + raise TypeError("invalid curated evidence projection") + try: + return CuratedEvidence.model_validate(raw) + except ValueError as error: + raise ValueError("invalid curated evidence projection") from error + + +_ENGLISH_LABELS = { + "purpose": "Purpose", "concept": "Concept", "tables": "Tables", "columns": "Columns", + "column": "Column", "value": "Value", "meaning": "Meaning", "input": "Input", + "output": "Output", "rule": "Rule", "definition": "Definition", "synonyms": "Synonyms", + "variants": "Variants", "question": "Question", "interpretation": "Interpretation", + "label": "Label", "url": "URL", "description": "Description", "provenance": "Provenance", +} +_ITALIAN_LABELS = { + "purpose": "Scopi", "concept": "Concetto", "tables": "Tabelle", "columns": "Colonne", + "column": "Colonna", "value": "Valore", "meaning": "Significato", "input": "Input", + "output": "Output", "rule": "Regola", "definition": "Definizione", "synonyms": "Sinonimi", + "variants": "Varianti", "question": "Domanda", "interpretation": "Interpretazione", + "label": "Etichetta", "url": "URL", "description": "Descrizione", "provenance": "Provenienza", +} +_ITALIAN_KIND_LABELS = { + "glossary": "Glossario", "domain": "Dominio", "enum": "Enum", "example": "Esempio", + "mapping": "Mappatura", "normalization": "Normalizzazione", "formula": "Formula", + "reference": "Riferimento", +} + + +def _labels(evidence: CuratedEvidence) -> Mapping[str, str]: + return _ITALIAN_LABELS if evidence.language.lower().startswith("it") else _ENGLISH_LABELS + + +def _kind_label(evidence: CuratedEvidence) -> str: + if evidence.language.lower().startswith("it"): + return _ITALIAN_KIND_LABELS[evidence.kind] + return evidence.kind.title() + + +def _scope_lines(evidence: CuratedEvidence, labels: Mapping[str, str]) -> list[str]: + lines: list[str] = [] + if evidence.applies_to.concepts: + lines.append(f"{labels['concept']}: " + ", ".join(evidence.applies_to.concepts)) + if evidence.applies_to.tables: + lines.append(f"{labels['tables']}: " + ", ".join(evidence.applies_to.tables)) + if evidence.applies_to.columns: + lines.append(f"{labels['columns']}: " + ", ".join(evidence.applies_to.columns)) + return lines + + +def _curated_fragment_content(evidence: CuratedEvidence) -> list[str]: + """Render semantic atoms without treating structured values as arbitrary text.""" + label = _kind_label(evidence) + labels = _labels(evidence) + common = [ + f"{label}: {evidence.title}", + f"{labels['purpose']}: " + ", ".join(evidence.purposes), + *_scope_lines(evidence, labels), + ] + payload = evidence.payload + if evidence.kind == "formula": + body = [ + f"{labels['concept']}: {payload.concept}", + f"{labels['columns']}: " + ", ".join(payload.columns), + f"SQL: {payload.sql}", + ] + return ["\n".join([*common, *body, f"{labels['provenance']}: {evidence.provenance.source_file}"])] + if evidence.kind == "enum": + return [ + "\n".join([ + *common, + f"{labels['column']}: {payload.column}", + f"{labels['value']}: {value}", + f"{labels['meaning']}: {meaning}", + f"{labels['provenance']}: {evidence.provenance.source_file}", + ]) + for value, meaning in sorted(payload.values.items()) + ] + if evidence.kind == "mapping": + body = [ + f"{labels['concept']}: {payload.concept}", + f"{labels['tables']}: " + ", ".join(payload.tables), + f"{labels['columns']}: " + ", ".join(payload.columns), + ] + elif evidence.kind == "normalization": + body = [ + f"{labels['input']}: {payload.input}", + f"{labels['output']}: {payload.output}", + f"{labels['rule']}: {payload.rule}", + ] + elif evidence.kind == "glossary": + body = [ + f"{labels['definition']}: {payload.definition}", + *( [f"{labels['synonyms']}: " + ", ".join(payload.synonyms)] if payload.synonyms else []), + *( [f"{labels['variants']}: " + ", ".join(payload.variants)] if payload.variants else []), + ] + elif evidence.kind == "domain": + body = [f"{labels['rule']}: {payload.rule}"] + elif evidence.kind == "example": + body = [f"{labels['question']}: {payload.question}", f"{labels['interpretation']}: {payload.interpretation}"] + elif evidence.kind == "reference": + body = [f"{labels['label']}: {payload.label}", f"{labels['url']}: {payload.url}", f"{labels['description']}: {payload.description}"] + else: # pragma: no cover - CuratedEvidence validates the finite kind set. + raise ValueError("unsupported curated evidence kind") + return ["\n".join([*common, *body, f"{labels['provenance']}: {evidence.provenance.source_file}"])] + + +def _curated_metadata(evidence: CuratedEvidence) -> dict: + return { + "evidence_id": evidence.id, + "evidence_kind": evidence.kind, + "purposes": list(evidence.purposes), + "scope": evidence.applies_to.model_dump(mode="json"), + "language": evidence.language, + "provenance": evidence.provenance.model_dump(mode="json"), + } + + def chunk(document: CanonicalDocument, policy: ChunkPolicy) -> list[CanonicalChunk]: """Split canonical text with stable character-count boundaries and identifiers.""" chunks: list[CanonicalChunk] = [] policy_fingerprint = _policy_fingerprint(policy) - for ordinal, content in enumerate(_contents(document.content, policy.max_chars)): + curated = _curated_evidence(document) + contents = ( + _curated_fragment_content(curated) + if curated is not None + else _contents(document.content, policy.max_chars) + ) + if any(len(content) > policy.max_chars for content in contents): + if curated is not None: + raise AtomicContentTooLargeError(curated) + raise ValueError("chunk content exceeds max_chars") + for ordinal, content in enumerate(contents): chunk_hash = f"sha256:{_hash(content)}" identifier = _hash( ":".join( @@ -78,6 +238,7 @@ def chunk(document: CanonicalDocument, policy: ChunkPolicy) -> list[CanonicalChu "document": document.model_dump(mode="json")["metadata"], "source_fingerprint": document.source_fingerprint, "title": document.title, + **(_curated_metadata(curated) if curated is not None else {}), }, ) ) diff --git a/harness/tht/evidence/corpus/models.py b/harness/tht/evidence/corpus/models.py index fe7ba324..69779bc5 100644 --- a/harness/tht/evidence/corpus/models.py +++ b/harness/tht/evidence/corpus/models.py @@ -8,6 +8,7 @@ from typing import Self from pydantic import BaseModel, ConfigDict, Field, JsonValue, field_validator, model_validator +from tht.evidence.canonical import EVIDENCE_KINDS, EVIDENCE_PURPOSES from tht.evidence.contracts import ( canonical_provenance_uri, normalize_aware_datetime, @@ -17,6 +18,11 @@ from tht.evidence.contracts import ( _NAMESPACED_ID = re.compile(r"^[a-z][a-z0-9_-]*:[A-Za-z0-9._:-]+$") _SHA256 = re.compile(r"^sha256:[0-9a-f]{64}$") +_EVIDENCE_ID = re.compile(r"^evidence:[a-z0-9]+(?:-[a-z0-9]+)*$") +_EVIDENCE_METADATA_KEYS = frozenset({ + "evidence_id", "evidence_kind", "purposes", "scope", "language", "provenance", +}) +_CHUNK_METADATA_KEYS = frozenset({"chunk_policy", "document", "source_fingerprint", "title"}) def _validate_namespaced_id(value: str) -> str: @@ -37,6 +43,36 @@ def _require_content_hash(content: str, content_hash: str) -> None: raise ValueError("content_hash must match the exact canonical UTF-8 content") +def _validate_evidence_metadata(metadata: Mapping[str, JsonValue]) -> None: + if not (set(metadata) & _EVIDENCE_METADATA_KEYS): + return + keys = set(metadata) + missing = _EVIDENCE_METADATA_KEYS - keys + unknown = keys - _EVIDENCE_METADATA_KEYS - _CHUNK_METADATA_KEYS + if missing or unknown: + raise ValueError("typed Evidence metadata must use the exact allowlisted keys") + if not isinstance(metadata["evidence_id"], str) or not _EVIDENCE_ID.fullmatch(metadata["evidence_id"]): + raise ValueError("typed Evidence metadata must contain a valid evidence_id") + if metadata["evidence_kind"] not in EVIDENCE_KINDS: + raise ValueError("typed Evidence metadata must contain a valid evidence_kind") + purposes = metadata["purposes"] + if not isinstance(purposes, list) or not purposes or any(value not in EVIDENCE_PURPOSES for value in purposes): + raise ValueError("typed Evidence metadata must contain public purposes") + scope = metadata["scope"] + if not isinstance(scope, dict) or set(scope) != {"concepts", "tables", "columns"}: + raise ValueError("typed Evidence metadata must contain canonical scope") + if any(not isinstance(scope[name], list) or any(not isinstance(value, str) for value in scope[name]) + for name in ("concepts", "tables", "columns")): + raise ValueError("typed Evidence metadata must contain canonical scope") + if not isinstance(metadata["language"], str) or not metadata["language"]: + raise ValueError("typed Evidence metadata must contain language") + provenance = metadata["provenance"] + if not isinstance(provenance, dict) or set(provenance) != { + "source_file", "source_sha256", "supporting_excerpts", + }: + raise ValueError("typed Evidence metadata must contain canonical provenance") + + class _CanonicalValue(BaseModel): model_config = ConfigDict( frozen=True, extra="forbid", validate_default=True, revalidate_instances="always" @@ -99,6 +135,7 @@ class CanonicalChunk(_WithMetadata): @model_validator(mode="after") def content_hash_matches(self) -> "CanonicalChunk": _require_content_hash(self.content, self.content_hash) + _validate_evidence_metadata(self.metadata) return self diff --git a/harness/tht/evidence/corpus/normalize.py b/harness/tht/evidence/corpus/normalize.py index 5076077a..3000d688 100644 --- a/harness/tht/evidence/corpus/normalize.py +++ b/harness/tht/evidence/corpus/normalize.py @@ -4,12 +4,15 @@ import hashlib import re import unicodedata from collections.abc import Mapping +from pathlib import Path +from urllib.parse import urlsplit import yaml from pydantic import JsonValue, TypeAdapter, ValidationError from yaml.events import AliasEvent from yaml.nodes import MappingNode +from tht.evidence.canonical import EVIDENCE_KINDS, parse_curated_markdown from tht.evidence.contracts import AcquiredDocument, canonical_provenance_uri from tht.evidence.corpus.models import CanonicalDocument @@ -108,6 +111,20 @@ def _frontmatter(text: str) -> tuple[dict[str, JsonValue], str]: return metadata, text[match.end() :] +def _is_curated_filesystem_document(acquired: AcquiredDocument) -> bool: + if urlsplit(acquired.source.uri).scheme != "file": + return False + path = Path(urlsplit(acquired.source.uri).path) + if path.suffix != ".md": + return False + parts = path.parts + try: + curated_index = parts.index("curated") + except ValueError: + return False + return len(parts) >= curated_index + 3 and parts[curated_index + 1] in EVIDENCE_KINDS + + def normalize(acquired: AcquiredDocument, pipeline_version: str) -> CanonicalDocument: """Normalize one transport result without I/O or implicit data loss.""" if not pipeline_version: @@ -115,7 +132,6 @@ def normalize(acquired: AcquiredDocument, pipeline_version: str) -> CanonicalDoc decoded = _decode(acquired) canonical = unicodedata.normalize("NFC", decoded.replace("\r\n", "\n").replace("\r", "\n")) - frontmatter, content = _frontmatter(canonical) source_uri = canonical_provenance_uri(acquired.source.uri) identity = f"{acquired.source.source_id}\n{source_uri}" media_type = (acquired.media_type or "text/plain").split(";", 1)[0].strip().lower() @@ -123,8 +139,21 @@ def normalize(acquired: AcquiredDocument, pipeline_version: str) -> CanonicalDoc "source": acquired.source.model_dump(mode="json")["metadata"], "acquisition": acquired.model_dump(mode="json")["metadata"], } - if frontmatter: - metadata["frontmatter"] = frontmatter + if _is_curated_filesystem_document(acquired): + try: + evidence = parse_curated_markdown(canonical, path=Path(urlsplit(acquired.source.uri).path)) + except ValueError as error: + raise PermanentNormalizationError("invalid_curated_evidence") from error + if evidence.review_items: + raise PermanentNormalizationError("curated_evidence_requires_review") + content = canonical + title = evidence.title + metadata["curated_evidence"] = evidence.model_dump(mode="json") + else: + frontmatter, content = _frontmatter(canonical) + title = str(frontmatter.get("title", "")) + if frontmatter: + metadata["frontmatter"] = frontmatter try: return CanonicalDocument( @@ -133,7 +162,7 @@ def normalize(acquired: AcquiredDocument, pipeline_version: str) -> CanonicalDoc source_uri=source_uri, source_fingerprint=acquired.source.fingerprint, content_hash=f"sha256:{_sha256(content)}", - title=str(frontmatter.get("title", "")), + title=title, content=content, media_type=media_type, modified_at=acquired.source.modified_at, @@ -141,6 +170,6 @@ def normalize(acquired: AcquiredDocument, pipeline_version: str) -> CanonicalDoc metadata=metadata, ) except ValidationError as error: - if frontmatter: + if "frontmatter" in metadata: raise PermanentNormalizationError("invalid_frontmatter") from error raise diff --git a/harness/tht/evidence/corpus/pipeline.py b/harness/tht/evidence/corpus/pipeline.py index 0cb78a03..319d4020 100644 --- a/harness/tht/evidence/corpus/pipeline.py +++ b/harness/tht/evidence/corpus/pipeline.py @@ -13,8 +13,9 @@ from datetime import UTC from pathlib import Path import tht.evidence.acquisition as evidence_acquisition +from tht.evidence.canonical import ReviewItem from tht.evidence.contracts import EvidenceSource, SourceObject, canonical_provenance_uri -from tht.evidence.corpus.chunk import ChunkPolicy, chunk +from tht.evidence.corpus.chunk import AtomicContentTooLargeError, ChunkPolicy, chunk from tht.evidence.corpus.models import CanonicalChunk, CanonicalDocument, CorpusManifest from tht.evidence.corpus.normalize import normalize from tht.evidence.corpus.store import CorpusStore @@ -51,6 +52,7 @@ class PipelineResult: manifest: CorpusManifest = field(repr=False) run_id: str | None = None resumed_from: str | None = None + review_items: tuple[ReviewItem, ...] = () def __repr__(self) -> str: counts = { @@ -85,6 +87,7 @@ class PipelineResult: "manifest_id": self.manifest.manifest_id, "run_id": self.run_id, "resumed_from": self.resumed_from, + "review_items": [item.model_dump(mode="json") for item in self.review_items], } @@ -466,7 +469,11 @@ class CorpusPipeline: evidence_acquisition.acquire(source, item), self.pipeline_version, )) documents.sort(key=lambda value: value.source_id) - chunks = [part for document in documents for part in chunk(document, self.chunk_policy)] + try: + chunks = [part for document in documents for part in chunk(document, self.chunk_policy)] + except AtomicContentTooLargeError as error: + write(context, "review-items.json", [error.review_item.model_dump(mode="json")]) + raise previous_generations = dict(previous.metadata.get("document_generations", {})) if previous else {} changed = set(plan["changed"]) generations = { @@ -667,6 +674,14 @@ class CorpusPipeline: ], after_stage_return=after_stage_return) run_dir = workspace_root / ".tht-jobs" / "evidence" / "runs" / report.run_id plan = json.loads((run_dir / "artifacts" / "plan.json").read_text()) + review_items_path = run_dir / "artifacts" / "review-items.json" + try: + review_items = tuple( + ReviewItem.model_validate(item) + for item in json.loads(review_items_path.read_text(encoding="utf-8")) + ) if review_items_path.exists() else () + except (OSError, ValueError) as error: + raise PipelineError("preprocessing review items are corrupt") from error if dry_run: manifest = self.store.active_manifest() or CorpusManifest(pipeline_version=self.pipeline_version) generation = None @@ -682,9 +697,9 @@ class CorpusPipeline: if manifest_path.exists() else CorpusManifest(pipeline_version=self.pipeline_version)) published = False return PipelineResult( - report.status, generation, published, tuple(plan["changed"]), + "blocked" if review_items else report.status, generation, published, tuple(plan["changed"]), tuple(plan["unchanged"]), tuple(plan["removed"]), manifest, - report.run_id, report.resumed_from, + report.run_id, report.resumed_from, review_items, ) def _run(self, *, dry_run: bool = False, resume: str | None = None) -> PipelineResult: @@ -778,6 +793,13 @@ class CorpusPipeline: ) self.store.publish(staged) self.gc(workspace_root=self.store.root.parent) + except AtomicContentTooLargeError as error: + self._compensate(generation, vector_written) + manifest = previous or CorpusManifest(pipeline_version=self.pipeline_version) + return PipelineResult( + "blocked", None, False, changed, unchanged, removed, manifest, + review_items=(error.review_item,), + ) except PipelineError: self._compensate(generation, vector_written) raise