97 lines
4.3 KiB
Python
97 lines
4.3 KiB
Python
"""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
|
|
from tht.ports.vector import VectorStore
|
|
|
|
|
|
def build_dwh(cfg: Config) -> DwhAdapter:
|
|
"""Build the DWH adapter selected by the validated workspace resource."""
|
|
resource = cfg.dwh
|
|
match resource.type:
|
|
case "postgres_direct":
|
|
return PostgresDwhAdapter(
|
|
resource.connection,
|
|
statement_timeout_ms=cfg.execution.statement_timeout_ms,
|
|
)
|
|
case "thoth_rest":
|
|
return ThothRestDwhAdapter(resource.database, resource.endpoint)
|
|
case other: # pragma: no cover - Pydantic's discriminator rejects this first.
|
|
raise ConfigError(f"Adapter DWH non supportato: {other}")
|
|
|
|
|
|
def build_vector_store(cfg: Config, *, require_write: bool = False) -> VectorStore:
|
|
"""Build the vector adapter, optionally requiring write capability."""
|
|
resource = cfg.vectors
|
|
if resource is None:
|
|
raise ConfigError("Risorsa vectors non configurata")
|
|
|
|
match resource.type:
|
|
case "qdrant":
|
|
return QdrantVectorStore(
|
|
base_url=resource.base_url,
|
|
collection=resource.collection,
|
|
workspace_id=cfg._workspace_id,
|
|
workspace_revision=cfg._workspace_revision,
|
|
expected_dimension=cfg.embeddings.dim if cfg.embeddings is not None else None,
|
|
)
|
|
case other: # pragma: no cover - Pydantic's discriminator rejects this first.
|
|
raise ConfigError(f"Adapter vector non supportato: {other}")
|
|
|
|
|
|
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
|
|
|
|
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_dwh", "build_evidence_sources", "build_vector_store"]
|