Files
ThothII/harness/tht/config.py
T

890 lines
34 KiB
Python

import json
import os
import re
import stat
import warnings
from ipaddress import ip_address
from pathlib import Path
from typing import Annotated, Any, Literal
from urllib.parse import urlparse
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
_ENV_RE = re.compile(r"\$\{([A-Za-z_][A-Za-z0-9_]*)\}")
class ConfigError(Exception):
"""Errore di configurazione, con messaggio leggibile per l'utente."""
def local_tht_home() -> Path:
"""Return the private local ThothII home, honoring the explicit override."""
return Path(os.environ.get("THT_HOME", "~/.thothii")).expanduser()
def _expand_env(value: Any) -> Any:
if isinstance(value, str):
def repl(m: re.Match) -> str:
var = m.group(1)
if var not in os.environ:
raise ConfigError(
f"Variabile d'ambiente non definita: {var} "
f"(definiscila nel file .env o nell'ambiente)"
)
return os.environ[var]
return _ENV_RE.sub(repl, value)
if isinstance(value, dict):
return {k: _expand_env(v) for k, v in value.items()}
if isinstance(value, list):
return [_expand_env(v) for v in value]
return value
_MAX_SIGNED_URL_FILE_BYTES = 1024 * 1024
def _resolve_evidence_secret_files(value: Any) -> Any:
"""Resolve only signed HTTP URL arrays, keeping their values out of public errors."""
if isinstance(value, dict):
resolved = {
key: _resolve_evidence_secret_files(item)
for key, item in value.items()
}
if resolved.get("type") == "s3" and any(
name in resolved for name in ("access_key", "secret_key", "session_token")
):
raise ConfigError("S3 Evidence credentials require *_file references")
if resolved.get("type") != "http" or "signed_urls_file" not in resolved:
return resolved
if "urls" in resolved:
raise ConfigError("HTTP signed URL file cannot be combined with urls")
provenance_urls = resolved.get("provenance_urls")
if (
not isinstance(provenance_urls, list)
or not provenance_urls
or any(not isinstance(item, str) or not item for item in provenance_urls)
):
raise ConfigError("HTTP signed URL file requires provenance_urls")
path_value = resolved.pop("signed_urls_file")
if not isinstance(path_value, str):
raise ConfigError("Invalid signed URL file reference")
path = Path(path_value)
try:
entry = path.lstat()
target = path.stat()
if stat.S_ISLNK(entry.st_mode) or not stat.S_ISREG(target.st_mode):
raise OSError
if target.st_size > _MAX_SIGNED_URL_FILE_BYTES:
raise OSError
with path.open("rb") as stream:
payload = stream.read(_MAX_SIGNED_URL_FILE_BYTES + 1)
if len(payload) > _MAX_SIGNED_URL_FILE_BYTES:
raise OSError
parsed = json.loads(payload.decode("utf-8"))
except (OSError, UnicodeError, json.JSONDecodeError):
raise ConfigError("Cannot read signed URL file") from None
if (
not isinstance(parsed, list)
or not parsed
or any(not isinstance(item, str) or not item for item in parsed)
):
raise ConfigError("Invalid signed URL file")
resolved["urls"] = parsed
return resolved
if isinstance(value, list):
return [_resolve_evidence_secret_files(item) for item in value]
return value
_MAX_SCALAR_SECRET_FILE_BYTES = 64 * 1024
def _resolve_secret_files(value: Any) -> Any:
if isinstance(value, dict):
resolved = {key: _resolve_secret_files(item) for key, item in value.items()}
for secret_name in ("password", "access_key", "secret_key", "session_token"):
file_name = f"{secret_name}_file"
if file_name not in resolved:
continue
if secret_name in resolved:
raise ConfigError(f"{secret_name} and {file_name} are mutually exclusive")
path = Path(resolved.pop(file_name))
try:
with path.open("rb") as stream:
payload = stream.read(_MAX_SCALAR_SECRET_FILE_BYTES + 1)
if len(payload) > _MAX_SCALAR_SECRET_FILE_BYTES:
raise OSError
secret = payload.decode("utf-8")
except (OSError, UnicodeError):
raise ConfigError(f"Cannot read secret file: {path}") from None
if not secret or any(char.isspace() for char in secret) or "\x00" in secret:
raise ConfigError(f"Invalid secret file: {path}")
resolved[secret_name] = secret
return resolved
if isinstance(value, list):
return [_resolve_secret_files(item) for item in value]
return value
class DatabaseConfig(BaseModel):
host: str = "localhost"
port: int = 5432
database: str
db_schema: str = Field(alias="schema")
user: str
password: str
# transport: `direct` (Postgres via SQLAlchemy) o `rest` (Supabase/PostgREST).
# In `rest` deve esistere la sezione `rest` (validato a livello di Config).
transport: Literal["direct", "rest"] = "direct"
model_config = {"populate_by_name": True}
class SessionPostgresConfig(DatabaseConfig):
"""TLS-verified direct session storage; its login must be granted thoth_sessions_runtime.
Provision the dedicated LOGIN role and membership out of band so deployment credentials
never appear in the versioned migration pack.
"""
sslmode: Literal["verify-ca", "verify-full"] = "verify-full"
sslrootcert: Path | None = None
class SessionStorageConfig(BaseModel):
type: Literal["postgres_direct"]
connection: SessionPostgresConfig
class RestConfig(BaseModel):
"""Accesso al DWH via Supabase/PostgREST. base_url es. https://host/dwh/ ."""
base_url: str
api_key: str
timeout: int = 30
connect_timeout: int = 5
ssl_ca: str | None = None # path al certificato CA (per server con CA interna)
class DatabaseIdentityConfig(BaseModel):
database: str
db_schema: str = Field(alias="schema")
model_config = {"populate_by_name": True}
class PostgresDwhConfig(BaseModel):
type: Literal["postgres_direct"]
connection: DatabaseConfig
class ThothRestDwhConfig(BaseModel):
type: Literal["thoth_rest"]
database: DatabaseIdentityConfig
endpoint: RestConfig
DwhResourceConfig = Annotated[
PostgresDwhConfig | ThothRestDwhConfig,
Field(discriminator="type"),
]
class PgvectorDirectConfig(BaseModel):
type: Literal["pgvector_direct"]
reader: DatabaseConfig | None = None
writer: DatabaseConfig | None = None
# Deprecated compatibility: a single direct connection historically meant read-only.
connection: DatabaseConfig | None = None
@model_validator(mode="after")
def validate_connections(self):
if self.reader is None and self.writer is None and self.connection is None:
raise ValueError("pgvector_direct requires a reader or writer connection")
return self
class ThothVectorHttpConfig(BaseModel):
type: Literal["thoth_vector_http"]
reader: RestConfig | None = None
writer: RestConfig | None = None
# Transitional direct loading path used by the server profile.
direct: DatabaseConfig | None = None
class QdrantConfig(BaseModel):
type: Literal["qdrant"]
base_url: str
collection: str = Field(min_length=1)
# Internal runtime policy. Registry-rendered configs must not create or alter
# the workspace-owned semantic collection; legacy configs retain compatibility.
collection_lifecycle: Literal["create_if_missing", "require_existing"] = "create_if_missing"
VectorResourceConfig = Annotated[
PgvectorDirectConfig | ThothVectorHttpConfig | QdrantConfig,
Field(discriminator="type"),
]
class PathsConfig(BaseModel):
artifacts: Path = Path("artifacts")
indexes: Path = Path("indexes")
sessions: Path = Path("sessions")
class RuntimeIdentityConfig(BaseModel):
workspace_id: str = Field(pattern=r"^[a-z][a-z0-9-]{2,62}$")
workspace_revision: str = Field(pattern=r"^[0-9a-f]{40}$")
source_identity: str | None = Field(
default=None,
pattern=r"^workspace://[a-z][a-z0-9-]{2,62}$",
)
@model_validator(mode="after")
def source_matches_workspace(self):
expected = f"workspace://{self.workspace_id}"
if self.source_identity is None:
self.source_identity = expected
elif self.source_identity != expected:
raise ValueError("source_identity must match workspace_id")
return self
class WorkspaceRoots(PathsConfig):
pass
class ExamplesConfig(BaseModel):
max_per_column: int = 10
class LshSkipConfig(BaseModel):
# DEPRECATO: l'euristica di lunghezza (skip_column vendored) non è più usata. La
# selezione delle colonne da indicizzare segue il principio di column eligibility
# (sezione `eligibility`). Mantenuto solo per compatibilità con tht.yaml esistenti.
max_total_chars: int = 50000
max_avg_length: int = 20
class LshConfig(BaseModel):
signature_size: int = 64
n_gram: int = 3
threshold: float = 0.5
max_values_per_column: int = 1000
skip: LshSkipConfig = LshSkipConfig()
class EligibilityConfig(BaseModel):
# Soglie del principio di column eligibility (testo ampio ignorato ovunque).
max_declared_len: int = 128 # char/varchar dichiarati <= soglia: eligible senza campionare
max_avg_length: int = 40 # fallback data-driven: lunghezza media valori campionati
max_sampled_len: int = 200 # fallback data-driven: lunghezza massima valore campionato
# Colonne di servizio sempre ignorate per nome (match case-insensitive), a prescindere
# dal tipo: metadati ETL/audit non analitici (es. timestamp di ultimo aggiornamento).
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)
model_config = {"extra": "forbid"}
class HttpEvidenceSourceConfig(BaseModel):
type: Literal["http"]
# Transport URLs are secret-bearing. Signed-file configurations retain only public,
# query-free provenance identities alongside the masked transport values.
urls: list[SecretStr] = Field(min_length=1)
provenance_urls: list[str] | None = Field(default=None, 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)
allow_private_hosts: bool = False
max_cache_bytes: int = Field(default=64 * 1024 * 1024, gt=0)
model_config = {"extra": "forbid"}
@model_validator(mode="after")
def validate_provenance_mapping(self):
transport_urls = [url.get_secret_value() for url in self.urls]
try:
canonical = [canonical_provenance_uri(url) for url in transport_urls]
except ValueError as exc:
raise ValueError("HTTP transport URL is invalid") from exc
if any(
urlparse(url).scheme not in ("http", "https") or not urlparse(url).hostname
for url in transport_urls
):
raise ValueError("HTTP transport URL must use http or https")
if self.provenance_urls is None:
if canonical != transport_urls or len(set(canonical)) != len(canonical):
raise ValueError("Public HTTP URLs must be canonical query-free identities")
return self
try:
provenance = [canonical_provenance_uri(url) for url in self.provenance_urls]
except ValueError as exc:
raise ValueError("HTTP provenance URL is invalid") from exc
if provenance != self.provenance_urls:
raise ValueError("HTTP provenance URLs must be canonical query-free identities")
if len(set(provenance)) != len(provenance):
raise ValueError("HTTP provenance URLs must not repeat")
if len(canonical) != len(provenance) or canonical != provenance:
raise ValueError("Signed HTTP URLs must map one-to-one to provenance URLs in order")
return self
def transport_urls(self) -> list[str]:
"""Expose secret transport values only at the adapter-construction boundary."""
return [url.get_secret_value() for url in self.urls]
class S3EvidenceSourceConfig(BaseModel):
type: Literal["s3"]
bucket: str = Field(min_length=1)
prefix: str = ""
endpoint_url: str | None = None
region: str | None = None
access_key: SecretStr | None = None
secret_key: SecretStr | None = None
session_token: SecretStr | None = None
trusted_endpoint: bool = False
allow_private_endpoint: bool = False
allow_insecure_endpoint: bool = False
max_bytes: int = Field(default=10 * 1024 * 1024, gt=0)
max_objects: int = Field(default=10_000, gt=0)
max_pages: int = Field(default=100, gt=0)
page_size: int = Field(default=1000, gt=0, le=1000)
model_config = {"extra": "forbid"}
@model_validator(mode="after")
def validate_static_credentials(self):
if (self.access_key is None) != (self.secret_key is None):
raise ValueError("S3 access_key and secret_key must be configured together")
if self.session_token is not None and self.access_key is None:
raise ValueError("S3 session_token requires static credentials")
return self
EvidenceSourceConfig = Annotated[
FilesystemEvidenceSourceConfig | HttpEvidenceSourceConfig | S3EvidenceSourceConfig,
Field(discriminator="type"),
]
class EvidenceSourcesConfig(BaseModel):
# 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
model_config = {"extra": "forbid"}
class EmbeddingsConfig(BaseModel):
provider: str = "ollama_internal"
base_url: str
model: str = "nomic-embed-text-v2-moe"
dim: int = Field(default=768, alias="dimensions")
batch_size: int = 32
timeout: int = 30
connect_timeout: int = 5
bin: str = "ollama"
start_cmd: list[str] | None = None
model_config = {"populate_by_name": True, "extra": "forbid"}
class VectorConfig(BaseModel):
max_chunk_chars: int = Field(default=4000, gt=0)
# ACTIVE plus the two most recent rollback generations by default.
retain_published_generations: int = Field(default=3, ge=1)
model_config = {"extra": "forbid"}
class SearchConfig(BaseModel):
rrf_k: int = 60
top_schema_tables: int = 12 # default `--top` per `tht search --kind schema` (n. tabelle)
schema_chunk_pool: int = 150 # chunk tabella/colonna fusi prima dell'aggregazione a tabella
class ExecutionConfig(BaseModel):
allow: list[str] = ["cte_test", "explain", "preview", "aggregate", "export"]
max_preview_rows: int = 10
statement_timeout_ms: int = 30000
warn_execution_ms: int = 5000
warn_plan_rows: int = 1_000_000
max_aggregate_cells: int = 20
max_export_rows: int = 100000
forbidden_functions: list[str] = [
"setval",
"nextval",
"pg_advisory_lock",
"pg_advisory_xact_lock",
"dblink",
"dblink_exec",
"pg_terminate_backend",
"pg_cancel_backend",
"lo_import",
"lo_export",
"pg_reload_conf",
]
class Config(BaseModel):
_workspace_id: str = PrivateAttr(default="default")
_workspace_revision: str | None = PrivateAttr(default=None)
_config_source: str = PrivateAttr(default="direct")
runtime_identity: RuntimeIdentityConfig | None = None
dwh: DwhResourceConfig
vectors: VectorResourceConfig | None = None
session_storage: SessionStorageConfig | None = None
roots: WorkspaceRoots = WorkspaceRoots()
# Compatibility views retained until all call sites consume typed resources.
database: DatabaseConfig
# Profilo dell'installazione, letto da THT_PROFILE (.env), non dallo yaml versionato.
# server: ricostruisce i derivati (artefatti, LSH, vettori schema nel vectordb).
# workstation: postazione locale che legge il vectordb via REST; gli upsert remoti
# richiedono vector_write_rest, mentre init/clear/rebuild restano solo-server.
profile: Literal["server", "workstation"] = "server"
# Language in which table/column descriptions and evidence are written. The skill
# instructions stay in English; only content/output follow this language. Default
# 'en' so Thoth is not bound to any customer's language.
language: str = "en"
paths: PathsConfig = PathsConfig()
examples: ExamplesConfig = ExamplesConfig()
lsh: LshConfig = LshConfig()
eligibility: EligibilityConfig = EligibilityConfig()
evidence: EvidenceSourcesConfig | None = None
embeddings: EmbeddingsConfig | None = None
vector_db: DatabaseConfig | None = None
vector: VectorConfig = VectorConfig()
search: SearchConfig = SearchConfig()
execution: ExecutionConfig = ExecutionConfig()
rest: RestConfig | None = None
# Endpoint REST dedicato per la similarity search del pgvector (rpc search_similar).
# Se presente, `tht search` legge via REST; altrimenti legge in diretto (dev/test).
vector_rest: RestConfig | None = None
# Endpoint REST dedicato alle scritture controllate del pgvector. E' opzionale e usa
# una API key separata dalla lettura; espone solo upsert/hash via RPC allowlist.
vector_write_rest: RestConfig | None = None
@model_validator(mode="before")
@classmethod
def accept_legacy_constructor_fields(cls, value: Any) -> Any:
if not isinstance(value, dict) or "dwh" in value:
return value
translated, _ = translate_legacy_config(value)
_populate_legacy_views(translated)
return translated
def workspace_id_from_path(path: Path) -> str:
"""Return the stable workspace identity, resolving deployment aliases first."""
return path.resolve().stem.lower().replace(".", "-").replace("_", "-")
def workspace_id_for_config(config: Config, path: Path) -> str:
"""Return the canonical runtime identity, falling back for legacy configs."""
if config.runtime_identity is not None:
return config.runtime_identity.workspace_id
return workspace_id_from_path(path)
def _format_validation_error(error: ValidationError) -> str:
messages = {
"missing": "required field",
"extra_forbidden": "unknown field",
"greater_than": "value must be greater than the configured bound",
"greater_than_equal": "value must meet the configured lower bound",
"less_than_equal": "value exceeds the configured upper bound",
"literal_error": "unsupported literal value",
"union_tag_invalid": "unsupported discriminator",
"union_tag_not_found": "missing discriminator",
"string_too_short": "string is too short",
"too_short": "collection is too short",
"value_error": "configuration value is invalid",
}
lines = []
for issue in error.errors(include_input=False, include_url=False):
location = ".".join(str(part) for part in issue.get("loc", ())) or "configuration"
error_type = str(issue.get("type", "validation_error"))
message = messages.get(error_type, "invalid configuration value")
if location == "runtime_identity" and error_type == "value_error":
message = "source_identity does not match workspace identity"
lines.append(f"{location} [{error_type}]: {message}")
return "\n".join(lines)
def _validate_raw_config_shape(raw: dict[str, Any], path: Path) -> None:
"""Reject unsafe YAML shapes before compatibility translation or sorting keys."""
seen: set[int] = set()
def walk(value: Any, location: str) -> None:
if isinstance(value, dict):
marker = id(value)
if marker in seen:
return
seen.add(marker)
for key, item in value.items():
if not isinstance(key, str):
raise ConfigError(
f"Configurazione non valida in {path}: mapping key at {location} "
"must be a string"
)
walk(item, f"{location}.{key}")
elif isinstance(value, list):
for index, item in enumerate(value):
walk(item, f"{location}[{index}]")
walk(raw, "configuration")
resources = raw.get("resources")
if "resources" in raw and not isinstance(resources, dict):
raise ConfigError(
f"Configurazione non valida in {path}: resources must be a mapping"
)
if isinstance(resources, dict):
for name in ("vector", "embeddings"):
if name in resources and not isinstance(resources[name], dict):
raise ConfigError(
f"Configurazione non valida in {path}: resources.{name} must be a mapping"
)
def load_config(path: Path) -> Config:
if not path.exists():
raise ConfigError(f"File di configurazione non trovato: {path}")
try:
raw = yaml.safe_load(path.read_text())
except yaml.YAMLError as exc:
raise ConfigError(f"Configurazione YAML non valida: {path}") from exc
if not isinstance(raw, dict):
raise ConfigError(f"Configurazione non valida (atteso un mapping YAML): {path}")
_validate_raw_config_shape(raw, path)
expanded = _resolve_secret_files(_resolve_evidence_secret_files(_expand_env(raw)))
_validate_internal_embedding_contract(expanded, path)
_validate_internal_vector_contract(expanded, path)
translated, used_legacy = translate_legacy_config(expanded)
_validate_vector_resource_consistency(expanded, translated, path)
vectors = translated.get("vectors")
# runtime_identity is the registry marker. The lifecycle is an internal
# runtime policy, never a descriptor-controlled option.
if (
isinstance(translated.get("runtime_identity"), dict)
and isinstance(vectors, dict)
and vectors.get("type") == "qdrant"
):
vectors["collection_lifecycle"] = "require_existing"
_populate_legacy_views(translated)
try:
cfg = Config.model_validate(translated)
except ValidationError as e:
details = _format_validation_error(e)
raise ConfigError(f"Configurazione non valida in {path}:\n{details}") from None
env_profile = os.environ.get("THT_PROFILE")
if env_profile is not None:
if env_profile not in ("server", "workstation"):
raise ConfigError(
f"THT_PROFILE non valido: {env_profile!r} (atteso 'server' o 'workstation')."
)
cfg = cfg.model_copy(update={"profile": env_profile})
if cfg.database.transport == "rest" and cfg.rest is None:
raise ConfigError(
f"transport: rest richiede la sezione `rest` (base_url, api_key) in {path}."
)
# THT_HOME is the local-user storage root. Retain THT_DATA_ROOT as the
# portable-deployment compatibility name until all callers use repositories.
data_root = os.environ.get("THT_HOME") or os.environ.get("THT_DATA_ROOT")
if data_root:
# Import locally: paths owns resolution, while ConfigError remains the public
# configuration exception callers already handle.
from tht.paths import resolve_workspace_paths
resolved = resolve_workspace_paths(path, cfg, Path(data_root))
cfg = cfg.model_copy(
update={
"paths": PathsConfig(
sessions=resolved.sessions,
artifacts=resolved.artifacts,
indexes=resolved.indexes,
)
}
)
elif not used_legacy:
# Modern `roots` replace `paths`; without a mounted data root retain the old
# working-directory-relative behavior used by local development.
cfg = cfg.model_copy(update={"paths": PathsConfig(**cfg.roots.model_dump())})
if used_legacy:
warnings.warn(
"DEPRECATION: legacy workspace resource keys are deprecated; "
"use dwh, vectors, and roots.",
FutureWarning,
stacklevel=2,
)
cfg._workspace_id = workspace_id_for_config(cfg, path)
cfg._workspace_revision = (
cfg.runtime_identity.workspace_revision
if cfg.runtime_identity is not None
else None
)
cfg._config_source = (
cfg.runtime_identity.source_identity
if cfg.runtime_identity is not None
else path.resolve().as_posix()
)
_validate_active_embeddings_config(cfg.embeddings, path)
_validate_active_vector_config(cfg.vectors, path)
return cfg
def _validate_internal_embedding_contract(raw: dict[str, Any], path: Path) -> None:
resources = raw.get("resources")
if not isinstance(resources, dict):
return
embeddings = resources.get("embeddings")
if not isinstance(embeddings, dict):
return
provider = embeddings.get("provider")
model = embeddings.get("model")
dimensions = embeddings.get("dimensions")
base_url = embeddings.get("base_url")
allowed = {"provider", "base_url", "model", "dimensions"}
unexpected = sorted(set(embeddings) - allowed)
if unexpected:
raise ConfigError(
f"Configurazione non valida in {path}:\n"
f"resources.embeddings non supporta: {', '.join(unexpected)}"
)
if provider != "ollama_internal":
raise ConfigError(
f"Configurazione non valida in {path}:\n"
"resources.embeddings.provider deve essere 'ollama_internal'"
)
if model != "qwen3-embedding:0.6b":
raise ConfigError(
f"Configurazione non valida in {path}:\n"
"resources.embeddings.model deve essere 'qwen3-embedding:0.6b'"
)
if dimensions != 1024:
raise ConfigError(
f"Configurazione non valida in {path}:\n"
"resources.embeddings.dimensions deve essere 1024"
)
if not _is_allowed_internal_embedding_url(base_url):
raise ConfigError(
f"Configurazione non valida in {path}:\n"
"resources.embeddings.base_url deve usare http://embedding:11434 "
"oppure un endpoint loopback di sviluppo su porta 11434"
)
def _normalize_qdrant_base_url(value: Any) -> str | None:
"""Normalize a URL, returning ``None`` for untrusted raw YAML values."""
if not isinstance(value, str):
return None
try:
parsed = urlparse(value)
hostname = (parsed.hostname or "").lower()
port = parsed.port
except (TypeError, ValueError):
return None
host = f"[{hostname}]" if ":" in hostname and not hostname.startswith("[") else hostname
netloc = f"{host}:{port}" if port is not None else host
return f"{parsed.scheme.lower()}://{netloc}{parsed.path.rstrip('/') or '/'}"
def _validate_vector_resource_consistency(raw: dict[str, Any], translated: dict[str, Any], path: Path) -> None:
"""Reject divergent top-level and compatibility Qdrant resource views."""
resources = raw.get("resources")
resource = resources.get("vector") if isinstance(resources, dict) else None
vectors = translated.get("vectors")
if not isinstance(resource, dict) or not isinstance(vectors, dict):
return
if vectors.get("type") != "qdrant" or resource.get("engine") != "qdrant":
raise ConfigError(
f"Configurazione non valida in {path}: vectors and resources.vector must both describe qdrant"
)
vector_url = _normalize_qdrant_base_url(vectors.get("base_url"))
resource_url = _normalize_qdrant_base_url(resource.get("base_url"))
vector_collection = vectors.get("collection")
resource_collection = resource.get("collection")
if (
vector_url is None
or resource_url is None
or not isinstance(vector_collection, str)
or not isinstance(resource_collection, str)
or vector_url != resource_url
or vector_collection != resource_collection
):
raise ConfigError(
f"Configurazione non valida in {path}: vectors and resources.vector disagree"
)
def _validate_internal_vector_contract(raw: dict[str, Any], path: Path) -> None:
resources = raw.get("resources")
if not isinstance(resources, dict):
return
vector = resources.get("vector")
if not isinstance(vector, dict):
return
engine = vector.get("engine")
base_url = vector.get("base_url")
collection = vector.get("collection")
allowed = {"engine", "base_url", "collection"}
unexpected = sorted(set(vector) - allowed)
if unexpected:
raise ConfigError(
f"Configurazione non valida in {path}:\n"
f"resources.vector non supporta: {', '.join(unexpected)}"
)
if engine != "qdrant":
raise ConfigError(
f"Configurazione non valida in {path}:\n"
"resources.vector.engine deve essere 'qdrant'"
)
if not isinstance(collection, str) or not collection:
raise ConfigError(
f"Configurazione non valida in {path}:\n"
"resources.vector.collection deve essere valorizzato"
)
if not _is_allowed_internal_qdrant_url(base_url):
raise ConfigError(
f"Configurazione non valida in {path}:\n"
"resources.vector.base_url deve usare http://qdrant:6333 "
"oppure un endpoint loopback di sviluppo su porta 6333"
)
def _validate_active_embeddings_config(
embeddings: "EmbeddingsConfig | None",
path: Path,
) -> None:
if embeddings is None:
return
if embeddings.provider != "ollama_internal":
raise ConfigError(
f"Configurazione non valida in {path}:\n"
"embeddings.provider deve essere 'ollama_internal'"
)
if embeddings.model != "qwen3-embedding:0.6b":
raise ConfigError(
f"Configurazione non valida in {path}:\n"
"embeddings.model deve essere 'qwen3-embedding:0.6b'"
)
if embeddings.dim != 1024:
raise ConfigError(
f"Configurazione non valida in {path}:\n"
"embeddings.dim deve essere 1024"
)
if not _is_allowed_internal_embedding_url(embeddings.base_url):
raise ConfigError(
f"Configurazione non valida in {path}:\n"
"embeddings.base_url deve usare http://embedding:11434 "
"oppure un endpoint loopback di sviluppo su porta 11434"
)
def _validate_active_vector_config(
vectors: "VectorResourceConfig | None",
path: Path,
) -> None:
if vectors is None or vectors.type != "qdrant":
return
if not _is_allowed_internal_qdrant_url(vectors.base_url):
raise ConfigError(
f"Configurazione non valida in {path}:\n"
"vectors.base_url deve usare http://qdrant:6333 "
"oppure un endpoint loopback di sviluppo su porta 6333"
)
def _is_allowed_internal_embedding_url(value: Any) -> bool:
if not isinstance(value, str):
return False
parsed = urlparse(value)
if parsed.scheme != "http" or not parsed.hostname or parsed.port != 11434:
return False
if parsed.params or parsed.query or parsed.fragment:
return False
if parsed.path not in ("", "/"):
return False
if parsed.hostname == "embedding":
return True
try:
host = ip_address(parsed.hostname)
except ValueError:
return parsed.hostname == "localhost"
return host.is_loopback
def _is_allowed_internal_qdrant_url(value: Any) -> bool:
if not isinstance(value, str):
return False
parsed = urlparse(value)
if parsed.scheme != "http" or not parsed.hostname or parsed.port != 6333:
return False
if parsed.params or parsed.query or parsed.fragment:
return False
if parsed.path not in ("", "/"):
return False
if parsed.hostname == "qdrant":
return True
try:
host = ip_address(parsed.hostname)
except ValueError:
return parsed.hostname == "localhost"
return host.is_loopback
def _populate_legacy_views(raw: dict[str, Any]) -> None:
"""Populate old Config attributes for command compatibility during migration."""
dwh = raw.get("dwh")
if "database" not in raw and isinstance(dwh, dict):
if dwh.get("type") == "postgres_direct":
raw["database"] = {**dwh["connection"], "transport": "direct"}
elif dwh.get("type") == "thoth_rest":
raw["database"] = {
**dwh["database"],
"user": "rest",
"password": "",
"transport": "rest",
}
raw["rest"] = dwh["endpoint"]
vectors = raw.get("vectors")
if isinstance(vectors, dict):
if vectors.get("type") == "pgvector_direct":
raw.setdefault(
"vector_db",
vectors.get("writer") or vectors.get("reader") or vectors.get("connection"),
)
elif vectors.get("type") == "thoth_vector_http":
raw.setdefault("vector_rest", vectors.get("reader"))
raw.setdefault("vector_write_rest", vectors.get("writer"))
raw.setdefault("vector_db", vectors.get("direct"))
raw.setdefault("paths", raw.get("roots", {}))