From 5e39cfa347c40e2ed1b2f4430d76e3dee6fdc96c Mon Sep 17 00:00:00 2001 From: mptyl Date: Sat, 8 Aug 2026 18:03:57 +0200 Subject: [PATCH] feat: index semantic records in qdrant --- harness/tests/test_adapter_factory.py | 46 +++- harness/tests/test_config_resources.py | 28 +++ harness/tests/test_memory_save_one.py | 49 +++- harness/tests/test_qdrant_vector_store.py | 3 + harness/tests/test_search_pack.py | 19 +- harness/tests/test_semantic_kind_isolation.py | 222 ++++++++++++++++++ harness/tht/adapters/factory.py | 10 +- harness/tht/adapters/vector/qdrant.py | 4 + harness/tht/cli/memory_cmd.py | 24 +- harness/tht/cli/vector_cmd.py | 44 +++- harness/tht/config.py | 79 ++++++- harness/tht/config_compat.py | 8 + harness/tht/vectorstore/records.py | 9 +- 13 files changed, 516 insertions(+), 29 deletions(-) create mode 100644 harness/tests/test_semantic_kind_isolation.py diff --git a/harness/tests/test_adapter_factory.py b/harness/tests/test_adapter_factory.py index a29d38bd..13fb5c54 100644 --- a/harness/tests/test_adapter_factory.py +++ b/harness/tests/test_adapter_factory.py @@ -1,8 +1,8 @@ import pytest from tht.adapters.dwh import PostgresDwhAdapter, ThothRestDwhAdapter -from tht.adapters.vector import PgVectorStore, ThothHttpVectorStore from tht.adapters.factory import build_dwh, build_vector_store +from tht.adapters.vector import PgVectorStore, QdrantVectorStore, ThothHttpVectorStore from tht.config import Config, ConfigError @@ -125,6 +125,50 @@ def test_factory_builds_writer_only_direct_vector_when_write_is_required(): assert store.capabilities.upsert is True +def test_factory_selects_qdrant_for_schema_v3_runtime(): + config = Config.model_validate( + { + "dwh": { + "type": "postgres_direct", + "connection": { + "host": "db", + "database": "analytics", + "schema": "mart", + "user": "reader", + "password": "secret", + }, + }, + "database": { + "host": "db", + "database": "analytics", + "schema": "mart", + "user": "reader", + "password": "secret", + "transport": "direct", + }, + "vectors": { + "type": "qdrant", + "base_url": "http://qdrant:6333", + "collection": "psd-clinical", + }, + "embeddings": { + "provider": "ollama_internal", + "base_url": "http://embedding:11434", + "model": "qwen3-embedding:0.6b", + "dim": 1024, + }, + } + ) + config._workspace_id = "psd-clinical" + config._workspace_revision = "a" * 40 + + store = build_vector_store(config, require_write=True) + + assert isinstance(store, QdrantVectorStore) + assert store.capabilities.search is True + assert store.capabilities.upsert is True + + def test_factory_reuses_legacy_direct_connection_for_server_writes_only(): server = _config(vector_type="pgvector_direct", writer=False) server.vectors.connection = server.vectors.reader diff --git a/harness/tests/test_config_resources.py b/harness/tests/test_config_resources.py index dcf327bf..df051a68 100644 --- a/harness/tests/test_config_resources.py +++ b/harness/tests/test_config_resources.py @@ -6,6 +6,7 @@ from tht.config import ( ConfigError, PgvectorDirectConfig, PostgresDwhConfig, + QdrantConfig, ThothRestDwhConfig, ThothVectorHttpConfig, load_config, @@ -237,6 +238,33 @@ resources: assert cfg.embeddings.dim == 1024 +def test_accepts_internal_qdrant_resource_contract(tmp_path): + workspace = tmp_path / "workspace.yaml" + workspace.write_text( + """ +dwh: + type: postgres_direct + connection: {database: analytics, schema: mart, user: reader, password: secret} +resources: + vector: + engine: qdrant + base_url: http://qdrant:6333 + collection: psd-clinical + embeddings: + provider: ollama_internal + base_url: http://embedding:11434 + model: qwen3-embedding:0.6b + dimensions: 1024 +""" + ) + + cfg = load_config(workspace) + + assert isinstance(cfg.vectors, QdrantConfig) + assert cfg.vectors.base_url == "http://qdrant:6333" + assert cfg.vectors.collection == "psd-clinical" + + @pytest.mark.parametrize( ("snippet", "pattern"), [ diff --git a/harness/tests/test_memory_save_one.py b/harness/tests/test_memory_save_one.py index bed7cce7..c4fb3d8d 100644 --- a/harness/tests/test_memory_save_one.py +++ b/harness/tests/test_memory_save_one.py @@ -7,19 +7,28 @@ pgvector as a one-row upsert. This test pins the pure core of that behavior: - the writer.upsert_records is called once with a single row - writer.sync is NEVER called (that is the full-resync path) """ -from datetime import datetime +from datetime import UTC, datetime from unittest.mock import MagicMock +from tht.adapters.vector.qdrant import point_id from tht.memory import MemoryRecord, memory_vector_record_for_decision, save_one_memory +from tht.vectorstore.records import qdrant_payload def _record(seq: int = 7, **kw) -> MemoryRecord: - base = dict( - id="mem-0007", ts=datetime(2025, 1, 1), session_id="s1", decision_seq=seq, - type="concept_clarified", subject="paziente attivo", - detail="flag_attivo = TRUE", rationale="r", - question_context="dammi i pazienti", tables=[], concepts=["paziente attivo"], - ) + base = { + "id": "mem-0007", + "ts": datetime(2025, 1, 1, tzinfo=UTC), + "session_id": "s1", + "decision_seq": seq, + "type": "concept_clarified", + "subject": "paziente attivo", + "detail": "flag_attivo = TRUE", + "rationale": "r", + "question_context": "dammi i pazienti", + "tables": [], + "concepts": ["paziente attivo"], + } base.update(kw) return MemoryRecord(**base) @@ -85,3 +94,29 @@ def test_save_one_uses_writer_key_for_upsert(): save_one_memory(records, decision_seq=7, store=writer, embedder=embedder) # one upsert call, single row, table=memory assert writer.upsert.call_count == 1 + + +def test_save_one_preserves_semantic_point_identity_fields(): + records = [_record(seq=7)] + writer = MagicMock() + writer.existing_hashes.return_value = {} + writer.upsert.return_value = 1 + embedder = MagicMock() + embedder.embed_documents.return_value = [[0.0] * 4] + + save_one_memory(records, decision_seq=7, store=writer, embedder=embedder) + + row = writer.upsert.call_args.args[1][0] + payload = qdrant_payload( + row.record, + content_hash=row.content_hash, + workspace_id="psd-clinical", + workspace_revision="a" * 40, + ) + + assert point_id("psd-clinical", "memory", row.record.id) == point_id( + "psd-clinical", "memory", "memory:mem-0007" + ) + assert payload["kind"] == "memory" + assert payload["workspace_id"] == "psd-clinical" + assert payload["workspace_revision"] == "a" * 40 diff --git a/harness/tests/test_qdrant_vector_store.py b/harness/tests/test_qdrant_vector_store.py index 902b643d..ab5789c4 100644 --- a/harness/tests/test_qdrant_vector_store.py +++ b/harness/tests/test_qdrant_vector_store.py @@ -166,6 +166,7 @@ def _store(fake: FakeQdrantHttp) -> QdrantVectorStore: base_url="http://qdrant:6333", collection="workspace-semantic", workspace_id="demo", + workspace_revision="a" * 40, expected_dimension=1024, request=fake.request, ) @@ -195,6 +196,7 @@ def test_upsert_creates_collection_and_keyword_indexes_idempotently(): "record_kind", "vector_generation", "workspace_id", + "workspace_revision", } @@ -239,6 +241,7 @@ def test_upsert_serializes_qdrant_point_payloads(record, semantic_kind): assert point["id"] == point_id("demo", semantic_kind, record.record.id) assert point["vector"] == record.embedding assert point["payload"]["workspace_id"] == "demo" + assert point["payload"]["workspace_revision"] == "a" * 40 assert point["payload"]["kind"] == semantic_kind assert point["payload"]["record_kind"] == record.record.kind assert point["payload"]["record_key"] == record.record.id diff --git a/harness/tests/test_search_pack.py b/harness/tests/test_search_pack.py index 8df7e047..cdb517eb 100644 --- a/harness/tests/test_search_pack.py +++ b/harness/tests/test_search_pack.py @@ -1,5 +1,5 @@ import json -from datetime import datetime +from datetime import UTC, datetime from types import SimpleNamespace from typer.testing import CliRunner @@ -7,8 +7,8 @@ from typer.testing import CliRunner from tht.cli import app from tht.config import load_config from tht.jobs.dwh_pipeline import DwhPreprocessPipeline, config_dwh_binding -from tht.ports.vector import VectorReadUnavailable from tht.mschema.models import ColumnPhysical, PhysicalSchema, TablePhysical +from tht.ports.vector import VectorReadUnavailable from tht.vectorstore.embeddings import EmbeddingsError @@ -22,7 +22,11 @@ class _FakeEmbedder: class _FakeSearcher: + def __init__(self): + self.calls = [] + def search(self, vec, top_n, kinds=None): + self.calls.append({"top_n": top_n, "kinds": kinds}) if kinds == ["solved_question"]: return [SimpleNamespace( kind="memory", ref="s-1", id="m1", title="q solved", @@ -47,7 +51,7 @@ class _FakeSearcher: def _workspace(tmp_path, with_session=None): physical = PhysicalSchema( - database="d", schema="s", introspected_at=datetime(2026, 1, 1), + database="d", schema="s", introspected_at=datetime(2026, 1, 1, tzinfo=UTC), tables={"fact_ablazione": TablePhysical( comment="Ablazioni", columns={"cod_paz": ColumnPhysical(type="bigint")})}, ) @@ -55,7 +59,7 @@ def _workspace(tmp_path, with_session=None): cfg.write_text( "database: {database: d, schema: s, user: u, password: p, transport: direct}\n" "vector_db: {database: v, schema: public, user: u, password: p}\n" - "embeddings: {base_url: 'http://localhost:11434', model: nomic-embed-text, dim: 8}\n" + "embeddings: {base_url: 'http://localhost:11434', model: qwen3-embedding:0.6b, dim: 1024}\n" f"paths: {{artifacts: {tmp_path/'artifacts'}, indexes: {tmp_path/'i'}, " f"sessions: {tmp_path/'sessions'}}}\n" ) @@ -91,10 +95,15 @@ def _patch(monkeypatch, embedder, searcher): def test_pack_single_embed_and_sections(tmp_path, monkeypatch): cfg = _workspace(tmp_path) emb = _FakeEmbedder() - _patch(monkeypatch, emb, _FakeSearcher()) + searcher = _FakeSearcher() + _patch(monkeypatch, emb, searcher) res = CliRunner().invoke(app, ["search", "pack", "quanti pazienti", "-c", str(cfg)]) assert res.exit_code == 0, res.output assert emb.calls == 1 # UN solo embedding per le tre ricerche + assert [call["kinds"] for call in searcher.calls] == [ + ["schema_table", "schema_column"], + ["solved_question"], + ] assert "fact_ablazione" in res.output and "Ablazioni" in res.output # Evidence is fail-closed until an ACTIVE corpus exists; legacy vector rows # must not leak into a new search pack. diff --git a/harness/tests/test_semantic_kind_isolation.py b/harness/tests/test_semantic_kind_isolation.py new file mode 100644 index 00000000..5f000a46 --- /dev/null +++ b/harness/tests/test_semantic_kind_isolation.py @@ -0,0 +1,222 @@ +from __future__ import annotations + +import hashlib +from dataclasses import dataclass +from datetime import UTC, datetime + +from tht.adapters.vector.qdrant import point_id +from tht.cli.vector_cmd import sync_canonical_records +from tht.corpus.chunk import ChunkPolicy +from tht.corpus.models import CanonicalChunk +from tht.corpus.pipeline import CorpusPipeline +from tht.corpus.store import CorpusStore +from tht.memory import MemoryRecord, save_one_memory +from tht.mschema.models import ( + Annotations, + ColumnPhysical, + PhysicalSchema, + TablePhysical, +) +from tht.ports.vector import VectorCapabilities, VectorHealth +from tht.vectorstore.records import qdrant_payload, schema_records + + +def _sha(content: str) -> str: + return f"sha256:{hashlib.sha256(content.encode('utf-8')).hexdigest()}" + + +class _Embedder: + def embed_documents(self, documents): + return [[float(index + 1)] * 4 for index, _ in enumerate(documents)] + + +@dataclass +class _Point: + point_id: str + payload: dict + embedding: list[float] + + +class FakeVectorStore: + def __init__(self, workspace_id="psd-clinical", workspace_revision=None): + self.workspace_id = workspace_id + self.workspace_revision = workspace_revision or "a" * 40 + self.points: dict[str, _Point] = {} + self.search_calls: list[dict] = [] + + @property + def capabilities(self): + return VectorCapabilities( + search=True, + existing_hashes=True, + upsert=True, + metadata_filter=True, + delete_generation=True, + list_evidence_generations=True, + ) + + def health(self): + return VectorHealth(ok=True) + + def search(self, collections, embedding, *, limit, kinds=None, metadata_filter=None): + self.search_calls.append( + { + "collections": collections, + "embedding": embedding, + "limit": limit, + "kinds": kinds, + "metadata_filter": metadata_filter, + } + ) + return [] + + def existing_hashes(self, collection, kinds): + allowed = set(kinds) + return { + point.payload["record_key"]: point.payload["content_hash"] + for point in self.points.values() + if point.payload["record_kind"] in allowed + } + + def upsert(self, collection, records): + for row in records: + semantic_kind = qdrant_payload( + row.record, + content_hash=row.content_hash, + workspace_id=self.workspace_id, + workspace_revision=self.workspace_revision, + )["kind"] + payload = qdrant_payload( + row.record, + content_hash=row.content_hash, + workspace_id=self.workspace_id, + workspace_revision=self.workspace_revision, + ) + self.points[point_id(self.workspace_id, semantic_kind, row.record.id)] = _Point( + point_id=point_id(self.workspace_id, semantic_kind, row.record.id), + payload=payload, + embedding=row.embedding, + ) + return len(records) + + def delete_generation(self, collection, generation, workspace_id): + doomed = [ + key + for key, point in self.points.items() + if point.payload.get("record_kind") == "evidence" + and point.payload.get("vector_generation") == generation + and point.payload.get("workspace_id") == workspace_id + ] + for key in doomed: + self.points.pop(key) + return len(doomed) + + def list_evidence_generations(self, collection, workspace_id): + return sorted( + { + point.payload["vector_generation"] + for point in self.points.values() + if point.payload.get("record_kind") == "evidence" + and point.payload.get("workspace_id") == workspace_id + } + ) + + +def _schema_records(): + return schema_records( + PhysicalSchema( + database="analytics", + schema="mart", + introspected_at=datetime.now(UTC), + tables={ + "fact_patient": TablePhysical( + comment="Patients", + columns={"id": ColumnPhysical(type="bigint", comment="pk")}, + ) + }, + ), + Annotations(), + ) + + +def _memory_records(): + return [ + MemoryRecord( + id="mem-0001", + ts=datetime(2026, 1, 1, tzinfo=UTC), + session_id="s1", + decision_seq=7, + type="concept_clarified", + subject="paziente attivo", + detail="flag_attivo = true", + rationale="r", + question_context="dammi i pazienti attivi", + tables=[], + concepts=["paziente attivo"], + ) + ] + + +def test_schema_and_memory_use_expected_semantic_kinds_and_shared_identity(): + store = FakeVectorStore() + embedder = _Embedder() + + schema_stats = sync_canonical_records( + "schema_records", + _schema_records(), + store=store, + embedder=embedder, + ) + memory_count = save_one_memory(_memory_records(), 7, store=store, embedder=embedder) + + assert schema_stats.added == 2 + assert memory_count == 1 + payloads = {point.payload["record_kind"]: point.payload for point in store.points.values()} + assert payloads["schema_table"]["kind"] == "schema" + assert payloads["schema_column"]["kind"] == "schema" + assert payloads["memory"]["kind"] == "memory" + assert {payload["workspace_id"] for payload in payloads.values()} == {"psd-clinical"} + assert {payload["workspace_revision"] for payload in payloads.values()} == {"a" * 40} + + +def test_corpus_vector_records_keep_exact_generation_and_retry_is_idempotent(tmp_path): + store = FakeVectorStore() + pipeline = CorpusPipeline( + store=CorpusStore(tmp_path / "corpus"), + sources=[], + embedder=None, + vector_store=store, + embedding_model="qwen3-embedding:0.6b", + embedding_dimensions=1024, + chunk_policy=ChunkPolicy(version="chunk-v1", max_chars=4000), + pipeline_version="evidence-v1", + workspace_id="psd-clinical", + ) + chunk = CanonicalChunk( + chunk_id="chunk:1", + document_id="doc:patient-guide", + ordinal=0, + content="Patient evidence", + content_hash=_sha("Patient evidence"), + source_uri="file:///tmp/patient-guide.md", + pipeline_version="evidence-v1", + ) + row = pipeline._vector_record( + chunk, + [0.1, 0.2, 0.3, 0.4], + "gen:" + "1" * 32, + "psd-clinical", + ) + + assert row.record.kind == "evidence" + assert row.record.metadata["vector_generation"] == "gen:" + "1" * 32 + + store.upsert("evidence", [row]) + store.upsert("evidence", [row]) + + assert len(store.points) == 1 + point = next(iter(store.points.values())) + assert point.payload["kind"] == "evidence" + assert point.payload["vector_generation"] == "gen:" + "1" * 32 + assert point.payload["workspace_id"] == "psd-clinical" + assert point.payload["workspace_revision"] == "a" * 40 diff --git a/harness/tht/adapters/factory.py b/harness/tht/adapters/factory.py index 60f9b59e..85e906f0 100644 --- a/harness/tht/adapters/factory.py +++ b/harness/tht/adapters/factory.py @@ -3,7 +3,7 @@ from tht.adapters.dwh import PostgresDwhAdapter, ThothRestDwhAdapter from tht.adapters.evidence import FilesystemEvidenceSource, HttpManifestEvidenceSource from tht.adapters.evidence.s3 import S3EvidenceSource -from tht.adapters.vector import PgVectorStore, ThothHttpVectorStore +from tht.adapters.vector import PgVectorStore, QdrantVectorStore, ThothHttpVectorStore from tht.config import Config, ConfigError from tht.db.connection import make_engine from tht.ports.dwh import DwhAdapter @@ -56,6 +56,14 @@ def build_vector_store(cfg: Config, *, require_write: bool = False) -> VectorSto VectorRestClient(resource.writer) if resource.writer is not None else None, expected_dimension=cfg.embeddings.dim if cfg.embeddings is not None else None, ) + case "qdrant": + return QdrantVectorStore( + base_url=resource.base_url, + collection=resource.collection, + workspace_id=cfg._workspace_id, + workspace_revision=cfg._workspace_revision, + expected_dimension=cfg.embeddings.dim if cfg.embeddings is not None else None, + ) case other: # pragma: no cover - Pydantic's discriminator rejects this first. raise ConfigError(f"Adapter vector non supportato: {other}") diff --git a/harness/tht/adapters/vector/qdrant.py b/harness/tht/adapters/vector/qdrant.py index dc691a98..5d65e2a1 100644 --- a/harness/tht/adapters/vector/qdrant.py +++ b/harness/tht/adapters/vector/qdrant.py @@ -32,6 +32,7 @@ _KEYWORD_INDEXES = ( "record_kind", "vector_generation", "workspace_id", + "workspace_revision", ) @@ -52,6 +53,7 @@ class QdrantVectorStore: base_url: str, collection: str, workspace_id: str, + workspace_revision: str | None = None, expected_dimension: int | None = None, request: Callable[..., object] | None = None, connect_timeout: float = 2.0, @@ -60,6 +62,7 @@ class QdrantVectorStore: self._base_url = base_url.rstrip("/") self._collection = collection self._workspace_id = workspace_id + self._workspace_revision = workspace_revision self._expected_dimension = expected_dimension self._request = request or requests.request self._timeout = (connect_timeout, read_timeout) @@ -198,6 +201,7 @@ class QdrantVectorStore: write_record.record, content_hash=write_record.content_hash, workspace_id=self._workspace_id, + workspace_revision=self._workspace_revision, ), } ) diff --git a/harness/tht/cli/memory_cmd.py b/harness/tht/cli/memory_cmd.py index 0f1ea79e..dada0df2 100644 --- a/harness/tht/cli/memory_cmd.py +++ b/harness/tht/cli/memory_cmd.py @@ -10,17 +10,18 @@ from pathlib import Path import typer from sqlalchemy.exc import OperationalError, ProgrammingError -from tht.cli.config_cmd import CONFIG_OPT -from tht.cli.schema_cmd import _load_config_or_exit 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 from tht.cli.session_cmd import load_snapshot_or_exit from tht.cli.vector_cmd import require_vector_cfg memory_app = typer.Typer(help="Review memory (registro canonico + indice pgvector)") +DECISION_OPT = typer.Option(None, "--decision", help="Seq da promuovere (ripetibile).") def registry_path(cfg) -> Path: @@ -29,18 +30,23 @@ def registry_path(cfg) -> Path: def _resync_memory(cfg): """Risincronizza l'indice pgvector col registro corrente (incrementale).""" - from tht.cli.vector_cmd import make_embedder, open_store + from tht.adapters.factory import build_vector_store + from tht.cli.vector_cmd import make_embedder, sync_canonical_records from tht.memory import load_registry, memory_vector_records records = memory_vector_records(load_registry(registry_path(cfg))) - store = open_store(cfg, "memory") - return store.sync(records, make_embedder(cfg.embeddings), kinds={"memory"}) + return sync_canonical_records( + "memory", + records, + store=build_vector_store(cfg, require_write=True), + embedder=make_embedder(cfg.embeddings), + ) @memory_app.command("promote") def promote_cmd( session: str = typer.Option(..., "--session"), - decision: list[int] = typer.Option(None, "--decision", help="Seq da promuovere (ripetibile)."), + decision: list[int] = DECISION_OPT, preview: bool = typer.Option(False, "--preview", help="Mostra i candidati in JSON, non scrive."), json_out: bool = typer.Option(False, "--json", help="Output JSON (per Pi)."), config: Path = CONFIG_OPT, @@ -55,7 +61,9 @@ def promote_cmd( if preview: from tht.memory import ( - MAX_PROMOTION_CANDIDATES, preview_promotions_snapshot, reusable_promotions_snapshot, + MAX_PROMOTION_CANDIDATES, + preview_promotions_snapshot, + reusable_promotions_snapshot, ) cand = preview_promotions_snapshot(snapshot, registry_path(cfg)) extra = len(reusable_promotions_snapshot(snapshot, registry_path(cfg))) - len(cand) @@ -304,8 +312,8 @@ def update_cmd( """Modifica i campi di merito di una memoria (provenienza immutabile).""" from typing import get_args - from tht.memory import MemoryNotFound, update_record from tht.decisions import DecisionType + from tht.memory import MemoryNotFound, update_record cfg = _load_config_or_exit(config) diff --git a/harness/tht/cli/vector_cmd.py b/harness/tht/cli/vector_cmd.py index f54546bc..ccde26ca 100644 --- a/harness/tht/cli/vector_cmd.py +++ b/harness/tht/cli/vector_cmd.py @@ -2,9 +2,15 @@ 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._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)") @@ -69,6 +75,31 @@ def open_searcher(cfg): 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, " @@ -88,7 +119,6 @@ def init_cmd( from sqlalchemy.exc import OperationalError from tht.vectorstore.embeddings import EmbeddingsError - from tht.vectorstore.reader import ALL_TABLES cfg = _load_config_or_exit(config) @@ -131,8 +161,12 @@ def index_schema_cmd(config: Path = CONFIG_OPT) -> None: physical = PhysicalSchema.from_yaml(phys_file) annotations = Annotations.from_yaml(annotations_path(cfg)) records = schema_records(physical, annotations) - store = open_store(cfg, "schema_records") - stats = store.sync( - records, make_embedder(cfg.embeddings), kinds={"schema_table", "schema_column"} + 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) diff --git a/harness/tht/config.py b/harness/tht/config.py index aedcef0a..331d22b0 100644 --- a/harness/tht/config.py +++ b/harness/tht/config.py @@ -152,8 +152,14 @@ class ThothVectorHttpConfig(BaseModel): direct: DatabaseConfig | None = None +class QdrantConfig(BaseModel): + type: Literal["qdrant"] + base_url: str + collection: str = Field(min_length=1) + + VectorResourceConfig = Annotated[ - PgvectorDirectConfig | ThothVectorHttpConfig, + PgvectorDirectConfig | ThothVectorHttpConfig | QdrantConfig, Field(discriminator="type"), ] @@ -397,6 +403,7 @@ def load_config(path: Path) -> Config: raise ConfigError(f"Configurazione non valida (atteso un mapping YAML): {path}") expanded = _resolve_secret_files(_expand_env(raw)) _validate_internal_embedding_contract(expanded, path) + _validate_internal_vector_contract(expanded, path) translated, used_legacy = translate_legacy_config(expanded) _populate_legacy_views(translated) try: @@ -455,6 +462,7 @@ def load_config(path: Path) -> Config: else path.resolve().as_posix() ) _validate_active_embeddings_config(cfg.embeddings, path) + _validate_active_vector_config(cfg.vectors, path) return cfg @@ -500,6 +508,42 @@ def _validate_internal_embedding_contract(raw: dict[str, Any], path: Path) -> No ) +def _validate_internal_vector_contract(raw: dict[str, Any], path: Path) -> None: + resources = raw.get("resources") + if not isinstance(resources, dict): + return + vector = resources.get("vector") + if not isinstance(vector, dict): + return + + engine = vector.get("engine") + base_url = vector.get("base_url") + collection = vector.get("collection") + allowed = {"engine", "base_url", "collection"} + unexpected = sorted(set(vector) - allowed) + if unexpected: + raise ConfigError( + f"Configurazione non valida in {path}:\n" + f"resources.vector non supporta: {', '.join(unexpected)}" + ) + if engine != "qdrant": + raise ConfigError( + f"Configurazione non valida in {path}:\n" + "resources.vector.engine deve essere 'qdrant'" + ) + if not isinstance(collection, str) or not collection: + raise ConfigError( + f"Configurazione non valida in {path}:\n" + "resources.vector.collection deve essere valorizzato" + ) + if not _is_allowed_internal_qdrant_url(base_url): + raise ConfigError( + f"Configurazione non valida in {path}:\n" + "resources.vector.base_url deve usare http://qdrant:6333 " + "oppure un endpoint loopback di sviluppo su porta 6333" + ) + + def _validate_active_embeddings_config( embeddings: "EmbeddingsConfig | None", path: Path, @@ -529,6 +573,20 @@ def _validate_active_embeddings_config( ) +def _validate_active_vector_config( + vectors: "VectorResourceConfig | None", + path: Path, +) -> None: + if vectors is None or vectors.type != "qdrant": + return + if not _is_allowed_internal_qdrant_url(vectors.base_url): + raise ConfigError( + f"Configurazione non valida in {path}:\n" + "vectors.base_url deve usare http://qdrant:6333 " + "oppure un endpoint loopback di sviluppo su porta 6333" + ) + + def _is_allowed_internal_embedding_url(value: Any) -> bool: if not isinstance(value, str): return False @@ -548,6 +606,25 @@ def _is_allowed_internal_embedding_url(value: Any) -> bool: return host.is_loopback +def _is_allowed_internal_qdrant_url(value: Any) -> bool: + if not isinstance(value, str): + return False + parsed = urlparse(value) + if parsed.scheme != "http" or not parsed.hostname or parsed.port != 6333: + return False + if parsed.params or parsed.query or parsed.fragment: + return False + if parsed.path not in ("", "/"): + return False + if parsed.hostname == "qdrant": + return True + try: + host = ip_address(parsed.hostname) + except ValueError: + return parsed.hostname == "localhost" + return host.is_loopback + + def _populate_legacy_views(raw: dict[str, Any]) -> None: """Populate old Config attributes for command compatibility during migration.""" dwh = raw.get("dwh") diff --git a/harness/tht/config_compat.py b/harness/tht/config_compat.py index 04dbede7..639566fe 100644 --- a/harness/tht/config_compat.py +++ b/harness/tht/config_compat.py @@ -33,6 +33,14 @@ def translate_legacy_config(raw: dict[str, Any]) -> tuple[dict[str, Any], bool]: translated["embeddings"] = embedding if "dimensions" in translated["embeddings"] and "dim" not in translated["embeddings"]: translated["embeddings"]["dim"] = translated["embeddings"].pop("dimensions") + if isinstance(resources, dict) and "vector" in resources and "vectors" not in translated: + vector = _as_mapping(resources.get("vector")) + if isinstance(vector, dict): + translated["vectors"] = { + "type": "qdrant", + "base_url": vector.get("base_url"), + "collection": vector.get("collection"), + } legacy = any(key in raw for key in _LEGACY_RESOURCE_KEYS) if not legacy: return translated, False diff --git a/harness/tht/vectorstore/records.py b/harness/tht/vectorstore/records.py index 792adc6c..d522374b 100644 --- a/harness/tht/vectorstore/records.py +++ b/harness/tht/vectorstore/records.py @@ -27,11 +27,18 @@ def qdrant_semantic_kind(kind: str) -> str: raise ValueError(f"Unsupported vector kind: {kind}") -def qdrant_payload(record: VectorRecord, *, content_hash: str, workspace_id: str) -> dict: +def qdrant_payload( + record: VectorRecord, + *, + content_hash: str, + workspace_id: str, + workspace_revision: str | None = None, +) -> dict: semantic_kind = qdrant_semantic_kind(record.kind) return { **record.metadata, "workspace_id": workspace_id, + **({"workspace_revision": workspace_revision} if workspace_revision else {}), "kind": semantic_kind, "record_kind": record.kind, "record_key": record.id,