Implement approved specification #32 and tickets #33-#37. Keep host authentication server-verified and pin session interaction language. Compile scoped base selectors for browser compatibility and retain full gutters during CSS pruning.
1018 lines
40 KiB
Python
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; session dialogue uses interaction_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", {}))
|