From 32cd2165eb03f8aca3a90ad0c80a337aa0eab98d Mon Sep 17 00:00:00 2001 From: mptyl Date: Sat, 8 Aug 2026 18:14:03 +0200 Subject: [PATCH] fix: support qdrant-only vector maintenance --- harness/tests/test_qdrant_cli_commands.py | 164 ++++++++++++++++++++++ harness/tests/test_solved_search_cli.py | 22 +-- harness/tht/adapters/vector/pgvector.py | 29 +++- harness/tht/adapters/vector/qdrant.py | 15 ++ harness/tht/adapters/vector/thoth_http.py | 20 ++- harness/tht/cli/memory_cmd.py | 12 +- harness/tht/cli/vector_cmd.py | 4 +- harness/tht/ports/vector.py | 2 + harness/tht/vectorstore/rest_client.py | 21 ++- 9 files changed, 261 insertions(+), 28 deletions(-) create mode 100644 harness/tests/test_qdrant_cli_commands.py diff --git a/harness/tests/test_qdrant_cli_commands.py b/harness/tests/test_qdrant_cli_commands.py new file mode 100644 index 00000000..b7dd0ca1 --- /dev/null +++ b/harness/tests/test_qdrant_cli_commands.py @@ -0,0 +1,164 @@ +from __future__ import annotations + +import json +from datetime import UTC, datetime +from pathlib import Path +from types import SimpleNamespace + +from typer.testing import CliRunner + +from tht.cli import app +from tht.memory import MemoryRecord, save_registry + + +class _FakeEmbedder: + def embed_documents(self, documents): + return [[0.1] * 4 for _ in documents] + + +class _FakeVectorStore: + def __init__(self): + self.upserts = [] + self.deleted = [] + + def existing_hashes(self, collection, kinds): + return {} + + def upsert(self, collection, records): + self.upserts.append((collection, records)) + return len(records) + + def delete_kinds(self, collection, kinds): + self.deleted.append((collection, list(kinds))) + return 3 + + +def _qdrant_runtime_config(tmp_path: Path) -> Path: + cfg = tmp_path / "workspace.yaml" + cfg.write_text( + f""" +runtime_identity: + workspace_id: psd-clinical + workspace_revision: {'a' * 40} +dwh: + type: postgres_direct + connection: {{database: analytics, schema: mart, user: reader, password: secret}} +vectors: + type: qdrant + base_url: http://qdrant:6333 + collection: psd-clinical +roots: + sessions: {tmp_path / 'sessions'} + artifacts: {tmp_path / 'artifacts'} + indexes: {tmp_path / 'indexes'} +embeddings: + provider: ollama_internal + base_url: http://embedding:11434 + model: qwen3-embedding:0.6b + dim: 1024 +""" + ) + return cfg + + +def _write_schema_artifacts(tmp_path: Path) -> None: + (tmp_path / "artifacts" / "mschema").mkdir(parents=True, exist_ok=True) + (tmp_path / "artifacts" / "mschema" / "physical.yaml").write_text( + """ +database: analytics +schema: mart +introspected_at: 2026-01-01T00:00:00+00:00 +tables: + fact_patient: + comment: Patients + columns: + id: + type: bigint +""" + ) + (tmp_path / "artifacts" / "mschema" / "annotations.yaml").write_text( + "tables: {}\n" + ) + + +def _memory_record() -> MemoryRecord: + 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_vector_index_schema_accepts_qdrant_only_runtime_config(tmp_path, monkeypatch): + cfg = _qdrant_runtime_config(tmp_path) + _write_schema_artifacts(tmp_path) + store = _FakeVectorStore() + + monkeypatch.setattr("tht.adapters.factory.build_vector_store", lambda cfg, require_write: store) + monkeypatch.setattr("tht.cli.vector_cmd.make_embedder", lambda _: _FakeEmbedder()) + + res = CliRunner().invoke(app, ["vector", "index-schema", "-c", str(cfg)]) + + assert res.exit_code == 0, res.output + assert store.upserts + + +def test_memory_promote_accepts_qdrant_only_runtime_config(tmp_path, monkeypatch): + cfg = _qdrant_runtime_config(tmp_path) + store = _FakeVectorStore() + promoted = [_memory_record()] + snapshot = SimpleNamespace(manifest=SimpleNamespace(id="s1")) + + monkeypatch.setattr("tht.cli.memory_cmd.load_snapshot_or_exit", lambda cfg, session: snapshot) + monkeypatch.setattr("tht.memory.promote_snapshot", lambda *args, **kwargs: promoted) + monkeypatch.setattr("tht.memory.load_registry", lambda path: promoted) + monkeypatch.setattr("tht.adapters.factory.build_vector_store", lambda cfg, require_write: store) + monkeypatch.setattr("tht.cli.vector_cmd.make_embedder", lambda _: _FakeEmbedder()) + + res = CliRunner().invoke( + app, + ["memory", "promote", "--session", "s1", "--decision", "7", "--json", "-c", str(cfg)], + ) + + assert res.exit_code == 0, res.output + assert json.loads(res.stdout)["indexed"] is True + assert store.upserts + + +def test_memory_index_accepts_qdrant_only_runtime_config(tmp_path, monkeypatch): + cfg = _qdrant_runtime_config(tmp_path) + store = _FakeVectorStore() + records = [_memory_record()] + + save_registry(records, tmp_path / "artifacts" / "memory" / "registry.jsonl") + monkeypatch.setattr("tht.adapters.factory.build_vector_store", lambda cfg, require_write: store) + monkeypatch.setattr("tht.cli.vector_cmd.make_embedder", lambda _: _FakeEmbedder()) + + res = CliRunner().invoke(app, ["memory", "index", "-c", str(cfg)]) + + assert res.exit_code == 0, res.output + assert "OK:" in res.output + assert store.upserts + + +def test_memory_clear_accepts_qdrant_only_runtime_config(tmp_path, monkeypatch): + cfg = _qdrant_runtime_config(tmp_path) + store = _FakeVectorStore() + records = [_memory_record()] + registry = tmp_path / "artifacts" / "memory" / "registry.jsonl" + save_registry(records, registry) + + monkeypatch.setattr("tht.adapters.factory.build_vector_store", lambda cfg, require_write: store) + + res = CliRunner().invoke(app, ["memory", "clear", "--yes", "-c", str(cfg)]) + + assert res.exit_code == 0, res.output + assert not registry.exists() diff --git a/harness/tests/test_solved_search_cli.py b/harness/tests/test_solved_search_cli.py index fe6f0018..d1b27f91 100644 --- a/harness/tests/test_solved_search_cli.py +++ b/harness/tests/test_solved_search_cli.py @@ -7,7 +7,7 @@ puro (`[]` in modalita' --json) ed exit 0, cosi' il modello prosegue senza exemplar. Il finalize-hook gestisce gia' lo stesso scenario in modo analogo. """ import json -from datetime import datetime +from datetime import UTC, datetime from typer.testing import CliRunner @@ -24,7 +24,7 @@ def _cfg(tmp_path): "database: {database: d, schema: s, user: u, password: p, transport: direct}\n" f"paths: {{artifacts: {tmp_path/'a'}, indexes: {tmp_path/'i'}, sessions: {tmp_path/'se'}}}\n" "vector_db: {database: v, schema: vectors, user: u, password: p, transport: direct}\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" ) return cfg @@ -103,15 +103,15 @@ def test_solved_search_json_maps_hit_metadata(tmp_path, monkeypatch): def test_memory_search_excludes_legacy_table_records(tmp_path, monkeypatch): records = [ - MemoryRecord( - id="mem-0001", ts=datetime(2026, 1, 1), session_id="s1", - decision_seq=1, type="table_promoted", subject="fact_pazienti", - ), - MemoryRecord( - id="mem-0002", ts=datetime(2026, 1, 1), session_id="s1", - decision_seq=2, type="concept_clarified", subject="paziente attivo", - detail="flag_attivo = TRUE", - ), + MemoryRecord( + id="mem-0001", ts=datetime(2026, 1, 1, tzinfo=UTC), session_id="s1", + decision_seq=1, type="table_promoted", subject="fact_pazienti", + ), + MemoryRecord( + id="mem-0002", ts=datetime(2026, 1, 1, tzinfo=UTC), session_id="s1", + decision_seq=2, type="concept_clarified", subject="paziente attivo", + detail="flag_attivo = TRUE", + ), ] cfg = _cfg(tmp_path) save_registry(records, tmp_path / "a" / "memory" / "registry.jsonl") diff --git a/harness/tht/adapters/vector/pgvector.py b/harness/tht/adapters/vector/pgvector.py index 6843d975..b4ed96d8 100644 --- a/harness/tht/adapters/vector/pgvector.py +++ b/harness/tht/adapters/vector/pgvector.py @@ -3,8 +3,10 @@ import json import re +from psycopg2 import Error as PsycopgError from psycopg2 import sql from sqlalchemy import Engine +from sqlalchemy.exc import SQLAlchemyError from tht.config import DatabaseConfig from tht.db.connection import make_engine @@ -19,7 +21,6 @@ from tht.ports.vector import ( ) from tht.vectorstore.store import VectorHit, hit_from_metadata - COLLECTION_KINDS = { "schema_records": {"schema_table", "schema_column"}, "evidence": {"evidence"}, @@ -243,7 +244,7 @@ class PgVectorStore: return True, None, dimensions finally: raw.close() - except Exception as exc: + except (AttributeError, TypeError, ValueError, PsycopgError, SQLAlchemyError) as exc: return False, f"vector database probe failed: {type(exc).__name__}", set() def health(self) -> VectorHealth: @@ -471,6 +472,30 @@ class PgVectorStore: if raw is not None: raw.close() + def delete_kinds(self, collection: str, kinds: list[str]) -> int: + _collection(self._schema, collection) + _validate_collection_kinds(collection, kinds) + raw = None + try: + raw = self._require_writer().raw_connection() + with raw.cursor() as cursor: + cursor.execute( + sql.SQL("DELETE FROM {} WHERE kind = ANY(%s)").format( + _collection(self._schema, collection) + ), + (kinds,), + ) + count = cursor.rowcount + raw.commit() + return count + except Exception as exc: + if raw is not None: + raw.rollback() + raise VectorWriteUnavailable("Vector kind cleanup unavailable") from exc + finally: + if raw is not None: + raw.close() + def list_evidence_generations(self, collection: str, workspace_id: str) -> list[str]: if collection != "evidence": raise VectorStoreError("Only exact Evidence generations may be listed") diff --git a/harness/tht/adapters/vector/qdrant.py b/harness/tht/adapters/vector/qdrant.py index 5d65e2a1..f8804eeb 100644 --- a/harness/tht/adapters/vector/qdrant.py +++ b/harness/tht/adapters/vector/qdrant.py @@ -212,6 +212,21 @@ class QdrantVectorStore: ) return len(records) + def delete_kinds(self, collection: str, kinds: list[str]) -> int: + _collection("vectors", collection) + _validate_collection_kinds(collection, kinds) + must = [ + *self._workspace_filter(), + {"key": "record_kind", "match": {"any": sorted(kinds)}}, + ] + before = len(self._scroll(must)) + self._call( + "POST", + f"/collections/{self._collection}/points/delete?wait=true", + {"filter": {"must": must}}, + ) + return before + def delete_generation(self, collection: str, generation: str, workspace_id: str) -> int: if collection != "evidence" or _GENERATION.fullmatch(generation) is None: raise VectorStoreError("Only exact Evidence generations may be deleted") diff --git a/harness/tht/adapters/vector/thoth_http.py b/harness/tht/adapters/vector/thoth_http.py index 5a93b3c3..7f768e60 100644 --- a/harness/tht/adapters/vector/thoth_http.py +++ b/harness/tht/adapters/vector/thoth_http.py @@ -2,6 +2,11 @@ import re +from tht.adapters.vector.pgvector import ( + _collection, + _validate_collection_kinds, + _validate_known_kinds, +) from tht.ports.vector import ( VectorCapabilities, VectorHealth, @@ -14,11 +19,6 @@ from tht.ports.vector import ( ) from tht.vectorstore.rest_client import VectorRestClient, VectorRestError from tht.vectorstore.store import hit_from_metadata -from tht.adapters.vector.pgvector import ( - _collection, - _validate_collection_kinds, - _validate_known_kinds, -) def _merge(hits: list[VectorHit], limit: int) -> list[VectorHit]: @@ -85,7 +85,7 @@ class ThothHttpVectorStore: return None, None, [] try: return True, None, client.list_tables() - except Exception as exc: + except (RuntimeError, VectorRestError) as exc: return False, str(exc), [] def search( @@ -155,6 +155,14 @@ class ThothHttpVectorStore: except VectorRestError as exc: raise VectorStoreError(str(exc)) from exc + def delete_kinds(self, collection: str, kinds: list[str]) -> int: + _collection("vectors", collection) + _validate_collection_kinds(collection, kinds) + try: + return self._require_writer().delete_kinds(collection, kinds) + except VectorRestError as exc: + raise VectorStoreError(str(exc)) from exc + def delete_generation(self, collection: str, generation: str, workspace_id: str) -> int: if collection != "evidence" or re.fullmatch(r"gen:[0-9a-f]{32}", generation) is None: raise VectorStoreError("Only exact Evidence generations may be deleted") diff --git a/harness/tht/cli/memory_cmd.py b/harness/tht/cli/memory_cmd.py index dada0df2..be818e25 100644 --- a/harness/tht/cli/memory_cmd.py +++ b/harness/tht/cli/memory_cmd.py @@ -43,6 +43,12 @@ def _resync_memory(cfg): ) +def clear_memory_index(cfg): + from tht.adapters.factory import build_vector_store + + return build_vector_store(cfg, require_write=True).delete_kinds("memory", ["memory"]) + + @memory_app.command("promote") def promote_cmd( session: str = typer.Option(..., "--session"), @@ -192,7 +198,6 @@ def clear_cmd( config: Path = CONFIG_OPT, ) -> None: """Cancella TUTTA la review memory: registro canonico + indice pgvector (kind=memory).""" - from tht.cli.vector_cmd import make_embedder, open_store, require_direct_vector_cfg from tht.memory import load_registry cfg = _load_config_or_exit(config) @@ -209,10 +214,7 @@ def clear_cmd( typer.secho("Annullato.", fg=typer.colors.YELLOW) raise typer.Exit(code=1) - # Indice pgvector: rimuove i record kind=memory (sync con insieme vuoto). - require_direct_vector_cfg(cfg) - store = open_store(cfg, "memory") - store.sync([], make_embedder(cfg.embeddings), kinds={"memory"}) + clear_memory_index(cfg) # Registro canonico. registry.unlink() diff --git a/harness/tht/cli/vector_cmd.py b/harness/tht/cli/vector_cmd.py index ccde26ca..3072160a 100644 --- a/harness/tht/cli/vector_cmd.py +++ b/harness/tht/cli/vector_cmd.py @@ -26,8 +26,8 @@ def require_vector_cfg(cfg): missing = [] if cfg.embeddings is None: missing.append("embeddings") - if cfg.vector_db is None and not has_vector_write_rest(cfg): - missing.append("vector_db o vector_write_rest") + 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)}.", diff --git a/harness/tht/ports/vector.py b/harness/tht/ports/vector.py index e2bdef21..7aaf5687 100644 --- a/harness/tht/ports/vector.py +++ b/harness/tht/ports/vector.py @@ -80,6 +80,8 @@ class VectorStore(Protocol): def upsert(self, collection: str, records: list[VectorWriteRecord]) -> int: ... + def delete_kinds(self, collection: str, kinds: list[str]) -> int: ... + def delete_generation(self, collection: str, generation: str, workspace_id: str) -> int: ... def list_evidence_generations(self, collection: str, workspace_id: str) -> list[str]: ... diff --git a/harness/tht/vectorstore/rest_client.py b/harness/tht/vectorstore/rest_client.py index dd7e1ac9..19edce5f 100644 --- a/harness/tht/vectorstore/rest_client.py +++ b/harness/tht/vectorstore/rest_client.py @@ -5,9 +5,10 @@ Endpoint dedicato (es. https://host/vector/v1/), distinto dal DWH. La lettura us Errori in italiano e azionabili, stile `rest/client.py`. """ -import requests import re +import requests + from tht.config import RestConfig @@ -46,7 +47,7 @@ class VectorRestClient: try: body = resp.json() detail = body.get("message") or body.get("details") or resp.text - except Exception: + except ValueError: detail = resp.text return f"Vector REST rpc {fn} → HTTP {resp.status_code}: {detail}" @@ -149,6 +150,22 @@ class VectorRestClient: return int(payload.get("deleted", 0)) return 0 + def delete_kinds(self, table_name: str, kinds: list[str]) -> int: + try: + payload = self._call( + "delete_vector_kinds", + {"table_name": table_name, "kinds": kinds}, + ) + except VectorRestError as error: + if "HTTP 404" in str(error): + raise VectorRestError( + "delete_vector_kinds RPC is unavailable; deploy the cleanup migration" + ) from None + raise + if isinstance(payload, dict): + return int(payload.get("deleted", 0)) + return 0 + def list_evidence_generations(self, table_name: str, workspace_id: str) -> list[str]: if re.fullmatch(r"[a-z][a-z0-9_-]{0,63}", workspace_id) is None: raise ValueError("workspace namespace must be canonical")