fix: harden task3 qdrant compatibility boundaries
This commit is contained in:
@@ -1,4 +1,3 @@
|
||||
import { execFileSync } from "node:child_process";
|
||||
import { mkdtempSync, rmSync, writeFileSync } from "node:fs";
|
||||
import { tmpdir } from "node:os";
|
||||
import { join } from "node:path";
|
||||
@@ -109,7 +108,7 @@ test("keeps create_if_missing for non-registry runtime renders", () => {
|
||||
});
|
||||
|
||||
test.each(["session", "maintenance"])(
|
||||
"renders require_existing vectors for %s acquisition",
|
||||
"pure renderer emits require_existing vectors for %s acquisition",
|
||||
(mode) => {
|
||||
const rendered = parse(renderRuntimeConfig(workspaceV3, directBindings, paths, {
|
||||
workspaceId: "psd-clinical", workspaceRevision: "a".repeat(40),
|
||||
@@ -124,36 +123,12 @@ test.each(["session", "maintenance"])(
|
||||
},
|
||||
);
|
||||
|
||||
test("session and maintenance renders are byte-identical and bind identically in harness", () => {
|
||||
test("pure session and maintenance renders are byte-identical", () => {
|
||||
const context = { workspaceId: "psd-clinical", workspaceRevision: "a".repeat(40) };
|
||||
const passwordFile = evidenceSecretFile("dwh-password", "not-a-canary");
|
||||
const bindings = {
|
||||
...directBindings,
|
||||
dwh: { ...directBindings.dwh, values: {
|
||||
...directBindings.dwh.values,
|
||||
THT_WS_PSD_CLINICAL_DWH_PASSWORD_FILE: passwordFile,
|
||||
} },
|
||||
};
|
||||
const session = renderRuntimeConfig(workspaceV3, bindings, paths, context, {}, semanticRuntime);
|
||||
const maintenance = renderRuntimeConfig(workspaceV3, bindings, paths, context, {}, semanticRuntime);
|
||||
const session = renderRuntimeConfig(workspaceV3, directBindings, paths, context, {}, semanticRuntime);
|
||||
const maintenance = renderRuntimeConfig(workspaceV3, directBindings, paths, context, {}, semanticRuntime);
|
||||
|
||||
expect(maintenance).toBe(session);
|
||||
|
||||
const script = `
|
||||
import json, sys, tempfile
|
||||
from pathlib import Path
|
||||
from tht.config import load_config
|
||||
from tht.jobs.dwh_pipeline import config_dwh_binding
|
||||
with tempfile.TemporaryDirectory() as root:
|
||||
path = Path(root) / "runtime.yaml"
|
||||
path.write_text(sys.stdin.read())
|
||||
print(json.dumps(config_dwh_binding(load_config(path)), sort_keys=True))
|
||||
`;
|
||||
const python = join(process.cwd(), "../harness/.venv/bin/python");
|
||||
const run = (yaml: string) => execFileSync(python, ["-c", script], {
|
||||
cwd: join(process.cwd(), "../harness"), input: yaml, encoding: "utf8",
|
||||
}).trim();
|
||||
expect(run(session)).toBe(run(maintenance));
|
||||
});
|
||||
|
||||
test("renders schema-v3 DWH REST without exposing secret contents", () => {
|
||||
|
||||
@@ -0,0 +1,163 @@
|
||||
"""Shared fake Qdrant HTTP boundary for adapter and CLI tests."""
|
||||
|
||||
import json
|
||||
|
||||
from tht.ports.vector import VectorWriteRecord
|
||||
from tht.vectorstore.records import VectorRecord
|
||||
|
||||
|
||||
class FakeResponse:
|
||||
def __init__(self, status_code: int, payload=None, text: str | None = None):
|
||||
self.status_code = status_code
|
||||
self._payload = payload
|
||||
self.text = text if text is not None else (
|
||||
"" if payload is None else json.dumps(payload)
|
||||
)
|
||||
|
||||
@property
|
||||
def ok(self) -> bool:
|
||||
return 200 <= self.status_code < 300
|
||||
|
||||
def json(self):
|
||||
if isinstance(self._payload, Exception):
|
||||
raise self._payload
|
||||
return self._payload
|
||||
|
||||
|
||||
class FakeQdrantHttp:
|
||||
def __init__(self, *, dimension=1024, distance="Cosine"):
|
||||
self.dimension = dimension
|
||||
self.distance = distance
|
||||
self.collection = None
|
||||
self.payload_indexes: set[str] = set()
|
||||
self.payload_index_types: dict[str, str] = {}
|
||||
self.points: dict[str, dict] = {}
|
||||
self.calls: list[tuple[str, str, dict | None]] = []
|
||||
self.fail_request: Exception | None = None
|
||||
self.malformed_query = False
|
||||
self.malformed_scroll = False
|
||||
self.scroll_pages: list[dict] | None = None
|
||||
|
||||
def request(self, method, url, *, json=None, timeout=None):
|
||||
self.calls.append((method, url, json))
|
||||
if self.fail_request is not None:
|
||||
raise self.fail_request
|
||||
|
||||
path = url.split("://", 1)[-1].split("/", 1)[-1]
|
||||
path = "/" + path.split("?", 1)[0]
|
||||
|
||||
if method == "GET" and path == "/collections/workspace-semantic":
|
||||
if self.collection is None:
|
||||
return FakeResponse(404, {"status": "error"})
|
||||
return FakeResponse(200, {
|
||||
"result": {
|
||||
"config": {
|
||||
"params": {
|
||||
"vectors": {"size": self.dimension, "distance": self.distance}
|
||||
}
|
||||
},
|
||||
"payload_schema": {
|
||||
field: {"data_type": self.payload_index_types.get(field, "keyword")} for field in sorted(self.payload_indexes)
|
||||
},
|
||||
}
|
||||
})
|
||||
|
||||
if method == "PUT" and path == "/collections/workspace-semantic":
|
||||
self.collection = json
|
||||
self.dimension = json["vectors"]["size"]
|
||||
self.distance = json["vectors"]["distance"]
|
||||
return FakeResponse(200, {"status": "ok"})
|
||||
|
||||
if method == "PUT" and path == "/collections/workspace-semantic/index":
|
||||
self.payload_indexes.add(json["field_name"])
|
||||
return FakeResponse(200, {"status": "ok"})
|
||||
|
||||
if method == "PUT" and path == "/collections/workspace-semantic/points":
|
||||
if self.collection is None:
|
||||
return FakeResponse(404, {"status": "error"})
|
||||
for point in json["points"]:
|
||||
self.points[point["id"]] = point
|
||||
return FakeResponse(200, {"result": {"status": "acknowledged"}})
|
||||
|
||||
if method == "POST" and path == "/collections/workspace-semantic/points/query":
|
||||
if self.malformed_query:
|
||||
return FakeResponse(200, {"result": {"points": "nope"}})
|
||||
wanted = _match_points(self.points.values(), json["filter"])
|
||||
scored = sorted(
|
||||
(
|
||||
{
|
||||
"id": point["id"],
|
||||
"score": point.get("score", 0.9),
|
||||
"payload": point["payload"],
|
||||
}
|
||||
for point in wanted
|
||||
),
|
||||
key=lambda point: (-point["score"], point["payload"]["record_key"]),
|
||||
)
|
||||
return FakeResponse(200, {"result": {"points": scored[: json["limit"]]}})
|
||||
|
||||
if method == "POST" and path == "/collections/workspace-semantic/points/scroll":
|
||||
if self.malformed_scroll:
|
||||
return FakeResponse(200, {"result": {"points": "bad"}})
|
||||
if self.scroll_pages is not None:
|
||||
offset = json.get("offset")
|
||||
for page in self.scroll_pages:
|
||||
if page["offset"] == offset:
|
||||
filtered = _match_points(page["points"], json["filter"])
|
||||
return FakeResponse(200, {
|
||||
"result": {
|
||||
"points": filtered,
|
||||
"next_page_offset": page["next_page_offset"],
|
||||
}
|
||||
})
|
||||
raise AssertionError(("unexpected offset", offset, self.scroll_pages))
|
||||
wanted = sorted(
|
||||
_match_points(self.points.values(), json["filter"]),
|
||||
key=lambda point: point["payload"]["record_key"],
|
||||
)
|
||||
return FakeResponse(200, {"result": {"points": wanted, "next_page_offset": None}})
|
||||
|
||||
if method == "POST" and path == "/collections/workspace-semantic/points/delete":
|
||||
doomed = [point["id"] for point in _match_points(self.points.values(), json["filter"])]
|
||||
for point_id_value in doomed:
|
||||
self.points.pop(point_id_value, None)
|
||||
return FakeResponse(200, {"result": {"status": "acknowledged"}})
|
||||
|
||||
raise AssertionError((method, path, json))
|
||||
|
||||
|
||||
def _match_points(points, flt):
|
||||
matches = []
|
||||
must = flt["must"]
|
||||
for point in points:
|
||||
payload = point["payload"]
|
||||
if all(_match_clause(payload, clause) for clause in must):
|
||||
matches.append(point)
|
||||
return matches
|
||||
|
||||
|
||||
def _match_clause(payload, clause):
|
||||
key = clause["key"]
|
||||
match = clause["match"]
|
||||
if "value" in match:
|
||||
return payload.get(key) == match["value"]
|
||||
if "any" in match:
|
||||
return payload.get(key) in set(match["any"])
|
||||
raise AssertionError(clause)
|
||||
|
||||
|
||||
def _write_record(record_id: str, kind: str, *, metadata=None):
|
||||
return VectorWriteRecord(
|
||||
record=VectorRecord(
|
||||
id=record_id,
|
||||
kind=kind,
|
||||
ref=f"ref:{record_id}",
|
||||
title=f"title:{record_id}",
|
||||
content=f"content:{record_id}",
|
||||
metadata=metadata or {},
|
||||
),
|
||||
embedding=[0.1] * 1024,
|
||||
content_hash="sha256:" + "a" * 64,
|
||||
)
|
||||
|
||||
|
||||
@@ -67,6 +67,15 @@ def test_factory_selects_dwh_adapter(dwh_type, adapter_type):
|
||||
assert isinstance(build_dwh(_config(dwh_type=dwh_type)), adapter_type)
|
||||
|
||||
|
||||
def test_factory_rejects_require_existing_without_embedding_dimension():
|
||||
config = _config()
|
||||
config.vectors.collection_lifecycle = "require_existing"
|
||||
config.embeddings = None
|
||||
|
||||
with pytest.raises(ConfigError, match="explicit positive embedding dimension"):
|
||||
build_vector_store(config)
|
||||
|
||||
|
||||
def test_factory_selects_qdrant_for_schema_v3_runtime():
|
||||
config = _config(dwh_type="postgres_direct")
|
||||
|
||||
|
||||
@@ -218,6 +218,53 @@ def test_vector_index_schema_json_is_single_document(monkeypatch, tmp_path):
|
||||
assert payload["counts"]["added"] == 2
|
||||
|
||||
|
||||
def test_vector_index_schema_json_maps_require_existing_delete_race(monkeypatch, tmp_path):
|
||||
cfg = _qdrant_runtime_config(tmp_path)
|
||||
_write_schema_artifacts(tmp_path)
|
||||
monkeypatch.setattr("tht.cli.vector_cmd.make_embedder", lambda _: SimpleNamespace(embed_documents=lambda docs: [[0.1] * 1024 for _ in docs]))
|
||||
keyword_indexes = {
|
||||
"content_hash", "document_id", "kind", "record_key", "record_kind",
|
||||
"vector_generation", "workspace_id", "workspace_revision",
|
||||
}
|
||||
deleted = False
|
||||
|
||||
class Response:
|
||||
def __init__(self, status_code, payload=None):
|
||||
self.status_code = status_code
|
||||
self.ok = status_code < 400
|
||||
self._payload = payload
|
||||
self.text = "" if payload is None else "{}"
|
||||
|
||||
def json(self):
|
||||
return self._payload
|
||||
|
||||
def request(method, url, **kwargs):
|
||||
nonlocal deleted
|
||||
if method == "POST" and url.endswith("/points/scroll"):
|
||||
return Response(200, {"result": {"points": [], "next_page_offset": None}})
|
||||
if method == "GET" and url.endswith("/collections/psd-clinical"):
|
||||
if deleted:
|
||||
return Response(404, {"status": {"error": "missing"}})
|
||||
response = Response(200, {"result": {
|
||||
"config": {"params": {"vectors": {"size": 1024, "distance": "Cosine"}}},
|
||||
"payload_schema": {
|
||||
key: {"data_type": "keyword"} for key in keyword_indexes
|
||||
},
|
||||
}})
|
||||
deleted = True
|
||||
return response
|
||||
if method == "PUT" and "/points?wait=true" in url:
|
||||
return Response(404, {"status": {"error": "missing"}})
|
||||
raise AssertionError((method, url))
|
||||
|
||||
monkeypatch.setattr("requests.request", request)
|
||||
response = CliRunner().invoke(app, ["vector", "index-schema", "--json", "-c", str(cfg)])
|
||||
|
||||
assert response.exit_code == 1
|
||||
assert response.stdout == '{"code":"semantic_index_incompatible","status":"failed"}\n'
|
||||
assert response.stderr == ""
|
||||
|
||||
|
||||
def test_vector_index_schema_json_failure_is_safe(monkeypatch, tmp_path):
|
||||
import json
|
||||
|
||||
|
||||
@@ -1,167 +1,11 @@
|
||||
import json
|
||||
from uuid import NAMESPACE_URL, uuid5
|
||||
|
||||
import pytest
|
||||
import requests
|
||||
from qdrant_test_helpers import FakeQdrantHttp, _write_record
|
||||
|
||||
from tht.adapters.vector.qdrant import QdrantVectorStore, point_id
|
||||
from tht.ports.vector import VectorStoreError, VectorWriteRecord
|
||||
from tht.vectorstore.records import VectorRecord
|
||||
|
||||
|
||||
class FakeResponse:
|
||||
def __init__(self, status_code: int, payload=None, text: str | None = None):
|
||||
self.status_code = status_code
|
||||
self._payload = payload
|
||||
self.text = text if text is not None else (
|
||||
"" if payload is None else json.dumps(payload)
|
||||
)
|
||||
|
||||
@property
|
||||
def ok(self) -> bool:
|
||||
return 200 <= self.status_code < 300
|
||||
|
||||
def json(self):
|
||||
if isinstance(self._payload, Exception):
|
||||
raise self._payload
|
||||
return self._payload
|
||||
|
||||
|
||||
class FakeQdrantHttp:
|
||||
def __init__(self, *, dimension=1024, distance="Cosine"):
|
||||
self.dimension = dimension
|
||||
self.distance = distance
|
||||
self.collection = None
|
||||
self.payload_indexes: set[str] = set()
|
||||
self.payload_index_types: dict[str, str] = {}
|
||||
self.points: dict[str, dict] = {}
|
||||
self.calls: list[tuple[str, str, dict | None]] = []
|
||||
self.fail_request: Exception | None = None
|
||||
self.malformed_query = False
|
||||
self.malformed_scroll = False
|
||||
self.scroll_pages: list[dict] | None = None
|
||||
|
||||
def request(self, method, url, *, json=None, timeout=None):
|
||||
self.calls.append((method, url, json))
|
||||
if self.fail_request is not None:
|
||||
raise self.fail_request
|
||||
|
||||
path = url.split("://", 1)[-1].split("/", 1)[-1]
|
||||
path = "/" + path.split("?", 1)[0]
|
||||
|
||||
if method == "GET" and path == "/collections/workspace-semantic":
|
||||
if self.collection is None:
|
||||
return FakeResponse(404, {"status": "error"})
|
||||
return FakeResponse(200, {
|
||||
"result": {
|
||||
"config": {
|
||||
"params": {
|
||||
"vectors": {"size": self.dimension, "distance": self.distance}
|
||||
}
|
||||
},
|
||||
"payload_schema": {
|
||||
field: {"data_type": self.payload_index_types.get(field, "keyword")} for field in sorted(self.payload_indexes)
|
||||
},
|
||||
}
|
||||
})
|
||||
|
||||
if method == "PUT" and path == "/collections/workspace-semantic":
|
||||
self.collection = json
|
||||
self.dimension = json["vectors"]["size"]
|
||||
self.distance = json["vectors"]["distance"]
|
||||
return FakeResponse(200, {"status": "ok"})
|
||||
|
||||
if method == "PUT" and path == "/collections/workspace-semantic/index":
|
||||
self.payload_indexes.add(json["field_name"])
|
||||
return FakeResponse(200, {"status": "ok"})
|
||||
|
||||
if method == "PUT" and path == "/collections/workspace-semantic/points":
|
||||
if self.collection is None:
|
||||
return FakeResponse(404, {"status": "error"})
|
||||
for point in json["points"]:
|
||||
self.points[point["id"]] = point
|
||||
return FakeResponse(200, {"result": {"status": "acknowledged"}})
|
||||
|
||||
if method == "POST" and path == "/collections/workspace-semantic/points/query":
|
||||
if self.malformed_query:
|
||||
return FakeResponse(200, {"result": {"points": "nope"}})
|
||||
wanted = _match_points(self.points.values(), json["filter"])
|
||||
scored = sorted(
|
||||
(
|
||||
{
|
||||
"id": point["id"],
|
||||
"score": point.get("score", 0.9),
|
||||
"payload": point["payload"],
|
||||
}
|
||||
for point in wanted
|
||||
),
|
||||
key=lambda point: (-point["score"], point["payload"]["record_key"]),
|
||||
)
|
||||
return FakeResponse(200, {"result": {"points": scored[: json["limit"]]}})
|
||||
|
||||
if method == "POST" and path == "/collections/workspace-semantic/points/scroll":
|
||||
if self.malformed_scroll:
|
||||
return FakeResponse(200, {"result": {"points": "bad"}})
|
||||
if self.scroll_pages is not None:
|
||||
offset = json.get("offset")
|
||||
for page in self.scroll_pages:
|
||||
if page["offset"] == offset:
|
||||
filtered = _match_points(page["points"], json["filter"])
|
||||
return FakeResponse(200, {
|
||||
"result": {
|
||||
"points": filtered,
|
||||
"next_page_offset": page["next_page_offset"],
|
||||
}
|
||||
})
|
||||
raise AssertionError(("unexpected offset", offset, self.scroll_pages))
|
||||
wanted = sorted(
|
||||
_match_points(self.points.values(), json["filter"]),
|
||||
key=lambda point: point["payload"]["record_key"],
|
||||
)
|
||||
return FakeResponse(200, {"result": {"points": wanted, "next_page_offset": None}})
|
||||
|
||||
if method == "POST" and path == "/collections/workspace-semantic/points/delete":
|
||||
doomed = [point["id"] for point in _match_points(self.points.values(), json["filter"])]
|
||||
for point_id_value in doomed:
|
||||
self.points.pop(point_id_value, None)
|
||||
return FakeResponse(200, {"result": {"status": "acknowledged"}})
|
||||
|
||||
raise AssertionError((method, path, json))
|
||||
|
||||
|
||||
def _match_points(points, flt):
|
||||
matches = []
|
||||
must = flt["must"]
|
||||
for point in points:
|
||||
payload = point["payload"]
|
||||
if all(_match_clause(payload, clause) for clause in must):
|
||||
matches.append(point)
|
||||
return matches
|
||||
|
||||
|
||||
def _match_clause(payload, clause):
|
||||
key = clause["key"]
|
||||
match = clause["match"]
|
||||
if "value" in match:
|
||||
return payload.get(key) == match["value"]
|
||||
if "any" in match:
|
||||
return payload.get(key) in set(match["any"])
|
||||
raise AssertionError(clause)
|
||||
|
||||
|
||||
def _write_record(record_id: str, kind: str, *, metadata=None):
|
||||
return VectorWriteRecord(
|
||||
record=VectorRecord(
|
||||
id=record_id,
|
||||
kind=kind,
|
||||
ref=f"ref:{record_id}",
|
||||
title=f"title:{record_id}",
|
||||
content=f"content:{record_id}",
|
||||
metadata=metadata or {},
|
||||
),
|
||||
embedding=[0.1] * 1024,
|
||||
content_hash="sha256:" + "a" * 64,
|
||||
)
|
||||
from tht.ports.vector import VectorStoreError
|
||||
|
||||
|
||||
def _store(fake: FakeQdrantHttp, *, collection_lifecycle="create_if_missing") -> QdrantVectorStore:
|
||||
@@ -183,6 +27,20 @@ def test_point_id_is_deterministic_uuidv5():
|
||||
|
||||
|
||||
|
||||
def test_require_existing_requires_an_explicit_embedding_dimension():
|
||||
fake = FakeQdrantHttp()
|
||||
with pytest.raises(ValueError, match="expected dimension"):
|
||||
QdrantVectorStore(
|
||||
base_url="http://qdrant:6333",
|
||||
collection="workspace-semantic",
|
||||
workspace_id="demo",
|
||||
expected_dimension=None,
|
||||
collection_lifecycle="require_existing",
|
||||
request=fake.request,
|
||||
)
|
||||
assert fake.calls == []
|
||||
|
||||
|
||||
def test_require_existing_refuses_missing_collection_without_mutations():
|
||||
fake = FakeQdrantHttp()
|
||||
store = _store(fake, collection_lifecycle="require_existing")
|
||||
@@ -255,8 +113,11 @@ def test_require_existing_write_fails_after_collection_is_deleted_without_recrea
|
||||
collection_lifecycle="require_existing", request=request,
|
||||
)
|
||||
|
||||
with pytest.raises(VectorStoreError):
|
||||
from tht.ports.vector import SemanticIndexIncompatibleError
|
||||
|
||||
with pytest.raises(SemanticIndexIncompatibleError) as caught:
|
||||
store.upsert("memory", [_write_record("memory:1", "memory")])
|
||||
assert caught.value.code == "semantic_index_incompatible"
|
||||
assert not [call for call in fake.calls if call[0] == "PUT" and call[1].endswith("/collections/workspace-semantic")]
|
||||
|
||||
|
||||
|
||||
@@ -460,7 +460,7 @@ def test_registry_rendered_session_and_maintenance_configs_require_existing_qdra
|
||||
assert yaml.safe_load(maintenance_path.read_text())["vectors"]["collection_lifecycle"] == "require_existing"
|
||||
assert session_path.read_bytes() == maintenance_path.read_bytes()
|
||||
|
||||
from test_qdrant_vector_store import FakeQdrantHttp, _write_record
|
||||
from qdrant_test_helpers import FakeQdrantHttp, _write_record
|
||||
|
||||
for cfg in (session_cfg, maintenance_cfg):
|
||||
fake = FakeQdrantHttp(dimension=384) if collection_state == "incompatible" else FakeQdrantHttp()
|
||||
@@ -481,6 +481,40 @@ def test_registry_rendered_session_and_maintenance_configs_bind_equally(tmp_path
|
||||
)
|
||||
|
||||
|
||||
def test_loader_rejects_split_brain_qdrant_compatibility_resources(tmp_path):
|
||||
values = _qdrant_config_yaml(tmp_path, registry=True)
|
||||
values["vectors"]["collection"] = "workspace-a"
|
||||
values["resources"] = {
|
||||
"vector": {
|
||||
"engine": "qdrant",
|
||||
"base_url": "http://qdrant:6333/",
|
||||
"collection": "workspace-b",
|
||||
}
|
||||
}
|
||||
path = tmp_path / "split-brain.yaml"
|
||||
path.write_text(yaml.safe_dump(values))
|
||||
|
||||
with pytest.raises(ConfigError, match="vectors.*resources.vector|disagree"):
|
||||
load_config(path)
|
||||
|
||||
|
||||
def test_loader_rejects_lifecycle_policy_in_compatibility_resource(tmp_path):
|
||||
values = _qdrant_config_yaml(tmp_path, registry=True)
|
||||
values["resources"] = {
|
||||
"vector": {
|
||||
"engine": "qdrant",
|
||||
"base_url": "http://qdrant:6333",
|
||||
"collection": "workspace-semantic",
|
||||
"collection_lifecycle": "require_existing",
|
||||
}
|
||||
}
|
||||
path = tmp_path / "lifecycle.yaml"
|
||||
path.write_text(yaml.safe_dump(values))
|
||||
|
||||
with pytest.raises(ConfigError, match="collection_lifecycle"):
|
||||
load_config(path)
|
||||
|
||||
|
||||
def test_legacy_runtime_config_keeps_create_capable_qdrant_default(tmp_path):
|
||||
path = tmp_path / "legacy.yaml"
|
||||
path.write_text(yaml.safe_dump(_qdrant_config_yaml(tmp_path, registry=False)))
|
||||
|
||||
@@ -32,12 +32,20 @@ def build_vector_store(cfg: Config, *, require_write: bool = False) -> VectorSto
|
||||
|
||||
match resource.type:
|
||||
case "qdrant":
|
||||
expected_dimension = cfg.embeddings.dim if cfg.embeddings is not None else None
|
||||
if resource.collection_lifecycle == "require_existing" and (
|
||||
expected_dimension is None or expected_dimension <= 0
|
||||
):
|
||||
raise ConfigError(
|
||||
"require_existing Qdrant vectors require an explicit positive embedding dimension"
|
||||
)
|
||||
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,
|
||||
expected_dimension=expected_dimension,
|
||||
expected_distance="Cosine",
|
||||
collection_lifecycle=resource.collection_lifecycle,
|
||||
)
|
||||
case other: # pragma: no cover - Pydantic's discriminator rejects this first.
|
||||
|
||||
@@ -13,9 +13,12 @@ from tht.adapters.vector._shared import (
|
||||
validate_known_kinds,
|
||||
)
|
||||
from tht.ports.vector import (
|
||||
SemanticIndexIncompatibleError,
|
||||
VectorCapabilities,
|
||||
VectorHealth,
|
||||
VectorResponseError,
|
||||
VectorStoreError,
|
||||
VectorTransportError,
|
||||
VectorWriteRecord,
|
||||
require_positive_limit,
|
||||
)
|
||||
@@ -55,6 +58,7 @@ class QdrantVectorStore:
|
||||
workspace_id: str,
|
||||
workspace_revision: str | None = None,
|
||||
expected_dimension: int | None = None,
|
||||
expected_distance: str | None = "Cosine",
|
||||
collection_lifecycle: str = "create_if_missing",
|
||||
request: Callable[..., object] | None = None,
|
||||
connect_timeout: float = 2.0,
|
||||
@@ -65,8 +69,16 @@ class QdrantVectorStore:
|
||||
self._workspace_id = workspace_id
|
||||
self._workspace_revision = workspace_revision
|
||||
self._expected_dimension = expected_dimension
|
||||
self._expected_distance = expected_distance
|
||||
if collection_lifecycle not in ("create_if_missing", "require_existing"):
|
||||
raise ValueError("Unsupported Qdrant collection lifecycle")
|
||||
if collection_lifecycle == "require_existing" and (
|
||||
expected_dimension is None or expected_dimension <= 0 or expected_distance is None
|
||||
):
|
||||
raise ValueError(
|
||||
"require_existing requires an explicit positive expected dimension "
|
||||
"and expected distance"
|
||||
)
|
||||
self._collection_lifecycle = collection_lifecycle
|
||||
self._request = request or requests.request
|
||||
self._timeout = (connect_timeout, read_timeout)
|
||||
@@ -100,9 +112,17 @@ class QdrantVectorStore:
|
||||
|
||||
dimension = info["config"]["params"]["vectors"]["size"]
|
||||
dimensions = (dimension,)
|
||||
compatible = (
|
||||
dimension_compatible = (
|
||||
None if self._expected_dimension is None else dimensions == (self._expected_dimension,)
|
||||
)
|
||||
observed_distance = info["config"]["params"]["vectors"].get("distance")
|
||||
distance_compatible = (
|
||||
None if self._expected_distance is None else observed_distance == self._expected_distance
|
||||
)
|
||||
compatible = (
|
||||
None if dimension_compatible is None and distance_compatible is None
|
||||
else dimension_compatible is not False and distance_compatible is not False
|
||||
)
|
||||
return VectorHealth(
|
||||
ok=compatible is not False,
|
||||
read_configured=True,
|
||||
@@ -161,7 +181,7 @@ class QdrantVectorStore:
|
||||
)
|
||||
points = response.get("result", {}).get("points")
|
||||
if not isinstance(points, list):
|
||||
raise VectorStoreError("Qdrant returned malformed query response")
|
||||
raise VectorResponseError("Qdrant returned malformed query response")
|
||||
hits = [self._hit_from_point(point) for point in points]
|
||||
return sorted(hits, key=lambda hit: (-hit.similarity, hit.id))[:limit]
|
||||
|
||||
@@ -179,11 +199,11 @@ class QdrantVectorStore:
|
||||
for point in points:
|
||||
payload = point.get("payload")
|
||||
if not isinstance(payload, dict):
|
||||
raise VectorStoreError("Qdrant returned malformed scroll response")
|
||||
raise VectorResponseError("Qdrant returned malformed scroll response")
|
||||
record_key = payload.get("record_key")
|
||||
content_hash = payload.get("content_hash")
|
||||
if not isinstance(record_key, str) or not isinstance(content_hash, str):
|
||||
raise VectorStoreError("Qdrant returned malformed scroll response")
|
||||
raise VectorResponseError("Qdrant returned malformed scroll response")
|
||||
hashes[record_key] = content_hash
|
||||
return hashes
|
||||
|
||||
@@ -207,11 +227,18 @@ class QdrantVectorStore:
|
||||
),
|
||||
}
|
||||
)
|
||||
self._call(
|
||||
"PUT",
|
||||
f"/collections/{self._collection}/points?wait=true",
|
||||
{"points": points},
|
||||
)
|
||||
try:
|
||||
self._call(
|
||||
"PUT",
|
||||
f"/collections/{self._collection}/points?wait=true",
|
||||
{"points": points},
|
||||
)
|
||||
except VectorTransportError as exc:
|
||||
if self._collection_lifecycle == "require_existing" and exc.status_code == 404:
|
||||
raise SemanticIndexIncompatibleError(
|
||||
"Qdrant collection disappeared during semantic index write"
|
||||
) from exc
|
||||
raise
|
||||
return len(records)
|
||||
|
||||
def delete_kinds(self, collection: str, kinds: list[str]) -> int:
|
||||
@@ -308,10 +335,10 @@ class QdrantVectorStore:
|
||||
def _ensure_collection(self, *, strict: bool) -> dict | None:
|
||||
response = self._call("GET", f"/collections/{self._collection}", None, allow_missing=True)
|
||||
if response is None:
|
||||
if self._collection_lifecycle == "require_existing":
|
||||
raise SemanticIndexIncompatibleError("Qdrant collection is missing")
|
||||
if not strict:
|
||||
raise VectorStoreError("Qdrant collection is missing")
|
||||
if self._collection_lifecycle == "require_existing":
|
||||
raise VectorStoreError("Qdrant collection configuration mismatch (semantic_index_incompatible)")
|
||||
self._call(
|
||||
"PUT",
|
||||
f"/collections/{self._collection}",
|
||||
@@ -327,24 +354,31 @@ class QdrantVectorStore:
|
||||
result = response.get("result") if isinstance(response, dict) else None
|
||||
config = result.get("config", {}).get("params", {}).get("vectors") if isinstance(result, dict) else None
|
||||
if not isinstance(config, dict):
|
||||
raise VectorStoreError("Qdrant returned malformed collection response")
|
||||
raise VectorResponseError("Qdrant returned malformed collection response")
|
||||
size = config.get("size")
|
||||
distance = config.get("distance")
|
||||
if (
|
||||
self._expected_dimension is not None
|
||||
and (size != self._expected_dimension or distance != "Cosine")
|
||||
self._expected_dimension is not None and size != self._expected_dimension
|
||||
) or (
|
||||
self._expected_distance is not None and distance != self._expected_distance
|
||||
):
|
||||
raise VectorStoreError("Qdrant collection configuration mismatch (semantic_index_incompatible)")
|
||||
if self._collection_lifecycle == "require_existing":
|
||||
raise SemanticIndexIncompatibleError(
|
||||
"Qdrant collection dimension or distance is incompatible"
|
||||
)
|
||||
raise VectorStoreError("Qdrant collection configuration mismatch")
|
||||
payload_schema = result.get("payload_schema")
|
||||
if not isinstance(payload_schema, dict):
|
||||
raise VectorStoreError("Qdrant returned malformed collection response")
|
||||
raise VectorResponseError("Qdrant returned malformed collection response")
|
||||
for field_name in _KEYWORD_INDEXES:
|
||||
field = payload_schema.get(field_name)
|
||||
if not isinstance(field, dict) or field.get("data_type") != "keyword":
|
||||
if self._collection_lifecycle == "require_existing":
|
||||
raise SemanticIndexIncompatibleError(
|
||||
"Qdrant collection payload indexes are incompatible"
|
||||
)
|
||||
if not strict:
|
||||
raise VectorStoreError("Qdrant collection payload indexes mismatch")
|
||||
if self._collection_lifecycle == "require_existing":
|
||||
raise VectorStoreError("Qdrant collection configuration mismatch (semantic_index_incompatible)")
|
||||
self._call(
|
||||
"PUT",
|
||||
f"/collections/{self._collection}/index",
|
||||
@@ -370,13 +404,13 @@ class QdrantVectorStore:
|
||||
result = response.get("result", {})
|
||||
page = result.get("points")
|
||||
if not isinstance(page, list):
|
||||
raise VectorStoreError("Qdrant returned malformed scroll response")
|
||||
raise VectorResponseError("Qdrant returned malformed scroll response")
|
||||
points.extend(page)
|
||||
next_page_offset = result.get("next_page_offset")
|
||||
if next_page_offset is None:
|
||||
return points
|
||||
if next_page_offset in seen_offsets:
|
||||
raise VectorStoreError("Qdrant returned malformed scroll response")
|
||||
raise VectorResponseError("Qdrant returned malformed scroll response")
|
||||
seen_offsets.add(next_page_offset)
|
||||
offset = next_page_offset
|
||||
|
||||
@@ -384,7 +418,7 @@ class QdrantVectorStore:
|
||||
payload = point.get("payload")
|
||||
score = point.get("score")
|
||||
if not isinstance(payload, dict) or not isinstance(score, (int, float)):
|
||||
raise VectorStoreError("Qdrant returned malformed query response")
|
||||
raise VectorResponseError("Qdrant returned malformed query response")
|
||||
return hit_from_metadata(float(score), payload)
|
||||
|
||||
def _call(self, method: str, path: str, payload: dict | None, allow_missing: bool = False) -> dict | None:
|
||||
@@ -396,19 +430,19 @@ class QdrantVectorStore:
|
||||
timeout=self._timeout,
|
||||
)
|
||||
except requests.RequestException as exc:
|
||||
raise VectorStoreError(_sanitize_exception(exc)) from exc
|
||||
raise VectorTransportError(_sanitize_exception(exc)) from exc
|
||||
if response.status_code == 404 and allow_missing:
|
||||
return None
|
||||
if not response.ok:
|
||||
raise VectorStoreError(f"Qdrant request failed: HTTP {response.status_code}")
|
||||
raise VectorTransportError(f"Qdrant request failed: HTTP {response.status_code}", status_code=response.status_code)
|
||||
if response.status_code == 204 or not getattr(response, "text", ""):
|
||||
return {}
|
||||
try:
|
||||
data = response.json()
|
||||
except Exception as exc:
|
||||
raise VectorStoreError("Qdrant returned malformed JSON response") from exc
|
||||
raise VectorResponseError("Qdrant returned malformed JSON response") from exc
|
||||
if not isinstance(data, dict):
|
||||
raise VectorStoreError("Qdrant returned malformed JSON response")
|
||||
raise VectorResponseError("Qdrant returned malformed JSON response")
|
||||
return data
|
||||
|
||||
|
||||
|
||||
@@ -18,7 +18,7 @@ from tht.cli.schema_cmd import (
|
||||
physical_path,
|
||||
)
|
||||
from tht.config import Config, ConfigError
|
||||
from tht.ports.vector import VectorWriteRecord
|
||||
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)")
|
||||
@@ -265,6 +265,12 @@ def index_schema_cmd(
|
||||
_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=(",", ":")))
|
||||
|
||||
+33
-11
@@ -550,17 +550,8 @@ def load_config(path: Path) -> Config:
|
||||
_validate_internal_embedding_contract(expanded, path)
|
||||
_validate_internal_vector_contract(expanded, path)
|
||||
translated, used_legacy = translate_legacy_config(expanded)
|
||||
_validate_vector_resource_consistency(expanded, translated, path)
|
||||
vectors = translated.get("vectors")
|
||||
resource_vector = expanded.get("resources", {}).get("vector") if isinstance(
|
||||
expanded.get("resources"), dict
|
||||
) else None
|
||||
if (
|
||||
isinstance(vectors, dict)
|
||||
and vectors.get("type") == "qdrant"
|
||||
and isinstance(resource_vector, dict)
|
||||
and "collection_lifecycle" in resource_vector
|
||||
):
|
||||
vectors["collection_lifecycle"] = resource_vector["collection_lifecycle"]
|
||||
# runtime_identity is the registry marker. The lifecycle is an internal
|
||||
# runtime policy, never a descriptor-controlled option.
|
||||
if (
|
||||
@@ -673,6 +664,37 @@ def _validate_internal_embedding_contract(raw: dict[str, Any], path: Path) -> No
|
||||
)
|
||||
|
||||
|
||||
def _normalize_qdrant_base_url(value: str) -> str:
|
||||
parsed = urlparse(value)
|
||||
hostname = (parsed.hostname or "").lower()
|
||||
try:
|
||||
port = parsed.port
|
||||
except ValueError:
|
||||
port = None
|
||||
host = f"[{hostname}]" if ":" in hostname and not hostname.startswith("[") else hostname
|
||||
netloc = f"{host}:{port}" if port is not None else host
|
||||
return f"{parsed.scheme.lower()}://{netloc}{parsed.path.rstrip('/') or '/'}"
|
||||
|
||||
|
||||
def _validate_vector_resource_consistency(raw: dict[str, Any], translated: dict[str, Any], path: Path) -> None:
|
||||
"""Reject divergent top-level and compatibility Qdrant resource views."""
|
||||
resources = raw.get("resources")
|
||||
resource = resources.get("vector") if isinstance(resources, dict) else None
|
||||
vectors = translated.get("vectors")
|
||||
if not isinstance(resource, dict) or not isinstance(vectors, dict):
|
||||
return
|
||||
if vectors.get("type") != "qdrant" or resource.get("engine") != "qdrant":
|
||||
raise ConfigError(
|
||||
f"Configurazione non valida in {path}: vectors and resources.vector must both describe qdrant"
|
||||
)
|
||||
if _normalize_qdrant_base_url(vectors.get("base_url", "")) != _normalize_qdrant_base_url(
|
||||
resource.get("base_url", "")
|
||||
) or vectors.get("collection") != resource.get("collection"):
|
||||
raise ConfigError(
|
||||
f"Configurazione non valida in {path}: vectors and resources.vector disagree"
|
||||
)
|
||||
|
||||
|
||||
def _validate_internal_vector_contract(raw: dict[str, Any], path: Path) -> None:
|
||||
resources = raw.get("resources")
|
||||
if not isinstance(resources, dict):
|
||||
@@ -684,7 +706,7 @@ def _validate_internal_vector_contract(raw: dict[str, Any], path: Path) -> None:
|
||||
engine = vector.get("engine")
|
||||
base_url = vector.get("base_url")
|
||||
collection = vector.get("collection")
|
||||
allowed = {"engine", "base_url", "collection", "collection_lifecycle"}
|
||||
allowed = {"engine", "base_url", "collection"}
|
||||
unexpected = sorted(set(vector) - allowed)
|
||||
if unexpected:
|
||||
raise ConfigError(
|
||||
|
||||
@@ -45,6 +45,27 @@ class VectorStoreError(Exception):
|
||||
"""Base error exposed by vector adapters."""
|
||||
|
||||
|
||||
class SemanticIndexIncompatibleError(VectorStoreError):
|
||||
"""The configured semantic collection cannot safely serve this workspace."""
|
||||
|
||||
code = "semantic_index_incompatible"
|
||||
|
||||
def __init__(self, message: str = "Qdrant semantic index is incompatible"):
|
||||
super().__init__(f"{self.code}: {message}")
|
||||
|
||||
|
||||
class VectorTransportError(VectorStoreError):
|
||||
"""Qdrant could not be reached or returned an HTTP failure."""
|
||||
|
||||
def __init__(self, message: str, *, status_code: int | None = None):
|
||||
super().__init__(message)
|
||||
self.status_code = status_code
|
||||
|
||||
|
||||
class VectorResponseError(VectorStoreError):
|
||||
"""Qdrant returned a malformed response."""
|
||||
|
||||
|
||||
class VectorWriteUnavailable(VectorStoreError):
|
||||
"""Raised when a deployment has no vector writer credential."""
|
||||
|
||||
@@ -86,13 +107,16 @@ class VectorStore(Protocol):
|
||||
|
||||
|
||||
__all__ = [
|
||||
"SemanticIndexIncompatibleError",
|
||||
"VectorCapabilities",
|
||||
"VectorHealth",
|
||||
"VectorHit",
|
||||
"VectorReadUnavailable",
|
||||
"VectorRecord",
|
||||
"VectorResponseError",
|
||||
"VectorStore",
|
||||
"VectorStoreError",
|
||||
"VectorTransportError",
|
||||
"VectorWriteRecord",
|
||||
"VectorWriteUnavailable",
|
||||
"require_positive_limit",
|
||||
|
||||
Reference in New Issue
Block a user