fix: harden qdrant payload and paging
This commit is contained in:
@@ -96,3 +96,75 @@ git diff --check
|
|||||||
`/Users/mp/.agents/skills/adversarial-review`; I performed a manual adversarial self-review
|
`/Users/mp/.agents/skills/adversarial-review`; I performed a manual adversarial self-review
|
||||||
instead
|
instead
|
||||||
- the focused suite still emits one pre-existing warning from `testcontainers.postgres`
|
- the focused suite still emits one pre-existing warning from `testcontainers.postgres`
|
||||||
|
|
||||||
|
## Fix round 1 — 2026-08-08
|
||||||
|
|
||||||
|
### Findings addressed
|
||||||
|
|
||||||
|
- IMPORTANT: metadata collisions could override canonical Qdrant payload identity fields and break
|
||||||
|
workspace isolation
|
||||||
|
- IMPORTANT: scroll-based operations only read the first page and did not follow
|
||||||
|
`next_page_offset`, making `existing_hashes`, `list_evidence_generations`, and delete counts
|
||||||
|
inexact beyond one page
|
||||||
|
|
||||||
|
### RED evidence
|
||||||
|
|
||||||
|
Command:
|
||||||
|
|
||||||
|
```bash
|
||||||
|
cd harness
|
||||||
|
./.venv/bin/pytest tests/test_qdrant_vector_store.py tests/test_vector_port_contract.py -q
|
||||||
|
```
|
||||||
|
|
||||||
|
Observed before the fix:
|
||||||
|
|
||||||
|
- exit code `1`
|
||||||
|
- `2 failed, 31 passed, 1 warning`
|
||||||
|
|
||||||
|
Representative failures:
|
||||||
|
|
||||||
|
- `assert payload["workspace_id"] == "demo"` failed because colliding `record.metadata`
|
||||||
|
overwrote canonical payload fields
|
||||||
|
- paginated scroll test missed later pages, so `existing_hashes` and generation cleanup counts
|
||||||
|
were incomplete
|
||||||
|
|
||||||
|
### GREEN evidence
|
||||||
|
|
||||||
|
Command:
|
||||||
|
|
||||||
|
```bash
|
||||||
|
cd harness
|
||||||
|
./.venv/bin/pytest tests/test_qdrant_vector_store.py tests/test_vector_port_contract.py -q
|
||||||
|
```
|
||||||
|
|
||||||
|
Observed after the fix:
|
||||||
|
|
||||||
|
- exit code `0`
|
||||||
|
- `33 passed, 1 warning`
|
||||||
|
|
||||||
|
Touched-file lint:
|
||||||
|
|
||||||
|
```bash
|
||||||
|
cd harness
|
||||||
|
./.venv/bin/ruff check tht/adapters/vector/qdrant.py tht/vectorstore/records.py \
|
||||||
|
tests/test_qdrant_vector_store.py tests/test_vector_port_contract.py
|
||||||
|
```
|
||||||
|
|
||||||
|
- exit code `0`
|
||||||
|
- `All checks passed!`
|
||||||
|
|
||||||
|
Patch hygiene:
|
||||||
|
|
||||||
|
```bash
|
||||||
|
git diff --check
|
||||||
|
```
|
||||||
|
|
||||||
|
- exit code `0`
|
||||||
|
|
||||||
|
### Minimal fix
|
||||||
|
|
||||||
|
- made `qdrant_payload` apply canonical fields after `record.metadata` so workspace ID, semantic
|
||||||
|
kind, original record kind, canonical record key, and content hash cannot be overridden by
|
||||||
|
metadata collisions
|
||||||
|
- paginated `_scroll` until `next_page_offset` is absent, sent the returned `offset` back on the
|
||||||
|
next request, and reject repeated offsets as malformed to avoid infinite loops
|
||||||
|
|||||||
@@ -38,6 +38,7 @@ class FakeQdrantHttp:
|
|||||||
self.fail_request: Exception | None = None
|
self.fail_request: Exception | None = None
|
||||||
self.malformed_query = False
|
self.malformed_query = False
|
||||||
self.malformed_scroll = False
|
self.malformed_scroll = False
|
||||||
|
self.scroll_pages: list[dict] | None = None
|
||||||
|
|
||||||
def request(self, method, url, *, json=None, timeout=None):
|
def request(self, method, url, *, json=None, timeout=None):
|
||||||
self.calls.append((method, url, json))
|
self.calls.append((method, url, json))
|
||||||
@@ -98,11 +99,23 @@ class FakeQdrantHttp:
|
|||||||
if method == "POST" and path == "/collections/workspace-semantic/points/scroll":
|
if method == "POST" and path == "/collections/workspace-semantic/points/scroll":
|
||||||
if self.malformed_scroll:
|
if self.malformed_scroll:
|
||||||
return FakeResponse(200, {"result": {"points": "bad"}})
|
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(
|
wanted = sorted(
|
||||||
_match_points(self.points.values(), json["filter"]),
|
_match_points(self.points.values(), json["filter"]),
|
||||||
key=lambda point: point["payload"]["record_key"],
|
key=lambda point: point["payload"]["record_key"],
|
||||||
)
|
)
|
||||||
return FakeResponse(200, {"result": {"points": wanted}})
|
return FakeResponse(200, {"result": {"points": wanted, "next_page_offset": None}})
|
||||||
|
|
||||||
if method == "POST" and path == "/collections/workspace-semantic/points/delete":
|
if method == "POST" and path == "/collections/workspace-semantic/points/delete":
|
||||||
doomed = [point["id"] for point in _match_points(self.points.values(), json["filter"])]
|
doomed = [point["id"] for point in _match_points(self.points.values(), json["filter"])]
|
||||||
@@ -310,3 +323,96 @@ def test_sanitizes_timeout_and_malformed_responses():
|
|||||||
fake.malformed_scroll = True
|
fake.malformed_scroll = True
|
||||||
with pytest.raises(VectorStoreError, match="Qdrant returned malformed scroll response"):
|
with pytest.raises(VectorStoreError, match="Qdrant returned malformed scroll response"):
|
||||||
store.existing_hashes("memory", ["memory"])
|
store.existing_hashes("memory", ["memory"])
|
||||||
|
|
||||||
|
|
||||||
|
def test_upsert_payload_keeps_canonical_identity_when_metadata_collides():
|
||||||
|
fake = FakeQdrantHttp()
|
||||||
|
store = _store(fake)
|
||||||
|
record = _write_record(
|
||||||
|
"memory:1",
|
||||||
|
"memory",
|
||||||
|
metadata={
|
||||||
|
"workspace_id": "evil",
|
||||||
|
"kind": "evil",
|
||||||
|
"record_kind": "evil",
|
||||||
|
"record_key": "evil",
|
||||||
|
"content_hash": "evil",
|
||||||
|
},
|
||||||
|
)
|
||||||
|
|
||||||
|
store.upsert("memory", [record])
|
||||||
|
|
||||||
|
payload = next(iter(fake.points.values()))["payload"]
|
||||||
|
assert payload["workspace_id"] == "demo"
|
||||||
|
assert payload["kind"] == "memory"
|
||||||
|
assert payload["record_kind"] == "memory"
|
||||||
|
assert payload["record_key"] == "memory:1"
|
||||||
|
assert payload["content_hash"] == "sha256:" + "a" * 64
|
||||||
|
|
||||||
|
|
||||||
|
def test_scroll_based_operations_paginate_until_next_page_offset_is_absent():
|
||||||
|
fake = FakeQdrantHttp()
|
||||||
|
generation_a = "gen:" + "1" * 32
|
||||||
|
generation_b = "gen:" + "2" * 32
|
||||||
|
fake.scroll_pages = [
|
||||||
|
{
|
||||||
|
"offset": None,
|
||||||
|
"points": [
|
||||||
|
{
|
||||||
|
"id": "p1",
|
||||||
|
"payload": {
|
||||||
|
"workspace_id": "demo",
|
||||||
|
"kind": "evidence",
|
||||||
|
"record_kind": "evidence",
|
||||||
|
"record_key": f"demo:{generation_a}:chunk:1",
|
||||||
|
"content_hash": "sha256:" + "a" * 64,
|
||||||
|
"vector_generation": generation_a,
|
||||||
|
},
|
||||||
|
}
|
||||||
|
],
|
||||||
|
"next_page_offset": "page-2",
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"offset": "page-2",
|
||||||
|
"points": [
|
||||||
|
{
|
||||||
|
"id": "p2",
|
||||||
|
"payload": {
|
||||||
|
"workspace_id": "demo",
|
||||||
|
"kind": "evidence",
|
||||||
|
"record_kind": "evidence",
|
||||||
|
"record_key": f"demo:{generation_a}:chunk:2",
|
||||||
|
"content_hash": "sha256:" + "b" * 64,
|
||||||
|
"vector_generation": generation_a,
|
||||||
|
},
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"id": "p3",
|
||||||
|
"payload": {
|
||||||
|
"workspace_id": "demo",
|
||||||
|
"kind": "evidence",
|
||||||
|
"record_kind": "evidence",
|
||||||
|
"record_key": f"demo:{generation_b}:chunk:3",
|
||||||
|
"content_hash": "sha256:" + "c" * 64,
|
||||||
|
"vector_generation": generation_b,
|
||||||
|
},
|
||||||
|
},
|
||||||
|
],
|
||||||
|
"next_page_offset": None,
|
||||||
|
},
|
||||||
|
]
|
||||||
|
store = _store(fake)
|
||||||
|
|
||||||
|
assert store.existing_hashes("evidence", ["evidence"]) == {
|
||||||
|
f"demo:{generation_a}:chunk:1": "sha256:" + "a" * 64,
|
||||||
|
f"demo:{generation_a}:chunk:2": "sha256:" + "b" * 64,
|
||||||
|
f"demo:{generation_b}:chunk:3": "sha256:" + "c" * 64,
|
||||||
|
}
|
||||||
|
assert store.list_evidence_generations("evidence", "demo") == [generation_a, generation_b]
|
||||||
|
assert store.delete_generation("evidence", generation_a, "demo") == 2
|
||||||
|
offsets = [
|
||||||
|
call[2].get("offset")
|
||||||
|
for call in fake.calls
|
||||||
|
if call[0] == "POST" and call[1].endswith("/points/scroll")
|
||||||
|
]
|
||||||
|
assert offsets[:2] == [None, "page-2"]
|
||||||
|
|||||||
@@ -318,15 +318,32 @@ class QdrantVectorStore:
|
|||||||
return result
|
return result
|
||||||
|
|
||||||
def _scroll(self, must: list[dict]) -> list[dict]:
|
def _scroll(self, must: list[dict]) -> list[dict]:
|
||||||
response = self._call(
|
points: list[dict] = []
|
||||||
"POST",
|
offset = None
|
||||||
f"/collections/{self._collection}/points/scroll",
|
seen_offsets = set()
|
||||||
{"with_payload": True, "limit": 10000, "filter": {"must": must}},
|
while True:
|
||||||
)
|
response = self._call(
|
||||||
points = response.get("result", {}).get("points")
|
"POST",
|
||||||
if not isinstance(points, list):
|
f"/collections/{self._collection}/points/scroll",
|
||||||
raise VectorStoreError("Qdrant returned malformed scroll response")
|
{
|
||||||
return points
|
"with_payload": True,
|
||||||
|
"limit": 10000,
|
||||||
|
"filter": {"must": must},
|
||||||
|
"offset": offset,
|
||||||
|
},
|
||||||
|
)
|
||||||
|
result = response.get("result", {})
|
||||||
|
page = result.get("points")
|
||||||
|
if not isinstance(page, list):
|
||||||
|
raise VectorStoreError("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")
|
||||||
|
seen_offsets.add(next_page_offset)
|
||||||
|
offset = next_page_offset
|
||||||
|
|
||||||
def _hit_from_point(self, point: dict) -> VectorHit:
|
def _hit_from_point(self, point: dict) -> VectorHit:
|
||||||
payload = point.get("payload")
|
payload = point.get("payload")
|
||||||
|
|||||||
@@ -30,6 +30,7 @@ def qdrant_semantic_kind(kind: str) -> str:
|
|||||||
def qdrant_payload(record: VectorRecord, *, content_hash: str, workspace_id: str) -> dict:
|
def qdrant_payload(record: VectorRecord, *, content_hash: str, workspace_id: str) -> dict:
|
||||||
semantic_kind = qdrant_semantic_kind(record.kind)
|
semantic_kind = qdrant_semantic_kind(record.kind)
|
||||||
return {
|
return {
|
||||||
|
**record.metadata,
|
||||||
"workspace_id": workspace_id,
|
"workspace_id": workspace_id,
|
||||||
"kind": semantic_kind,
|
"kind": semantic_kind,
|
||||||
"record_kind": record.kind,
|
"record_kind": record.kind,
|
||||||
@@ -38,7 +39,6 @@ def qdrant_payload(record: VectorRecord, *, content_hash: str, workspace_id: str
|
|||||||
"title": record.title,
|
"title": record.title,
|
||||||
"content": record.content,
|
"content": record.content,
|
||||||
"content_hash": content_hash,
|
"content_hash": content_hash,
|
||||||
**record.metadata,
|
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user