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

316 lines
11 KiB
Python

import logging
from collections.abc import Mapping
from pathlib import Path
from typing import Literal, TypedDict
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,
_load_schema_config,
annotations_path,
physical_path,
)
from tht.config import Config, ConfigError
from tht.ports.vector import SemanticIndexIncompatibleError, VectorWriteRecord
from tht.vectorstore.store import SyncStats, content_hash
vector_app = typer.Typer(help="Indice semantico Qdrant (derivato, rigenerabile)")
logger = logging.getLogger(__name__)
def _require_writer_capability(*, workspace_id: str | None = None, revision: str | None = None) -> None:
# Authorization is unconditional: an environment marker is attacker-controlled
# and must never turn a mutating direct invocation into an authorized child.
from tht.workspace_writer_lock import require_workspace_writer_capability
require_workspace_writer_capability(workspace_id=workspace_id, revision=revision)
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):
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, config=None):
_require_writer_capability(workspace_id=getattr(getattr(config, "runtime_identity", None), "workspace_id", None), revision=getattr(getattr(config, "runtime_identity", None), "workspace_revision", None))
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: SyncStats | Mapping[str, int]) -> None:
if isinstance(stats, SyncStats):
counts = {
"added": stats.added,
"updated": stats.updated,
"deleted": stats.deleted,
"unchanged": stats.unchanged,
}
else:
counts = stats
typer.secho(
f"OK: {counts['added']} nuovi, {counts['updated']} aggiornati, "
f"{counts['deleted']} rimossi, {counts['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:
"""Verifica il runtime Qdrant e la raggiungibilita' dell'embedder configurato."""
from tht.adapters.factory import build_vector_store
from tht.vectorstore.embeddings import EmbeddingsError
cfg = _load_config_or_exit(config)
require_server_profile(cfg, "vector init")
require_vector_cfg(cfg)
health = build_vector_store(cfg, require_write=True).health()
if not health.ok:
typer.secho(
f"ERRORE runtime vettoriale: {health.detail or 'Qdrant non raggiungibile o incompatibile'}",
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: runtime Qdrant pronto per la collezione {cfg.vectors.collection}",
fg=typer.colors.GREEN,
)
class _MachineVectorError(Exception):
"""Expected vector CLI failure with a stable machine-readable code."""
def __init__(self, code: str):
self.code = code
super().__init__(code)
class IndexCounts(TypedDict):
added: int
updated: int
deleted: int
unchanged: int
class IndexSchemaResult(TypedDict):
status: Literal["succeeded", "failed"]
code: Literal["ok"]
counts: IndexCounts
def _vector_cfg_or_error(cfg: Config) -> None:
missing = []
if cfg.embeddings is None:
missing.append("embeddings")
if cfg.vectors is None:
missing.append("vectors")
if missing:
raise _MachineVectorError("vector_configuration_missing")
def _vector_write_or_error(cfg: Config) -> None:
if cfg.profile == "workstation" and not has_vector_write_rest(cfg):
raise _MachineVectorError("vector_write_not_allowed")
def _load_schema_artifacts(cfg: Config):
import yaml
from pydantic import ValidationError
from tht.mschema.models import Annotations, PhysicalSchema
phys_file = physical_path(cfg)
if not phys_file.exists():
raise _MachineVectorError("physical_schema_missing")
try:
return (
PhysicalSchema.from_yaml(phys_file),
Annotations.from_yaml(annotations_path(cfg)),
)
except (OSError, UnicodeError, TypeError, ValueError, yaml.YAMLError, ValidationError):
raise _MachineVectorError("schema_artifacts_invalid") from None
def index_schema_data(
config: Config | Path, *, suppress_legacy_warning: bool = False, physical=None, annotations=None
) -> IndexSchemaResult:
"""Synchronize schema records and return a bounded machine result."""
from tht.mschema.models import Annotations, PhysicalSchema
from tht.vectorstore.records import schema_records
if isinstance(config, Path):
try:
cfg = _load_schema_config(config, suppress_legacy_warning=suppress_legacy_warning)
except ConfigError:
raise _MachineVectorError("invalid_configuration") from None
else:
cfg = config
_vector_write_or_error(cfg)
_vector_cfg_or_error(cfg)
if physical is None:
phys_file = physical_path(cfg)
if not phys_file.exists():
raise _MachineVectorError("physical_schema_missing")
import yaml
from pydantic import ValidationError
try:
if physical is None:
physical = PhysicalSchema.from_yaml(phys_file)
if annotations is None:
annotations = Annotations.from_yaml(annotations_path(cfg))
except (OSError, UnicodeError, TypeError, ValueError, yaml.YAMLError, ValidationError):
raise _MachineVectorError("schema_artifacts_invalid") from None
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), config=cfg,
)
return {
"status": "succeeded",
"code": "ok",
"counts": {
"added": stats.added,
"updated": stats.updated,
"deleted": stats.deleted,
"unchanged": stats.unchanged,
},
}
@vector_app.command("index-schema")
def index_schema_cmd(
config: Path = CONFIG_OPT,
json_output: bool = typer.Option(False, "--json", help="Emetti JSON puro su stdout."),
) -> None:
"""Embedda e sincronizza i record schema (tabelle e colonne) nel semantic store."""
import json
if json_output:
try:
cfg = _load_schema_config(config, suppress_legacy_warning=True)
except ConfigError:
typer.echo(json.dumps({"status": "failed", "code": "invalid_configuration"}, sort_keys=True, separators=(",", ":")))
raise typer.Exit(code=1) from None
else:
cfg = _load_config_or_exit(config)
try:
# Authorization and configuration are checked before touching any artifacts.
_vector_write_or_error(cfg)
_vector_cfg_or_error(cfg)
physical, annotations = _load_schema_artifacts(cfg)
payload = index_schema_data(cfg, physical=physical, annotations=annotations)
except SemanticIndexIncompatibleError:
if json_output:
typer.echo(json.dumps({"status": "failed", "code": "semantic_index_incompatible"}, sort_keys=True, separators=(",", ":")))
raise typer.Exit(code=1) from None
typer.secho("ERRORE: indice semantico incompatibile.", fg=typer.colors.RED, err=True)
raise typer.Exit(code=1) from None
except _MachineVectorError as error:
if json_output:
typer.echo(json.dumps({"status": "failed", "code": error.code}, sort_keys=True, separators=(",", ":")))
raise typer.Exit(code=1) from None
_render_index_schema_error(error, cfg)
except Exception:
if json_output:
typer.echo(json.dumps({"status": "failed", "code": "schema_index_failed"}, sort_keys=True, separators=(",", ":")))
raise typer.Exit(code=1) from None
logger.exception("Schema indexing failed")
typer.secho("ERRORE: impossibile indicizzare lo schema.", fg=typer.colors.RED, err=True)
raise typer.Exit(code=1) from None
if json_output:
typer.echo(json.dumps(payload, sort_keys=True, separators=(",", ":")))
return
_print_stats(payload["counts"])
def _render_index_schema_error(error: _MachineVectorError, cfg) -> None:
"""Render expected failures without changing the old human CLI messages."""
if error.code == "physical_schema_missing":
typer.secho(
f"ERRORE: {physical_path(cfg)} non trovato. Esegui prima `tht schema introspect`.",
fg=typer.colors.RED,
err=True,
)
elif error.code == "vector_configuration_missing":
require_vector_cfg(cfg)
elif error.code == "vector_write_not_allowed":
require_vector_write_allowed(cfg, "vector index-schema")
elif error.code == "schema_artifacts_invalid":
typer.secho("ERRORE: schema artifacts non validi.", fg=typer.colors.RED, err=True)
else:
typer.secho("ERRORE: impossibile indicizzare lo schema.", fg=typer.colors.RED, err=True)
raise typer.Exit(code=1) from None