From 44f1efa5ba9be77bf188285bdf3227a11f44b3cd Mon Sep 17 00:00:00 2001 From: mptyl Date: Mon, 24 Aug 2026 02:25:11 +0200 Subject: [PATCH] refactor(evidence): migrate acquisition and preprocessing (#30) --- harness/tests/test_corpus_pipeline.py | 23 ++ .../tests/test_evidence_facade_contract.py | 64 +++++ harness/tests/test_preprocess_cli.py | 7 +- harness/tht/adapters/evidence/filesystem.py | 2 +- harness/tht/adapters/evidence/http.py | 2 +- harness/tht/adapters/evidence/s3.py | 2 +- harness/tht/adapters/factory.py | 51 +--- harness/tht/cli/preprocess_cmd.py | 16 +- harness/tht/config.py | 2 +- harness/tht/corpus/models.py | 2 +- harness/tht/corpus/normalize.py | 2 +- harness/tht/corpus/pipeline.py | 13 +- harness/tht/evidence/__init__.py | 58 +++- harness/tht/evidence/acquisition.py | 18 ++ harness/tht/evidence/contracts.py | 230 ++++++++++++++++ harness/tht/evidence/preprocessing.py | 44 +++ harness/tht/evidence/sources.py | 66 +++++ harness/tht/ports/evidence.py | 255 ++---------------- 18 files changed, 554 insertions(+), 303 deletions(-) create mode 100644 harness/tht/evidence/acquisition.py create mode 100644 harness/tht/evidence/contracts.py create mode 100644 harness/tht/evidence/preprocessing.py create mode 100644 harness/tht/evidence/sources.py diff --git a/harness/tests/test_corpus_pipeline.py b/harness/tests/test_corpus_pipeline.py index 2e62499d..572938ef 100644 --- a/harness/tests/test_corpus_pipeline.py +++ b/harness/tests/test_corpus_pipeline.py @@ -116,6 +116,29 @@ def pipeline(tmp_path, source, *, embedder=None, vectors=None, model="model-a", ) +def test_pipeline_routes_source_io_through_evidence_facade(tmp_path, monkeypatch): + import tht.evidence.acquisition as evidence_acquisition + + calls = [] + + def discover(source): + calls.append(("discover", source)) + return source.discover() + + def acquire(source, source_item): + calls.append(("acquire", source_item.source_id)) + return source.acquire(source_item) + + monkeypatch.setattr(evidence_acquisition, "discover", discover) + monkeypatch.setattr(evidence_acquisition, "acquire", acquire) + source = Source([(item("one", "a"), "body")]) + + result = pipeline(tmp_path, source).run() + + assert result.status == "succeeded" + assert calls == [("discover", source), ("acquire", "fs:one")] + + def test_retention_bounds_generations_and_purges_vectors_after_publish(tmp_path): vectors = Vectors() generations = [] diff --git a/harness/tests/test_evidence_facade_contract.py b/harness/tests/test_evidence_facade_contract.py index 73d5eae9..aac845aa 100644 --- a/harness/tests/test_evidence_facade_contract.py +++ b/harness/tests/test_evidence_facade_contract.py @@ -1,5 +1,6 @@ from datetime import UTC, datetime import hashlib +import inspect from types import SimpleNamespace import pytest @@ -10,7 +11,9 @@ from tht.decisions import DecisionRecord from tht.evidence import ( acquire, active_searcher, + build_preprocessing_pipeline, build_retrieval_entries, + build_sources, discover, project_session, resolve_citation, @@ -87,6 +90,67 @@ def test_acquisition_facade_preserves_classified_errors(): assert captured.value.details == {"operation": "download"} +def test_source_factory_preserves_legacy_first_order_and_filesystem_configuration(tmp_path): + from tht.adapters.factory import build_evidence_sources + + legacy_root = tmp_path / "legacy" + configured_root = tmp_path / "configured" + (legacy_root / "evidence").mkdir(parents=True) + configured_root.mkdir() + cfg = SimpleNamespace(evidence=SimpleNamespace( + source_root=legacy_root, + evidence_dir="evidence", + sources=[SimpleNamespace( + type="filesystem", + root=configured_root, + patterns=("*.md",), + max_bytes=1024, + )], + )) + + legacy = build_evidence_sources(cfg) + current = build_sources(cfg.evidence) + + assert [type(source) for source in current] == [type(source) for source in legacy] + assert [source.root for source in current] == [ + (legacy_root / "evidence").resolve(), + configured_root.resolve(), + ] + assert current[1].patterns == legacy[1].patterns == ("*.md",) + assert current[1].max_bytes == legacy[1].max_bytes == 1024 + + +def test_preprocessing_factory_forwards_only_evidence_pipeline_dependencies(monkeypatch): + captured = {} + + class FakePipeline: + def __init__(self, **kwargs): + captured.update(kwargs) + + monkeypatch.setattr("tht.corpus.pipeline.CorpusPipeline", FakePipeline) + dependencies = { + "store": object(), + "sources": [object()], + "embedder": object(), + "vector_store": object(), + "embedding_model": "model", + "embedding_dimensions": 3, + "chunk_policy": object(), + "pipeline_version": "evidence-v1", + "retain_published_generations": 2, + "workspace_id": None, + } + + pipeline = build_preprocessing_pipeline(**dependencies) + + assert isinstance(pipeline, FakePipeline) + assert captured == dependencies + assert all( + parameter.kind is not inspect.Parameter.VAR_KEYWORD + for parameter in inspect.signature(build_preprocessing_pipeline).parameters.values() + ) + + def _active_config(tmp_path): store = CorpusStore(tmp_path / "corpus") generation = store.stage( diff --git a/harness/tests/test_preprocess_cli.py b/harness/tests/test_preprocess_cli.py index 792111eb..7bba5994 100644 --- a/harness/tests/test_preprocess_cli.py +++ b/harness/tests/test_preprocess_cli.py @@ -247,10 +247,13 @@ def test_run_from_config_uses_runtime_identity_workspace_id(monkeypatch, tmp_pat calls["run_as_job"] = kwargs return SimpleNamespace(model_dump=lambda mode=None: {"status": "succeeded"}) - monkeypatch.setattr("tht.adapters.factory.build_evidence_sources", lambda cfg: []) + monkeypatch.setattr("tht.evidence.build_sources", lambda cfg: []) + monkeypatch.setattr( + "tht.evidence.build_preprocessing_pipeline", + lambda **kwargs: FakePipeline(**kwargs), + ) monkeypatch.setattr("tht.adapters.factory.build_vector_store", lambda cfg, require_write: object()) monkeypatch.setattr("tht.cli.vector_cmd.make_embedder", lambda cfg: object()) - monkeypatch.setattr("tht.corpus.pipeline.CorpusPipeline", FakePipeline) command.run_from_config(config) diff --git a/harness/tht/adapters/evidence/filesystem.py b/harness/tht/adapters/evidence/filesystem.py index 8eca8d1a..ae415607 100644 --- a/harness/tht/adapters/evidence/filesystem.py +++ b/harness/tht/adapters/evidence/filesystem.py @@ -7,7 +7,7 @@ from datetime import UTC, datetime from pathlib import Path, PurePosixPath from urllib.parse import unquote, urlsplit -from tht.ports.evidence import ( +from tht.evidence.contracts import ( AcquiredDocument, EvidenceSourceError, EvidenceSourceErrorCategory, diff --git a/harness/tht/adapters/evidence/http.py b/harness/tht/adapters/evidence/http.py index 9c895693..ff423d22 100644 --- a/harness/tht/adapters/evidence/http.py +++ b/harness/tht/adapters/evidence/http.py @@ -10,7 +10,7 @@ from urllib.parse import urljoin, urlsplit import requests -from tht.ports.evidence import ( +from tht.evidence.contracts import ( AcquiredDocument, EvidenceSourceError, EvidenceSourceErrorCategory, diff --git a/harness/tht/adapters/evidence/s3.py b/harness/tht/adapters/evidence/s3.py index e0a45fb7..cbbfee90 100644 --- a/harness/tht/adapters/evidence/s3.py +++ b/harness/tht/adapters/evidence/s3.py @@ -6,7 +6,7 @@ import re from datetime import UTC, datetime from urllib.parse import quote, urlsplit -from tht.ports.evidence import ( +from tht.evidence.contracts import ( AcquiredDocument, EvidenceSourceError, EvidenceSourceErrorCategory, SourceObject, ) diff --git a/harness/tht/adapters/factory.py b/harness/tht/adapters/factory.py index a9f52936..faa465b2 100644 --- a/harness/tht/adapters/factory.py +++ b/harness/tht/adapters/factory.py @@ -1,8 +1,6 @@ """Central construction of deployment-specific adapters.""" from tht.adapters.dwh import PostgresDwhAdapter, ThothRestDwhAdapter -from tht.adapters.evidence import FilesystemEvidenceSource, HttpManifestEvidenceSource -from tht.adapters.evidence.s3 import S3EvidenceSource from tht.adapters.vector import QdrantVectorStore from tht.config import Config, ConfigError from tht.ports.dwh import DwhAdapter @@ -45,53 +43,10 @@ def build_vector_store(cfg: Config, *, require_write: bool = False) -> VectorSto def build_evidence_sources(cfg: Config): - """Build configured Evidence sources, including the legacy curated filesystem tree.""" - evidence = cfg.evidence - if evidence is None: - return [] - sources = [] - if evidence.source_root is not None: - sources.append(FilesystemEvidenceSource(evidence.source_root / evidence.evidence_dir)) - for resource in evidence.sources: - match resource.type: - case "filesystem": - sources.append( - FilesystemEvidenceSource( - resource.root, - patterns=resource.patterns, - max_bytes=resource.max_bytes, - ) - ) - case "http": - sources.append( - HttpManifestEvidenceSource( - resource.transport_urls(), - connect_timeout=resource.connect_timeout, - read_timeout=resource.read_timeout, - max_bytes=resource.max_bytes, - max_redirects=resource.max_redirects, - allow_private_hosts=resource.allow_private_hosts, - max_cache_bytes=resource.max_cache_bytes, - ) - ) - case "s3": - def secret(value): - return value.get_secret_value() if value is not None else None + """Compatibility shim for callers not yet migrated to ``tht.evidence``.""" + from tht.evidence import build_sources - sources.append(S3EvidenceSource( - bucket=resource.bucket, prefix=resource.prefix, - endpoint_url=resource.endpoint_url, region=resource.region, - access_key=secret(resource.access_key), secret_key=secret(resource.secret_key), - session_token=secret(resource.session_token), - trusted_endpoint=resource.trusted_endpoint, - allow_private_endpoint=resource.allow_private_endpoint, - allow_insecure_endpoint=resource.allow_insecure_endpoint, - max_bytes=resource.max_bytes, max_objects=resource.max_objects, - max_pages=resource.max_pages, page_size=resource.page_size, - )) - case other: # pragma: no cover - Pydantic rejects unsupported discriminators. - raise ConfigError(f"Adapter evidence non supportato: {other}") - return sources + return build_sources(cfg.evidence) __all__ = ["build_dwh", "build_evidence_sources", "build_vector_store"] diff --git a/harness/tht/cli/preprocess_cmd.py b/harness/tht/cli/preprocess_cmd.py index 3bd008c6..aa6f1238 100644 --- a/harness/tht/cli/preprocess_cmd.py +++ b/harness/tht/cli/preprocess_cmd.py @@ -93,19 +93,19 @@ def _parse_dwh_steps(value: str) -> tuple[str, ...]: def run_from_config(config: Path, *, dry_run: bool = False, resume: str | None = None): - from tht.adapters.factory import build_evidence_sources, build_vector_store + 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.corpus.chunk import ChunkPolicy - from tht.corpus.pipeline import CorpusPipeline from tht.corpus.store import CorpusStore + from tht.evidence import build_preprocessing_pipeline, build_sources cfg = _load_config_or_exit(config) if cfg.embeddings is None: raise RuntimeError("embeddings are not configured") corpus_root = cfg.paths.artifacts.parent / "corpus" - pipeline = CorpusPipeline( - store=CorpusStore(corpus_root), sources=build_evidence_sources(cfg), + 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), embedding_model=cfg.embeddings.model, embedding_dimensions=cfg.embeddings.dim, @@ -127,19 +127,19 @@ def run_from_config(config: Path, *, dry_run: bool = False, resume: str | None = def gc_from_config(config: Path, *, dry_run: bool = False): - from tht.adapters.factory import build_evidence_sources, build_vector_store + 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.corpus.chunk import ChunkPolicy - from tht.corpus.pipeline import CorpusPipeline from tht.corpus.store import CorpusStore + from tht.evidence import build_preprocessing_pipeline, build_sources cfg = _load_config_or_exit(config) if cfg.embeddings is None: raise RuntimeError("embeddings are not configured") corpus_root = cfg.paths.artifacts.parent / "corpus" - pipeline = CorpusPipeline( - store=CorpusStore(corpus_root), sources=build_evidence_sources(cfg), + 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), embedding_model=cfg.embeddings.model, embedding_dimensions=cfg.embeddings.dim, chunk_policy=ChunkPolicy(version="chunk-v1", max_chars=cfg.vector.max_chunk_chars), diff --git a/harness/tht/config.py b/harness/tht/config.py index a0182ec3..93ac3d47 100644 --- a/harness/tht/config.py +++ b/harness/tht/config.py @@ -13,7 +13,7 @@ import yaml from pydantic import BaseModel, Field, PrivateAttr, SecretStr, ValidationError, model_validator from tht.config_compat import translate_legacy_config -from tht.ports.evidence import canonical_provenance_uri +from tht.evidence.contracts import canonical_provenance_uri _ENV_RE = re.compile(r"\$\{([A-Za-z_][A-Za-z0-9_]*)\}") diff --git a/harness/tht/corpus/models.py b/harness/tht/corpus/models.py index 4c4f93f5..d398124c 100644 --- a/harness/tht/corpus/models.py +++ b/harness/tht/corpus/models.py @@ -8,7 +8,7 @@ from typing import Self from pydantic import BaseModel, ConfigDict, Field, JsonValue, field_validator, model_validator -from tht.ports.evidence import ( +from tht.evidence.contracts import ( canonical_provenance_uri, normalize_aware_datetime, validate_namespaced_value, diff --git a/harness/tht/corpus/normalize.py b/harness/tht/corpus/normalize.py index 54076826..22d4ac21 100644 --- a/harness/tht/corpus/normalize.py +++ b/harness/tht/corpus/normalize.py @@ -11,7 +11,7 @@ from yaml.events import AliasEvent from yaml.nodes import MappingNode from tht.corpus.models import CanonicalDocument -from tht.ports.evidence import AcquiredDocument, canonical_provenance_uri +from tht.evidence.contracts import AcquiredDocument, canonical_provenance_uri MAX_DOCUMENT_BYTES = 10 * 1024 * 1024 diff --git a/harness/tht/corpus/pipeline.py b/harness/tht/corpus/pipeline.py index 625a5ff1..7a55012d 100644 --- a/harness/tht/corpus/pipeline.py +++ b/harness/tht/corpus/pipeline.py @@ -15,7 +15,8 @@ from tht.corpus.chunk import ChunkPolicy, chunk from tht.corpus.models import CanonicalChunk, CanonicalDocument, CorpusManifest from tht.corpus.normalize import normalize from tht.corpus.store import CorpusStore -from tht.ports.evidence import EvidenceSource, SourceObject, canonical_provenance_uri +import tht.evidence.acquisition as evidence_acquisition +from tht.evidence.contracts import EvidenceSource, SourceObject, canonical_provenance_uri from tht.ports.vector import VectorStore, VectorWriteRecord from tht.vectorstore.records import VectorRecord from tht.jobs.models import JobSpec @@ -228,7 +229,7 @@ class CorpusPipeline: discovered = [] seen = set() for source in self.sources: - for item in source.discover(): + for item in evidence_acquisition.discover(source): if item.source_id in seen: raise PipelineError("duplicate Evidence source identity") seen.add(item.source_id) @@ -455,7 +456,9 @@ class CorpusPipeline: documents = [prior[source_id] for source_id in plan["unchanged"]] for source_id in plan["changed"]: source, item = source_by_id[source_id] - documents.append(normalize(source.acquire(item), self.pipeline_version)) + documents.append(normalize( + 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)] previous_generations = dict(previous.metadata.get("document_generations", {})) if previous else {} @@ -710,7 +713,9 @@ class CorpusPipeline: try: for source, item in discovered: if item.source_id in changed_set: - documents.append(normalize(source.acquire(item), self.pipeline_version)) + documents.append(normalize( + evidence_acquisition.acquire(source, item), self.pipeline_version, + )) documents.sort(key=lambda document: document.source_id) chunks: list[CanonicalChunk] = [] for document in documents: diff --git a/harness/tht/evidence/__init__.py b/harness/tht/evidence/__init__.py index f5e4cce4..d2756b57 100644 --- a/harness/tht/evidence/__init__.py +++ b/harness/tht/evidence/__init__.py @@ -7,33 +7,69 @@ and exceptions; they do not define a cross-domain service protocol. from __future__ import annotations -from collections.abc import Iterable from pathlib import Path from typing import TYPE_CHECKING -from tht.ports.evidence import ( +from tht.evidence.acquisition import acquire, discover +from tht.evidence.contracts import ( AcquiredDocument, EvidenceSource, EvidenceSourceError, EvidenceSourceErrorCategory, SourceObject, + canonical_provenance_uri, + normalize_aware_datetime, + validate_namespaced_value, + validate_safe_metadata, ) if TYPE_CHECKING: + from tht.config import EvidenceSourcesConfig + from tht.corpus.chunk import ChunkPolicy + from tht.corpus.pipeline import CorpusPipeline from tht.corpus.store import CorpusStore from tht.decisions import DecisionRecord + from tht.evidence.preprocessing import EvidenceEmbedder + from tht.ports.vector import VectorStore from tht.search.evidence import ActiveEvidenceSearcher from tht.session.models import SchemaLinking -def discover(source: EvidenceSource) -> Iterable[SourceObject]: - """Discover source objects without changing source-defined laziness or ordering.""" - return source.discover() +def build_sources(evidence: "EvidenceSourcesConfig | None") -> list[EvidenceSource]: + """Build configured Evidence source adapters in the existing deterministic order.""" + from tht.evidence.sources import build_sources as build_configured_sources + + return build_configured_sources(evidence) -def acquire(source: EvidenceSource, item: SourceObject) -> AcquiredDocument: - """Acquire one discovered object, preserving the source's classified failures.""" - return source.acquire(item) +def build_preprocessing_pipeline( + *, + store: "CorpusStore", + sources: list[EvidenceSource], + embedder: "EvidenceEmbedder", + 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, +) -> "CorpusPipeline": + """Construct the Evidence preprocessing use case from core-owned infrastructure.""" + from tht.evidence.preprocessing import build_preprocessing_pipeline as build_pipeline + + return build_pipeline( + store=store, + sources=sources, + embedder=embedder, + vector_store=vector_store, + embedding_model=embedding_model, + embedding_dimensions=embedding_dimensions, + chunk_policy=chunk_policy, + pipeline_version=pipeline_version, + retain_published_generations=retain_published_generations, + workspace_id=workspace_id, + ) def active_searcher( @@ -135,9 +171,15 @@ __all__ = [ "SourceObject", "acquire", "active_searcher", + "build_preprocessing_pipeline", "build_retrieval_entries", + "build_sources", + "canonical_provenance_uri", "discover", + "normalize_aware_datetime", "project_session", "resolve_citation", "validate_corpus_workspace", + "validate_namespaced_value", + "validate_safe_metadata", ] diff --git a/harness/tht/evidence/acquisition.py b/harness/tht/evidence/acquisition.py new file mode 100644 index 00000000..a4471c11 --- /dev/null +++ b/harness/tht/evidence/acquisition.py @@ -0,0 +1,18 @@ +"""Evidence acquisition operations over the credential-free source contract.""" + +from collections.abc import Iterable + +from tht.evidence.contracts import AcquiredDocument, EvidenceSource, SourceObject + + +def discover(source: EvidenceSource) -> Iterable[SourceObject]: + """Discover source objects without changing source-defined laziness or ordering.""" + return source.discover() + + +def acquire(source: EvidenceSource, item: SourceObject) -> AcquiredDocument: + """Acquire one discovered object, preserving the source's classified failures.""" + return source.acquire(item) + + +__all__ = ["acquire", "discover"] diff --git a/harness/tht/evidence/contracts.py b/harness/tht/evidence/contracts.py new file mode 100644 index 00000000..73bea547 --- /dev/null +++ b/harness/tht/evidence/contracts.py @@ -0,0 +1,230 @@ +"""Credential-free contracts for discovering and acquiring Evidence objects.""" + +import re +from collections.abc import Iterable, Mapping, Sequence +from datetime import UTC, datetime +from enum import Enum +from typing import Protocol, Self, runtime_checkable +from urllib.parse import parse_qsl, urlsplit, urlunsplit + +from pydantic import BaseModel, ConfigDict, Field, JsonValue, TypeAdapter, field_validator + + +class FrozenDict(dict): + """A JSON-serializable dict whose mutation operations are disabled.""" + + def _immutable(self, *args, **kwargs): + raise TypeError("frozen JSON metadata cannot be mutated") + + __delitem__ = _immutable + __ior__ = _immutable + __setitem__ = _immutable + clear = _immutable + pop = _immutable + popitem = _immutable + setdefault = _immutable + update = _immutable + + +class FrozenList(list): + """A JSON-serializable list whose mutation operations are disabled.""" + + def _immutable(self, *args, **kwargs): + raise TypeError("frozen JSON metadata cannot be mutated") + + __delitem__ = _immutable + __iadd__ = _immutable + __imul__ = _immutable + __setitem__ = _immutable + append = _immutable + clear = _immutable + extend = _immutable + insert = _immutable + pop = _immutable + remove = _immutable + reverse = _immutable + sort = _immutable + + +_CAMEL_BOUNDARY = re.compile(r"(?<=[a-z0-9])(?=[A-Z])") +_SEPARATORS = re.compile(r"[^a-z0-9]+") +_NAMESPACED_VALUE = re.compile(r"^[a-z][a-z0-9_-]*:[A-Za-z0-9._:-]+$") +_CREDENTIAL_KEYS = { + "apikey", + "authorization", + "authtoken", + "bearertoken", + "clientsecret", + "credential", + "credentials", + "password", + "passwd", + "privatekey", + "refreshtoken", + "sessioncookie", + "xapikey", + "accesstoken", +} +_JSON_METADATA = TypeAdapter(dict[str, JsonValue]) + + +def _normalize_key(key: str) -> str: + return _SEPARATORS.sub("", _CAMEL_BOUNDARY.sub("_", key).lower()) + + +def _is_credential_key(key: str) -> bool: + return _normalize_key(key) in _CREDENTIAL_KEYS + + +def _reject_credentials(value, path: str = "metadata") -> None: + if isinstance(value, Mapping): + for key, child in value.items(): + if _is_credential_key(str(key)): + raise ValueError(f"credential-like metadata key is not allowed: {path}.{key}") + _reject_credentials(child, f"{path}.{key}") + elif isinstance(value, Sequence) and not isinstance(value, (str, bytes, bytearray)): + for index, child in enumerate(value): + _reject_credentials(child, f"{path}[{index}]") + + +def freeze_json(value): + """Recursively freeze a Pydantic-validated JSON value without changing its JSON shape.""" + if isinstance(value, Mapping): + return FrozenDict({str(key): freeze_json(child) for key, child in value.items()}) + if isinstance(value, Sequence) and not isinstance(value, (str, bytes, bytearray)): + return FrozenList(freeze_json(child) for child in value) + return value + + +def validate_safe_metadata(value: dict[str, JsonValue]) -> FrozenDict: + _reject_credentials(value) + return freeze_json(value) + + +def validate_canonical_uri(value: str) -> str: + try: + parsed = urlsplit(value) + _ = parsed.port + except ValueError as error: + raise ValueError("invalid canonical URI") from error + if not parsed.scheme: + raise ValueError("canonical URI must include a scheme") + if parsed.username is not None or parsed.password is not None: + raise ValueError("canonical URI must not contain credentials in userinfo") + for key, _ in parse_qsl(parsed.query, keep_blank_values=True): + if _is_credential_key(key): + raise ValueError("canonical URI must not contain credentials in query parameters") + return value + + +def canonical_provenance_uri(value: str) -> str: + """Return only stable URI identity; transport query/fragment data is never provenance.""" + try: + parsed = urlsplit(value) + _ = parsed.port + except ValueError as error: + raise ValueError("invalid canonical URI") from error + if not parsed.scheme: + raise ValueError("canonical URI must include a scheme") + if parsed.username is not None or parsed.password is not None: + raise ValueError("canonical URI must not contain credentials in userinfo") + return urlunsplit((parsed.scheme, parsed.netloc, parsed.path, "", "")) + + +def normalize_aware_datetime(value: datetime | None) -> datetime | None: + if value is None: + return None + if value.tzinfo is None or value.utcoffset() is None: + raise ValueError("datetime must be timezone-aware") + return value.astimezone(UTC) + + +def validate_namespaced_value(value: str) -> str: + if not _NAMESPACED_VALUE.fullmatch(value): + raise ValueError("value must be namespaced as ':'") + return value + + +class _EvidenceValue(BaseModel): + model_config = ConfigDict( + frozen=True, + extra="forbid", + revalidate_instances="always", + validate_default=True, + ser_json_bytes="base64", + val_json_bytes="base64", + ) + + def model_copy(self, *, update: Mapping[str, object] | None = None, deep: bool = False) -> Self: + """Copy through validation; Pydantic's unchecked update-copy is unsafe for contracts.""" + data = self.model_dump(round_trip=True) + if update: + data.update(update) + return type(self).model_validate(data) + + +class SourceObject(_EvidenceValue): + source_id: str = Field(min_length=1) + uri: str = Field(min_length=1) + fingerprint: str = Field(min_length=1) + modified_at: datetime | None = None + metadata: dict[str, JsonValue] = Field(default_factory=dict) + + _source_id = field_validator("source_id")(validate_namespaced_value) + _fingerprint = field_validator("fingerprint")(validate_namespaced_value) + _safe_uri = field_validator("uri")(validate_canonical_uri) + _aware_modified_at = field_validator("modified_at")(normalize_aware_datetime) + _frozen_metadata = field_validator("metadata")(validate_safe_metadata) + + +class AcquiredDocument(_EvidenceValue): + """Transport result; bytes use explicit base64 encoding in JSON mode.""" + + source: SourceObject + content: bytes + media_type: str | None = None + acquired_at: datetime | None = None + metadata: dict[str, JsonValue] = Field(default_factory=dict) + + _aware_acquired_at = field_validator("acquired_at")(normalize_aware_datetime) + _frozen_metadata = field_validator("metadata")(validate_safe_metadata) + + +class EvidenceSourceErrorCategory(str, Enum): + TRANSIENT = "transient" + PERMANENT = "permanent" + + +class EvidenceSourceError(Exception): + """Classified source failure with credential-free structured diagnostics.""" + + def __init__( + self, + _message: str, + *, + category: EvidenceSourceErrorCategory, + details: dict[str, JsonValue] | None = None, + ) -> None: + super().__init__("evidence source operation failed") + object.__setattr__(self, "category", EvidenceSourceErrorCategory(category)) + object.__setattr__( + self, + "details", + validate_safe_metadata(_JSON_METADATA.validate_python(details or {})), + ) + + def __setattr__(self, name: str, value) -> None: + if name in {"args", "category", "details"} and hasattr(self, name): + raise AttributeError(f"{name} is immutable") + super().__setattr__(name, value) + + @property + def retryable(self) -> bool: + return self.category is EvidenceSourceErrorCategory.TRANSIENT + + +@runtime_checkable +class EvidenceSource(Protocol): + def discover(self) -> Iterable[SourceObject]: ... + + def acquire(self, item: SourceObject) -> AcquiredDocument: ... diff --git a/harness/tht/evidence/preprocessing.py b/harness/tht/evidence/preprocessing.py new file mode 100644 index 00000000..e1bf29e2 --- /dev/null +++ b/harness/tht/evidence/preprocessing.py @@ -0,0 +1,44 @@ +"""Explicit construction boundary for Evidence preprocessing.""" + +from typing import Protocol + +from tht.corpus.chunk import ChunkPolicy +from tht.corpus.pipeline import CorpusPipeline +from tht.corpus.store import CorpusStore +from tht.evidence.contracts import EvidenceSource +from tht.ports.vector import VectorStore + + +class EvidenceEmbedder(Protocol): + def embed_documents(self, texts: list[str]) -> list[list[float]]: ... + + +def build_preprocessing_pipeline( + *, + store: CorpusStore, + sources: list[EvidenceSource], + embedder: EvidenceEmbedder, + 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, +) -> CorpusPipeline: + """Construct preprocessing from the bounded infrastructure supplied by core.""" + return CorpusPipeline( + store=store, + sources=sources, + embedder=embedder, + vector_store=vector_store, + embedding_model=embedding_model, + embedding_dimensions=embedding_dimensions, + chunk_policy=chunk_policy, + pipeline_version=pipeline_version, + retain_published_generations=retain_published_generations, + workspace_id=workspace_id, + ) + + +__all__ = ["EvidenceEmbedder", "build_preprocessing_pipeline"] diff --git a/harness/tht/evidence/sources.py b/harness/tht/evidence/sources.py new file mode 100644 index 00000000..4c2a6c87 --- /dev/null +++ b/harness/tht/evidence/sources.py @@ -0,0 +1,66 @@ +"""Construction of configured Evidence source adapters.""" + +from tht.adapters.evidence import FilesystemEvidenceSource, HttpManifestEvidenceSource +from tht.adapters.evidence.s3 import S3EvidenceSource +from tht.config import ConfigError, EvidenceSourcesConfig +from tht.evidence.contracts import EvidenceSource + + +def build_sources(evidence: EvidenceSourcesConfig | None) -> list[EvidenceSource]: + """Build configured Evidence adapters in the existing deterministic order.""" + if evidence is None: + return [] + sources: list[EvidenceSource] = [] + if evidence.source_root is not None: + sources.append(FilesystemEvidenceSource(evidence.source_root / evidence.evidence_dir)) + for resource in evidence.sources: + match resource.type: + case "filesystem": + sources.append( + FilesystemEvidenceSource( + resource.root, + patterns=resource.patterns, + max_bytes=resource.max_bytes, + ) + ) + case "http": + sources.append( + HttpManifestEvidenceSource( + resource.transport_urls(), + connect_timeout=resource.connect_timeout, + read_timeout=resource.read_timeout, + max_bytes=resource.max_bytes, + max_redirects=resource.max_redirects, + allow_private_hosts=resource.allow_private_hosts, + max_cache_bytes=resource.max_cache_bytes, + ) + ) + case "s3": + + def secret(value): + return value.get_secret_value() if value is not None else None + + sources.append( + S3EvidenceSource( + bucket=resource.bucket, + prefix=resource.prefix, + endpoint_url=resource.endpoint_url, + region=resource.region, + access_key=secret(resource.access_key), + secret_key=secret(resource.secret_key), + session_token=secret(resource.session_token), + trusted_endpoint=resource.trusted_endpoint, + allow_private_endpoint=resource.allow_private_endpoint, + allow_insecure_endpoint=resource.allow_insecure_endpoint, + max_bytes=resource.max_bytes, + max_objects=resource.max_objects, + max_pages=resource.max_pages, + page_size=resource.page_size, + ) + ) + case other: # pragma: no cover - Pydantic rejects unsupported discriminators. + raise ConfigError(f"Adapter evidence non supportato: {other}") + return sources + + +__all__ = ["build_sources"] diff --git a/harness/tht/ports/evidence.py b/harness/tht/ports/evidence.py index 0f736b0a..a56a1a6b 100644 --- a/harness/tht/ports/evidence.py +++ b/harness/tht/ports/evidence.py @@ -1,230 +1,31 @@ -"""Credential-free port for discovering and acquiring Evidence objects.""" +"""Legacy import path for Evidence contracts. -import re -from collections.abc import Iterable, Mapping, Sequence -from datetime import UTC, datetime -from enum import Enum -from typing import Protocol, Self, runtime_checkable -from urllib.parse import parse_qsl, urlsplit, urlunsplit +New production code imports the leaf contracts from ``tht.evidence.contracts``. This one-way +compatibility shim remains only until the old Evidence layout is removed. +""" -from pydantic import BaseModel, ConfigDict, Field, JsonValue, TypeAdapter, field_validator +from tht.evidence.contracts import ( + AcquiredDocument, + EvidenceSource, + EvidenceSourceError, + EvidenceSourceErrorCategory, + SourceObject, + canonical_provenance_uri, + normalize_aware_datetime, + validate_canonical_uri, + validate_namespaced_value, + validate_safe_metadata, +) - -class FrozenDict(dict): - """A JSON-serializable dict whose mutation operations are disabled.""" - - def _immutable(self, *args, **kwargs): - raise TypeError("frozen JSON metadata cannot be mutated") - - __delitem__ = _immutable - __ior__ = _immutable - __setitem__ = _immutable - clear = _immutable - pop = _immutable - popitem = _immutable - setdefault = _immutable - update = _immutable - - -class FrozenList(list): - """A JSON-serializable list whose mutation operations are disabled.""" - - def _immutable(self, *args, **kwargs): - raise TypeError("frozen JSON metadata cannot be mutated") - - __delitem__ = _immutable - __iadd__ = _immutable - __imul__ = _immutable - __setitem__ = _immutable - append = _immutable - clear = _immutable - extend = _immutable - insert = _immutable - pop = _immutable - remove = _immutable - reverse = _immutable - sort = _immutable - - -_CAMEL_BOUNDARY = re.compile(r"(?<=[a-z0-9])(?=[A-Z])") -_SEPARATORS = re.compile(r"[^a-z0-9]+") -_NAMESPACED_VALUE = re.compile(r"^[a-z][a-z0-9_-]*:[A-Za-z0-9._:-]+$") -_CREDENTIAL_KEYS = { - "apikey", - "authorization", - "authtoken", - "bearertoken", - "clientsecret", - "credential", - "credentials", - "password", - "passwd", - "privatekey", - "refreshtoken", - "sessioncookie", - "xapikey", - "accesstoken", -} -_JSON_METADATA = TypeAdapter(dict[str, JsonValue]) - - -def _normalize_key(key: str) -> str: - return _SEPARATORS.sub("", _CAMEL_BOUNDARY.sub("_", key).lower()) - - -def _is_credential_key(key: str) -> bool: - return _normalize_key(key) in _CREDENTIAL_KEYS - - -def _reject_credentials(value, path: str = "metadata") -> None: - if isinstance(value, Mapping): - for key, child in value.items(): - if _is_credential_key(str(key)): - raise ValueError(f"credential-like metadata key is not allowed: {path}.{key}") - _reject_credentials(child, f"{path}.{key}") - elif isinstance(value, Sequence) and not isinstance(value, (str, bytes, bytearray)): - for index, child in enumerate(value): - _reject_credentials(child, f"{path}[{index}]") - - -def freeze_json(value): - """Recursively freeze a Pydantic-validated JSON value without changing its JSON shape.""" - if isinstance(value, Mapping): - return FrozenDict({str(key): freeze_json(child) for key, child in value.items()}) - if isinstance(value, Sequence) and not isinstance(value, (str, bytes, bytearray)): - return FrozenList(freeze_json(child) for child in value) - return value - - -def validate_safe_metadata(value: dict[str, JsonValue]) -> FrozenDict: - _reject_credentials(value) - return freeze_json(value) - - -def validate_canonical_uri(value: str) -> str: - try: - parsed = urlsplit(value) - _ = parsed.port - except ValueError as error: - raise ValueError("invalid canonical URI") from error - if not parsed.scheme: - raise ValueError("canonical URI must include a scheme") - if parsed.username is not None or parsed.password is not None: - raise ValueError("canonical URI must not contain credentials in userinfo") - for key, _ in parse_qsl(parsed.query, keep_blank_values=True): - if _is_credential_key(key): - raise ValueError("canonical URI must not contain credentials in query parameters") - return value - - -def canonical_provenance_uri(value: str) -> str: - """Return only stable URI identity; transport query/fragment data is never provenance.""" - try: - parsed = urlsplit(value) - _ = parsed.port - except ValueError as error: - raise ValueError("invalid canonical URI") from error - if not parsed.scheme: - raise ValueError("canonical URI must include a scheme") - if parsed.username is not None or parsed.password is not None: - raise ValueError("canonical URI must not contain credentials in userinfo") - return urlunsplit((parsed.scheme, parsed.netloc, parsed.path, "", "")) - - -def normalize_aware_datetime(value: datetime | None) -> datetime | None: - if value is None: - return None - if value.tzinfo is None or value.utcoffset() is None: - raise ValueError("datetime must be timezone-aware") - return value.astimezone(UTC) - - -def validate_namespaced_value(value: str) -> str: - if not _NAMESPACED_VALUE.fullmatch(value): - raise ValueError("value must be namespaced as ':'") - return value - - -class _EvidenceValue(BaseModel): - model_config = ConfigDict( - frozen=True, - extra="forbid", - revalidate_instances="always", - validate_default=True, - ser_json_bytes="base64", - val_json_bytes="base64", - ) - - def model_copy(self, *, update: Mapping[str, object] | None = None, deep: bool = False) -> Self: - """Copy through validation; Pydantic's unchecked update-copy is unsafe for contracts.""" - data = self.model_dump(round_trip=True) - if update: - data.update(update) - return type(self).model_validate(data) - - -class SourceObject(_EvidenceValue): - source_id: str = Field(min_length=1) - uri: str = Field(min_length=1) - fingerprint: str = Field(min_length=1) - modified_at: datetime | None = None - metadata: dict[str, JsonValue] = Field(default_factory=dict) - - _source_id = field_validator("source_id")(validate_namespaced_value) - _fingerprint = field_validator("fingerprint")(validate_namespaced_value) - _safe_uri = field_validator("uri")(validate_canonical_uri) - _aware_modified_at = field_validator("modified_at")(normalize_aware_datetime) - _frozen_metadata = field_validator("metadata")(validate_safe_metadata) - - -class AcquiredDocument(_EvidenceValue): - """Transport result; bytes use explicit base64 encoding in JSON mode.""" - - source: SourceObject - content: bytes - media_type: str | None = None - acquired_at: datetime | None = None - metadata: dict[str, JsonValue] = Field(default_factory=dict) - - _aware_acquired_at = field_validator("acquired_at")(normalize_aware_datetime) - _frozen_metadata = field_validator("metadata")(validate_safe_metadata) - - -class EvidenceSourceErrorCategory(str, Enum): - TRANSIENT = "transient" - PERMANENT = "permanent" - - -class EvidenceSourceError(Exception): - """Classified source failure with credential-free structured diagnostics.""" - - def __init__( - self, - _message: str, - *, - category: EvidenceSourceErrorCategory, - details: dict[str, JsonValue] | None = None, - ) -> None: - super().__init__("evidence source operation failed") - object.__setattr__(self, "category", EvidenceSourceErrorCategory(category)) - object.__setattr__( - self, - "details", - validate_safe_metadata(_JSON_METADATA.validate_python(details or {})), - ) - - def __setattr__(self, name: str, value) -> None: - if name in {"args", "category", "details"} and hasattr(self, name): - raise AttributeError(f"{name} is immutable") - super().__setattr__(name, value) - - @property - def retryable(self) -> bool: - return self.category is EvidenceSourceErrorCategory.TRANSIENT - - -@runtime_checkable -class EvidenceSource(Protocol): - def discover(self) -> Iterable[SourceObject]: ... - - def acquire(self, item: SourceObject) -> AcquiredDocument: ... +__all__ = [ + "AcquiredDocument", + "EvidenceSource", + "EvidenceSourceError", + "EvidenceSourceErrorCategory", + "SourceObject", + "canonical_provenance_uri", + "normalize_aware_datetime", + "validate_canonical_uri", + "validate_namespaced_value", + "validate_safe_metadata", +]