1016 lines
40 KiB
Python
1016 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):
|
|
# 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", {}))
|