diff --git a/harness/tests/test_qdrant_vector_store.py b/harness/tests/test_qdrant_vector_store.py index fe36104d..676897e3 100644 --- a/harness/tests/test_qdrant_vector_store.py +++ b/harness/tests/test_qdrant_vector_store.py @@ -33,6 +33,7 @@ class FakeQdrantHttp: 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 @@ -59,7 +60,7 @@ class FakeQdrantHttp: } }, "payload_schema": { - field: {"data_type": "keyword"} for field in sorted(self.payload_indexes) + field: {"data_type": self.payload_index_types.get(field, "keyword")} for field in sorted(self.payload_indexes) }, } }) @@ -75,6 +76,8 @@ class FakeQdrantHttp: 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"}}) @@ -161,13 +164,14 @@ def _write_record(record_id: str, kind: str, *, metadata=None): ) -def _store(fake: FakeQdrantHttp) -> QdrantVectorStore: +def _store(fake: FakeQdrantHttp, *, collection_lifecycle="create_if_missing") -> QdrantVectorStore: return QdrantVectorStore( base_url="http://qdrant:6333", collection="workspace-semantic", workspace_id="demo", workspace_revision="a" * 40, expected_dimension=1024, + collection_lifecycle=collection_lifecycle, request=fake.request, ) @@ -178,6 +182,84 @@ def test_point_id_is_deterministic_uuidv5(): ) + +def test_require_existing_refuses_missing_collection_without_mutations(): + fake = FakeQdrantHttp() + store = _store(fake, collection_lifecycle="require_existing") + + with pytest.raises(VectorStoreError, match="semantic_index_incompatible"): + store.upsert("memory", [_write_record("memory:1", "memory")]) + + assert [call for call in fake.calls if call[0] == "PUT"] == [] + + +@pytest.mark.parametrize( + ("dimension", "distance", "indexes", "index_types"), + [(384, "Cosine", set(), {}), + (1024, "Dot", set(), {}), + (1024, "Cosine", {"content_hash"}, {}), + (1024, "Cosine", { + "content_hash", "document_id", "kind", "record_key", "record_kind", + "vector_generation", "workspace_id", "workspace_revision", + }, {"kind": "integer"})], +) +def test_require_existing_refuses_incompatible_collection_without_mutations( + dimension, distance, indexes, index_types +): + fake = FakeQdrantHttp(dimension=dimension, distance=distance) + fake.collection = {"vectors": {"size": dimension, "distance": distance}} + fake.payload_indexes = indexes + fake.payload_index_types = index_types + store = _store(fake, collection_lifecycle="require_existing") + + with pytest.raises(VectorStoreError, match="semantic_index_incompatible"): + store.upsert("memory", [_write_record("memory:1", "memory")]) + + assert [call for call in fake.calls if call[0] == "PUT"] == [] + + +def test_require_existing_writes_compatible_collection_without_lifecycle_mutations(): + 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", + } + store = _store(fake, collection_lifecycle="require_existing") + + assert store.upsert("memory", [_write_record("memory:1", "memory")]) == 1 + assert not [call for call in fake.calls if call[0] == "PUT" and call[1].endswith("/index")] + + +def test_require_existing_write_fails_after_collection_is_deleted_without_recreating(): + 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 + response = original_request(method, url, **kwargs) + if method == "GET" and url.endswith("/collections/workspace-semantic") and not deleted: + deleted = True + fake.collection = None + return response + + 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): + store.upsert("memory", [_write_record("memory:1", "memory")]) + assert not [call for call in fake.calls if call[0] == "PUT" and call[1].endswith("/collections/workspace-semantic")] + + def test_upsert_creates_collection_and_keyword_indexes_idempotently(): fake = FakeQdrantHttp() store = _store(fake) diff --git a/harness/tests/test_registry_evidence_config.py b/harness/tests/test_registry_evidence_config.py index f972974b..718650ef 100644 --- a/harness/tests/test_registry_evidence_config.py +++ b/harness/tests/test_registry_evidence_config.py @@ -10,6 +10,7 @@ from tht.adapters.evidence import HttpManifestEvidenceSource from tht.adapters.factory import build_evidence_sources from tht.cli import app from tht.config import ConfigError, load_config +from tht.jobs.dwh_pipeline import config_dwh_binding SIGNED_CANARY = "SIGNED-CANARY-QUERY" ACCESS_CANARY = "ACCESS-CANARY" @@ -350,3 +351,59 @@ def test_validation_repr_cli_and_exception_output_never_disclose_transport_secre assert result.exit_code == 0 assert_no_canaries(result.stdout) assert_no_canaries(result.stderr) + + +def _qdrant_config_yaml(tmp_path, *, registry: bool) -> dict: + value = { + "dwh": { + "type": "postgres_direct", + "connection": { + "database": "analytics", "schema": "public", "user": "reader", + "password": "not-a-canary", + }, + }, + "vectors": { + "type": "qdrant", "base_url": "http://localhost:6333", + "collection": "workspace-semantic", + }, + "embeddings": { + "base_url": "http://localhost:11434", "model": "qwen3-embedding:0.6b", "dim": 1024, + }, + "evidence": {"source_root": str(tmp_path / "evidence")}, + } + if registry: + value["runtime_identity"] = { + "workspace_id": "demo-workspace", "workspace_revision": "a" * 40, + } + return value + + +@pytest.mark.parametrize("mode", ["session", "maintenance"]) +def test_registry_runtime_configs_require_existing_qdrant_collection(tmp_path, mode): + path = tmp_path / f"{mode}.yaml" + path.write_text(yaml.safe_dump(_qdrant_config_yaml(tmp_path, registry=True))) + + cfg = load_config(path) + + assert cfg.vectors.collection_lifecycle == "require_existing" + + +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))) + + cfg = load_config(path) + + assert cfg.vectors.collection_lifecycle == "create_if_missing" + + +def test_registry_session_and_maintenance_bindings_are_equal(tmp_path): + values = _qdrant_config_yaml(tmp_path, registry=True) + session_path = tmp_path / "session.yaml" + maintenance_path = tmp_path / "maintenance.yaml" + session_path.write_text(yaml.safe_dump(values)) + maintenance_path.write_text(yaml.safe_dump(values)) + + assert config_dwh_binding(load_config(session_path)) == config_dwh_binding( + load_config(maintenance_path) + ) diff --git a/harness/tht/adapters/factory.py b/harness/tht/adapters/factory.py index 3d594ed7..a9f52936 100644 --- a/harness/tht/adapters/factory.py +++ b/harness/tht/adapters/factory.py @@ -38,6 +38,7 @@ def build_vector_store(cfg: Config, *, require_write: bool = False) -> VectorSto workspace_id=cfg._workspace_id, workspace_revision=cfg._workspace_revision, expected_dimension=cfg.embeddings.dim if cfg.embeddings is not None else None, + collection_lifecycle=resource.collection_lifecycle, ) 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 f95611f6..f99359d9 100644 --- a/harness/tht/adapters/vector/qdrant.py +++ b/harness/tht/adapters/vector/qdrant.py @@ -55,6 +55,7 @@ class QdrantVectorStore: workspace_id: str, workspace_revision: str | None = None, expected_dimension: int | None = None, + collection_lifecycle: str = "create_if_missing", request: Callable[..., object] | None = None, connect_timeout: float = 2.0, read_timeout: float = 10.0, @@ -64,6 +65,9 @@ class QdrantVectorStore: self._workspace_id = workspace_id self._workspace_revision = workspace_revision self._expected_dimension = expected_dimension + if collection_lifecycle not in ("create_if_missing", "require_existing"): + raise ValueError("Unsupported Qdrant collection lifecycle") + self._collection_lifecycle = collection_lifecycle self._request = request or requests.request self._timeout = (connect_timeout, read_timeout) @@ -306,6 +310,8 @@ class QdrantVectorStore: if response is None: 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}", @@ -328,11 +334,17 @@ class QdrantVectorStore: self._expected_dimension is not None and (size != self._expected_dimension or distance != "Cosine") ): - raise VectorStoreError("Qdrant collection configuration mismatch") + raise VectorStoreError("Qdrant collection configuration mismatch (semantic_index_incompatible)") + payload_schema = result.get("payload_schema") + if not isinstance(payload_schema, dict): + raise VectorStoreError("Qdrant returned malformed collection response") for field_name in _KEYWORD_INDEXES: - if field_name not in result.get("payload_schema", {}): + field = payload_schema.get(field_name) + if not isinstance(field, dict) or field.get("data_type") != "keyword": 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", diff --git a/harness/tht/config.py b/harness/tht/config.py index 6916d13d..b0debe03 100644 --- a/harness/tht/config.py +++ b/harness/tht/config.py @@ -222,6 +222,9 @@ class QdrantConfig(BaseModel): type: Literal["qdrant"] base_url: str collection: str = Field(min_length=1) + # Internal runtime policy. Registry-rendered configs must not create or alter + # the workspace-owned semantic collection; legacy configs retain compatibility. + collection_lifecycle: Literal["create_if_missing", "require_existing"] = "create_if_missing" VectorResourceConfig = Annotated[ @@ -547,6 +550,25 @@ 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) + 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 ( + isinstance(translated.get("runtime_identity"), dict) + and isinstance(vectors, dict) + and vectors.get("type") == "qdrant" + ): + vectors["collection_lifecycle"] = "require_existing" _populate_legacy_views(translated) try: cfg = Config.model_validate(translated) @@ -662,7 +684,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"} + allowed = {"engine", "base_url", "collection", "collection_lifecycle"} unexpected = sorted(set(vector) - allowed) if unexpected: raise ConfigError(