Files
ThothII/harness/tht/cli/vector_cmd.py
T

173 lines
5.9 KiB
Python

from pathlib import Path
import typer
from tht.cli._guards import (
has_vector_write_rest,
require_server_profile,
require_vector_write_allowed,
)
from tht.cli.config_cmd import CONFIG_OPT
from tht.cli.schema_cmd import _load_config_or_exit, annotations_path, physical_path
from tht.ports.vector import VectorWriteRecord
from tht.vectorstore.store import SyncStats, content_hash
vector_app = typer.Typer(help="Indice semantico pgvector (derivato, rigenerabile)")
def make_embedder(embeddings_cfg):
"""Factory del client embeddings (monkeypatchabile nei test)."""
from tht.vectorstore.embeddings import OllamaEmbeddings
return OllamaEmbeddings(embeddings_cfg)
def require_vector_cfg(cfg):
missing = []
if cfg.embeddings is None:
missing.append("embeddings")
if cfg.vectors is None and cfg.vector_db is None and not has_vector_write_rest(cfg):
missing.append("vectors o vector_db o vector_write_rest")
if missing:
typer.secho(
f"ERRORE: sezioni mancanti nel workspace yaml: {', '.join(missing)}.",
fg=typer.colors.RED, err=True,
)
raise typer.Exit(code=1)
def require_direct_vector_cfg(cfg):
missing = [k for k in ("vector_db", "embeddings") if getattr(cfg, k) is None]
if missing:
typer.secho(
f"ERRORE: sezioni mancanti nel workspace yaml: {', '.join(missing)}.",
fg=typer.colors.RED, err=True,
)
raise typer.Exit(code=1)
def open_store(cfg, table: str):
"""Writer table-scoped per il LOADING.
Sul server preferisce la connessione diretta. In profilo workstation usa `vector_write_rest`
se configurato, con upsert remoto non distruttivo.
"""
from tht.adapters.factory import build_vector_loader
return build_vector_loader(cfg, table)
def open_searcher(cfg):
"""Searcher per la LETTURA (similarity search): via REST se `vector_rest` è configurato,
altrimenti connessione diretta (dev/test)."""
from tht.adapters.factory import build_vector_store
from tht.vectorstore.reader import tables_for_kinds
store = build_vector_store(cfg)
class AdapterSearcher:
def search(self, query_vec, top_n=10, kinds=None, metadata_filter=None):
return store.search(
tables_for_kinds(kinds), query_vec, limit=top_n, kinds=kinds,
metadata_filter=metadata_filter,
)
return AdapterSearcher()
def sync_canonical_records(collection, records, *, store, embedder):
kinds = sorted({record.kind for record in records})
existing = store.existing_hashes(collection, kinds)
pending = []
stats = SyncStats()
changed = []
for record in records:
hashed = content_hash(record.content)
current = existing.get(record.id)
if current == hashed:
stats.unchanged += 1
continue
changed.append((record, hashed, current is None))
if changed:
embeddings = embedder.embed_documents([record.content for record, *_ in changed])
for (record, hashed, is_added), embedding in zip(changed, embeddings, strict=True):
pending.append(VectorWriteRecord(record=record, embedding=embedding, content_hash=hashed))
if is_added:
stats.added += 1
else:
stats.updated += 1
store.upsert(collection, pending)
return stats
def _print_stats(stats) -> None:
typer.secho(
f"OK: {stats.added} nuovi, {stats.updated} aggiornati, "
f"{stats.deleted} rimossi, {stats.unchanged} invariati",
fg=typer.colors.GREEN,
)
@vector_app.command("init")
def init_cmd(
config: Path = CONFIG_OPT,
skip_ollama_check: bool = typer.Option(
False, "--skip-ollama-check", help="Non verificare la raggiungibilita' di Ollama."
),
) -> None:
"""Crea schema e tabella pgvector (idempotente) e verifica le connessioni."""
from sqlalchemy.exc import OperationalError
from tht.vectorstore.embeddings import EmbeddingsError
from tht.vectorstore.reader import ALL_TABLES
cfg = _load_config_or_exit(config)
require_server_profile(cfg, "vector init")
require_direct_vector_cfg(cfg)
try:
for table in ALL_TABLES:
open_store(cfg, table).init_schema()
except OperationalError as e:
typer.secho(f"ERRORE connessione pgvector: {e.orig}", fg=typer.colors.RED, err=True)
raise typer.Exit(code=1)
if not skip_ollama_check:
try:
make_embedder(cfg.embeddings).embed_query("ping")
except EmbeddingsError as e:
typer.secho(f"ERRORE: {e}", fg=typer.colors.RED, err=True)
raise typer.Exit(code=1)
typer.secho(
f"OK: schema {cfg.vector_db.db_schema} pronto (tabelle: {', '.join(ALL_TABLES)}) su "
f"{cfg.vector_db.host}:{cfg.vector_db.port}", fg=typer.colors.GREEN,
)
@vector_app.command("index-schema")
def index_schema_cmd(config: Path = CONFIG_OPT) -> None:
"""Embedda e sincronizza i record schema (tabelle e colonne) da mschema."""
from tht.mschema.models import Annotations, PhysicalSchema
from tht.vectorstore.records import schema_records
cfg = _load_config_or_exit(config)
require_vector_write_allowed(cfg, "vector index-schema")
require_vector_cfg(cfg)
phys_file = physical_path(cfg)
if not phys_file.exists():
typer.secho(
f"ERRORE: {phys_file} non trovato. Esegui prima `tht schema introspect`.",
fg=typer.colors.RED, err=True,
)
raise typer.Exit(code=1)
physical = PhysicalSchema.from_yaml(phys_file)
annotations = Annotations.from_yaml(annotations_path(cfg))
records = schema_records(physical, annotations)
from tht.adapters.factory import build_vector_store
stats = sync_canonical_records(
"schema_records",
records,
store=build_vector_store(cfg, require_write=True),
embedder=make_embedder(cfg.embeddings),
)
_print_stats(stats)