import hashlib import json import os import re import stat import sys 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 _read_runtime_fd(fd: int, label: str, expected_mode: int = 0o400) -> tuple[bytes, os.stat_result]: try: info = os.fstat(fd) if (not stat.S_ISREG(info.st_mode) or info.st_nlink != 1 or stat.S_IMODE(info.st_mode) != expected_mode or info.st_uid != os.getuid()): raise OSError("unsafe runtime descriptor") os.lseek(fd, 0, os.SEEK_SET) chunks: list[bytes] = [] total = 0 while chunk := os.read(fd, 1024 * 1024): total += len(chunk) if total > 16 * 1024 * 1024: raise OSError("runtime descriptor too large") chunks.append(chunk) return b"".join(chunks), info except OSError as exc: raise ConfigError(f"File runtime {label} non attendibile") from exc def _strict_runtime_manifest(raw: object) -> dict[str, object]: required = { "version", "workspace_id", "workspace_revision", "descriptor_git_blob", "descriptor_sha256", "descriptor_dev", "descriptor_ino", "config_sha256", "config_dwh_binding", "config_dev", "config_ino", "config_size", "config_mode", "config_uid", "config_nlink", "directory_identities", } if not isinstance(raw, dict) or set(raw) != required or raw.get("version") != 1: raise ConfigError("Manifest runtime non valido") if (not isinstance(raw.get("workspace_id"), str) or not re.fullmatch(r"[a-z][a-z0-9-]{2,62}", raw["workspace_id"]) or not isinstance(raw.get("workspace_revision"), str) or not re.fullmatch(r"[0-9a-f]{40}", raw["workspace_revision"]) or not isinstance(raw.get("descriptor_git_blob"), str) or not re.fullmatch(r"[0-9a-f]{40}", raw["descriptor_git_blob"])): raise ConfigError("Manifest runtime non valido") if not isinstance(raw.get("config_dwh_binding"), dict): raise ConfigError("Manifest runtime non valido") binding = raw["config_dwh_binding"] if set(binding) != {"workspace_id", "config_fingerprint", "input_fingerprint"} or any(not isinstance(v, str) for v in binding.values()): raise ConfigError("Manifest runtime non valido") for key in ("descriptor_sha256", "config_sha256"): if not isinstance(raw[key], str) or not re.fullmatch(r"[0-9a-f]{64}", raw[key]): raise ConfigError("Manifest runtime non valido") for key in ("descriptor_dev", "descriptor_ino", "config_dev", "config_ino", "config_size", "config_uid", "config_nlink"): if not isinstance(raw[key], str) or not raw[key].isdigit(): raise ConfigError("Manifest runtime non valido") identities = raw.get("directory_identities") if (not isinstance(identities, list) or not identities or any(not isinstance(item, dict) or set(item) != {"path", "dev", "ino", "mode", "uid"} or not isinstance(item["path"], str) or not item["path"].startswith("/") or any(not isinstance(item[key], str) or not item[key].isdigit() for key in ("dev", "ino", "mode", "uid")) or item["mode"] == "0" for item in identities)): raise ConfigError("Manifest runtime non valido") if len({item["path"] for item in identities}) != len(identities): raise ConfigError("Manifest runtime non valido") if raw["config_mode"] != "400": raise ConfigError("Manifest runtime non valido") return raw def _canonical_runtime_path(path: Path) -> Path: value = str(path) if value == "/tmp" or value.startswith("/tmp/"): return Path("/private" + value) if value == "/var" or value.startswith("/var/"): return Path("/private" + value) return path def _runtime_directory_identities(paths: list[Path]) -> list[dict[str, str]]: """Open/stat every directory component and return its current identity chain.""" out: list[dict[str, str]] = [] seen: set[str] = set() for path in paths: canonical = _canonical_runtime_path(path) parts = list(canonical.parts) if not parts or parts[0] != "/": raise OSError("runtime config path is invalid") current = "/" components = ["/"] + parts[1:] for component in components: if component != "/": current = current.rstrip("/") + "/" + component # Re-open shared prefixes for each destination branch too: a # replacement between config-parent and manifest-parent traversal # must not be hidden by de-duplication. # lstat before and fstat after open closes the stat/open replacement # window for each component, including canonical destination parents. parent = os.open("/", os.O_RDONLY | getattr(os, "O_DIRECTORY", 0)) try: for name in [part for part in Path(current).parts[1:-1]]: nxt = _open_runtime_component(parent, name) os.close(parent); parent = nxt if current == "/": fd = os.dup(parent) else: name = Path(current).name fd = _open_runtime_component(parent, name) try: info = os.fstat(fd) if not stat.S_ISDIR(info.st_mode) or info.st_nlink < 1: raise OSError("unsafe runtime directory") identity = {"path": current, "dev": str(info.st_dev), "ino": str(info.st_ino), "mode": format(stat.S_IMODE(info.st_mode), "o"), "uid": str(info.st_uid)} if current in seen: previous = next(item for item in out if item["path"] == current) if identity != previous: raise OSError("runtime config directory changed") else: out.append(identity) seen.add(current) finally: os.close(fd) finally: os.close(parent) return out def _verify_runtime_directory_identities(config_path: Path, manifest_path: Path, expected: object) -> None: current = _runtime_directory_identities([config_path.parent, manifest_path.parent]) if current != expected: raise ConfigError("Destinazione config runtime modificata") def _open_runtime_component(parent: int, name: str) -> int: before = os.stat(name, dir_fd=parent, follow_symlinks=False) if stat.S_ISLNK(before.st_mode): raise OSError("runtime config path contains a symlink") flags = os.O_RDONLY | getattr(os, "O_DIRECTORY", 0) | os.O_NOFOLLOW fd = os.open(name, flags, dir_fd=parent) after = os.fstat(fd) if (before.st_dev != after.st_dev or before.st_ino != after.st_ino or not stat.S_ISDIR(after.st_mode)): os.close(fd) raise OSError("runtime config path changed during open") return fd def _open_runtime_file(path: Path, expected_mode: int) -> int: if not path.is_absolute(): raise OSError("runtime config path must be absolute") parts = list(path.parts) if sys.platform == "darwin" and len(parts) > 1 and parts[1] in ("var", "tmp"): parts = ["/", "private", *parts[1:]] if not parts or parts[0] != "/" or any(part in ("", ".", "..") or "/" in part for part in parts[1:]): raise OSError("runtime config path is invalid") current = os.open("/", os.O_RDONLY | getattr(os, "O_DIRECTORY", 0)) try: for component in parts[1:-1]: nxt = _open_runtime_component(current, component) os.close(current) current = nxt fd = os.open(parts[-1], os.O_RDONLY | os.O_NOFOLLOW, dir_fd=current) info = os.fstat(fd) if (not stat.S_ISREG(info.st_mode) or info.st_nlink != 1 or stat.S_IMODE(info.st_mode) != expected_mode or info.st_uid != os.getuid()): os.close(fd) raise OSError("unsafe runtime file") return fd finally: os.close(current) def _runtime_manifest_path(config_path: Path) -> Path: if config_path.name == "" or config_path.suffix != ".yaml" or config_path.parent.name != "runtime-config": raise OSError("runtime config path is not canonical") if not re.fullmatch(r"[0-9a-f]{40}", config_path.stem): raise OSError("runtime config path is not canonical") return config_path.parent.parent / "runtime-config-manifests" / f"{config_path.stem}.json" def load_config(path: Path) -> Config: # Registry leases authenticate the canonical pathname through a durable manifest # digest. The config and manifest are opened component-by-component; FD 3/4 are # reserved for the maintenance writer/root ABI. runtime_manifest: dict[str, object] | None = None expected_manifest = os.environ.get("THT_RUNTIME_CONFIG_MANIFEST_SHA256") runtime_fd = os.environ.get("THT_CONFIG_FD") manifest_fd = os.environ.get("THT_CONFIG_MANIFEST_FD") legacy_expected = os.environ.get("THT_CONFIG_MANIFEST_SHA256") if expected_manifest is not None: if not re.fullmatch(r"[0-9a-f]{64}", expected_manifest): raise ConfigError("Handoff runtime incompleto") config_fd = manifest_fd_local = None try: config_fd = _open_runtime_file(path, 0o400) manifest_fd_local = _open_runtime_file(_runtime_manifest_path(path), 0o600) config_bytes, config_info = _read_runtime_fd(config_fd, "config") manifest_bytes, manifest_info = _read_runtime_fd(manifest_fd_local, "manifest", 0o600) if hashlib.sha256(manifest_bytes).hexdigest() != expected_manifest: raise ConfigError("Manifest runtime modificato") runtime_manifest = _strict_runtime_manifest(json.loads(manifest_bytes.decode("utf-8"))) _verify_runtime_directory_identities(path, _runtime_manifest_path(path), runtime_manifest["directory_identities"]) if (runtime_manifest["config_sha256"] != hashlib.sha256(config_bytes).hexdigest() or int(runtime_manifest["config_dev"]) != config_info.st_dev or int(runtime_manifest["config_ino"]) != config_info.st_ino or int(runtime_manifest["config_size"]) != config_info.st_size or int(runtime_manifest["config_uid"]) != config_info.st_uid or int(runtime_manifest["config_nlink"]) != config_info.st_nlink or config_info.st_dev == manifest_info.st_dev and config_info.st_ino == manifest_info.st_ino): raise ConfigError("Identità config runtime non valida") source_text = config_bytes.decode("utf-8") except (OSError, UnicodeError, ValueError, json.JSONDecodeError) as exc: if isinstance(exc, ConfigError): raise raise ConfigError("File di configurazione runtime non attendibile") from exc finally: if config_fd is not None: os.close(config_fd) if manifest_fd_local is not None: os.close(manifest_fd_local) elif runtime_fd is not None or manifest_fd is not None or legacy_expected is not None: # Compatibility for direct /dev/fd callers. New backend leases never use it. if runtime_fd is None or manifest_fd is None or legacy_expected is None or not re.fullmatch(r"[0-9a-f]{64}", legacy_expected): raise ConfigError("Handoff runtime incompleto") try: config_bytes, config_info = _read_runtime_fd(int(runtime_fd), "config") manifest_bytes, manifest_info = _read_runtime_fd(int(manifest_fd), "manifest", 0o600) if hashlib.sha256(manifest_bytes).hexdigest() != legacy_expected: raise ConfigError("Manifest runtime modificato") runtime_manifest = _strict_runtime_manifest(json.loads(manifest_bytes.decode("utf-8"))) _verify_runtime_directory_identities(path, _runtime_manifest_path(path), runtime_manifest["directory_identities"]) if (runtime_manifest["config_sha256"] != hashlib.sha256(config_bytes).hexdigest() or int(runtime_manifest["config_dev"]) != config_info.st_dev or int(runtime_manifest["config_ino"]) != config_info.st_ino or int(runtime_manifest["config_size"]) != config_info.st_size or int(runtime_manifest["config_uid"]) != config_info.st_uid or int(runtime_manifest["config_nlink"]) != config_info.st_nlink or config_info.st_dev == manifest_info.st_dev and config_info.st_ino == manifest_info.st_ino): raise ConfigError("Identità config runtime non valida") source_text = config_bytes.decode("utf-8") except (OSError, UnicodeError, ValueError, json.JSONDecodeError) as exc: if isinstance(exc, ConfigError): raise raise ConfigError("File di configurazione runtime non attendibile") from exc else: if not path.exists(): raise ConfigError(f"File di configurazione non trovato: {path}") try: source_text = path.read_text() except OSError as exc: raise ConfigError(f"File di configurazione non trovato: {path}") from exc try: raw = yaml.safe_load(source_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) if runtime_manifest is not None: from tht.jobs.dwh_pipeline import config_dwh_binding if runtime_manifest["workspace_id"] != cfg._workspace_id or runtime_manifest["workspace_revision"] != cfg._workspace_revision: raise ConfigError("Identità workspace runtime non valida") if config_dwh_binding(cfg) != runtime_manifest["config_dwh_binding"]: raise ConfigError("Binding DWH runtime modificato") 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", {}))