fix: prevent P2 from owning Qdrant lifecycle
This commit is contained in:
@@ -33,6 +33,7 @@ class FakeQdrantHttp:
|
|||||||
self.distance = distance
|
self.distance = distance
|
||||||
self.collection = None
|
self.collection = None
|
||||||
self.payload_indexes: set[str] = set()
|
self.payload_indexes: set[str] = set()
|
||||||
|
self.payload_index_types: dict[str, str] = {}
|
||||||
self.points: dict[str, dict] = {}
|
self.points: dict[str, dict] = {}
|
||||||
self.calls: list[tuple[str, str, dict | None]] = []
|
self.calls: list[tuple[str, str, dict | None]] = []
|
||||||
self.fail_request: Exception | None = None
|
self.fail_request: Exception | None = None
|
||||||
@@ -59,7 +60,7 @@ class FakeQdrantHttp:
|
|||||||
}
|
}
|
||||||
},
|
},
|
||||||
"payload_schema": {
|
"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"})
|
return FakeResponse(200, {"status": "ok"})
|
||||||
|
|
||||||
if method == "PUT" and path == "/collections/workspace-semantic/points":
|
if method == "PUT" and path == "/collections/workspace-semantic/points":
|
||||||
|
if self.collection is None:
|
||||||
|
return FakeResponse(404, {"status": "error"})
|
||||||
for point in json["points"]:
|
for point in json["points"]:
|
||||||
self.points[point["id"]] = point
|
self.points[point["id"]] = point
|
||||||
return FakeResponse(200, {"result": {"status": "acknowledged"}})
|
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(
|
return QdrantVectorStore(
|
||||||
base_url="http://qdrant:6333",
|
base_url="http://qdrant:6333",
|
||||||
collection="workspace-semantic",
|
collection="workspace-semantic",
|
||||||
workspace_id="demo",
|
workspace_id="demo",
|
||||||
workspace_revision="a" * 40,
|
workspace_revision="a" * 40,
|
||||||
expected_dimension=1024,
|
expected_dimension=1024,
|
||||||
|
collection_lifecycle=collection_lifecycle,
|
||||||
request=fake.request,
|
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():
|
def test_upsert_creates_collection_and_keyword_indexes_idempotently():
|
||||||
fake = FakeQdrantHttp()
|
fake = FakeQdrantHttp()
|
||||||
store = _store(fake)
|
store = _store(fake)
|
||||||
|
|||||||
@@ -10,6 +10,7 @@ from tht.adapters.evidence import HttpManifestEvidenceSource
|
|||||||
from tht.adapters.factory import build_evidence_sources
|
from tht.adapters.factory import build_evidence_sources
|
||||||
from tht.cli import app
|
from tht.cli import app
|
||||||
from tht.config import ConfigError, load_config
|
from tht.config import ConfigError, load_config
|
||||||
|
from tht.jobs.dwh_pipeline import config_dwh_binding
|
||||||
|
|
||||||
SIGNED_CANARY = "SIGNED-CANARY-QUERY"
|
SIGNED_CANARY = "SIGNED-CANARY-QUERY"
|
||||||
ACCESS_CANARY = "ACCESS-CANARY"
|
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 result.exit_code == 0
|
||||||
assert_no_canaries(result.stdout)
|
assert_no_canaries(result.stdout)
|
||||||
assert_no_canaries(result.stderr)
|
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)
|
||||||
|
)
|
||||||
|
|||||||
@@ -38,6 +38,7 @@ def build_vector_store(cfg: Config, *, require_write: bool = False) -> VectorSto
|
|||||||
workspace_id=cfg._workspace_id,
|
workspace_id=cfg._workspace_id,
|
||||||
workspace_revision=cfg._workspace_revision,
|
workspace_revision=cfg._workspace_revision,
|
||||||
expected_dimension=cfg.embeddings.dim if cfg.embeddings is not None else None,
|
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.
|
case other: # pragma: no cover - Pydantic's discriminator rejects this first.
|
||||||
raise ConfigError(f"Adapter vector non supportato: {other}")
|
raise ConfigError(f"Adapter vector non supportato: {other}")
|
||||||
|
|||||||
@@ -55,6 +55,7 @@ class QdrantVectorStore:
|
|||||||
workspace_id: str,
|
workspace_id: str,
|
||||||
workspace_revision: str | None = None,
|
workspace_revision: str | None = None,
|
||||||
expected_dimension: int | None = None,
|
expected_dimension: int | None = None,
|
||||||
|
collection_lifecycle: str = "create_if_missing",
|
||||||
request: Callable[..., object] | None = None,
|
request: Callable[..., object] | None = None,
|
||||||
connect_timeout: float = 2.0,
|
connect_timeout: float = 2.0,
|
||||||
read_timeout: float = 10.0,
|
read_timeout: float = 10.0,
|
||||||
@@ -64,6 +65,9 @@ class QdrantVectorStore:
|
|||||||
self._workspace_id = workspace_id
|
self._workspace_id = workspace_id
|
||||||
self._workspace_revision = workspace_revision
|
self._workspace_revision = workspace_revision
|
||||||
self._expected_dimension = expected_dimension
|
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._request = request or requests.request
|
||||||
self._timeout = (connect_timeout, read_timeout)
|
self._timeout = (connect_timeout, read_timeout)
|
||||||
|
|
||||||
@@ -306,6 +310,8 @@ class QdrantVectorStore:
|
|||||||
if response is None:
|
if response is None:
|
||||||
if not strict:
|
if not strict:
|
||||||
raise VectorStoreError("Qdrant collection is missing")
|
raise VectorStoreError("Qdrant collection is missing")
|
||||||
|
if self._collection_lifecycle == "require_existing":
|
||||||
|
raise VectorStoreError("Qdrant collection configuration mismatch (semantic_index_incompatible)")
|
||||||
self._call(
|
self._call(
|
||||||
"PUT",
|
"PUT",
|
||||||
f"/collections/{self._collection}",
|
f"/collections/{self._collection}",
|
||||||
@@ -328,11 +334,17 @@ class QdrantVectorStore:
|
|||||||
self._expected_dimension is not None
|
self._expected_dimension is not None
|
||||||
and (size != self._expected_dimension or distance != "Cosine")
|
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:
|
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:
|
if not strict:
|
||||||
raise VectorStoreError("Qdrant collection payload indexes mismatch")
|
raise VectorStoreError("Qdrant collection payload indexes mismatch")
|
||||||
|
if self._collection_lifecycle == "require_existing":
|
||||||
|
raise VectorStoreError("Qdrant collection configuration mismatch (semantic_index_incompatible)")
|
||||||
self._call(
|
self._call(
|
||||||
"PUT",
|
"PUT",
|
||||||
f"/collections/{self._collection}/index",
|
f"/collections/{self._collection}/index",
|
||||||
|
|||||||
+23
-1
@@ -222,6 +222,9 @@ class QdrantConfig(BaseModel):
|
|||||||
type: Literal["qdrant"]
|
type: Literal["qdrant"]
|
||||||
base_url: str
|
base_url: str
|
||||||
collection: str = Field(min_length=1)
|
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[
|
VectorResourceConfig = Annotated[
|
||||||
@@ -547,6 +550,25 @@ def load_config(path: Path) -> Config:
|
|||||||
_validate_internal_embedding_contract(expanded, path)
|
_validate_internal_embedding_contract(expanded, path)
|
||||||
_validate_internal_vector_contract(expanded, path)
|
_validate_internal_vector_contract(expanded, path)
|
||||||
translated, used_legacy = translate_legacy_config(expanded)
|
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)
|
_populate_legacy_views(translated)
|
||||||
try:
|
try:
|
||||||
cfg = Config.model_validate(translated)
|
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")
|
engine = vector.get("engine")
|
||||||
base_url = vector.get("base_url")
|
base_url = vector.get("base_url")
|
||||||
collection = vector.get("collection")
|
collection = vector.get("collection")
|
||||||
allowed = {"engine", "base_url", "collection"}
|
allowed = {"engine", "base_url", "collection", "collection_lifecycle"}
|
||||||
unexpected = sorted(set(vector) - allowed)
|
unexpected = sorted(set(vector) - allowed)
|
||||||
if unexpected:
|
if unexpected:
|
||||||
raise ConfigError(
|
raise ConfigError(
|
||||||
|
|||||||
Reference in New Issue
Block a user