import hashlib import json from pathlib import Path import typer from tht.cli._guards import require_vector_write_allowed from tht.cli.config_cmd import CONFIG_OPT from tht.cli.schema_cmd import _load_config_or_exit from tht.ports.vector import VectorWriteRecord from tht.vectorstore.store import SyncStats, content_hash vector_app = typer.Typer(help="Indice semantico Qdrant (derivato, rigenerabile)") def _artifact_digest(path: Path) -> str: try: contents = path.read_bytes() except FileNotFoundError: # An absent annotations file is valid (empty curation); digest the empty artifact. contents = b"" return "sha256:" + hashlib.sha256(contents).hexdigest() def _emit_json(payload: dict) -> None: typer.echo(json.dumps(payload, ensure_ascii=False, sort_keys=True)) 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: missing.append("vectors") 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_searcher(cfg): """Searcher per la lettura semantic search sul runtime vettoriale attivo.""" 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, query_text=None, query_language=None, ): return store.search( tables_for_kinds(kinds), query_vec, limit=top_n, kinds=kinds, metadata_filter=metadata_filter, query_text=query_text, query_language=query_language, ) 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("index-schema") def index_schema_cmd( config: Path = CONFIG_OPT, json_output: bool = typer.Option(False, "--json"), ) -> None: """Replace the schema slice from the PostgreSQL Catalog projection.""" from tht.ports.vector import VectorStoreError cfg = _load_config_or_exit(config) require_vector_write_allowed(cfg, "vector index-schema") require_vector_cfg(cfg) try: stats, counts, snapshot_path = index_catalog_schema(cfg) except VectorStoreError as exc: code = str(exc) error = "semantic index incompatible" if code == "semantic_index_incompatible" else "schema indexing failed" if json_output: _emit_json({ "code": code, "error": error, "operation": "index_schema", "schemaVersion": 1, "status": "failed", "workspaceId": cfg._workspace_id, "workspaceRevision": cfg._workspace_revision, }) else: typer.secho(f"ERRORE: {error}", fg=typer.colors.RED, err=True) raise typer.Exit(code=1) from None if json_output: _emit_json({ "artifactIdentities": [ {"digest": _artifact_digest(snapshot_path), "kind": "catalog_metadata_snapshot"}, ], "code": "ok", "collection": cfg.vectors.collections["reference"], "counts": { "added": stats.added, "columns": counts["columns"], "deleted": stats.deleted, "records": counts["records"], "relationships": counts["relationships"], "tables": counts["tables"], "unchanged": stats.unchanged, "updated": stats.updated, }, "operation": "index_schema", "schemaVersion": 1, "status": "succeeded", "workspaceId": cfg._workspace_id, "workspaceRevision": cfg._workspace_revision, }) return _print_stats(stats) def index_catalog_schema(cfg): from tht.adapters.factory import build_vector_store from tht.mschema.catalog_snapshot import load_catalog_metadata_snapshot from tht.vectorstore.records import catalog_schema_records snapshot_path = cfg.paths.catalog_metadata_snapshot if snapshot_path is None: raise ValueError("catalog metadata snapshot is not configured") snapshot = load_catalog_metadata_snapshot(snapshot_path, cfg._workspace_id) records = catalog_schema_records(snapshot) store = build_vector_store(cfg, require_write=True) deleted = store.delete_kinds( "schema_records", ["schema_table", "schema_column", "schema_relationship"] ) stats = sync_canonical_records( "schema_records", records, store=store, embedder=make_embedder(cfg.embeddings), ) stats.deleted = deleted return stats, { "tables": len(snapshot.tables), "columns": sum(len(table.columns) for table in snapshot.tables), "relationships": len(snapshot.relationships), "records": len(records), }, snapshot_path