fix: close task3 qdrant and config validation gaps
This commit is contained in:
@@ -107,21 +107,18 @@ test("keeps create_if_missing for non-registry runtime renders", () => {
|
||||
});
|
||||
});
|
||||
|
||||
test.each(["session", "maintenance"])(
|
||||
"pure renderer emits require_existing vectors for %s acquisition",
|
||||
(mode) => {
|
||||
const rendered = parse(renderRuntimeConfig(workspaceV3, directBindings, paths, {
|
||||
workspaceId: "psd-clinical", workspaceRevision: "a".repeat(40),
|
||||
}, {}, semanticRuntime));
|
||||
test("pure renderer emits require_existing vectors for registry runtime", () => {
|
||||
const rendered = parse(renderRuntimeConfig(workspaceV3, directBindings, paths, {
|
||||
workspaceId: "psd-clinical", workspaceRevision: "a".repeat(40),
|
||||
}, {}, semanticRuntime));
|
||||
|
||||
expect(rendered.vectors).toEqual({
|
||||
type: "qdrant",
|
||||
base_url: "http://qdrant:6333",
|
||||
collection: "psd-clinical",
|
||||
collection_lifecycle: "require_existing",
|
||||
});
|
||||
},
|
||||
);
|
||||
expect(rendered.vectors).toEqual({
|
||||
type: "qdrant",
|
||||
base_url: "http://qdrant:6333",
|
||||
collection: "psd-clinical",
|
||||
collection_lifecycle: "require_existing",
|
||||
});
|
||||
});
|
||||
|
||||
test("pure session and maintenance renders are byte-identical", () => {
|
||||
const context = { workspaceId: "psd-clinical", workspaceRevision: "a".repeat(40) };
|
||||
|
||||
@@ -44,5 +44,6 @@ testpaths = ["tests"]
|
||||
markers = [
|
||||
"l0: testcontainers tests (need Docker, run locally)",
|
||||
"l2: end-to-end tests requiring real GLM 5.2 + remote DB (skipped when .env incomplete)",
|
||||
"integration: cross-runtime integration tests requiring repository-local toolchains",
|
||||
]
|
||||
addopts = "-m 'not l2'" # L0 runs by default (Docker present); L2 opt-in
|
||||
|
||||
@@ -16,6 +16,17 @@ class _FakeEmbedder:
|
||||
return [[0.1] * 4 for _ in documents]
|
||||
|
||||
|
||||
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
|
||||
|
||||
|
||||
class _FakeVectorStore:
|
||||
def __init__(self):
|
||||
self.upserts = []
|
||||
@@ -218,6 +229,55 @@ def test_vector_index_schema_json_is_single_document(monkeypatch, tmp_path):
|
||||
assert payload["counts"]["added"] == 2
|
||||
|
||||
|
||||
def test_vector_index_schema_json_maps_initial_missing_collection(monkeypatch, tmp_path):
|
||||
cfg = _qdrant_runtime_config(tmp_path)
|
||||
_write_schema_artifacts(tmp_path)
|
||||
calls = []
|
||||
|
||||
def request(method, url, **kwargs):
|
||||
calls.append((method, url))
|
||||
if method == "GET" and url.endswith("/collections/psd-clinical"):
|
||||
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 == ""
|
||||
assert not [call for call in calls if call[0] == "PUT"]
|
||||
assert not [call for call in calls if call[1].endswith("/points/scroll")]
|
||||
|
||||
|
||||
def test_vector_index_schema_json_rejects_incompatible_empty_collection(monkeypatch, tmp_path):
|
||||
cfg = _qdrant_runtime_config(tmp_path)
|
||||
_write_schema_artifacts(tmp_path)
|
||||
keyword_indexes = {
|
||||
"content_hash", "document_id", "kind", "record_key", "record_kind",
|
||||
"vector_generation", "workspace_id", "workspace_revision",
|
||||
}
|
||||
calls = []
|
||||
|
||||
def request(method, url, **kwargs):
|
||||
calls.append((method, url))
|
||||
if method == "GET" and url.endswith("/collections/psd-clinical"):
|
||||
return _Response(200, {"result": {
|
||||
"config": {"params": {"vectors": {"size": 384, "distance": "Cosine"}}},
|
||||
"payload_schema": {key: {"data_type": "keyword"} for key in keyword_indexes},
|
||||
}})
|
||||
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 == ""
|
||||
assert not [call for call in calls if call[0] == "PUT"]
|
||||
assert not [call for call in calls if call[1].endswith("/points/scroll")]
|
||||
|
||||
|
||||
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)
|
||||
|
||||
@@ -2,7 +2,7 @@ from uuid import NAMESPACE_URL, uuid5
|
||||
|
||||
import pytest
|
||||
import requests
|
||||
from qdrant_test_helpers import FakeQdrantHttp, _write_record
|
||||
from qdrant_test_helpers import FakeQdrantHttp, FakeResponse, _write_record
|
||||
|
||||
from tht.adapters.vector.qdrant import QdrantVectorStore, point_id
|
||||
from tht.ports.vector import VectorStoreError
|
||||
@@ -51,6 +51,59 @@ def test_require_existing_refuses_missing_collection_without_mutations():
|
||||
assert [call for call in fake.calls if call[0] == "PUT"] == []
|
||||
|
||||
|
||||
def test_require_existing_preflights_before_existing_hash_scroll():
|
||||
fake = FakeQdrantHttp()
|
||||
original_request = fake.request
|
||||
|
||||
def request(method, url, **kwargs):
|
||||
if method == "POST" and url.endswith("/points/scroll"):
|
||||
return FakeResponse(404, {"status": "error"})
|
||||
return original_request(method, url, **kwargs)
|
||||
|
||||
store = QdrantVectorStore(
|
||||
base_url="http://qdrant:6333", collection="workspace-semantic", workspace_id="demo",
|
||||
workspace_revision="a" * 40, expected_dimension=1024,
|
||||
collection_lifecycle="require_existing", request=request,
|
||||
)
|
||||
|
||||
with pytest.raises(VectorStoreError, match="semantic_index_incompatible"):
|
||||
store.existing_hashes("memory", ["memory"])
|
||||
|
||||
assert [call for call in fake.calls if call[1].endswith("/points/scroll")] == []
|
||||
assert [call for call in fake.calls if call[0] == "PUT"] == []
|
||||
|
||||
|
||||
def test_require_existing_delete_maps_collection_404_after_preflight():
|
||||
fake = FakeQdrantHttp()
|
||||
fake.collection = {"vectors": {"size": 1024, "distance": "Cosine"}}
|
||||
fake.payload_indexes = {
|
||||
"content_hash", "document_id", "kind", "record_key", "record_kind",
|
||||
"vector_generation", "workspace_id", "workspace_revision",
|
||||
}
|
||||
original_request = fake.request
|
||||
deleted = False
|
||||
|
||||
def request(method, url, **kwargs):
|
||||
nonlocal deleted
|
||||
if method == "GET" and url.endswith("/collections/workspace-semantic") and not deleted:
|
||||
response = original_request(method, url, **kwargs)
|
||||
deleted = True
|
||||
fake.collection = None
|
||||
return response
|
||||
if method == "POST" and url.endswith("/points/delete?wait=true"):
|
||||
original_request(method, url, **kwargs)
|
||||
return FakeResponse(404, {"status": "error"})
|
||||
return original_request(method, url, **kwargs)
|
||||
|
||||
store = QdrantVectorStore(
|
||||
base_url="http://qdrant:6333", collection="workspace-semantic", workspace_id="demo",
|
||||
workspace_revision="a" * 40, expected_dimension=1024,
|
||||
collection_lifecycle="require_existing", request=request,
|
||||
)
|
||||
|
||||
with pytest.raises(VectorStoreError, match="semantic_index_incompatible"):
|
||||
store.delete_kinds("memory", ["memory"])
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("dimension", "distance", "indexes", "index_types"),
|
||||
[(384, "Cosine", set(), {}),
|
||||
|
||||
@@ -360,6 +360,9 @@ def _render_registry_runtime_config(tmp_path) -> str:
|
||||
"""Render the production schema-v3 runtime rather than duplicating its YAML."""
|
||||
tmp_path.mkdir(parents=True, exist_ok=True)
|
||||
backend = Path(__file__).resolve().parents[2] / "backend"
|
||||
tsx = backend / "node_modules/.bin/tsx"
|
||||
if not tsx.exists():
|
||||
pytest.skip("cross-runtime integration requires backend/node_modules/.bin/tsx")
|
||||
password_file = tmp_path / "dwh-password"
|
||||
password_file.write_text("not-a-canary")
|
||||
script = r"""
|
||||
@@ -407,7 +410,7 @@ process.stdout.write(renderRuntimeConfig(workspace, bindings, {
|
||||
script_path = tmp_path / "render-runtime.mts"
|
||||
script_path.write_text(script)
|
||||
result = subprocess.run(
|
||||
[str(backend / "node_modules/.bin/tsx"), str(script_path), str(password_file)],
|
||||
[str(tsx), str(script_path), str(password_file)],
|
||||
cwd=backend, check=True, capture_output=True, text=True,
|
||||
)
|
||||
return result.stdout
|
||||
@@ -446,8 +449,9 @@ def _qdrant_config_yaml(tmp_path, *, registry: bool) -> dict:
|
||||
return value
|
||||
|
||||
|
||||
@pytest.mark.integration
|
||||
@pytest.mark.parametrize("collection_state", ["missing", "incompatible"])
|
||||
def test_registry_rendered_session_and_maintenance_configs_require_existing_qdrant_collection(
|
||||
def test_registry_render_chain_produces_require_existing_qdrant_configs(
|
||||
tmp_path, monkeypatch, collection_state,
|
||||
):
|
||||
session_path, maintenance_path = _render_registry_runtime_configs(tmp_path)
|
||||
@@ -498,6 +502,49 @@ def test_loader_rejects_split_brain_qdrant_compatibility_resources(tmp_path):
|
||||
load_config(path)
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("field", "value"),
|
||||
[
|
||||
("base_url", {}),
|
||||
("base_url", []),
|
||||
("base_url", None),
|
||||
("base_url", 6333),
|
||||
("collection", {}),
|
||||
("collection", []),
|
||||
("collection", None),
|
||||
("collection", 7),
|
||||
],
|
||||
)
|
||||
def test_raw_qdrant_consistency_rejects_untrusted_top_level_values(tmp_path, field, value):
|
||||
values = _qdrant_config_yaml(tmp_path, registry=True)
|
||||
values["vectors"][field] = value
|
||||
values["resources"] = {
|
||||
"vector": {
|
||||
"engine": "qdrant",
|
||||
"base_url": "http://localhost:6333",
|
||||
"collection": "workspace-semantic",
|
||||
}
|
||||
}
|
||||
path = tmp_path / "invalid-qdrant.yaml"
|
||||
path.write_text(yaml.safe_dump(values))
|
||||
|
||||
with pytest.raises(ConfigError):
|
||||
load_config(path)
|
||||
|
||||
config_result = CliRunner().invoke(app, ["config", "check", "--config", str(path)])
|
||||
assert config_result.exit_code == 1
|
||||
assert "Traceback" not in config_result.stderr
|
||||
|
||||
vector_result = CliRunner().invoke(
|
||||
app, ["vector", "index-schema", "--json", "-c", str(path)]
|
||||
)
|
||||
assert vector_result.exit_code == 1
|
||||
assert vector_result.stderr == ""
|
||||
assert json.loads(vector_result.stdout) == {
|
||||
"status": "failed", "code": "invalid_configuration"
|
||||
}
|
||||
|
||||
|
||||
def test_loader_rejects_lifecycle_policy_in_compatibility_resource(tmp_path):
|
||||
values = _qdrant_config_yaml(tmp_path, registry=True)
|
||||
values["resources"] = {
|
||||
|
||||
@@ -188,6 +188,11 @@ class QdrantVectorStore:
|
||||
def existing_hashes(self, collection: str, kinds: list[str]) -> dict[str, str]:
|
||||
validate_collection(collection)
|
||||
validate_collection_kinds(collection, kinds)
|
||||
if self._collection_lifecycle == "require_existing":
|
||||
# Reconcile only after proving the exact existing collection contract.
|
||||
# This keeps an initial 404 typed and prevents an incompatible empty
|
||||
# collection from appearing healthy merely because there are no records.
|
||||
self._ensure_collection(strict=False)
|
||||
points = self._scroll(
|
||||
[
|
||||
*self._workspace_filter(),
|
||||
@@ -209,7 +214,6 @@ class QdrantVectorStore:
|
||||
|
||||
def upsert(self, collection: str, records: list[VectorWriteRecord]) -> int:
|
||||
validate_collection(collection)
|
||||
self._ensure_collection(strict=True)
|
||||
points = []
|
||||
for write_record in records:
|
||||
validate_collection_kinds(collection, [write_record.record.kind])
|
||||
@@ -228,6 +232,9 @@ class QdrantVectorStore:
|
||||
}
|
||||
)
|
||||
try:
|
||||
# Keep the preflight immediately adjacent to the mutation: a registry
|
||||
# collection may disappear after reconciliation has read its hashes.
|
||||
self._ensure_collection(strict=True)
|
||||
self._call(
|
||||
"PUT",
|
||||
f"/collections/{self._collection}/points?wait=true",
|
||||
@@ -249,12 +256,23 @@ class QdrantVectorStore:
|
||||
self._semantic_kind_filter(kinds),
|
||||
{"key": "record_kind", "match": {"any": sorted(kinds)}},
|
||||
]
|
||||
if self._collection_lifecycle == "require_existing":
|
||||
self._ensure_collection(strict=False)
|
||||
before = len(self._scroll(must))
|
||||
self._call(
|
||||
"POST",
|
||||
f"/collections/{self._collection}/points/delete?wait=true",
|
||||
{"filter": {"must": must}},
|
||||
)
|
||||
if self._collection_lifecycle == "require_existing":
|
||||
self._ensure_collection(strict=False)
|
||||
try:
|
||||
self._call(
|
||||
"POST",
|
||||
f"/collections/{self._collection}/points/delete?wait=true",
|
||||
{"filter": {"must": must}},
|
||||
)
|
||||
except VectorTransportError as exc:
|
||||
if self._collection_lifecycle == "require_existing" and exc.status_code == 404:
|
||||
raise SemanticIndexIncompatibleError(
|
||||
"Qdrant collection disappeared during semantic index deletion"
|
||||
) from exc
|
||||
raise
|
||||
return before
|
||||
|
||||
def delete_generation(self, collection: str, generation: str, workspace_id: str) -> int:
|
||||
@@ -269,14 +287,25 @@ class QdrantVectorStore:
|
||||
{"key": "record_kind", "match": {"any": ["evidence"]}},
|
||||
{"key": "vector_generation", "match": {"value": generation}},
|
||||
]
|
||||
if self._collection_lifecycle == "require_existing":
|
||||
self._ensure_collection(strict=False)
|
||||
before = len(
|
||||
self._scroll(must)
|
||||
)
|
||||
self._call(
|
||||
"POST",
|
||||
f"/collections/{self._collection}/points/delete?wait=true",
|
||||
{"filter": {"must": must}},
|
||||
)
|
||||
if self._collection_lifecycle == "require_existing":
|
||||
self._ensure_collection(strict=False)
|
||||
try:
|
||||
self._call(
|
||||
"POST",
|
||||
f"/collections/{self._collection}/points/delete?wait=true",
|
||||
{"filter": {"must": must}},
|
||||
)
|
||||
except VectorTransportError as exc:
|
||||
if self._collection_lifecycle == "require_existing" and exc.status_code == 404:
|
||||
raise SemanticIndexIncompatibleError(
|
||||
"Qdrant collection disappeared during semantic index deletion"
|
||||
) from exc
|
||||
raise
|
||||
return before
|
||||
|
||||
def list_evidence_generations(self, collection: str, workspace_id: str) -> list[str]:
|
||||
|
||||
+20
-8
@@ -664,13 +664,16 @@ 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()
|
||||
def _normalize_qdrant_base_url(value: Any) -> str | None:
|
||||
"""Normalize a URL, returning ``None`` for untrusted raw YAML values."""
|
||||
if not isinstance(value, str):
|
||||
return None
|
||||
try:
|
||||
parsed = urlparse(value)
|
||||
hostname = (parsed.hostname or "").lower()
|
||||
port = parsed.port
|
||||
except ValueError:
|
||||
port = None
|
||||
except (TypeError, ValueError):
|
||||
return 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 '/'}"
|
||||
@@ -687,9 +690,18 @@ def _validate_vector_resource_consistency(raw: dict[str, Any], translated: dict[
|
||||
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"):
|
||||
vector_url = _normalize_qdrant_base_url(vectors.get("base_url"))
|
||||
resource_url = _normalize_qdrant_base_url(resource.get("base_url"))
|
||||
vector_collection = vectors.get("collection")
|
||||
resource_collection = resource.get("collection")
|
||||
if (
|
||||
vector_url is None
|
||||
or resource_url is None
|
||||
or not isinstance(vector_collection, str)
|
||||
or not isinstance(resource_collection, str)
|
||||
or vector_url != resource_url
|
||||
or vector_collection != resource_collection
|
||||
):
|
||||
raise ConfigError(
|
||||
f"Configurazione non valida in {path}: vectors and resources.vector disagree"
|
||||
)
|
||||
|
||||
Reference in New Issue
Block a user