feat(evidence): add filesystem and HTTP sources

This commit is contained in:
2026-07-12 03:23:42 +02:00
parent 293d96e1a6
commit ffd683c587
9 changed files with 633 additions and 3 deletions
@@ -0,0 +1,6 @@
"""Evidence source adapter implementations."""
from tht.adapters.evidence.filesystem import FilesystemEvidenceSource
from tht.adapters.evidence.http import HttpManifestEvidenceSource
__all__ = ["FilesystemEvidenceSource", "HttpManifestEvidenceSource"]
+124
View File
@@ -0,0 +1,124 @@
"""Contained, deterministic filesystem Evidence source."""
import hashlib
from datetime import UTC, datetime
from pathlib import Path
from tht.ports.evidence import (
AcquiredDocument,
EvidenceSourceError,
EvidenceSourceErrorCategory,
SourceObject,
)
class FilesystemEvidenceSource:
def __init__(
self,
root: Path | str,
*,
patterns: tuple[str, ...] | list[str] = ("**/*.md",),
max_bytes: int = 10 * 1024 * 1024,
) -> None:
if max_bytes < 1:
raise ValueError("max_bytes must be positive")
if not patterns or any(not pattern for pattern in patterns):
raise ValueError("at least one non-empty discovery pattern is required")
try:
self.root = Path(root).expanduser().resolve(strict=True)
except OSError as error:
raise ValueError("filesystem evidence root is unavailable") from error
if not self.root.is_dir():
raise ValueError("filesystem evidence root must be a directory")
self.patterns = tuple(patterns)
self.max_bytes = max_bytes
def _contained(self, path: Path) -> Path:
try:
resolved = path.resolve(strict=True)
resolved.relative_to(self.root)
except (OSError, ValueError) as error:
raise EvidenceSourceError(
"unsafe filesystem object",
category=EvidenceSourceErrorCategory.PERMANENT,
details={"operation": "path_validation"},
) from error
if not resolved.is_file():
raise EvidenceSourceError(
"unsupported filesystem object",
category=EvidenceSourceErrorCategory.PERMANENT,
details={"operation": "path_validation"},
)
return resolved
def _read(self, path: Path) -> bytes:
try:
if path.stat().st_size > self.max_bytes:
raise EvidenceSourceError(
"filesystem object exceeds configured limit",
category=EvidenceSourceErrorCategory.PERMANENT,
details={"operation": "read", "limit_bytes": self.max_bytes},
)
with path.open("rb") as stream:
content = stream.read(self.max_bytes + 1)
except EvidenceSourceError:
raise
except OSError as error:
raise EvidenceSourceError(
"filesystem read failed",
category=EvidenceSourceErrorCategory.TRANSIENT,
details={"operation": "read"},
) from error
if len(content) > self.max_bytes:
raise EvidenceSourceError(
"filesystem object exceeds configured limit",
category=EvidenceSourceErrorCategory.PERMANENT,
details={"operation": "read", "limit_bytes": self.max_bytes},
)
return content
def _item(self, path: Path, content: bytes) -> SourceObject:
relative = path.relative_to(self.root).as_posix()
digest = hashlib.sha256(content).hexdigest()
stable_id = hashlib.sha256(relative.encode()).hexdigest()
modified = datetime.fromtimestamp(path.stat().st_mtime, tz=UTC)
return SourceObject(
source_id=f"filesystem:{stable_id}",
uri=path.as_uri(),
fingerprint=f"sha256:{digest}",
modified_at=modified,
metadata={"relative_path": relative},
)
def discover(self):
candidates = {path for pattern in self.patterns for path in self.root.glob(pattern)}
for candidate in sorted(candidates, key=lambda path: path.as_posix()):
path = self._contained(candidate)
content = self._read(path)
yield self._item(path, content)
def acquire(self, item: SourceObject) -> AcquiredDocument:
if not item.uri.startswith("file:"):
raise EvidenceSourceError(
"object does not belong to filesystem source",
category=EvidenceSourceErrorCategory.PERMANENT,
details={"operation": "acquire"},
)
from urllib.parse import unquote, urlsplit
parsed = urlsplit(item.uri)
path = self._contained(Path(unquote(parsed.path)))
content = self._read(path)
expected = self._item(path, content)
if item.source_id != expected.source_id:
raise EvidenceSourceError(
"object does not belong to filesystem source",
category=EvidenceSourceErrorCategory.PERMANENT,
details={"operation": "acquire"},
)
return AcquiredDocument(
source=expected,
content=content,
media_type="text/markdown" if path.suffix.lower() == ".md" else None,
acquired_at=datetime.now(UTC),
)
+202
View File
@@ -0,0 +1,202 @@
"""Explicit-manifest HTTP Evidence source with bounded streaming reads."""
import hashlib
import ipaddress
import socket
from datetime import UTC, datetime
from email.utils import parsedate_to_datetime
from urllib.parse import urljoin, urlsplit
import requests
from tht.ports.evidence import (
AcquiredDocument,
EvidenceSourceError,
EvidenceSourceErrorCategory,
SourceObject,
canonical_provenance_uri,
)
class HttpManifestEvidenceSource:
def __init__(
self,
urls: list[str] | tuple[str, ...],
*,
connect_timeout: float = 5,
read_timeout: float = 30,
max_bytes: int = 10 * 1024 * 1024,
max_redirects: int = 5,
) -> None:
if not urls:
raise ValueError("HTTP evidence manifest must contain at least one URL")
if connect_timeout <= 0 or read_timeout <= 0 or max_bytes < 1 or max_redirects < 0:
raise ValueError("HTTP evidence limits must be positive")
self._transport_by_uri: dict[str, str] = {}
for url in urls:
parsed = urlsplit(url)
if parsed.scheme not in {"http", "https"} or not parsed.hostname:
raise ValueError("HTTP evidence URLs must use http or https")
if parsed.username is not None or parsed.password is not None:
raise ValueError("HTTP evidence URLs must not contain userinfo credentials")
provenance = canonical_provenance_uri(url)
if provenance in self._transport_by_uri:
raise ValueError("HTTP evidence manifest contains duplicate canonical provenance")
self._transport_by_uri[provenance] = url
self.connect_timeout = connect_timeout
self.read_timeout = read_timeout
self.max_bytes = max_bytes
self.max_redirects = max_redirects
self._session = requests.Session()
self._cache: dict[str, AcquiredDocument] = {}
def __repr__(self) -> str:
return f"HttpManifestEvidenceSource(objects={len(self._transport_by_uri)})"
@staticmethod
def _source_id(uri: str) -> str:
return f"http:{hashlib.sha256(uri.encode()).hexdigest()}"
@staticmethod
def _reject_private_redirect(url: str) -> None:
parsed = urlsplit(url)
if parsed.scheme not in {"http", "https"} or not parsed.hostname:
raise EvidenceSourceError(
"redirect uses unsupported destination",
category=EvidenceSourceErrorCategory.PERMANENT,
details={"operation": "redirect"},
)
try:
addresses = {row[4][0] for row in socket.getaddrinfo(parsed.hostname, parsed.port)}
except OSError as error:
raise EvidenceSourceError(
"redirect destination resolution failed",
category=EvidenceSourceErrorCategory.TRANSIENT,
details={"operation": "redirect_resolution"},
) from error
if any(not ipaddress.ip_address(address).is_global for address in addresses):
raise EvidenceSourceError(
"redirect to private destination is forbidden",
category=EvidenceSourceErrorCategory.PERMANENT,
details={"operation": "redirect"},
)
@staticmethod
def _status_category(status: int) -> EvidenceSourceErrorCategory:
if status in {408, 425, 429} or 500 <= status <= 599:
return EvidenceSourceErrorCategory.TRANSIENT
return EvidenceSourceErrorCategory.PERMANENT
def _download(self, transport_url: str, provenance: str) -> AcquiredDocument:
current = transport_url
try:
for redirect_count in range(self.max_redirects + 1):
response = self._session.get(
current,
stream=True,
allow_redirects=False,
timeout=(self.connect_timeout, self.read_timeout),
)
if response.is_redirect:
response.close()
if redirect_count == self.max_redirects:
raise EvidenceSourceError(
"too many redirects",
category=EvidenceSourceErrorCategory.PERMANENT,
details={"operation": "redirect"},
)
destination = urljoin(current, response.headers.get("Location", ""))
self._reject_private_redirect(destination)
current = destination
continue
if not 200 <= response.status_code <= 299:
status = response.status_code
response.close()
raise EvidenceSourceError(
"HTTP status failure",
category=self._status_category(status),
details={"operation": "download", "status": status},
)
length = response.headers.get("Content-Length")
if length is not None and int(length) > self.max_bytes:
response.close()
raise EvidenceSourceError(
"HTTP object exceeds configured limit",
category=EvidenceSourceErrorCategory.PERMANENT,
details={"operation": "download", "limit_bytes": self.max_bytes},
)
content = bytearray()
for chunk in response.iter_content(chunk_size=min(64 * 1024, self.max_bytes + 1)):
content.extend(chunk)
if len(content) > self.max_bytes:
response.close()
raise EvidenceSourceError(
"HTTP object exceeds configured limit",
category=EvidenceSourceErrorCategory.PERMANENT,
details={"operation": "download", "limit_bytes": self.max_bytes},
)
etag = response.headers.get("ETag")
last_modified = response.headers.get("Last-Modified")
media_type = response.headers.get("Content-Type", "").split(";", 1)[0] or None
response.close()
break
except EvidenceSourceError:
raise
except (requests.Timeout, requests.ConnectionError) as error:
raise EvidenceSourceError(
"HTTP transport unavailable",
category=EvidenceSourceErrorCategory.TRANSIENT,
details={"operation": "download"},
) from error
except (requests.RequestException, ValueError) as error:
raise EvidenceSourceError(
"HTTP acquisition failed",
category=EvidenceSourceErrorCategory.PERMANENT,
details={"operation": "download"},
) from error
modified_at = None
if etag:
# Entity tags are commonly quoted; the port's stable-value grammar is deliberately
# narrower, so preserve the opaque validator through a deterministic digest.
fingerprint = f"etag:{hashlib.sha256(etag.encode()).hexdigest()}"
elif last_modified:
try:
modified_at = parsedate_to_datetime(last_modified).astimezone(UTC)
fingerprint = f"last-modified:{int(modified_at.timestamp())}"
except (TypeError, ValueError, OverflowError):
fingerprint = f"sha256:{hashlib.sha256(content).hexdigest()}"
else:
fingerprint = f"sha256:{hashlib.sha256(content).hexdigest()}"
item = SourceObject(
source_id=self._source_id(provenance),
uri=provenance,
fingerprint=fingerprint,
modified_at=modified_at,
)
return AcquiredDocument(
source=item,
content=bytes(content),
media_type=media_type,
acquired_at=datetime.now(UTC),
)
def discover(self):
for provenance in sorted(self._transport_by_uri):
document = self._download(self._transport_by_uri[provenance], provenance)
self._cache[document.source.source_id] = document
yield document.source
def acquire(self, item: SourceObject) -> AcquiredDocument:
expected_id = self._source_id(item.uri)
transport = self._transport_by_uri.get(item.uri)
if transport is None or item.source_id != expected_id:
raise EvidenceSourceError(
"object does not belong to HTTP source",
category=EvidenceSourceErrorCategory.PERMANENT,
details={"operation": "acquire"},
)
cached = self._cache.get(item.source_id)
if cached is not None and cached.source.fingerprint == item.fingerprint:
return cached
return self._download(transport, item.uri)
+35 -1
View File
@@ -1,6 +1,7 @@
"""Central construction of deployment-specific adapters."""
from tht.adapters.dwh import PostgresDwhAdapter, ThothRestDwhAdapter
from tht.adapters.evidence import FilesystemEvidenceSource, HttpManifestEvidenceSource
from tht.adapters.vector import PgVectorStore, ThothHttpVectorStore
from tht.config import Config, ConfigError
from tht.db.connection import make_engine
@@ -86,4 +87,37 @@ def build_vector_loader(cfg: Config, collection: str):
)
__all__ = ["build_dwh", "build_vector_loader", "build_vector_store"]
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(
[url.get_secret_value() for url in resource.urls],
connect_timeout=resource.connect_timeout,
read_timeout=resource.read_timeout,
max_bytes=resource.max_bytes,
max_redirects=resource.max_redirects,
)
)
case other: # pragma: no cover - Pydantic rejects unsupported discriminators.
raise ConfigError(f"Adapter evidence non supportato: {other}")
return sources
__all__ = ["build_dwh", "build_evidence_sources", "build_vector_loader", "build_vector_store"]
+34 -2
View File
@@ -5,7 +5,7 @@ from pathlib import Path
from typing import Annotated, Any, Literal
import yaml
from pydantic import BaseModel, Field, model_validator, ValidationError
from pydantic import BaseModel, Field, SecretStr, model_validator, ValidationError
from tht.config_compat import translate_legacy_config
@@ -172,13 +172,45 @@ class EligibilityConfig(BaseModel):
ignore_columns: list[str] = ["etl_last_update"]
class FilesystemEvidenceSourceConfig(BaseModel):
type: Literal["filesystem"]
root: Path
patterns: list[str] = ["**/*.md"]
max_bytes: int = Field(default=10 * 1024 * 1024, gt=0)
class HttpEvidenceSourceConfig(BaseModel):
type: Literal["http"]
# Manifest URLs may contain signed query parameters. Treat the complete transport URL as
# secret-bearing configuration; adapters derive a query-free provenance URI from it.
urls: list[SecretStr] = Field(min_length=1)
connect_timeout: float = Field(default=5, gt=0)
read_timeout: float = Field(default=30, gt=0)
max_bytes: int = Field(default=10 * 1024 * 1024, gt=0)
max_redirects: int = Field(default=5, ge=0)
EvidenceSourceConfig = Annotated[
FilesystemEvidenceSourceConfig | HttpEvidenceSourceConfig,
Field(discriminator="type"),
]
class EvidenceSourcesConfig(BaseModel):
source_root: Path
# Legacy curated-tree configuration remains accepted during migration.
source_root: Path | None = None
# cartella curata a mano nell'ETL (relativa a source_root): unica fonte delle
# evidence. Niente piu' estrazione automatica dalle schede tabella: i documenti
# qui dentro sono gia' evidence pronte (frontmatter + corpo), scelte e arricchite
# dall'autore ETL e organizzate in sottocartelle per dominio.
evidence_dir: str = "evidence"
sources: list[EvidenceSourceConfig] = []
@model_validator(mode="after")
def require_a_source(self):
if self.source_root is None and not self.sources:
raise ValueError("evidence requires source_root or sources")
return self
class EmbeddingsConfig(BaseModel):