refactor(harness): renaming prodotto tht (Onda -1)

Thoth (tht) è il prodotto, PSD è il cliente. Nessun riferimento al contesto
clinico nel codice.

Rinomine:
- comando+package nsp→tht (dir nsp/→tht/, 46 import, pyproject entry point)
- gate nsp-gate.js→tht-gate.js (+ rewrite token, relayIfNspFails→relayIfThtFails)
- workspace chirone.{example,test}.yaml→tht.{example,test}.yaml (generici)
- env THOTH_→THT_ (19 var) + NSP_ stragglers (NSP_HARNESS_ROOT, NSP_SESSION)
- commenti/docstring chirone/psdwp3/policlinico neutralizzati ('the reference
  implementation', 'the DWH')

Aggiunto [tool.setuptools.packages.find] include=['tht*'] (necessario: l'auto-
discovery rompeva con tht/ + workspaces/ come top-level multipli).

.env operatore aggiornato in-place (prefissi THT_, valori preservati, gitignored).

Verifica: pytest 109 passed, npm test 14 pass, tht phase meta --json OK, zero
residui nsp/THOTH_/NSP_/chirone nel package.
This commit is contained in:
2026-06-27 10:33:16 +02:00
parent 0dcc0246dc
commit fc5fbe6b65
72 changed files with 309 additions and 306 deletions
+1
View File
@@ -0,0 +1 @@
__version__ = "0.1.0"
+41
View File
@@ -0,0 +1,41 @@
"""tht CLI -- Typer app `tht`.
Minimal skeleton for now: registers `phase_app` (the gate needs `tht phase meta --json`).
Other command groups (db, schema, lsh, vector, memory, session, decision, sql, cte,
evidence, search) are wired in their porting tasks, once their backend deps land
(A9 ports db/mschema/_guards; B1 ports vectorstore; B3 ports search/lshindex/sampling).
"""
from __future__ import annotations
import typer
from tht import __version__
app = typer.Typer(
help="tht (Thoth) -- kit operativo per generazione SQL assistito da Pi",
no_args_is_help=True,
)
def _version_callback(value: bool) -> None:
if value:
typer.echo(__version__)
raise typer.Exit()
@app.callback()
def main(
version: bool = typer.Option(
False,
"--version",
callback=_version_callback,
is_eager=True,
help="Mostra la versione ed esce.",
),
) -> None:
pass
from tht.cli.phase_cmd import phase_app # noqa: E402
app.add_typer(phase_app, name="phase")
+40
View File
@@ -0,0 +1,40 @@
import typer
def require_server_profile(cfg, command: str) -> None:
"""Rifiuta i comandi di scrittura vectordb sul profilo workstation (exit 4).
Va chiamata subito dopo il caricamento della config e PRIMA di aprire qualunque
connessione, cosi' su workstation non si tenta mai la connessione diretta al vectordb.
"""
if cfg.profile == "workstation":
typer.secho(
f"ERRORE: `{command}` e' un comando solo-server (scrive nel vectordb centrale). "
f"Sulla postazione locale (profile: workstation) il vectordb si LEGGE via REST, "
f"non si ricostruisce. Esegui questo comando sul server di produzione "
f"(profile: server).",
fg=typer.colors.RED, err=True,
)
raise typer.Exit(code=4)
def has_vector_write_rest(cfg) -> bool:
return cfg.vector_write_rest is not None and bool(cfg.vector_write_rest.api_key.strip())
def require_vector_write_allowed(cfg, command: str) -> None:
"""Permette scritture vectordb da workstation solo con API key REST writer.
Senza `vector_write_rest`, la workstation resta read-only e i comandi di indexing sono
eseguibili solo sul server con connessione diretta al vectordb.
"""
if cfg.profile == "workstation" and not has_vector_write_rest(cfg):
typer.secho(
f"ERRORE: `{command}` e' un comando solo-server se manca `vector_write_rest`: "
f"scrive nel vectordb centrale. Sulla postazione locale serve la sezione "
f"`vector_write_rest` con una API key di upsert; in alternativa esegui il "
f"comando sul server di produzione (profile: server).",
fg=typer.colors.RED, err=True,
)
raise typer.Exit(code=4)
+51
View File
@@ -0,0 +1,51 @@
"""tht phase -- workflow HITL gate commands.
`meta` is the single source the gate (tht-gate.js) reads workflow facts from,
killing the JS/Python PHASE_NAMES drift (F2). advance/reopen/show land with the
session/db ports (A9) -- they need current_phase + load_session_or_exit which
depend on the not-yet-ported session store helpers.
"""
from __future__ import annotations
import json
import typer
from tht.workflow import load_workflow
phase_app = typer.Typer(help="Fase del workflow HITL (gate di avanzamento/ritorno)")
@phase_app.command("meta")
def meta_cmd(
as_json: bool = typer.Option(
False,
"--json",
help="Emette i metadati del workflow come JSON (usato dall'estensione gate).",
),
) -> None:
"""Stampa i metadati del workflow derivati da workflow.yaml (F2)."""
wf = load_workflow()
if not as_json:
typer.echo(f"max_phase: {wf.max_phase}")
for p in wf.phases:
typer.echo(f" {p.id} ({p.num}/{wf.max_phase}) {p.name} advance={p.advance}")
return
typer.echo(
json.dumps(
{
"schema_version": wf.schema_version,
"max_phase": wf.max_phase,
"phases": [
{
"num": p.num,
"id": p.id,
"name": p.name,
"advance": p.advance,
"artifacts_out": p.artifacts_out,
}
for p in wf.phases
],
}
)
)
+183
View File
@@ -0,0 +1,183 @@
import os
import re
from pathlib import Path
from typing import Any, Literal
import yaml
from pydantic import BaseModel, Field, ValidationError
_ENV_RE = re.compile(r"\$\{([A-Za-z_][A-Za-z0-9_]*)\}")
class ConfigError(Exception):
"""Errore di configurazione, con messaggio leggibile per l'utente."""
def _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
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 RestConfig(BaseModel):
"""Accesso al DWH via Supabase/PostgREST. base_url es. https://host/dwh/ ."""
base_url: str
api_key: str
timeout: int = 30
ssl_ca: str | None = None # path al certificato CA (per server con CA interna)
class PathsConfig(BaseModel):
artifacts: Path = Path("artifacts")
indexes: Path = Path("indexes")
sessions: Path = Path("sessions")
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 EvidenceSourcesConfig(BaseModel):
source_root: Path
# 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"
class EmbeddingsConfig(BaseModel):
base_url: str
model: str = "nomic-embed-text-v2-moe"
dim: int = 768
batch_size: int = 32
timeout: int = 120
class VectorConfig(BaseModel):
max_chunk_chars: int = 4000
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):
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"
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
def load_config(path: Path) -> Config:
if not path.exists():
raise ConfigError(f"File di configurazione non trovato: {path}")
raw = yaml.safe_load(path.read_text())
if not isinstance(raw, dict):
raise ConfigError(f"Configurazione non valida (atteso un mapping YAML): {path}")
try:
cfg = Config.model_validate(_expand_env(raw))
except ValidationError as e:
raise ConfigError(f"Configurazione non valida in {path}:\n{e}") from e
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}."
)
return cfg
View File
+37
View File
@@ -0,0 +1,37 @@
from sqlalchemy import Engine, create_engine, text
from tht.config import DatabaseConfig
def make_engine(cfg: DatabaseConfig) -> Engine:
url = (
f"postgresql+psycopg2://{cfg.user}:{cfg.password}"
f"@{cfg.host}:{cfg.port}/{cfg.database}"
)
return create_engine(url, echo=False)
def ping(engine: Engine) -> None:
with engine.connect() as conn:
conn.execute(text("SELECT 1"))
def writable_tables(engine: Engine, schema: str) -> list[str]:
"""Tabelle dello schema su cui l'utente corrente ha privilegi di scrittura."""
q = text("""
SELECT c.relname
FROM pg_class c
JOIN pg_namespace n ON n.oid = c.relnamespace
WHERE n.nspname = :schema
AND c.relkind IN ('r', 'p')
AND has_table_privilege(current_user, c.oid, 'INSERT, UPDATE, DELETE')
ORDER BY c.relname
""")
with engine.connect() as conn:
return [row[0] for row in conn.execute(q, {"schema": schema})]
def can_create_in_schema(engine: Engine, schema: str) -> bool:
q = text("SELECT has_schema_privilege(current_user, :schema, 'CREATE')")
with engine.connect() as conn:
return bool(conn.execute(q, {"schema": schema}).scalar())
+114
View File
@@ -0,0 +1,114 @@
"""Recupero della catena di certificati presentata da un endpoint HTTPS.
Serve al setup di una postazione *workstation* dietro una CA interna: scarica la
catena TLS del server REST e la salva in un bundle PEM da puntare con `THT_SSL_CA`
(consumato da `requests` via `verify=`). NON installa nulla nel trust store dell'OS.
"""
from __future__ import annotations
import _ssl
import socket
import ssl
import tempfile
from pathlib import Path
from urllib.parse import urlsplit
# Encoding atteso da Certificate.public_bytes() per la catena TLS non verificata.
_PEM_ENCODING = getattr(_ssl, "ENCODING_PEM", 1)
class CaFetchError(Exception):
"""Errore azionabile durante il recupero della catena CA."""
def parse_host_port(base_url: str) -> tuple[str, int]:
"""Estrae (host, port) da un URL REST https. Porta di default 443."""
parts = urlsplit(base_url)
if parts.scheme != "https":
raise CaFetchError(
f"URL non https: {base_url!r}. Il recupero CA ha senso solo su HTTPS."
)
if not parts.hostname:
raise CaFetchError(f"Host mancante nell'URL: {base_url!r}.")
return parts.hostname, parts.port or 443
def fetch_chain_pem(host: str, port: int = 443, timeout: int = 30) -> list[str]:
"""Restituisce la catena di certificati presentata da host:port come lista di PEM.
L'handshake è volutamente *non verificato* (CERT_NONE): stiamo recuperando la catena
per poter poi *stabilire* la fiducia, non per fidarci adesso. La verifica vera avviene
in seguito quando `THT_SSL_CA` punta al bundle salvato (es. `tht db ping`).
"""
ctx = ssl.SSLContext(ssl.PROTOCOL_TLS_CLIENT)
ctx.check_hostname = False
ctx.verify_mode = ssl.CERT_NONE
try:
with socket.create_connection((host, port), timeout=timeout) as sock:
with ctx.wrap_socket(sock, server_hostname=host) as tls:
certs = _unverified_chain(tls)
except (OSError, ssl.SSLError) as e:
raise CaFetchError(
f"Impossibile connettersi a {host}:{port} per recuperare i certificati: {e}"
) from e
if not certs:
raise CaFetchError(
f"Nessun certificato presentato da {host}:{port}. "
f"In alternativa, estrai la catena a mano con: "
f"openssl s_client -showcerts -connect {host}:{port} -servername {host}"
)
return [_to_pem(c) for c in certs]
def _unverified_chain(tls: ssl.SSLSocket) -> list:
"""Catena presentata dal server. Metodo pubblico su Python >= 3.13, API interna su 3.12."""
public = getattr(tls, "get_unverified_chain", None)
if public is not None:
return list(public() or [])
sslobj = getattr(tls, "_sslobj", None)
getter = getattr(sslobj, "get_unverified_chain", None) if sslobj is not None else None
if getter is None:
raise CaFetchError(
"Questa versione di Python non espone la catena TLS. "
"Estrai la catena a mano con `openssl s_client -showcerts`."
)
return list(getter() or [])
def describe_pem(pem: str) -> str:
"""Riassunto leggibile (subject / issuer) di un certificato PEM, best-effort.
Serve a far riconoscere all'utente la CA interna attesa (verifica out-of-band).
Restituisce "" se il certificato non è decodificabile.
"""
try:
with tempfile.NamedTemporaryFile("w", suffix=".pem", delete=False) as fh:
fh.write(pem)
tmp = fh.name
try:
info = _ssl._test_decode_cert(tmp)
finally:
Path(tmp).unlink(missing_ok=True)
except (OSError, ssl.SSLError, ValueError):
return ""
subject = _name(info.get("subject"))
issuer = _name(info.get("issuer"))
return f"subject={subject} issuer={issuer}"
def _name(rdns) -> str:
"""Estrae il CN (o l'intero RDN) da una struttura subject/issuer di _test_decode_cert."""
if not rdns:
return "?"
parts = {k: v for rdn in rdns for (k, v) in rdn}
return parts.get("commonName") or ", ".join(f"{k}={v}" for k, v in parts.items())
def _to_pem(cert) -> str:
"""Converte un certificato (_ssl.Certificate o DER bytes) in PEM."""
if isinstance(cert, (bytes, bytearray)):
return ssl.DER_cert_to_PEM_cert(bytes(cert))
pem = cert.public_bytes(_PEM_ENCODING)
return pem if isinstance(pem, str) else pem.decode("ascii")
+202
View File
@@ -0,0 +1,202 @@
from datetime import UTC, datetime
from sqlalchemy import Engine, text
from tht.mschema.models import (
ColumnPhysical,
ForeignKey,
Index,
PhysicalSchema,
TablePhysical,
)
# Query adattate da thoth_sqldb2 (Apache 2.0) — vedi src/tht/vendor/VENDORED.md.
_TABLES_Q = text("""
SELECT c.relname AS table_name,
COALESCE(d.description, '') AS comment,
GREATEST(c.reltuples::bigint, 0) AS row_count
FROM pg_class c
JOIN pg_namespace n ON n.oid = c.relnamespace
LEFT JOIN pg_description d ON d.objoid = c.oid AND d.objsubid = 0
WHERE c.relkind IN ('r', 'p') AND n.nspname = :schema
ORDER BY c.relname
""")
_COLUMNS_Q = text("""
SELECT a.attname AS column_name,
format_type(a.atttypid, a.atttypmod) AS data_type,
(NOT a.attnotnull) AS is_nullable,
pg_get_expr(d.adbin, d.adrelid) AS column_default,
COALESCE(pgd.description, '') AS comment,
(ty.typtype = 'e') AS is_enum,
EXISTS (
SELECT 1 FROM pg_index i
WHERE i.indrelid = c.oid AND i.indisprimary AND a.attnum = ANY (i.indkey)
) AS is_pk
FROM pg_class c
JOIN pg_namespace n ON n.oid = c.relnamespace
JOIN pg_attribute a ON a.attrelid = c.oid
JOIN pg_type ty ON ty.oid = a.atttypid
LEFT JOIN pg_attrdef d ON d.adrelid = c.oid AND d.adnum = a.attnum
LEFT JOIN pg_description pgd ON pgd.objoid = c.oid AND pgd.objsubid = a.attnum
WHERE c.relname = :table_name AND n.nspname = :schema
AND a.attnum > 0 AND NOT a.attisdropped
ORDER BY a.attnum
""")
_FOREIGN_KEYS_Q = text("""
SELECT con.conname AS constraint_name,
rel.relname AS source_table,
a.attname AS source_column,
frel.relname AS target_table,
fa.attname AS target_column,
src.ord
FROM pg_constraint con
JOIN pg_class rel ON rel.oid = con.conrelid
JOIN pg_namespace ns ON ns.oid = rel.relnamespace
JOIN pg_class frel ON frel.oid = con.confrelid
JOIN unnest(con.conkey) WITH ORDINALITY AS src(attnum, ord) ON true
JOIN pg_attribute a ON a.attrelid = con.conrelid AND a.attnum = src.attnum
JOIN unnest(con.confkey) WITH ORDINALITY AS dst(attnum, ord) ON dst.ord = src.ord
JOIN pg_attribute fa ON fa.attrelid = con.confrelid AND fa.attnum = dst.attnum
WHERE con.contype = 'f' AND ns.nspname = :schema
ORDER BY rel.relname, con.conname, src.ord
""")
_INDEXES_Q = text("""
SELECT i.relname AS index_name,
t.relname AS table_name,
ix.indisunique AS is_unique,
ix.indisprimary AS is_primary,
am.amname AS index_type,
array_agg(a.attname ORDER BY a.attnum) AS columns
FROM pg_index ix
JOIN pg_class i ON i.oid = ix.indexrelid
JOIN pg_class t ON t.oid = ix.indrelid
JOIN pg_namespace n ON n.oid = t.relnamespace
JOIN pg_am am ON am.oid = i.relam
JOIN pg_attribute a ON a.attrelid = t.oid AND a.attnum = ANY (ix.indkey)
WHERE n.nspname = :schema
GROUP BY i.relname, t.relname, ix.indisunique, ix.indisprimary, am.amname
ORDER BY t.relname, i.relname
""")
class IntrospectionError(Exception):
pass
def introspect(engine: Engine, database: str, schema: str) -> PhysicalSchema:
with engine.connect() as conn:
exists = conn.execute(
text("SELECT 1 FROM pg_namespace WHERE nspname = :schema"), {"schema": schema}
).scalar()
if not exists:
raise IntrospectionError(f"Schema inesistente: {schema}")
tables: dict[str, TablePhysical] = {}
for trow in conn.execute(_TABLES_Q, {"schema": schema}):
columns: dict[str, ColumnPhysical] = {}
for crow in conn.execute(
_COLUMNS_Q, {"table_name": trow.table_name, "schema": schema}
):
columns[crow.column_name] = ColumnPhysical(
type=crow.data_type,
nullable=bool(crow.is_nullable),
pk=bool(crow.is_pk),
default=crow.column_default,
comment=crow.comment,
is_enum=bool(crow.is_enum),
)
tables[trow.table_name] = TablePhysical(
comment=trow.comment, row_count=trow.row_count, columns=columns
)
# FK raggruppate per (tabella, constraint), ordinate per posizione
grouped: dict[tuple[str, str], ForeignKey] = {}
for row in conn.execute(_FOREIGN_KEYS_Q, {"schema": schema}):
key = (row.source_table, row.constraint_name)
fk = grouped.setdefault(
key,
ForeignKey(
columns=[], ref_table=row.target_table, ref_columns=[],
name=row.constraint_name,
),
)
fk.columns.append(row.source_column)
fk.ref_columns.append(row.target_column)
for (table_name, _), fk in grouped.items():
if table_name in tables:
tables[table_name].foreign_keys.append(fk)
for row in conn.execute(_INDEXES_Q, {"schema": schema}):
if row.table_name in tables:
tables[row.table_name].indexes.append(
Index(
name=row.index_name,
columns=list(row.columns),
unique=bool(row.is_unique),
primary=bool(row.is_primary),
type=row.index_type,
)
)
return PhysicalSchema(
database=database,
schema=schema,
introspected_at=datetime.now(UTC),
tables=tables,
)
def introspect_rest(client, database: str, schema: str) -> PhysicalSchema:
"""Introspezione via REST (rpc `list_tables`/`table_columns`/`table_comments`/
`table_foreign_keys`). Limiti rispetto al diretto: niente indici (nessun rpc) e
`is_enum` non disponibile (default False)."""
tables: dict[str, TablePhysical] = {}
for trow in client.list_tables(schema):
if trow.get("type") != "TABLE":
continue # le viste sono fuori scope (come l'introspezione diretta)
table_name = trow["table"]
col_comments = {
c["name"]: (c.get("comment") or "")
for c in client.table_comments(schema, table_name)
if c.get("object") == "COLUMN"
}
columns: dict[str, ColumnPhysical] = {}
for crow in client.table_columns(schema, table_name):
name = crow["column"]
columns[name] = ColumnPhysical(
type=crow["type"],
nullable=bool(crow.get("nullable", True)),
pk=bool(crow.get("pk", False)),
default=crow.get("default"),
comment=col_comments.get(name, ""),
)
foreign_keys: list[ForeignKey] = []
for fk in client.table_foreign_keys(schema, table_name):
foreign_keys.append(
ForeignKey(
columns=fk.get("columns") or [fk["column"]],
ref_table=fk.get("ref_table") or fk["target_table"],
ref_columns=fk.get("ref_columns") or [fk["target_column"]],
name=fk.get("name", ""),
)
)
tables[table_name] = TablePhysical(
comment=trow.get("comment") or "",
row_count=int(trow.get("rows") or 0),
columns=columns,
foreign_keys=foreign_keys,
)
return PhysicalSchema(
database=database,
schema=schema,
introspected_at=datetime.now(UTC),
tables=tables,
)
+160
View File
@@ -0,0 +1,160 @@
import logging
from dataclasses import dataclass
from sqlalchemy import Engine, text
from tht.config import ExamplesConfig, LshConfig
from tht.mschema.models import Annotations, PhysicalSchema
logger = logging.getLogger(__name__)
TEXT_TYPE_PREFIXES = ("text", "varchar", "character", "char")
def is_text_type(pg_type: str) -> bool:
return pg_type.lower().startswith(TEXT_TYPE_PREFIXES)
def add_examples(engine: Engine, physical: PhysicalSchema, cfg: ExamplesConfig) -> None:
"""Campiona i valori distinti piu' frequenti delle colonne testuali (in-place)."""
schema = physical.db_schema
with engine.connect() as conn:
for table_name, table in physical.tables.items():
for column_name, column in table.columns.items():
if not is_text_type(column.type):
continue
q = text(f'''
SELECT "{column_name}" FROM (
SELECT "{column_name}", count(*) AS _freq
FROM "{schema}"."{table_name}"
WHERE "{column_name}" IS NOT NULL AND length("{column_name}") > 0
GROUP BY "{column_name}"
ORDER BY _freq DESC
LIMIT :lim
) AS sub
''')
try:
rows = conn.execute(q, {"lim": cfg.max_per_column}).fetchall()
except Exception as e: # colonna non leggibile: si salta, non si interrompe
logger.warning("Campionamento saltato per %s.%s: %s", table_name, column_name, e)
continue
column.examples = [str(r[0]) for r in rows]
def add_examples_rest(client, physical: PhysicalSchema, cfg: ExamplesConfig) -> None:
"""Variante REST di add_examples: valori più frequenti via rpc `top_values`."""
schema = physical.db_schema
for table_name, table in physical.tables.items():
for column_name, column in table.columns.items():
if not is_text_type(column.type):
continue
rows = client.top_values(schema, table_name, column_name, cfg.max_per_column)
column.examples = [str(r["value"]) for r in rows if r["value"] not in (None, "")]
@dataclass
class SkippedColumn:
table: str
column: str
reason: str
@dataclass
class TruncatedColumn:
table: str
column: str
indexed: int # quanti valori (i più frequenti) sono stati indicizzati
def unique_values_for_lsh(
engine: Engine,
physical: PhysicalSchema,
cfg: LshConfig,
annotations: Annotations | None = None,
) -> tuple[dict[str, dict[str, list[str]]], list[SkippedColumn], list[TruncatedColumn]]:
"""Valori delle colonne testuali *eligible* per l'indice LSH.
Indicizza solo colonne con eligibilità effettiva True (le `wide_text` sono escluse:
vedi principio di column eligibility). Estrae i valori distinti *più frequenti*
(ORDER BY frequenza); se superano `max_values_per_column` la colonna è troncata e
segnalata (mai tagliata in silenzio).
"""
from tht.mschema.eligibility import effective_eligibility
annotations = annotations or Annotations()
schema = physical.db_schema
values: dict[str, dict[str, list[str]]] = {}
skipped: list[SkippedColumn] = []
truncated: list[TruncatedColumn] = []
with engine.connect() as conn:
for table_name, table in physical.tables.items():
table_ann = annotations.tables.get(table_name)
for column_name, column in table.columns.items():
if not is_text_type(column.type):
continue
ann_col = table_ann.columns.get(column_name) if table_ann else None
if not effective_eligibility(column, ann_col)[0]:
continue
q = text(f'''
SELECT "{column_name}" FROM (
SELECT "{column_name}", count(*) AS _freq
FROM "{schema}"."{table_name}"
WHERE "{column_name}" IS NOT NULL AND length("{column_name}") > 0
GROUP BY "{column_name}"
ORDER BY _freq DESC, "{column_name}"
LIMIT :lim
) AS sub
''')
try:
rows = conn.execute(q, {"lim": cfg.max_values_per_column}).fetchall()
except Exception as e:
skipped.append(SkippedColumn(table_name, column_name, f"errore: {e}"))
continue
vals = [str(r[0]) for r in rows]
if not vals:
continue
values.setdefault(table_name, {})[column_name] = vals
if len(vals) >= cfg.max_values_per_column:
truncated.append(TruncatedColumn(table_name, column_name, len(vals)))
return values, skipped, truncated
def unique_values_for_lsh_rest(
client,
physical: PhysicalSchema,
cfg: LshConfig,
annotations: Annotations | None = None,
) -> tuple[dict[str, dict[str, list[str]]], list[SkippedColumn], list[TruncatedColumn]]:
"""Variante REST di unique_values_for_lsh: valori più frequenti via rpc `top_values`.
Stessa logica di eligibility e di segnalazione del troncamento del transport diretto.
"""
from tht.mschema.eligibility import effective_eligibility
annotations = annotations or Annotations()
schema = physical.db_schema
values: dict[str, dict[str, list[str]]] = {}
skipped: list[SkippedColumn] = []
truncated: list[TruncatedColumn] = []
for table_name, table in physical.tables.items():
table_ann = annotations.tables.get(table_name)
for column_name, column in table.columns.items():
if not is_text_type(column.type):
continue
ann_col = table_ann.columns.get(column_name) if table_ann else None
if not effective_eligibility(column, ann_col)[0]:
continue
try:
rows = client.top_values(
schema, table_name, column_name, cfg.max_values_per_column
)
except Exception as e:
skipped.append(SkippedColumn(table_name, column_name, f"errore: {e}"))
continue
vals = [str(r["value"]) for r in rows if r["value"] not in (None, "")]
if not vals:
continue
values.setdefault(table_name, {})[column_name] = vals
if len(vals) >= cfg.max_values_per_column:
truncated.append(TruncatedColumn(table_name, column_name, len(vals)))
return values, skipped, truncated
+94
View File
@@ -0,0 +1,94 @@
from datetime import UTC, datetime
from pathlib import Path
from typing import Literal
from pydantic import BaseModel
DECISIONS_FILE = "review_decisions.jsonl"
# 22 tipi di the reference implementation (verified leggendo session/decisions.py) + 1 nuovo (D15):
# `decision_retracted` per il rollback a granularità step (ritira una decisione
# senza cancellarne la riga dal log di audit; effective_decisions la onora).
DecisionType = Literal[
"concept_clarified",
"question_rewritten",
"table_promoted",
"table_excluded",
"column_corrected",
"join_modified",
"evidence_accepted",
"evidence_rejected",
"ambiguity_open",
"memory_rejected",
"cte_approved",
"cte_corrected",
"cte_rejected",
"sql_revised",
"sql_approved",
"sql_rejected",
"phase_approved",
"phase_auto_approved",
"phase_reopened",
"phase_skipped",
"datamart_requested",
"datamart_declined",
# D15: marker di ritrazione. subject = "phase:N", retracts = decision_seq ritirata.
# Resta nel log di audit (append-only); effective_decisions() la esclude dalla vista.
"decision_retracted",
# D14a: valore citato nella domanda ancorato a una o piu' colonne. subject =
# "phase:4", detail = il valore (es. "ablazione"), rationale = la/e colonna/e scelta/e
# dal reviewer (aggregate_lsh_multi le espone tutte senza collassare al miglior match).
"value_grounded",
# D14b: formula di concetto approvata/rifiutata dal reviewer. subject = "phase:4",
# detail = il concetto (es. "fascia pediatrica"), rationale = la/e colonna/e o il motivo.
# retrieve_formula restituisce i candidati; queste decisioni registrano la scelta.
"concept_formula_approved",
"concept_formula_rejected",
]
class DecisionRecord(BaseModel):
seq: int
ts: datetime
type: DecisionType
subject: str
detail: str = ""
rationale: str = ""
# D15: se type == "decision_retracted", indica quale seq viene ritirata.
retracts: int | None = None
def list_decisions(session_dir: Path) -> list[DecisionRecord]:
path = session_dir / DECISIONS_FILE
if not path.exists():
return []
return [
DecisionRecord.model_validate_json(line)
for line in path.read_text().splitlines()
if line.strip()
]
def append_decision(
session_dir: Path,
*,
type: str,
subject: str,
detail: str = "",
rationale: str = "",
retracts: int | None = None,
) -> DecisionRecord:
record = DecisionRecord(
seq=len(list_decisions(session_dir)) + 1,
ts=datetime.now(UTC),
type=type,
subject=subject,
detail=detail,
rationale=rationale,
retracts=retracts,
)
path = session_dir / DECISIONS_FILE
path.parent.mkdir(parents=True, exist_ok=True)
with path.open("a") as f:
f.write(record.model_dump_json() + "\n")
return record
View File
+100
View File
@@ -0,0 +1,100 @@
"""SQL concept->formula evidence store (spec D14b, §4.7).
A concept (e.g. 'fascia pediatrica', 'ablazione') maps to a reusable SQL formula
(a CASE WHEN ...) that derives it from physical columns. These are reviewable
units: the gate surfaces a candidate formula, the reviewer approves or rejects it
(decision types concept_formula_approved / concept_formula_rejected), and approved
formulas travel with the schema-linking artifact.
Storage: one file per formula, frontmatter YAML + SQL body (same shape as
EvidenceDoc.parse). Directory layout: <root>/formulas/<slug>-<n>.sql.md.
Retrieve is by concept (may return several, e.g. competing drafts vs reviewed).
"""
from __future__ import annotations
import re
from pathlib import Path
from typing import Literal
import yaml
from pydantic import BaseModel
FORMULAS_SUBDIR = "formulas"
_SUFFIX_RE = re.compile(r"^(.*?)-(\d+)\.sql\.md$")
class ConceptFormula(BaseModel):
concept: str
columns: list[str] = []
sql: str
status: Literal["draft", "reviewed"] = "draft"
sources: list[str] = []
@property
def _slug(self) -> str:
"""ASCII slug for the filename (matches textutil.slugify shape)."""
import unicodedata
text = unicodedata.normalize("NFKD", self.concept).encode("ascii", "ignore").decode()
return re.sub(r"[^a-z0-9_]+", "-", text.lower()).strip("-") or "formula"
def dump(self) -> str:
meta = self.model_dump(exclude={"sql"}, mode="json")
fm = yaml.safe_dump(meta, sort_keys=False, allow_unicode=True)
return f"---\n{fm}---\n{self.sql}\n"
@classmethod
def parse(cls, text: str) -> "ConceptFormula":
if not text.startswith("---\n"):
raise ValueError("frontmatter mancante (atteso '---\\n' iniziale)")
try:
_, fm, body = text.split("---\n", 2)
except ValueError as e:
raise ValueError("frontmatter malformato") from e
meta = yaml.safe_load(fm)
if not isinstance(meta, dict):
raise ValueError("frontmatter non valido")
return cls.model_validate({**meta, "sql": body.strip("\n")})
def _next_path(root: Path, slug: str) -> Path:
"""First free <slug>-<n>.sql.md path under root (n starts at 1)."""
root.mkdir(parents=True, exist_ok=True)
existing = sorted(root.glob(f"{slug}-*.sql.md"))
n = 0
for p in existing:
m = _SUFFIX_RE.match(p.name)
if m:
n = max(n, int(m.group(2)))
return root / f"{slug}-{n + 1}.sql.md"
def save_formula(root: Path | str, formula: ConceptFormula) -> Path:
"""Persist a single concept->formula unit under <root>/formulas/. Returns the
written path. Append-only: each save writes a new file (so competing drafts and
reviewed versions coexist until a curator prunes)."""
root = Path(root)
formulas_dir = root / FORMULAS_SUBDIR
path = _next_path(formulas_dir, formula._slug)
path.parent.mkdir(parents=True, exist_ok=True)
path.write_text(formula.dump())
return path
def retrieve_formula(root: Path | str, concept: str) -> list[ConceptFormula]:
"""All formulas for `concept` under <root>/formulas/. Empty list if none (or if
the dir is absent). Multiple results mean competing drafts/versions for the same
concept -- the caller (gate) lets the reviewer pick."""
root = Path(root)
formulas_dir = root / FORMULAS_SUBDIR
if not formulas_dir.is_dir():
return []
out: list[ConceptFormula] = []
for f in sorted(formulas_dir.glob("*.sql.md")):
try:
formula = ConceptFormula.parse(f.read_text())
except ValueError:
continue # malformed file: skip, don't crash retrieval
if formula.concept == concept:
out.append(formula)
return out
+62
View File
@@ -0,0 +1,62 @@
from pathlib import Path
from typing import Literal
import yaml
from pydantic import BaseModel, ValidationError
class EvidenceError(Exception):
pass
class EvidenceDoc(BaseModel):
id: str
title: str
# tier e status sono opzionali: la sola presenza di un documento basta a
# vettorizzarlo, quindi l'autore ETL non e' obbligato a compilarli.
tier: Literal["structural", "concept"] = "structural"
status: Literal["auto", "draft", "reviewed"] = "reviewed"
sources: list[str] = []
tables: list[str] = []
concepts: list[str] = []
body: str = ""
path: Path | None = None # valorizzato al load, escluso dal dump
@classmethod
def parse(cls, text: str, path: Path | None = None) -> "EvidenceDoc":
if not text.startswith("---\n"):
raise EvidenceError(f"frontmatter mancante in {path or '<testo>'}")
try:
_, fm, body = text.split("---\n", 2)
except ValueError as e:
raise EvidenceError(f"frontmatter malformato in {path or '<testo>'}") from e
meta = yaml.safe_load(fm)
if not isinstance(meta, dict):
raise EvidenceError(f"frontmatter non valido in {path or '<testo>'}")
try:
return cls.model_validate({**meta, "body": body.strip("\n"), "path": path})
except ValidationError as e:
raise EvidenceError(f"evidence non valida in {path or '<testo>'}:\n{e}") from e
def dump(self) -> str:
meta = self.model_dump(exclude={"body", "path"}, mode="json")
fm = yaml.safe_dump(meta, sort_keys=False, allow_unicode=True)
return f"---\n{fm}---\n{self.body}\n"
def save(self, path: Path) -> None:
path.parent.mkdir(parents=True, exist_ok=True)
path.write_text(self.dump())
self.path = path
def load_evidence_dir(root: Path) -> list[EvidenceDoc]:
"""Carica ricorsivamente tutte le evidence sotto `root`, preservando la
gerarchia per dominio. I file README (di sola navigazione) sono ignorati."""
docs: list[EvidenceDoc] = []
if not root.is_dir():
return docs
for f in sorted(root.rglob("*.md")):
if f.name.upper().startswith("README"):
continue
docs.append(EvidenceDoc.parse(f.read_text(), path=f))
return docs
+257
View File
@@ -0,0 +1,257 @@
import os
import re
import tempfile
from datetime import UTC, datetime
from pathlib import Path
from pydantic import BaseModel
from tht.decisions import DecisionRecord, DecisionType, list_decisions
from tht.session.models import SessionManifest
from tht.vectorstore.records import VectorRecord
class MemoryRecord(BaseModel):
id: str
ts: datetime
session_id: str
decision_seq: int
type: DecisionType
subject: str
detail: str = ""
rationale: str = ""
question_context: str = ""
tables: list[str] = []
concepts: list[str] = []
_MEM_ID_RE = re.compile(r"\bmem-\d{4,}\b")
def decided_memory_ids(decisions: list[DecisionRecord]) -> set[str]:
"""Id delle memorie gia' DECISE nella sessione: rifiutate (decisione
`memory_rejected`, id nel subject) o applicate (citate come `mem-XXXX` nel
rationale, per convenzione della checklist memorie). Servono a
`tht memory search --session` per non riproporre cio' che il reviewer ha
gia' deciso (fix: memorie scartate riproposte)."""
out: set[str] = set()
for d in decisions:
if d.type == "memory_rejected" and d.subject:
out.add(d.subject)
out.update(_MEM_ID_RE.findall(d.rationale or ""))
return out
def load_registry(registry_path: Path) -> list[MemoryRecord]:
if not registry_path.exists():
return []
return [
MemoryRecord.model_validate_json(line)
for line in registry_path.read_text().splitlines()
if line.strip()
]
def save_registry(records: list[MemoryRecord], path: Path) -> None:
"""Riscrive il registro in modo atomico (tmp nella stessa dir + os.replace)."""
path.parent.mkdir(parents=True, exist_ok=True)
fd, tmp = tempfile.mkstemp(dir=path.parent, prefix=".registry-", suffix=".tmp")
try:
with os.fdopen(fd, "w") as f:
for r in records:
f.write(r.model_dump_json() + "\n")
os.replace(tmp, path)
except BaseException:
if os.path.exists(tmp):
os.unlink(tmp)
raise
class MemoryNotFound(Exception):
"""Sollevata quando un id memoria non esiste nel registro."""
EDITABLE_FIELDS = frozenset(
{"subject", "type", "detail", "rationale", "question_context", "tables", "concepts"}
)
def update_record(path: Path, mem_id: str, fields: dict) -> MemoryRecord:
records = load_registry(path)
for i, r in enumerate(records):
if r.id == mem_id:
data = r.model_dump()
data.update({k: v for k, v in fields.items() if k in EDITABLE_FIELDS})
records[i] = MemoryRecord.model_validate(data)
save_registry(records, path)
return records[i]
raise MemoryNotFound(mem_id)
def delete_record(path: Path, mem_id: str) -> MemoryRecord:
records = load_registry(path)
for i, r in enumerate(records):
if r.id == mem_id:
removed = records.pop(i)
save_registry(records, path)
return removed
raise MemoryNotFound(mem_id)
def _default_tables(decision: DecisionRecord) -> list[str]:
if decision.type in ("table_promoted", "table_excluded"):
return [decision.subject]
if decision.type in ("column_corrected", "join_modified") and "." in decision.subject:
return [decision.subject.split(".")[0]]
return []
def _default_concepts(decision: DecisionRecord) -> list[str]:
if decision.type == "concept_clarified":
return [decision.subject]
return []
def _question_context(decisions: list[DecisionRecord], manifest: SessionManifest) -> str:
rewritten = [d for d in decisions if d.type == "question_rewritten"]
return rewritten[-1].detail if rewritten else manifest.question
def _next_id_num(existing: list[MemoryRecord]) -> int:
nums = [int(r.id.split("-", 1)[1]) for r in existing if r.id.startswith("mem-")]
return (max(nums) + 1) if nums else 1
# Solo questi tipi di decisione sono concetti riusabili in altre generazioni
# (scelta reviewer): il resto e' query-specifico e non va proposto in promozione.
REUSABLE_TYPES = frozenset({"concept_clarified", "table_promoted", "table_excluded"})
MAX_PROMOTION_CANDIDATES = 5
def _compute_promotions(
session_dir: Path, manifest: SessionManifest, *,
seqs: list[int] | None, existing: list[MemoryRecord],
) -> list[MemoryRecord]:
decisions = list_decisions(session_dir)
selected = decisions if seqs is None else [d for d in decisions if d.seq in seqs]
already = {(r.session_id, r.decision_seq) for r in existing}
n = _next_id_num(existing)
context = _question_context(decisions, manifest)
out: list[MemoryRecord] = []
for d in selected:
if (manifest.id, d.seq) in already:
continue
out.append(
MemoryRecord(
id=f"mem-{n:04d}", ts=datetime.now(UTC), session_id=manifest.id,
decision_seq=d.seq, type=d.type, subject=d.subject, detail=d.detail,
rationale=d.rationale, question_context=context,
tables=_default_tables(d), concepts=_default_concepts(d),
)
)
n += 1
return out
def promote(
session_dir: Path, manifest: SessionManifest, *,
seqs: list[int] | None, registry_path: Path,
) -> list[MemoryRecord]:
"""Copia le decisioni indicate (tutte se seqs=None) nel registro globale.
Salta quelle gia' promosse (chiave: session_id + decision_seq)."""
existing = load_registry(registry_path)
promoted = _compute_promotions(session_dir, manifest, seqs=seqs, existing=existing)
if promoted:
save_registry(existing + promoted, registry_path)
return promoted
def reusable_promotions(
session_dir: Path, manifest: SessionManifest, registry_path: Path
) -> list[MemoryRecord]:
"""Candidati riusabili (tipi in REUSABLE_TYPES) non ancora promossi, SENZA
cap: il chiamante applica MAX_PROMOTION_CANDIDATES e segnala il troncamento."""
cand = _compute_promotions(
session_dir, manifest, seqs=None, existing=load_registry(registry_path)
)
return [c for c in cand if c.type in REUSABLE_TYPES]
def preview_promotions(
session_dir: Path, manifest: SessionManifest, registry_path: Path
) -> list[MemoryRecord]:
"""Candidati promuovibili: riusabili e troncati a MAX_PROMOTION_CANDIDATES."""
return reusable_promotions(session_dir, manifest, registry_path)[
:MAX_PROMOTION_CANDIDATES
]
def memory_vector_records(records: list[MemoryRecord]) -> list[VectorRecord]:
out: list[VectorRecord] = []
for r in records:
lines = [
f"Decisione {r.type}: {r.subject}",
r.detail,
r.rationale,
f"Domanda di contesto: {r.question_context}",
]
if r.tables:
lines.append("Tabelle: " + ", ".join(r.tables))
if r.concepts:
lines.append("Concetti: " + ", ".join(r.concepts))
out.append(
VectorRecord(
id=f"memory:{r.id}", kind="memory", ref=r.id,
title=f"{r.type}: {r.subject}",
content="\n".join(filter(None, lines)),
metadata={
"type": r.type, "session_id": r.session_id,
"tables": r.tables, "concepts": r.concepts,
},
)
)
return out
def memory_vector_record_for_decision(
records: list[MemoryRecord], decision_seq: int
) -> VectorRecord | None:
"""The single VectorRecord for `decision_seq`, or None if no memory matches.
D11 save-one builds only this one record (not the full memory_vector_records
list) so the remote upsert is a single row.
"""
match = [r for r in records if r.decision_seq == decision_seq]
if not match:
return None
return memory_vector_records(match)[0]
def save_one_memory(
records: list[MemoryRecord], decision_seq: int, *, writer, embedder
) -> int:
"""Targeted one-row upsert of a promoted decision to pgvector via the writer key
(spec D11). This is NOT a full vectorstore resync: it embeds and pushes a single
record, so a workstation with a writer key can publish one memory without
rebuilding the index. Returns the upsert count (0 if no record matched).
`writer` is a VectorRestClient (writer key); `embedder` an embeddings client.
The destructive cleanup (sync's delete-stale step) is intentionally absent: it
remains a server-side-only operation via the direct vectordb connection.
"""
from tht.vectorstore.rest_writer import pack_metadata
from tht.vectorstore.store import content_hash
record = memory_vector_record_for_decision(records, decision_seq)
if record is None:
return 0
embedding = embedder.embed_documents([record.content])[0]
row = {
"record_key": record.id,
"kind": record.kind,
"content_hash": content_hash(record.content),
"metadata": pack_metadata(record),
"embedding": embedding,
}
return writer.upsert_records("memory", [row])
View File
+93
View File
@@ -0,0 +1,93 @@
"""Classificazione di column eligibility (principio trasversale Thoth).
Vedi docs/superpowers/specs/2026-06-13-tht-column-eligibility-principle.md.
Il testo ampio (lettere di dimissione, note, anamnesi) è ignorato ovunque; i dati
provengono solo da numerici, enum, temporali, booleani e testo breve.
"""
import re
from tht.config import EligibilityConfig
from tht.db.sampling import is_text_type
from tht.mschema.models import ColumnAnnotation, ColumnPhysical, PhysicalSchema
_NUMERIC_PREFIXES = (
"smallint", "integer", "bigint", "numeric", "decimal", "real", "double", "money",
)
_TEMPORAL_PREFIXES = ("date", "time", "timestamp", "interval")
_LEN_RE = re.compile(r"\((\d+)\)")
def _declared_len(pg_type: str) -> int | None:
m = _LEN_RE.search(pg_type)
return int(m.group(1)) if m else None
def classify_column(
pg_type: str,
is_enum: bool,
sampled_avg: float | None,
sampled_max: int | None,
cfg: EligibilityConfig,
) -> tuple[bool, str]:
"""Classifica una colonna come (eligible, reason). Funzione pura, senza I/O."""
t = pg_type.strip().lower()
if t.endswith("[]"):
return False, "wide_text"
if is_enum:
return True, "enum"
if t.startswith(_NUMERIC_PREFIXES):
return True, "numeric"
if t.startswith("boolean"):
return True, "boolean"
if t.startswith(_TEMPORAL_PREFIXES):
return True, "temporal"
if t.startswith("uuid"):
return True, "code"
if is_text_type(t):
declared = _declared_len(t)
if declared is not None and declared <= cfg.max_declared_len:
return True, "short_text"
# bound grande o text/varchar non vincolato: decide il dato campionato
if sampled_avg is None or sampled_max is None:
return False, "wide_text"
if sampled_avg <= cfg.max_avg_length and sampled_max <= cfg.max_sampled_len:
return True, "short_text"
return False, "wide_text"
return False, "wide_text"
def _sampled_stats(examples: list[str]) -> tuple[float | None, int | None]:
if not examples:
return None, None
lengths = [len(v) for v in examples]
return sum(lengths) / len(lengths), max(lengths)
def classify_all(physical: PhysicalSchema, cfg: EligibilityConfig) -> None:
"""Assegna eligible/eligibility_reason a ogni colonna (in-place) e azzera gli
examples delle colonne ignored. Da chiamare DOPO add_examples."""
ignore_by_name = {name.lower() for name in cfg.ignore_columns}
for table in physical.tables.values():
for column_name, column in table.columns.items():
if column_name.lower() in ignore_by_name:
# colonna di servizio (ETL/audit): ignorata a prescindere dal tipo
column.eligible = False
column.eligibility_reason = "ignored_by_name"
column.examples = []
continue
avg, mx = _sampled_stats(column.examples)
eligible, reason = classify_column(column.type, column.is_enum, avg, mx, cfg)
column.eligible = eligible
column.eligibility_reason = reason
if not eligible:
column.examples = []
def effective_eligibility(
column: ColumnPhysical, annotation: ColumnAnnotation | None
) -> tuple[bool, str]:
"""Eligibilità effettiva: l'override in annotations.yaml vince sul fisico."""
if annotation is not None and annotation.eligible is not None:
return annotation.eligible, "override"
return column.eligible, column.eligibility_reason
+15
View File
@@ -0,0 +1,15 @@
from tht.mschema.models import Annotations, PhysicalSchema
def find_orphans(physical: PhysicalSchema, annotations: Annotations) -> list[str]:
"""Annotazioni che puntano a oggetti spariti dal fisico. Non le rimuove mai."""
orphans: list[str] = []
for table_name, table_ann in annotations.tables.items():
table = physical.tables.get(table_name)
if table is None:
orphans.append(table_name)
continue
for column_name in table_ann.columns:
if column_name not in table.columns:
orphans.append(f"{table_name}.{column_name}")
return orphans
+90
View File
@@ -0,0 +1,90 @@
from datetime import datetime
from pathlib import Path
from typing import Self
import yaml
from pydantic import BaseModel, Field
class _YamlModel(BaseModel):
def to_yaml(self, path: Path) -> None:
path.parent.mkdir(parents=True, exist_ok=True)
data = self.model_dump(by_alias=True, mode="json", exclude_defaults=False)
path.write_text(
yaml.safe_dump(data, sort_keys=False, allow_unicode=True, width=120)
)
@classmethod
def from_yaml(cls, path: Path) -> Self:
raw = yaml.safe_load(path.read_text())
return cls.model_validate(raw)
class ColumnPhysical(BaseModel):
type: str
nullable: bool = True
pk: bool = False
default: str | None = None
comment: str = ""
examples: list[str] = []
is_enum: bool = False
eligible: bool = True
eligibility_reason: str = ""
class ForeignKey(BaseModel):
columns: list[str]
ref_table: str
ref_columns: list[str]
name: str = ""
class Index(BaseModel):
name: str
columns: list[str]
unique: bool = False
primary: bool = False
type: str = "btree"
class TablePhysical(BaseModel):
comment: str = ""
row_count: int = 0 # stima da pg_class.reltuples
columns: dict[str, ColumnPhysical] = {}
foreign_keys: list[ForeignKey] = []
indexes: list[Index] = []
class PhysicalSchema(_YamlModel):
database: str
db_schema: str = Field(alias="schema")
introspected_at: datetime
tables: dict[str, TablePhysical] = {}
model_config = {"populate_by_name": True}
class ColumnAnnotation(BaseModel):
description: str = ""
synonyms: list[str] = []
concepts: list[str] = []
evidence: list[str] = []
notes: str = ""
eligible: bool | None = None
class TableAnnotation(BaseModel):
description: str = ""
concepts: list[str] = []
notes: str = ""
columns: dict[str, ColumnAnnotation] = {}
class Annotations(_YamlModel):
tables: dict[str, TableAnnotation] = {}
@classmethod
def from_yaml(cls, path: Path) -> "Annotations":
if not path.exists():
return cls()
return super().from_yaml(path)
+141
View File
@@ -0,0 +1,141 @@
from typing import Any
from tht.mschema.eligibility import effective_eligibility
from tht.mschema.models import Annotations, ColumnAnnotation, PhysicalSchema
MAX_EXAMPLES_IN_PROMPT = 5
def _ann_col(annotations: Annotations, table: str, column: str) -> ColumnAnnotation | None:
ann = annotations.tables.get(table)
if ann is None:
return None
return ann.columns.get(column)
def _table_description(physical: PhysicalSchema, annotations: Annotations, table: str) -> str:
ann = annotations.tables.get(table)
if ann and ann.description:
return ann.description
return physical.tables[table].comment
def _column_description(
physical: PhysicalSchema, annotations: Annotations, table: str, column: str
) -> str:
ann = annotations.tables.get(table)
if ann and column in ann.columns and ann.columns[column].description:
return ann.columns[column].description
return physical.tables[table].columns[column].comment
def to_mschema_text(
physical: PhysicalSchema,
annotations: Annotations | None = None,
tables: list[str] | None = None,
) -> str:
"""Serializzazione testuale in stile ThothAI (【Schema】/【Foreign keys】)."""
annotations = annotations or Annotations()
selected = [t for t in physical.tables if tables is None or t in tables]
lines: list[str] = ["【Schema】"]
fk_lines: list[str] = []
for table_name in selected:
table = physical.tables[table_name]
desc = _table_description(physical, annotations, table_name)
if desc:
lines.append(f"-- {desc}")
lines.append(f"CREATE TABLE {table_name} (")
for column_name, column in table.columns.items():
if not effective_eligibility(column, _ann_col(annotations, table_name, column_name))[0]:
continue
line = f" {column_name} {column.type.upper()}"
if column.pk:
line += " -- PRIMARY KEY"
lines.append(line)
cdesc = _column_description(physical, annotations, table_name, column_name)
if cdesc:
lines.append(f" -- {cdesc}")
if column.examples:
shown = ", ".join(column.examples[:MAX_EXAMPLES_IN_PROMPT])
lines.append(f" -- Examples: {shown}")
lines.append(");")
for fk in table.foreign_keys:
for src, dst in zip(fk.columns, fk.ref_columns):
fk_lines.append(f"{table_name}.{src}={fk.ref_table}.{dst}")
lines.extend(["", "【Foreign keys】", *fk_lines])
return "\n".join(lines)
def to_schema_dict(
physical: PhysicalSchema, annotations: Annotations | None = None
) -> dict[str, Any]:
"""Vista compatibile con le logiche AV-SQL (schema_dict)."""
annotations = annotations or Annotations()
out: dict[str, Any] = {}
for table_name, table in physical.tables.items():
cols = [
c
for c in table.columns
if effective_eligibility(table.columns[c], _ann_col(annotations, table_name, c))[0]
]
out[table_name] = {
"columns_name": cols,
"columns_type": [table.columns[c].type for c in cols],
"columns_description": [
_column_description(physical, annotations, table_name, c) for c in cols
],
"example_values": [table.columns[c].examples for c in cols],
"table_to_tablefullname": f"{physical.db_schema}.{table_name}",
"primary_keys": [c for c in cols if table.columns[c].pk],
"foreign_keys": [
{"columns": fk.columns, "ref_table": fk.ref_table, "ref_columns": fk.ref_columns}
for fk in table.foreign_keys
],
}
return out
def to_markdown(physical: PhysicalSchema, annotations: Annotations | None = None) -> str:
"""Report leggibile per il reviewer."""
annotations = annotations or Annotations()
lines = [
f"# Schema {physical.db_schema} ({physical.database})",
"",
f"Introspezione: {physical.introspected_at.isoformat()} — "
f"{len(physical.tables)} tabelle",
]
for table_name, table in physical.tables.items():
lines += ["", f"## {table_name}", ""]
desc = _table_description(physical, annotations, table_name)
if desc:
lines += [desc, ""]
lines += [
f"Righe (stima): {table.row_count}",
"",
"| Colonna | Tipo | Null | PK | Descrizione | Esempi |",
"|---|---|---|---|---|---|",
]
for column_name, column in table.columns.items():
eligible, reason = effective_eligibility(
column, _ann_col(annotations, table_name, column_name)
)
cdesc = _column_description(physical, annotations, table_name, column_name)
if not eligible:
lines.append(
f"| ~~{column_name}~~ | {column.type} | "
f"{'sì' if column.nullable else 'no'} | "
f"{'sì' if column.pk else ''} | {cdesc} | _ignorata: {reason}_ |"
)
continue
examples = ", ".join(column.examples[:3])
lines.append(
f"| {column_name} | {column.type} | {'sì' if column.nullable else 'no'} "
f"| {'sì' if column.pk else ''} | {cdesc} | {examples} |"
)
if table.foreign_keys:
lines += ["", "Foreign keys:"]
for fk in table.foreign_keys:
lines.append(
f"- ({', '.join(fk.columns)}) → {fk.ref_table} ({', '.join(fk.ref_columns)})"
)
return "\n".join(lines)
+204
View File
@@ -0,0 +1,204 @@
"""Phase machinery -- data-driven + effective_decisions (spec D15, F2, §4.8).
This is the single most important architectural fix vs the reference implementation: ALL helpers
consult effective_decisions() instead of raw list_decisions(), so the reopen-aware
view is consistent everywhere (fixes the bug where approved_ctes / advance_problems /
build_evidence conflated stale pre-reopen decisions with new ones).
Strada 2 (decisa in A5): effective_decisions + ladder if-phase-N che la consulta,
MAX_PHASE/PHASE_NAMES letti da workflow.yaml. L'evaluator generico dei prerequisites
di workflow.yaml (F2 pieno) entra in un secondo momento.
"""
from __future__ import annotations
import json
from pathlib import Path
from pydantic import ValidationError
from tht.decisions import DecisionRecord, list_decisions
from tht.session.models import SchemaLinking
from tht.workflow import load_workflow
def _phase_num(subject: str) -> int | None:
"""subject nel formato 'phase:N' -> N, oppure None."""
if not subject.startswith("phase:"):
return None
try:
return int(subject.split(":", 1)[1])
except ValueError:
return None
def _audit_excluding_retracted(session_dir: Path) -> list[DecisionRecord]:
"""Tutto il ledger (append-only) tranne le decisioni ritirate e i marker di ritrazione.
Base per il fold di current_phase: il guard 'n == cur' del fold e' gia' reopen-aware
(una phase_approved:N dopo un reopen a M<N non fa avanzare perche' cur!=N)."""
all_d = list_decisions(session_dir)
retracted_seqs = {
d.retracts for d in all_d if d.type == "decision_retracted" and d.retracts is not None
}
return [
d
for d in all_d
if d.seq not in retracted_seqs and d.type != "decision_retracted"
]
def current_phase(session_dir: Path) -> int:
"""Fase corrente come fold cronologico sull'audit (con ritirate escluse).
cur parte da 1; ogni phase_approved/phase_auto_approved per la fase CORRENTE avanza
(guard 'n == cur' -- gia' reopen-aware: dopo un reopen a M, le vecchie approvazioni
di N>M non fanno avanzare finche' non si riapprova in ordine); phase_reopened torna
indietro. Terminale: max_phase + 1.
"""
wf = load_workflow()
max_plus_one = wf.max_phase + 1
cur = 1
for d in _audit_excluding_retracted(session_dir):
n = _phase_num(d.subject)
if n is None:
continue
if d.type in ("phase_approved", "phase_auto_approved") and n == cur:
cur = min(cur + 1, max_plus_one)
elif d.type == "phase_reopened":
cur = max(1, min(cur, n))
return cur
def effective_decisions(session_dir: Path) -> list[DecisionRecord]:
"""La vista canonica 'effective as of pointer'. TUTTI gli helper non-fold devono usare questa.
Semantica: una decisione e' effective se appartiene a una fase <= current_phase.
Una decisione di fase 7 (es. sql_approved) e' stale quando current_phase=4 dopo un
rollback a F4, anche se fisicamente appare nel ledger.
Il fold di current_phase gestisce le *approvazioni* via guard 'n == cur'; qui
applichiamo la stessa nozione alle decisioni *sostanziali* (table_promoted, sql_approved,
cte_approved, ...): contano solo se la loro fase e' <= quella corrente.
Inoltre esclude le decisioni ritirate (decision_retracted) e i marker stessi.
"""
cur = current_phase(session_dir)
out: list[DecisionRecord] = []
for d in _audit_excluding_retracted(session_dir):
n = _phase_num(d.subject)
# decisioni senza subject di fase (es. concept_clarified puo' avere subject libero):
# le ammettiamo (non sono legate a una fase specifica da filtrare).
if n is None or n <= cur:
out.append(d)
return out
# --- AUTO-ADVANCE -----------------------------------------------------------
_AUTO_ADVANCE_PHASES = frozenset({2, 6}) # F2 Memorie, F6 CTE
_BOUNDARY_TYPES = frozenset({"phase_approved", "phase_auto_approved", "phase_reopened"})
_META_TYPES = frozenset(
{"phase_approved", "phase_auto_approved", "phase_reopened", "phase_skipped"}
)
def substantive_count_current_phase(session_dir: Path) -> int:
"""Numero di decisioni sostanziali dall'ultimo confine di fase (vista effective)."""
decs = effective_decisions(session_dir)
start = 0
for i, d in enumerate(decs):
if d.type in _BOUNDARY_TYPES:
start = i + 1
return sum(1 for d in decs[start:] if d.type not in _META_TYPES)
def auto_advance_eligible(session_dir: Path) -> bool:
"""Vero sse la fase corrente puo' auto-avanzare (zero decisioni sostanziali + prereq ok)."""
cur = current_phase(session_dir)
if cur not in _AUTO_ADVANCE_PHASES:
return False
if substantive_count_current_phase(session_dir) > 0:
return False
return not advance_problems(session_dir, cur)
# --- CTE helpers (consultano effective_decisions) ---------------------------
CTE_PLAN_FILE = "cte_plan.json"
def cte_plan(session_dir: Path) -> list[str]:
path = session_dir / CTE_PLAN_FILE
if not path.exists():
return []
return json.loads(path.read_text())
def approved_ctes(session_dir: Path) -> set[str]:
"""Insieme dei CTE approvati, dalla vista effective (esclude stale post-reopen)."""
return {d.subject for d in effective_decisions(session_dir) if d.type == "cte_approved"}
def next_cte(session_dir: Path) -> str | None:
approved = approved_ctes(session_dir)
for name in cte_plan(session_dir):
if name not in approved:
return name
return None
# --- advance_problems (ladder if-phase-N che consulta effective_decisions) ---
def _has_decision(session_dir: Path, type_: str) -> bool:
return any(d.type == type_ for d in effective_decisions(session_dir))
def _has_decision_subject(session_dir: Path, type_: str, subject: str) -> bool:
return any(
d.type == type_ and d.subject == subject
for d in effective_decisions(session_dir)
)
def advance_problems(session_dir: Path, phase: int) -> list[str]:
"""Prerequisiti minimi per chiudere `phase` (lista vuota = ok).
Ladder if-phase-N (Strada 2): la logica specifica resta, ma ogni lettura passa per
effective_decisions (fix D15). I prerequisiti sono anche documentati in workflow.yaml;
l'evaluator generico (F2 pieno) entra in un secondo momento.
"""
problems: list[str] = []
if phase == 3 and not _has_decision(session_dir, "question_rewritten"):
problems.append("manca la decisione question_rewritten (Fase 3)")
if phase == 5:
path = session_dir / "schema_linking.json"
if not path.exists():
problems.append("schema_linking.json assente (Fase 5)")
else:
try:
SchemaLinking.model_validate(json.loads(path.read_text()))
except (json.JSONDecodeError, ValidationError) as e:
problems.append(f"schema_linking.json non valido (Fase 5): {e}")
if phase == 6 and not _has_decision_subject(session_dir, "phase_skipped", "phase:6"):
plan = cte_plan(session_dir)
if not plan:
problems.append(
"Fase 6: nessun piano CTE (cte_plan.json) e nessun salto esplicito. "
"Approva un piano (reviewer_confirm kind:'cte_plan') oppure salta la Fase 6 "
"registrando una decisione phase_skipped subject phase:6."
)
else:
nc = next_cte(session_dir)
if nc is not None:
problems.append(f"CTE non ancora approvato: {nc} (Fase 6)")
if phase == 7 and not _has_decision(session_dir, "sql_approved"):
problems.append("manca la decisione sql_approved (Fase 7)")
if phase == 8 and not any(
d.type in ("datamart_requested", "datamart_declined")
for d in effective_decisions(session_dir)
):
problems.append(
"Fase 8: nessuna risposta sulla generazione dbt del datamart "
"(manca una decisione datamart_requested o datamart_declined)."
)
return problems
View File
+113
View File
@@ -0,0 +1,113 @@
"""Client per il the DWH esposto via Supabase/PostgREST.
Incapsula i 10 rpc verificati contro produzione. Errori sempre in italiano e azionabili
(stile `vectorstore/embeddings.py`). Lettura sola: l'API ammette solo SELECT/WITH.
"""
import requests
from tht.config import RestConfig
class RestError(Exception):
"""Errore di accesso al DWH via REST, con messaggio leggibile per il reviewer."""
class RestClient:
def __init__(self, cfg: RestConfig):
self.cfg = cfg
self._base = cfg.base_url.rstrip("/")
# -- trasporto -----------------------------------------------------------------
def _post(self, fn: str, args: dict) -> requests.Response:
url = f"{self._base}/rpc/{fn}"
verify: bool | str = self.cfg.ssl_ca if self.cfg.ssl_ca else True
try:
return requests.post(
url,
json=args,
headers={"X-API-Key": self.cfg.api_key},
timeout=self.cfg.timeout,
verify=verify,
)
except requests.RequestException as e:
raise RestError(
f"DWH REST non raggiungibile su {self.cfg.base_url} (rpc {fn}): {e}"
) from e
def _error_msg(self, fn: str, resp: requests.Response) -> str:
try:
body = resp.json()
detail = body.get("message") or body.get("details") or resp.text
except Exception:
detail = resp.text
return f"DWH REST rpc {fn} → HTTP {resp.status_code}: {detail}"
def _call(self, fn: str, args: dict):
resp = self._post(fn, args)
if not resp.ok:
raise RestError(self._error_msg(fn, resp))
if resp.status_code == 204 or not resp.text:
return None
return resp.json()
# -- rpc -----------------------------------------------------------------------
def ping(self) -> dict:
return self._call("ping", {})
def run_query(self, query_text: str) -> list[dict]:
"""Esegue una SELECT/WITH e restituisce le righe (read-only lato server)."""
return self._call("run_query", {"query_text": query_text}) or []
def explain_query(self, query_text: str) -> list[str]:
"""EXPLAIN di una SELECT/WITH: ritorna le righe testuali del piano."""
rows = self._call("explain_query", {"query_text": query_text}) or []
return [r["line"] for r in rows]
def validate_select(self, query_text: str) -> bool:
"""True se la query è una SELECT/WITH accettata, False se è una write."""
resp = self._post("_validate_select", {"query_text": query_text})
if resp.status_code == 204:
return True
if resp.status_code == 400:
return False
raise RestError(self._error_msg("_validate_select", resp))
def list_tables(self, schema_name: str) -> list[dict]:
return self._call("list_tables", {"schema_name": schema_name}) or []
def table_columns(self, schema_name: str, table_name: str) -> list[dict]:
return self._call(
"table_columns", {"schema_name": schema_name, "table_name": table_name}
) or []
def table_comments(self, schema_name: str, table_name: str) -> list[dict]:
return self._call(
"table_comments", {"schema_name": schema_name, "table_name": table_name}
) or []
def table_foreign_keys(self, schema_name: str, table_name: str) -> list[dict]:
return self._call(
"table_foreign_keys", {"schema_name": schema_name, "table_name": table_name}
) or []
def top_values(
self, schema_name: str, table_name: str, column_name: str, max_values: int
) -> list[dict]:
return self._call(
"top_values",
{
"schema_name": schema_name,
"table_name": table_name,
"column_name": column_name,
"max_values": max_values,
},
) or []
def column_stats(self, schema_name: str, table_name: str, column_name: str) -> dict:
return self._call(
"column_stats",
{"schema_name": schema_name, "table_name": table_name, "column_name": column_name},
)
+131
View File
@@ -0,0 +1,131 @@
from pydantic import BaseModel
from tht.vectorstore.store import VectorStore
class SearchResult(BaseModel):
key: str # column:<t>.<c> | table:<t> | evidence:<id>
label: str
kind: str # schema_column | schema_table | evidence | values
signals: dict # {"lsh": {"rank","score","value"?}, "vector": {"rank","score"}}
rrf: float
status: str = "" # status evidence, se applicabile
content: str = "" # testo matchato (per --explain)
def rrf_fuse(rankings: dict[str, list[tuple[str, float]]], k: int) -> dict[str, dict]:
"""rankings: nome_segnale -> lista (key, raw_score) gia' ordinata per rilevanza.
Ritorna key -> {"rrf": float, "signals": {segnale: {"rank", "score"}}}."""
fused: dict[str, dict] = {}
for signal, ranked in rankings.items():
for rank, (key, score) in enumerate(ranked, start=1):
entry = fused.setdefault(key, {"rrf": 0.0, "signals": {}})
entry["rrf"] += 1.0 / (k + rank)
entry["signals"][signal] = {"rank": rank, "score": round(score, 4)}
return fused
def _aggregate_lsh(lsh_hits: list[tuple[str, str, str, float]]) -> list[tuple[str, float, str]]:
"""Aggrega i match LSH per tabella.colonna tenendo il migliore: (key, score, value)."""
best: dict[str, tuple[float, str]] = {}
for table, column, value, score in lsh_hits:
key = f"column:{table}.{column}"
if key not in best or score > best[key][0]:
best[key] = (score, value)
ordered = sorted(best.items(), key=lambda kv: kv[1][0], reverse=True)
return [(key, score, value) for key, (score, value) in ordered]
def aggregate_lsh_multi(hits: list[dict]) -> dict[str, list[dict]]:
"""Value grounding (spec D14a): group LSH hits by table, keeping EVERY column
where the value appears -- NOT collapsed to a single best column.
The old _aggregate_lsh collapsed matches to one column per table.column key,
hiding alternative groundings (e.g. 'ablazione' matching both a boolean flag
and a free-text patologia field). This function exposes all of them so the
value-grounding widget can let the reviewer choose which column(s) anchor a
cited value.
hits: list of {table, column, value, score}.
Returns: {table -> [{column, value, score}, ...]}, each table's columns ordered
by score desc; within one (table, column) the best-scored value is kept.
"""
best: dict[tuple[str, str], dict] = {}
for h in hits:
key = (h["table"], h["column"])
cur = best.get(key)
if cur is None or h["score"] > cur["score"]:
best[key] = {"column": h["column"], "value": h["value"], "score": h["score"]}
grouped: dict[str, list[dict]] = {}
for (table, _), row in best.items():
grouped.setdefault(table, []).append(row)
for table in grouped:
grouped[table].sort(key=lambda r: r["score"], reverse=True)
return grouped
def _vector_key(hit) -> str:
if hit.kind == "schema_column":
return f"column:{hit.ref}"
if hit.kind == "schema_table":
return f"table:{hit.ref}"
return f"evidence:{hit.id}"
def schema_tables(results: list["SearchResult"], top_tables: int) -> list[tuple[str, float]]:
"""Aggrega i risultati schema a livello di tabella per lo schema-linking: ogni chunk
(tabella o colonna) contribuisce alla sua tabella tenendo il miglior RRF. Ritorna le
prime top_tables tabelle, ordinate per RRF desc (poi nome). I chunk non-schema sono
ignorati."""
best: dict[str, float] = {}
for r in results:
if r.kind not in ("schema_table", "schema_column"):
continue
table = r.key.split(":", 1)[1].split(".", 1)[0]
if table not in best or r.rrf > best[table]:
best[table] = r.rrf
ordered = sorted(best.items(), key=lambda kv: (-kv[1], kv[0]))
return ordered[:top_tables]
def combined_search(
keyword: str,
*,
lsh_hits: list[tuple[str, str, str, float]] | None,
store: VectorStore,
embedder,
top: int,
rrf_k: int,
kinds: list[str] | None,
) -> list[SearchResult]:
"""Fonde LSH (valori di campo) e pgvector con Reciprocal Rank Fusion."""
rankings: dict[str, list[tuple[str, float]]] = {}
lsh_values: dict[str, str] = {}
if lsh_hits:
aggregated = _aggregate_lsh(lsh_hits)
rankings["lsh"] = [(key, score) for key, score, _ in aggregated]
lsh_values = {key: value for key, _, value in aggregated}
vector_hits = store.search(embedder.embed_query(keyword), top_n=top * 2, kinds=kinds)
rankings["vector"] = [(_vector_key(h), h.similarity) for h in vector_hits]
by_key = {_vector_key(h): h for h in vector_hits}
fused = rrf_fuse(rankings, k=rrf_k)
results: list[SearchResult] = []
for key, data in fused.items():
hit = by_key.get(key)
if "lsh" in data["signals"] and key in lsh_values:
data["signals"]["lsh"]["value"] = lsh_values[key]
results.append(
SearchResult(
key=key,
label=hit.title if hit else key.removeprefix("column:"),
kind=hit.kind if hit else "values",
signals=data["signals"],
rrf=data["rrf"],
status=(hit.metadata.get("status", "") if hit else ""),
content=(hit.content if hit else lsh_values.get(key, "")),
)
)
results.sort(key=lambda r: r.rrf, reverse=True)
return results[:top]
View File
+39
View File
@@ -0,0 +1,39 @@
from pathlib import Path
from tht.decisions import DecisionRecord
from tht.session.models import SchemaLinking
def _find_evidence_file(evidence_root: Path, evidence_id: str) -> str:
for match in evidence_root.rglob(f"{evidence_id}.md"):
return str(match)
return ""
def build_evidence_entries(
decisions: list[DecisionRecord],
linking: SchemaLinking,
evidence_root: Path,
) -> list[dict]:
"""Elenco {id, file, esito, decision_seq}: evidence citate nello schema linking
(esito 'usata') e decisioni esplicite del reviewer (accettata/scartata, che
prevalgono sul linking)."""
entries: dict[str, dict] = {}
for candidate in linking.candidates:
for evidence_id in candidate.evidence:
entries.setdefault(evidence_id, {
"id": evidence_id,
"file": _find_evidence_file(evidence_root, evidence_id),
"esito": "usata",
"decision_seq": candidate.decision_seq,
})
for d in decisions:
if d.type not in ("evidence_accepted", "evidence_rejected"):
continue
entries[d.subject] = {
"id": d.subject,
"file": _find_evidence_file(evidence_root, d.subject),
"esito": "accettata" if d.type == "evidence_accepted" else "scartata",
"decision_seq": d.seq,
}
return list(entries.values())
+64
View File
@@ -0,0 +1,64 @@
from datetime import datetime
from typing import Literal
from pydantic import BaseModel, Field, ConfigDict
# Stub locale di _YamlModel (in the reference implementation vive in mschema/models.py).
# Qui serve solo come base con populate_by_name; mschema/ sara' portato nel Task A9.
class _YamlModel(BaseModel):
model_config = ConfigDict(populate_by_name=True)
class SessionManifest(_YamlModel):
id: str
created_at: datetime
status: Literal["open", "closed", "finalized"] = "open"
question: str
database: str
db_schema: str = Field(alias="schema")
# D12/D15: autore della sessione (auth) e versione del workflow usato.
author: str | None = None
summary: str | None = None
updated_at: datetime | None = None
updated_by: str | None = None
schema_version: int | None = None
class Candidate(BaseModel):
kind: Literal["table", "column"]
name: str
signals: dict = {}
evidence: list[str] = []
decision: Literal["promoted", "excluded", "pending"] = "pending"
decision_seq: int | None = None
# D14a: valori citati nella domanda ancorati a questa colonna/tabella.
grounded_values: list[dict] = []
class Join(BaseModel):
from_: str = Field(alias="from")
to: str
source: str = ""
decision: Literal["promoted", "excluded", "pending"] = "promoted"
decision_seq: int | None = None
model_config = {"populate_by_name": True}
class ExcludedItem(BaseModel):
kind: Literal["table", "column"]
name: str
decision_seq: int | None = None
class SchemaLinking(BaseModel):
question: str
candidates: list[Candidate] = []
joins: list[Join] = []
excluded: list[ExcludedItem] = []
open_questions: list[str] = []
# D14b: formule di concetto approvate, parte dello schema-linking.
concept_formulas: list[dict] = []
model_config = {"extra": "forbid"}
+84
View File
@@ -0,0 +1,84 @@
from datetime import UTC, datetime
from pathlib import Path
from tht.config import DatabaseConfig
from tht.session.models import SessionManifest
from tht.textutil import slugify
MANIFEST = "session_manifest.yaml"
MAX_SLUG_CHARS = 40
class SessionError(Exception):
pass
def render_question_md(question: str, assumptions: list[str] | None = None) -> str:
"""Rende question.md in modo deterministico: domanda + assunzioni opzionali.
Unica fonte di formattazione per question.md (riusata da create_session e
set_question), così il gate non deve costruire markdown a mano.
"""
body = f"# Domanda\n\n{question.strip()}\n"
items = [a.strip() for a in (assumptions or []) if a.strip()]
if items:
body += "\n## Assunzioni\n\n" + "".join(f"- {a}\n" for a in items)
return body
def _new_id(question: str, sessions_root: Path, stamp: str) -> str:
base = f"{stamp}-{slugify(question)[:MAX_SLUG_CHARS].rstrip('-')}"
candidate, n = base, 1
while (sessions_root / candidate).exists():
n += 1
candidate = f"{base}-{n}"
return candidate
def create_session(
question: str, db: DatabaseConfig, sessions_root: Path
) -> SessionManifest:
now = datetime.now(UTC)
# stamp con ora/min/sec: identifica univocamente sessioni dello stesso giorno
# sulla stessa domanda. Il contatore -n resta come rete per collisioni nello
# stesso secondo.
session_id = _new_id(question, sessions_root, now.strftime("%Y-%m-%d-%H%M%S"))
manifest = SessionManifest(
id=session_id, created_at=now, question=question,
database=db.database, schema=db.db_schema,
)
session_dir = sessions_root / session_id
manifest.to_yaml(session_dir / MANIFEST)
(session_dir / "question.md").write_text(render_question_md(question))
return manifest
def set_question(
session_id: str,
question: str,
assumptions: list[str],
sessions_root: Path,
) -> Path:
"""Riscrive question.md (Fase 3) in modo deterministico, senza edit tool.
Valida l'esistenza della sessione (SessionError se assente) e ritorna il
path scritto. Non tocca il manifest né le decisioni.
"""
load_session(session_id, sessions_root)
path = sessions_root / session_id / "question.md"
path.write_text(render_question_md(question, assumptions))
return path
def load_session(session_id: str, sessions_root: Path) -> SessionManifest:
path = sessions_root / session_id / MANIFEST
if not path.exists():
raise SessionError(f"Sessione non trovata: {session_id} (atteso {path})")
return SessionManifest.from_yaml(path)
def close_session(session_id: str, sessions_root: Path) -> SessionManifest:
manifest = load_session(session_id, sessions_root)
manifest.status = "closed"
manifest.to_yaml(sessions_root / session_id / MANIFEST)
return manifest
+69
View File
@@ -0,0 +1,69 @@
"""Per-step task document generator (spec D16, §4.9).
Emette un singolo documento compatto per fase/step, derivato dagli artefatti precedenti
e dalla vista effective delle decisioni, con byte budget enforced (target <20k token
per un modello 35B/<200k). MAI incorpora physical.yaml (~190k token, fatale).
D15+D16 complementari: il task doc e' generato dalla vista effective_decisions, quindi
post-rollback riflette automaticamente lo stato corretto (le decisioni stale di fasi
> current_phase sono escluse).
"""
from __future__ import annotations
from dataclasses import dataclass
from pathlib import Path
from tht.phase import effective_decisions
from tht.workflow import load_workflow
MAX_BODY_BYTES = 80_000 # ~20k token (target per task document di una fase)
@dataclass
class TaskDoc:
phase: int
body: str
byte_budget_ok: bool
def generate_task_doc(
session_dir: Path | str,
phase: int,
promoted_tables: list[str] | None = None,
) -> TaskDoc:
"""Genera il documento di task per la fase `phase`.
Contenuto (compatti, mai artefatti integrali fatali):
- Domanda (question.md) se presente.
- Schema linking (schema_linking.json) solo da fase >= 4.
- Brief delle decisioni effective (esclude stale post-rollback, esclude ritirate).
- Header del task con il numero/nome della fase.
"""
session_dir = Path(session_dir)
parts: list[str] = []
q = session_dir / "question.md"
if q.exists():
parts.append("## Domanda\n" + q.read_text())
sl = session_dir / "schema_linking.json"
if sl.exists() and phase >= 4:
parts.append("## Schema linking (deciso)\n```json\n" + sl.read_text() + "\n```")
# Brief decisioni effective (D15-aware)
eff = effective_decisions(session_dir)
if eff:
lines = [f"- {d.type} | {d.subject} | {d.detail}" for d in eff]
parts.append("## Decisioni effettive (effective)\n" + "\n".join(lines))
# Header fase
try:
wf = load_workflow()
name = wf.phase_name(phase)
header = f"## Task: fase {phase} ({name})"
except Exception:
header = f"## Task: fase {phase}"
parts.append(header)
body = "\n\n".join(parts)
return TaskDoc(phase=phase, body=body, byte_budget_ok=len(body.encode()) <= MAX_BODY_BYTES)
+54
View File
@@ -0,0 +1,54 @@
"""Artifact teardown on rollback (spec D15, §4.8).
teardown_to_phase cancella ogni artefatto la cui fase produttrice > target, usando la
mappa artifacts_out di workflow.yaml. Risolve il bug latente di the reference implementation: dopo un
re-derive con piano CTE diverso, i vecchi ctes/*.sql orfani restavano su disco e
bloccavano finalize (che itera glob('*.sql') esigendo che ognuno sia testato).
Da chiamare insieme all'append di phase_reopened per mantenere lo stato coerente
(l'artefatto su disco e il ledger effective devono allinearsi -- vedi phase.py).
"""
from __future__ import annotations
import shutil
from dataclasses import dataclass, field
from pathlib import Path
from tht.workflow import load_workflow
@dataclass
class TeardownReport:
target_phase: int
deleted_files: list[str] = field(default_factory=list)
def teardown_to_phase(session_dir: str | Path, target_phase: int) -> TeardownReport:
"""Cancella gli artefatti delle fasi > target_phase. Ritorna il report dei cancellati.
- File artefatto (es. 'schema_linking.json'): unlink se esiste.
- Directory artefatto (es. 'ctes/'): rimuove ricorsivamente (con tutti i .sql orfani).
- Artefatti di fase <= target: preservati (sono lavoro valido).
- Artefatti mancanti: noop (sessione nuova).
"""
session_dir = Path(session_dir)
wf = load_workflow()
report = TeardownReport(target_phase=target_phase)
for phase in wf.phases:
if phase.num <= target_phase:
continue
for artifact in phase.artifacts_out:
target = session_dir / artifact.rstrip("/")
is_dir = artifact.endswith("/")
if is_dir:
if target.exists() and target.is_dir():
# registra ogni file prima di rimuovere (utile per audit/debug orfani)
for f in sorted(target.glob("*")):
if f.is_file():
report.deleted_files.append(f.name)
shutil.rmtree(target)
else:
if target.exists():
target.unlink()
report.deleted_files.append(artifact)
return report
+8
View File
@@ -0,0 +1,8 @@
import re
import unicodedata
def slugify(text: str) -> str:
"""Slug ASCII minuscolo; preserva gli underscore (nomi tabella)."""
text = unicodedata.normalize("NFKD", text).encode("ascii", "ignore").decode()
return re.sub(r"[^a-z0-9_]+", "-", text.lower()).strip("-")
View File
+50
View File
@@ -0,0 +1,50 @@
import requests
from tht.config import EmbeddingsConfig
DOC_PREFIX = "search_document: "
QUERY_PREFIX = "search_query: "
class EmbeddingsError(Exception):
pass
class OllamaEmbeddings:
"""Client embeddings via Ollama. Applica i prefissi di task richiesti da nomic v2:
ometterli degrada il retrieval in modo silenzioso."""
def __init__(self, cfg: EmbeddingsConfig):
self.cfg = cfg
def _embed(self, texts: list[str]) -> list[list[float]]:
url = f"{self.cfg.base_url.rstrip('/')}/api/embed"
out: list[list[float]] = []
for i in range(0, len(texts), self.cfg.batch_size):
batch = texts[i : i + self.cfg.batch_size]
try:
resp = requests.post(
url, json={"model": self.cfg.model, "input": batch},
timeout=self.cfg.timeout,
)
resp.raise_for_status()
except requests.RequestException as e:
raise EmbeddingsError(
f"Ollama non raggiungibile su {self.cfg.base_url} "
f"(modello {self.cfg.model}): {e}"
) from e
embeddings = resp.json().get("embeddings", [])
for v in embeddings:
if len(v) != self.cfg.dim:
raise EmbeddingsError(
f"dimensione embedding inattesa: {len(v)} != {self.cfg.dim} "
f"(modello {self.cfg.model})"
)
out.extend(embeddings)
return out
def embed_documents(self, texts: list[str]) -> list[list[float]]:
return self._embed([DOC_PREFIX + t for t in texts])
def embed_query(self, text: str) -> list[float]:
return self._embed([QUERY_PREFIX + text])[0]
+67
View File
@@ -0,0 +1,67 @@
"""Lettura del pgvector dietro un'unica interfaccia `.search(query_vec, top_n, kinds)`, così
`search.combined_search` resta agnostico al transport. Due implementazioni:
- `RestSearcher` → produzione: similarity search via REST (`search_similar`).
- `DirectSearcher` → dev/test: connessione diretta a Postgres/pgvector.
Entrambe mappano i `kind` sulle tabelle per-dominio dello schema `vectors`.
"""
from sqlalchemy import Engine
from tht.vectorstore.rest_client import VectorRestClient
from tht.vectorstore.store import VectorHit, VectorStore, hit_from_metadata
# kind Thoth → tabella dello schema `vectors`.
KIND_TO_TABLE = {
"schema_table": "schema_records",
"schema_column": "schema_records",
"evidence": "evidence",
"memory": "memory",
}
ALL_TABLES = ["schema_records", "evidence", "memory"]
def tables_for_kinds(kinds: list[str] | None) -> list[str]:
"""Tabelle da interrogare per i kind richiesti (tutte se kinds è vuoto/None)."""
if not kinds:
return list(ALL_TABLES)
return sorted({KIND_TO_TABLE[k] for k in kinds if k in KIND_TO_TABLE})
def _merge(hits: list[VectorHit], top_n: int) -> list[VectorHit]:
return sorted(hits, key=lambda h: h.similarity, reverse=True)[:top_n]
class RestSearcher:
"""Similarity search via REST: una chiamata `search_similar` per tabella, poi fusione."""
def __init__(self, client: VectorRestClient):
self.client = client
def search(
self, query_vec: list[float], top_n: int = 10, kinds: list[str] | None = None
) -> list[VectorHit]:
hits: list[VectorHit] = []
for table in tables_for_kinds(kinds):
for row in self.client.search_similar(table, query_vec, top_n):
hits.append(hit_from_metadata(row.get("similarity", 0.0), row.get("metadata")))
return _merge(hits, top_n)
class DirectSearcher:
"""Similarity search diretta su Postgres/pgvector, interrogando le tabelle per-dominio."""
def __init__(self, engine: Engine, schema: str = "vectors", dim: int = 768):
self.engine = engine
self.schema = schema
self.dim = dim
def search(
self, query_vec: list[float], top_n: int = 10, kinds: list[str] | None = None
) -> list[VectorHit]:
hits: list[VectorHit] = []
for table in tables_for_kinds(kinds):
store = VectorStore(self.engine, schema=self.schema, table=table, dim=self.dim)
hits.extend(store.search(query_vec, top_n=top_n))
return _merge(hits, top_n)
+101
View File
@@ -0,0 +1,101 @@
import re
from pydantic import BaseModel
from tht.evidence.model import EvidenceDoc
from tht.mschema.models import Annotations, PhysicalSchema
MAX_EXAMPLES_IN_RECORD = 5
class VectorRecord(BaseModel):
id: str
kind: str # evidence | schema_table | schema_column
ref: str # file/chiave canonica di provenienza
title: str
content: str
metadata: dict = {}
def split_markdown(text: str, max_chars: int) -> list[str]:
"""Spezza un markdown: intero se sta nel limite, altrimenti per heading '##',
e in ultima istanza per accumulo greedy di righe."""
if len(text) <= max_chars:
return [text]
parts = re.split(r"(?=^## )", text, flags=re.MULTILINE)
chunks: list[str] = []
for part in parts:
part = part.strip("\n")
if not part:
continue
if len(part) <= max_chars:
chunks.append(part)
continue
current: list[str] = []
size = 0
for line in part.splitlines():
if size + len(line) > max_chars and current:
chunks.append("\n".join(current))
current, size = [], 0
current.append(line)
size += len(line) + 1
if current:
chunks.append("\n".join(current))
return chunks
def evidence_records(docs: list[EvidenceDoc], max_chunk_chars: int) -> list[VectorRecord]:
"""Record per tutte le evidence presenti: la sola presenza basta a indicizzarle."""
records: list[VectorRecord] = []
for doc in docs:
content = f"{doc.title}\n\n{doc.body}"
for i, chunk in enumerate(split_markdown(content, max_chunk_chars)):
records.append(
VectorRecord(
id=f"evidence:{doc.id}:{i}",
kind="evidence",
ref=str(doc.path) if doc.path else doc.id,
title=doc.title,
content=chunk,
metadata={
"status": doc.status, "tier": doc.tier,
"tables": doc.tables, "concepts": doc.concepts,
},
)
)
return records
def schema_records(physical: PhysicalSchema, annotations: Annotations) -> list[VectorRecord]:
"""Un record per tabella e uno per colonna, da mschema (physical + annotations)."""
records: list[VectorRecord] = []
for table_name, table in physical.tables.items():
ann_t = annotations.tables.get(table_name)
t_desc = (ann_t.description if ann_t and ann_t.description else table.comment)
t_concepts = ann_t.concepts if ann_t else []
lines = [f"Tabella {table_name}", t_desc]
if t_concepts:
lines.append("Concetti: " + ", ".join(t_concepts))
lines.append("Colonne: " + ", ".join(table.columns))
records.append(
VectorRecord(
id=f"schema_table:{table_name}", kind="schema_table", ref=table_name,
title=table_name, content="\n".join(filter(None, lines)),
)
)
for column_name, column in table.columns.items():
ann_c = ann_t.columns.get(column_name) if ann_t else None
c_desc = (ann_c.description if ann_c and ann_c.description else column.comment)
lines = [f"Colonna {table_name}.{column_name} ({column.type})", c_desc]
if ann_c and ann_c.synonyms:
lines.append("Sinonimi: " + ", ".join(ann_c.synonyms))
if column.examples:
lines.append("Esempi: " + ", ".join(column.examples[:MAX_EXAMPLES_IN_RECORD]))
records.append(
VectorRecord(
id=f"schema_column:{table_name}.{column_name}", kind="schema_column",
ref=f"{table_name}.{column_name}", title=f"{table_name}.{column_name}",
content="\n".join(filter(None, lines)),
)
)
return records
+104
View File
@@ -0,0 +1,104 @@
"""Client per la similarity search del pgvector esposta via Supabase/PostgREST.
Endpoint dedicato (es. https://host/vector/v1/), distinto dal DWH. La lettura usa
`search_similar`; la scrittura remota usa RPC allowlist con una API key separata.
Errori in italiano e azionabili, stile `rest/client.py`.
"""
import requests
from tht.config import RestConfig
class VectorRestError(Exception):
"""Errore di accesso al vector store via REST, con messaggio leggibile per il reviewer."""
class VectorRestClient:
def __init__(self, cfg: RestConfig):
self.cfg = cfg
self._base = cfg.base_url.rstrip("/")
@property
def api_key(self) -> str:
"""The REST API key for this client (spec D11: reader and writer carry
distinct keys against the same endpoint)."""
return self.cfg.api_key
def _post(self, fn: str, args: dict) -> requests.Response:
url = f"{self._base}/rpc/{fn}"
verify: bool | str = self.cfg.ssl_ca if self.cfg.ssl_ca else True
try:
return requests.post(
url,
json=args,
headers={"X-API-Key": self.cfg.api_key},
timeout=self.cfg.timeout,
verify=verify,
)
except requests.RequestException as e:
raise VectorRestError(
f"Vector REST non raggiungibile su {self.cfg.base_url} (rpc {fn}): {e}"
) from e
def _error_msg(self, fn: str, resp: requests.Response) -> str:
try:
body = resp.json()
detail = body.get("message") or body.get("details") or resp.text
except Exception:
detail = resp.text
return f"Vector REST rpc {fn} → HTTP {resp.status_code}: {detail}"
def _call(self, fn: str, args: dict):
resp = self._post(fn, args)
if not resp.ok:
raise VectorRestError(self._error_msg(fn, resp))
if resp.status_code == 204 or not resp.text:
return None
return resp.json()
def search_similar(
self, table_name: str, query_embedding: list[float], limit_count: int
) -> list[dict]:
"""Ricerca per similarità coseno su `vectors.<table_name>`: ritorna le righe
`{id, similarity, metadata}` ordinate per similarity decrescente."""
return self._call(
"search_similar",
{
"query_embedding": query_embedding,
"limit_count": limit_count,
"table_name": table_name,
},
) or []
def list_tables(self) -> list[dict]:
"""Tabelle vettoriali disponibili: `{table_name, vector_dimensions, …}`."""
return self._call("list_tables", {}) or []
def existing_hashes(self, table_name: str, kinds: list[str]) -> dict[str, str]:
"""Hash correnti per sync incrementale su una tabella vector allowlisted.
RPC attesa: `existing_vector_hashes(table_name, kinds)` -> righe
`{record_key, content_hash}`.
"""
rows = self._call(
"existing_vector_hashes",
{"table_name": table_name, "kinds": kinds},
) or []
return {row["record_key"]: row["content_hash"] for row in rows}
def upsert_records(self, table_name: str, rows: list[dict]) -> int:
"""Upsert controllato di record vettoriali già embeddati.
RPC attesa: `upsert_vector_records(table_name, rows)` -> `{upserted: N}` o righe.
Non espone delete/clear: il cleanup distruttivo resta solo-server.
"""
payload = self._call(
"upsert_vector_records",
{"table_name": table_name, "rows": rows},
)
if payload is None:
return len(rows)
if isinstance(payload, dict):
return int(payload.get("upserted", len(rows)))
return len(payload) if isinstance(payload, list) else len(rows)
+90
View File
@@ -0,0 +1,90 @@
"""Scrittura controllata del pgvector via REST.
Usata dalle postazioni remote solo quando e' configurata una seconda API key di scrittura.
Mantiene l'upsert incrementale del VectorStore diretto, ma non esegue delete/clear: le
operazioni distruttive restano solo-server via connessione Postgres diretta.
"""
from tht.vectorstore.records import VectorRecord
from tht.vectorstore.rest_client import VectorRestClient
from tht.vectorstore.store import SyncStats, content_hash
KIND_TO_TABLE = {
"schema_table": "schema_records",
"schema_column": "schema_records",
"evidence": "evidence",
"memory": "memory",
}
TABLE_TO_KINDS = {
"schema_records": {"schema_table", "schema_column"},
"evidence": {"evidence"},
"memory": {"memory"},
}
def pack_metadata(record: VectorRecord) -> dict:
"""Impacchetta nel metadata tutta la semantica letta poi da `search_similar`."""
return {
"kind": record.kind,
"ref": record.ref,
"record_key": record.id,
"title": record.title,
"content": record.content,
**record.metadata,
}
class RestVectorWriter:
"""Writer table-scoped via RPC REST allowlist.
Il metodo `sync` e' volutamente upsert-only: aggiorna/aggiunge record, conta gli stale,
ma non li elimina. Per cleanup completo usare i comandi server-side con `vector_db`.
"""
def __init__(self, client: VectorRestClient, table: str):
if table not in TABLE_TO_KINDS:
raise ValueError(f"Tabella vector non supportata per scrittura REST: {table}")
self.client = client
self.table = table
def existing_hashes(self, kinds: set[str]) -> dict[str, str]:
allowed = TABLE_TO_KINDS[self.table]
bad = kinds - allowed
if bad:
raise ValueError(
f"Kind non ammessi per vectors.{self.table}: {', '.join(sorted(bad))}"
)
return self.client.existing_hashes(self.table, sorted(kinds))
def sync(self, records: list[VectorRecord], embedder, kinds: set[str]) -> SyncStats:
stats = SyncStats()
existing = self.existing_hashes(kinds)
to_embed: list[VectorRecord] = []
for record in records:
h = content_hash(record.content)
if record.id not in existing:
to_embed.append(record)
stats.added += 1
elif existing[record.id] != h:
to_embed.append(record)
stats.updated += 1
else:
stats.unchanged += 1
stats.deleted = 0
vectors = embedder.embed_documents([r.content for r in to_embed]) if to_embed else []
rows = [
{
"record_key": record.id,
"kind": record.kind,
"content_hash": content_hash(record.content),
"metadata": pack_metadata(record),
"embedding": vector,
}
for record, vector in zip(to_embed, vectors)
]
if rows:
self.client.upsert_records(self.table, rows)
# Gli stale non vengono cancellati in REST writer: restano responsabilita' server-side.
return stats
+191
View File
@@ -0,0 +1,191 @@
import hashlib
import json
from dataclasses import dataclass
from sqlalchemy import Engine, text
from tht.vectorstore.records import VectorRecord
def content_hash(content: str) -> str:
return hashlib.sha256(content.encode()).hexdigest()
def _to_vector_literal(vec: list[float]) -> str:
return "[" + ",".join(f"{x:.8f}" for x in vec) + "]"
@dataclass
class SyncStats:
added: int = 0
updated: int = 0
deleted: int = 0
unchanged: int = 0
@dataclass
class VectorHit:
id: str
kind: str
ref: str
title: str
content: str
metadata: dict
similarity: float
def hit_from_metadata(similarity: float, metadata: dict | None) -> VectorHit:
"""Ricostruisce un VectorHit dal solo `metadata` (più la similarity). È l'unico modo
disponibile leggendo via REST (`search_similar` ritorna id/similarity/metadata), e viene
usato anche dalla lettura diretta per avere un'unica logica. Tollerante: usa default sui
campi assenti (es. metadata estranei della tabella fake remota)."""
md = metadata or {}
return VectorHit(
id=md.get("record_key", ""),
kind=md.get("kind", ""),
ref=md.get("ref", ""),
title=md.get("title", ""),
content=md.get("content", ""),
metadata=md,
similarity=float(similarity),
)
class VectorStore:
"""Tabella pgvector table-scoped: scrittura diretta (loading) su una tabella dello schema
`vectors`. La lettura via REST avviene su `search_similar`; questo store serve al loading e
alla lettura diretta (dev/test). Il contratto della tabella remota richiede `id` (BIGSERIAL),
`embedding vector(N)` e `metadata jsonb`; le colonne extra (`record_key`, `kind`,
`content_hash`) servono solo al loader e non sono esposte dalla REST."""
def __init__(
self, engine: Engine, schema: str = "vectors", table: str = "records", dim: int = 768
):
self.engine = engine
self.schema = schema
self.dim = dim
self._table = f"{schema}.{table}"
def init_schema(self) -> None:
with self.engine.begin() as conn:
conn.execute(text("CREATE EXTENSION IF NOT EXISTS vector"))
conn.execute(text(f"CREATE SCHEMA IF NOT EXISTS {self.schema}"))
conn.execute(text(f"""
CREATE TABLE IF NOT EXISTS {self._table} (
id bigserial PRIMARY KEY,
record_key text UNIQUE NOT NULL,
kind text NOT NULL,
content_hash text NOT NULL,
metadata jsonb NOT NULL DEFAULT '{{}}',
embedding vector({self.dim}) NOT NULL,
indexed_at timestamptz NOT NULL DEFAULT now()
)
"""))
conn.execute(text(
f"CREATE INDEX IF NOT EXISTS {self._idx('embedding')} ON {self._table} "
f"USING hnsw (embedding vector_cosine_ops) WITH (m = 16, ef_construction = 200)"
))
conn.execute(text(
f"CREATE INDEX IF NOT EXISTS {self._idx('kind')} ON {self._table} (kind)"
))
# GRANT al ruolo di sola lettura della REST, solo se esiste (assente in test/locale).
conn.execute(text(f"""
DO $$ BEGIN
IF EXISTS (SELECT 1 FROM pg_roles WHERE rolname = 'vector_reader') THEN
EXECUTE 'GRANT SELECT ON {self._table} TO vector_reader';
END IF;
END $$;
"""))
def _idx(self, suffix: str) -> str:
return f"{self._table.replace('.', '_')}_{suffix}_idx"
def clear(self) -> None:
with self.engine.begin() as conn:
conn.execute(text(f"DELETE FROM {self._table}"))
def existing_hashes(self, kinds: set[str]) -> dict[str, str]:
q = text(
f"SELECT record_key, content_hash FROM {self._table} WHERE kind = ANY(:kinds)"
)
with self.engine.connect() as conn:
return dict(conn.execute(q, {"kinds": list(kinds)}).fetchall())
def sync(self, records: list[VectorRecord], embedder, kinds: set[str]) -> SyncStats:
"""Allinea l'indice ai record correnti (per i kind dati): embedda solo il nuovo
o il modificato, elimina cio' che non esiste piu'."""
stats = SyncStats()
existing = self.existing_hashes(kinds)
current_ids = {r.id for r in records}
to_embed: list[VectorRecord] = []
for r in records:
h = content_hash(r.content)
if r.id not in existing:
to_embed.append(r)
stats.added += 1
elif existing[r.id] != h:
to_embed.append(r)
stats.updated += 1
else:
stats.unchanged += 1
vectors = embedder.embed_documents([r.content for r in to_embed]) if to_embed else []
upsert = text(f"""
INSERT INTO {self._table}
(record_key, kind, content_hash, metadata, embedding)
VALUES
(:record_key, :kind, :content_hash, CAST(:metadata AS jsonb),
CAST(:embedding AS vector))
ON CONFLICT (record_key) DO UPDATE SET
kind = EXCLUDED.kind, content_hash = EXCLUDED.content_hash,
metadata = EXCLUDED.metadata, embedding = EXCLUDED.embedding,
indexed_at = now()
""")
stale = [i for i in existing if i not in current_ids]
with self.engine.begin() as conn:
for r, vec in zip(to_embed, vectors):
conn.execute(upsert, {
"record_key": r.id, "kind": r.kind,
"content_hash": content_hash(r.content),
"metadata": json.dumps(_pack_metadata(r)),
"embedding": _to_vector_literal(vec),
})
if stale:
conn.execute(
text(f"DELETE FROM {self._table} WHERE record_key = ANY(:ids)"),
{"ids": stale},
)
stats.deleted = len(stale)
return stats
def search(
self, query_vec: list[float], top_n: int = 10, kinds: list[str] | None = None
) -> list[VectorHit]:
where = "WHERE kind = ANY(:kinds)" if kinds else ""
q = text(f"""
SELECT metadata, 1 - (embedding <=> CAST(:q AS vector)) AS similarity
FROM {self._table}
{where}
ORDER BY embedding <=> CAST(:q AS vector)
LIMIT :top_n
""")
params: dict = {"q": _to_vector_literal(query_vec), "top_n": top_n}
if kinds:
params["kinds"] = kinds
with self.engine.connect() as conn:
rows = conn.execute(q, params).fetchall()
return [hit_from_metadata(r.similarity, r.metadata) for r in rows]
def _pack_metadata(r: VectorRecord) -> dict:
"""Impacchetta nel `metadata` (unica colonna letta via REST) tutta la semantica Thoth."""
return {
"kind": r.kind,
"ref": r.ref,
"record_key": r.id,
"title": r.title,
"content": r.content,
**r.metadata,
}
+107
View File
@@ -0,0 +1,107 @@
"""Reads workflow.yaml -- the SINGLE source of workflow truth (spec F2, §5.3).
phase.py, the gate (tht-gate.js), and the skill all read from here.
No more duplicated constants (the JS/Python drift bug in the reference implementation -- PHASE_NAMES
truncated to 7 in JS -- is structurally impossible because there is one source).
Edit workflow.yaml to change the workflow: add/reorder/merge/skip phases.
"""
from __future__ import annotations
from dataclasses import dataclass, field
from pathlib import Path
from typing import Any
import yaml
_WF_PATH = Path(__file__).resolve().parent.parent / "workflow.yaml"
@dataclass
class PhaseSpec:
id: str
num: int
name: str
advance: str
prerequisites: list[Any]
artifacts_out: list[str] = field(default_factory=list)
@dataclass
class Workflow:
schema_version: int
phases: list[PhaseSpec]
_decision_min_map: dict[str, int] = field(default_factory=dict)
@property
def max_phase(self) -> int:
return len(self.phases)
def phase_by_num(self, n: int) -> PhaseSpec:
return self.phases[n - 1]
def phase_name(self, n: int) -> str:
if 1 <= n <= self.max_phase:
return self.phase_by_num(n).name
return "?"
def decision_min_phase(self, decision_type: str) -> int:
"""A decision type's min phase = the earliest phase whose prerequisites
reference it (via decision_exists / decision_subject_exists). Defaults to 1."""
return self._decision_min_map.get(decision_type, 1)
def _collect_decision_mins(phases: list[PhaseSpec]) -> dict[str, int]:
"""Scan prerequisites for decision_exists / decision_subject_exists mentions.
Supports both forms:
- decision_exists: <type> (scalar)
- decision_exists: [<type>, ...] (list, first element is the type)
- decision_subject_exists: [<type>, <subject>] (list, first element is the type)
"""
mins: dict[str, int] = {}
def scan(node: Any, phase_num: int) -> None:
if isinstance(node, dict):
for key, value in node.items():
if key in ("decision_exists", "decision_subject_exists"):
if isinstance(value, list) and value:
dtype = value[0]
elif isinstance(value, str):
dtype = value
else:
continue
if isinstance(dtype, str):
if dtype not in mins or phase_num < mins[dtype]:
mins[dtype] = phase_num
else:
scan(value, phase_num)
elif isinstance(node, list):
for item in node:
scan(item, phase_num)
for p in phases:
scan(p.prerequisites, p.num)
return mins
def load_workflow(path: Path | str = _WF_PATH) -> Workflow:
path = Path(path)
raw = yaml.safe_load(path.read_text())
phases: list[PhaseSpec] = []
for i, p in enumerate(raw["phases"], start=1):
phases.append(
PhaseSpec(
id=p["id"],
num=i,
name=p["name"],
advance=p["advance"],
prerequisites=p.get("prerequisites", []),
artifacts_out=p.get("artifacts_out", []),
)
)
return Workflow(
schema_version=raw.get("schema_version", 1),
phases=phases,
_decision_min_map=_collect_decision_mins(phases),
)
+31
View File
@@ -0,0 +1,31 @@
"""Workspace YAML loading — the single boundary for workspace configuration (spec D3).
Reads workspaces/<name>.yaml, expands ${VAR} from env, validates via the Config model
(ported from the reference implementation). Future migration to a DB store would replace only this module.
La struttura YAML rispecchia esattamente tht/config.py:
database + rest (DWH), vector_rest + vector_write_rest (pgvector, doppia key top-level),
vector_db (loading diretto, server-only), embeddings, evidence, execution.
"""
from __future__ import annotations
from pathlib import Path
from tht.config import Config, ConfigError, load_config
class WorkspaceError(Exception):
"""Errore di caricamento del workspace (file mancante, env var non definita, yaml invalido)."""
def load_workspace(path: str | Path) -> Config:
"""Carica e valida un workspace YAML.
Espande ${VAR} dall'ambiente; se una variabile referenziata manca, raises WorkspaceError.
Delega a load_config (portato da the reference implementation) per la validazione del modello Config.
"""
path = Path(path)
try:
return load_config(path)
except ConfigError as e:
raise WorkspaceError(str(e)) from e