193 lines
6.4 KiB
Python
193 lines
6.4 KiB
Python
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
|