Files
ThothII/harness/tht/config.py
T
Codex 82e2c91f42
Publish documentation / publish (push) Successful in 1m27s
feat: implement memory and evidence administration with guided repairs
Add PostgreSQL-backed memory, editable evidence with source review and activation, and human-approved archive repairs across the harness, API, and UI. Include migrations, deployment support, regression coverage, and validation documentation.

Refresh permissions from validated session roles so existing administrator logins can access newly deployed archive management features.
2026-09-10 10:31:34 +02:00

1018 lines
40 KiB
Python

import hashlib
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.evidence.contracts import canonical_provenance_uri
_ENV_RE = re.compile(r"\$\{([A-Za-z_][A-Za-z0-9_]*)\}")
def _canonical_fingerprint(value: str) -> str:
return "sha256:" + hashlib.sha256(value.encode("utf-8")).hexdigest()
def canonical_effective_config_document(cfg) -> dict:
"""Non-secret effective DWH/preprocessing configuration (versioned, P3).
Mirrors backend/src/workspaces/effective-config.ts: only fields that determine whether a
prepared DWH generation is reusable. Deliberately excludes session_storage, runtime_identity,
evidence, memory, search, execution, and every credential value.
"""
dwh: dict = {"engine": "postgres"}
dwh_cfg = getattr(cfg, "dwh", None)
if dwh_cfg is None:
raise ConfigError("DWH configuration is unavailable; cannot canonicalize effective config")
transport = getattr(dwh_cfg, "type", None)
if transport == "postgres_direct":
conn = dwh_cfg.connection
dwh["database"] = conn.database
dwh["schema"] = getattr(conn, "db_schema", None) or getattr(conn, "schema", None) or conn.database
dwh["transport"] = "postgres_direct"
dwh["host"] = conn.host
dwh["port"] = conn.port
dwh["user"] = conn.user
elif transport == "thoth_rest":
dwh["database"] = dwh_cfg.database.database
dwh["schema"] = dwh_cfg.database.db_schema
dwh["transport"] = "rest_api"
dwh["baseUrl"] = dwh_cfg.endpoint.base_url
else:
raise ConfigError("unsupported DWH transport in canonical effective config")
vectors = getattr(cfg, "vectors", None)
collections = getattr(vectors, "collections", None) if vectors is not None else None
if not collections or set(collections) != {"reference", "memory"}:
raise ConfigError("vector configuration is unavailable; cannot canonicalize effective config")
embeddings = getattr(cfg, "embeddings", None)
model = getattr(embeddings, "model", None) if embeddings is not None else None
embed_dim = getattr(embeddings, "dim", None) if embeddings is not None else None
if not model or not embed_dim:
raise ConfigError("embedding configuration is unavailable; cannot canonicalize effective config")
return {
"schemaVersion": 3,
"dwh": dwh,
"vector": {
"collections": {
"reference": collections["reference"],
"memory": collections["memory"],
},
"dimensions": int(embed_dim),
"distance": "cosine",
},
"embedding": {
"id": getattr(embeddings, "id", None) or f"ollama/{model}",
"model": model,
"dimensions": int(embed_dim),
},
"roots": {
"artifacts": str(getattr(cfg.paths, "artifacts", Path("artifacts"))),
"indexes": str(getattr(cfg.paths, "indexes", Path("indexes"))),
},
}
def canonical_effective_config_json(cfg) -> str:
return json.dumps(canonical_effective_config_document(cfg), separators=(",", ":"), ensure_ascii=False)
def effective_config_identity(workspace_id: str, cfg) -> str:
digest = hashlib.sha256(canonical_effective_config_json(cfg).encode("utf-8")).hexdigest()
return f"workspace://{workspace_id}@v1:{digest}"
def effective_config_fingerprint(cfg) -> str:
return _canonical_fingerprint(canonical_effective_config_json(cfg))
def effective_config_input_fingerprint(workspace_id: str, cfg) -> str:
return _canonical_fingerprint(effective_config_identity(workspace_id, cfg))
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 | None = None
# Installazione-local file path: il contenuto non entra mai in repo/descriptor/config
# renderizzata; viene letto solo a runtime dal processo che esegue il comando.
api_key_file: str | None = None
timeout: int = 30
connect_timeout: int = 5
ssl_ca: str | None = None # path al certificato CA (per server con CA interna)
@model_validator(mode="after")
def resolve_api_key(self) -> "RestConfig":
if self.api_key is None and self.api_key_file is not None:
path = Path(self.api_key_file)
if not path.is_file() or path.is_symlink():
raise ValueError("api_key_file is unavailable")
value = path.read_text(encoding="utf-8").strip()
if not value:
raise ValueError("api_key_file is empty")
self.api_key = value
if self.api_key is None:
raise ValueError("api_key or api_key_file is required")
return self
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
collections: dict[Literal["reference", "memory"], str]
collection_lifecycle: Literal["self_heal", "require_existing"] = "self_heal"
@model_validator(mode="before")
@classmethod
def accept_single_collection_test_fixture(cls, value: Any) -> Any:
if isinstance(value, dict) and "collections" not in value and "collection" in value:
translated = dict(value)
collection = translated.pop("collection")
translated["collections"] = {
"reference": collection,
"memory": f"{collection}-memory",
}
return translated
return value
@model_validator(mode="after")
def validate_collections(self):
if set(self.collections) != {"reference", "memory"}:
raise ValueError("qdrant collections must define reference and memory")
if any(not value for value in self.collections.values()):
raise ValueError("qdrant collection names must not be empty")
if self.collections["reference"] == self.collections["memory"]:
raise ValueError("qdrant reference and memory collections must be distinct")
return self
VectorResourceConfig = Annotated[
PgvectorDirectConfig | ThothVectorHttpConfig | QdrantConfig,
Field(discriminator="type"),
]
class PathsConfig(BaseModel):
artifacts: Path = Path("artifacts")
indexes: Path = Path("indexes")
sessions: Path = Path("sessions")
# Explicit workspace-global memory root (P3). When absent, legacy `artifacts/memory` is used
# only through the documented migration path.
memory: Path | None = None
# Revision-qualified curated FK annotations root (P5). When absent, legacy
# `artifacts/mschema/annotations.yaml` remains the annotations source.
annotations_root: Path | None = None
# Runtime-only, backend-derived effective relationship snapshot. When present,
# it is the exclusive FK source; it is not an authored workspace artifact.
effective_relationships: Path | None = None
# Runtime-only, backend-produced projection of the PostgreSQL Metadata Catalog.
catalog_metadata_snapshot: Path | None = None
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):
# Persistent curator checkout; only its activated snapshot is used by preprocessing.
local_archive_root: Path | None = None
# Version 2 is the materialized source/curated authoring layout.
schema_version: Literal[1, 2] = 1
# 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
id: str | None = None
model: str
dim: int = Field(alias="dimensions")
batch_size: int = 16
timeout: int = 300
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 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}")
expanded = _resolve_secret_files(_resolve_evidence_secret_files(_expand_env(raw)))
_apply_installation_embedding_projection(expanded, path)
_validate_internal_embedding_contract(expanded, path)
_validate_internal_vector_contract(expanded, path)
translated, used_legacy = translate_legacy_config(expanded)
_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,
memory=cfg.paths.memory,
annotations_root=cfg.paths.annotations_root,
effective_relationships=cfg.paths.effective_relationships,
catalog_metadata_snapshot=cfg.paths.catalog_metadata_snapshot,
)
}
)
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 _apply_installation_embedding_projection(raw: dict[str, Any], path: Path) -> None:
projected = {
"id": os.environ.get("THT_INTERNAL_EMBEDDING_ID"),
"model": os.environ.get("THT_INTERNAL_EMBEDDING_MODEL"),
"dimensions": os.environ.get("THT_INTERNAL_EMBEDDING_DIMENSIONS"),
}
if all(value is None for value in projected.values()):
return
if any(value is None for value in projected.values()):
raise ConfigError(
f"Configurazione non valida in {path}:\n"
"la proiezione embedding dell'installazione è incompleta"
)
embedding_id = projected["id"]
model = projected["model"]
try:
dimensions = int(projected["dimensions"] or "")
except ValueError:
dimensions = 0
if embedding_id != f"ollama/{model}" or dimensions <= 0:
raise ConfigError(
f"Configurazione non valida in {path}:\n"
"la proiezione embedding dell'installazione non è canonica"
)
sections: list[tuple[dict[str, Any], str]] = []
embeddings = raw.get("embeddings")
if isinstance(embeddings, dict):
sections.append((embeddings, "dim"))
resources = raw.get("resources")
if isinstance(resources, dict) and isinstance(resources.get("embeddings"), dict):
sections.append((resources["embeddings"], "dimensions"))
for section, dimension_key in sections:
declared = {
"id": section.get("id"),
"model": section.get("model"),
"dimensions": section.get(dimension_key),
}
expected = {"id": embedding_id, "model": model, "dimensions": dimensions}
for key, value in declared.items():
if value is not None and str(value) != str(expected[key]):
raise ConfigError(
f"Configurazione non valida in {path}:\n"
f"{key} è proprietà dell'installazione e diverge dalla proiezione attiva"
)
section["id"] = embedding_id
section["model"] = model
section[dimension_key] = dimensions
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")
embedding_id = embeddings.get("id")
dimensions = embeddings.get("dimensions")
base_url = embeddings.get("base_url")
allowed = {"provider", "base_url", "id", "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 not isinstance(model, str) or not model:
raise ConfigError(
f"Configurazione non valida in {path}:\n"
"resources.embeddings.model deve essere valorizzato"
)
try:
parsed_dimensions = int(dimensions)
except (TypeError, ValueError):
parsed_dimensions = 0
if isinstance(dimensions, bool) or parsed_dimensions <= 0:
raise ConfigError(
f"Configurazione non valida in {path}:\n"
"resources.embeddings.dimensions deve essere un intero positivo"
)
if embedding_id is not None and embedding_id != f"ollama/{model}":
raise ConfigError(
f"Configurazione non valida in {path}:\n"
"resources.embeddings.id deve essere l'identità canonica ollama/<model>"
)
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 _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")
collections = vector.get("collections")
lifecycle = vector.get("collection_lifecycle")
allowed = {"engine", "base_url", "collection", "collections", "collection_lifecycle"}
if lifecycle is not None and lifecycle not in ("self_heal", "require_existing"):
raise ConfigError(
f"Configurazione non valida in {path}:\n"
"resources.vector.collection_lifecycle deve essere 'self_heal' o 'require_existing'"
)
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'"
)
valid_collections = (
isinstance(collections, dict)
and set(collections) == {"reference", "memory"}
and all(isinstance(value, str) and value for value in collections.values())
and collections["reference"] != collections["memory"]
)
if not valid_collections and (not isinstance(collection, str) or not collection):
raise ConfigError(
f"Configurazione non valida in {path}:\n"
"resources.vector.collections deve definire reference e memory"
)
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 not embeddings.model:
raise ConfigError(
f"Configurazione non valida in {path}:\n"
"embeddings.model deve essere valorizzato"
)
if embeddings.dim <= 0:
raise ConfigError(
f"Configurazione non valida in {path}:\n"
"embeddings.dim deve essere un intero positivo"
)
if embeddings.id is not None and embeddings.id != f"ollama/{embeddings.model}":
raise ConfigError(
f"Configurazione non valida in {path}:\n"
"embeddings.id deve essere l'identità canonica ollama/<model>"
)
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", {}))