refactor(evidence): migrate acquisition and preprocessing (#30)

This commit is contained in:
2026-08-24 02:25:11 +02:00
parent d5e78febd3
commit 44f1efa5ba
18 changed files with 554 additions and 303 deletions
+23
View File
@@ -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 = []
@@ -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(
+5 -2
View File
@@ -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)
+1 -1
View File
@@ -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,
+1 -1
View File
@@ -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,
+1 -1
View File
@@ -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,
)
+3 -48
View File
@@ -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"]
+8 -8
View File
@@ -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),
+1 -1
View File
@@ -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_]*)\}")
+1 -1
View File
@@ -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,
+1 -1
View File
@@ -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
+9 -4
View File
@@ -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:
+50 -8
View File
@@ -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",
]
+18
View File
@@ -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"]
+230
View File
@@ -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 '<kind>:<stable-value>'")
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: ...
+44
View File
@@ -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"]
+66
View File
@@ -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"]
+28 -227
View File
@@ -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 '<kind>:<stable-value>'")
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",
]