fix: support qdrant-only vector maintenance

This commit is contained in:
2026-08-08 18:14:03 +02:00
parent 5e39cfa347
commit 32cd2165eb
9 changed files with 261 additions and 28 deletions
+164
View File
@@ -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()
+11 -11
View File
@@ -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")
+27 -2
View File
@@ -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")
+15
View File
@@ -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")
+14 -6
View File
@@ -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")
+7 -5
View File
@@ -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()
+2 -2
View File
@@ -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)}.",
+2
View File
@@ -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]: ...
+19 -2
View File
@@ -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")