Publish documentation / publish (push) Successful in 1m27s
Add PostgreSQL-backed memory, editable evidence with source review and activation, and human-approved archive repairs across the harness, API, and UI. Include migrations, deployment support, regression coverage, and validation documentation. Refresh permissions from validated session roles so existing administrator logins can access newly deployed archive management features.
1035 lines
39 KiB
Python
1035 lines
39 KiB
Python
import json
|
|
from uuid import NAMESPACE_URL, uuid5
|
|
|
|
import pytest
|
|
import requests
|
|
|
|
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.sparse_vectors: dict[str, dict] | None = None
|
|
self.payload_indexes: set[str] = set()
|
|
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
|
|
self.drop_collection_on_points = False
|
|
|
|
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},
|
|
**(
|
|
{"sparse_vectors": self.sparse_vectors}
|
|
if self.sparse_vectors is not None else {}
|
|
),
|
|
}
|
|
},
|
|
"payload_schema": {
|
|
field: {"data_type": "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"]
|
|
self.sparse_vectors = json.get("sparse_vectors")
|
|
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.drop_collection_on_points:
|
|
self.collection = 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"}})
|
|
filter_value = json["filter"] if "filter" in json else json["prefetch"][0]["filter"]
|
|
wanted = _match_points(self.points.values(), filter_value)
|
|
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, sparse_text=None, sparse_language=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,
|
|
sparse_text=sparse_text,
|
|
sparse_language=sparse_language,
|
|
)
|
|
|
|
|
|
def _store(
|
|
fake: FakeQdrantHttp, *, collection_lifecycle: str = "self_heal"
|
|
) -> 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,
|
|
)
|
|
|
|
|
|
def test_explicit_collections_route_reference_and_memory_independently():
|
|
calls = []
|
|
|
|
def request(method, url, *, json=None, timeout=None):
|
|
calls.append((method, url, json))
|
|
if method == "GET" and url.endswith("/collections/demo"):
|
|
return FakeResponse(404, {"status": "error"})
|
|
if method == "GET":
|
|
return FakeResponse(200, {
|
|
"result": {
|
|
"config": {"params": {"vectors": {"size": 1024, "distance": "Cosine"}}},
|
|
"payload_schema": {
|
|
field: {"data_type": "keyword"}
|
|
for field in (
|
|
"content_hash", "document_id", "kind", "record_key", "record_kind",
|
|
"vector_generation", "workspace_id", "workspace_revision",
|
|
)
|
|
},
|
|
}
|
|
})
|
|
if method == "POST" and url.endswith("/points/query"):
|
|
return FakeResponse(200, {"result": {"points": []}})
|
|
if method == "PUT" and "/points" in url:
|
|
return FakeResponse(200, {"result": {"status": "acknowledged"}})
|
|
raise AssertionError((method, url, json))
|
|
|
|
store = QdrantVectorStore(
|
|
base_url="http://qdrant:6333",
|
|
collections={"reference": "demo-reference", "memory": "demo-memory"},
|
|
workspace_id="demo",
|
|
workspace_revision="a" * 40,
|
|
expected_dimension=1024,
|
|
request=request,
|
|
)
|
|
store.search(["schema_records"], [0.2] * 1024, limit=1)
|
|
store.search(["memory"], [0.2] * 1024, limit=1)
|
|
store.upsert("memory", [_write_record("memory:1", "memory")])
|
|
|
|
urls = [url for _method, url, _payload in calls]
|
|
assert any("/collections/demo-reference/points/query" in url for url in urls)
|
|
assert any("/collections/demo-memory/points/query" in url for url in urls)
|
|
assert any("/collections/demo-memory/points" in url for url in urls)
|
|
assert not any("/collections/demo-reference/points?" in url for url in urls)
|
|
|
|
|
|
def test_clear_drops_reference_without_deleting_memory_collection():
|
|
calls = []
|
|
|
|
def request(method, url, *, json=None, timeout=None):
|
|
calls.append((method, url))
|
|
if method == "GET" and url.endswith("/collections/demo"):
|
|
return FakeResponse(404, {"status": "error"})
|
|
if method == "GET" and url.endswith("/collections/demo-reference"):
|
|
return FakeResponse(200, {"result": {}})
|
|
if method == "DELETE" and url.endswith("/collections/demo-reference"):
|
|
return FakeResponse(200, {"status": "ok"})
|
|
raise AssertionError((method, url, json))
|
|
|
|
store = QdrantVectorStore(
|
|
base_url="http://qdrant:6333",
|
|
collections={"reference": "demo-reference", "memory": "demo-memory"},
|
|
workspace_id="demo",
|
|
expected_dimension=1024,
|
|
request=request,
|
|
)
|
|
|
|
assert store.clear_reference() is True
|
|
assert ("DELETE", "http://qdrant:6333/collections/demo-reference") in calls
|
|
assert not any(method == "DELETE" and url.endswith("/collections/demo-memory") for method, url in calls)
|
|
|
|
|
|
def test_clear_does_not_read_or_import_legacy_memory():
|
|
calls = []
|
|
|
|
def request(method, url, *, json=None, timeout=None):
|
|
calls.append((method, url))
|
|
if method == "GET" and url.endswith("/collections/demo-reference"):
|
|
return FakeResponse(200, {"result": {}})
|
|
if method == "DELETE" and url.endswith("/collections/demo-reference"):
|
|
return FakeResponse(200, {"status": "ok"})
|
|
raise AssertionError((method, url, json))
|
|
|
|
store = QdrantVectorStore(
|
|
base_url="http://qdrant:6333",
|
|
collections={"reference": "demo-reference", "memory": "demo-memory"},
|
|
workspace_id="demo", expected_dimension=1024, request=request,
|
|
)
|
|
assert store.clear_reference() is True
|
|
assert calls == [
|
|
("GET", "http://qdrant:6333/collections/demo-reference"),
|
|
("DELETE", "http://qdrant:6333/collections/demo-reference"),
|
|
]
|
|
|
|
|
|
def _ready_collection_with_bm25(fake: FakeQdrantHttp) -> None:
|
|
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",
|
|
}
|
|
fake.sparse_vectors = {"bm25": {"modifier": "idf"}}
|
|
|
|
|
|
def test_point_id_is_deterministic_uuidv5():
|
|
assert point_id("demo", "memory", "memory:1") == str(
|
|
uuid5(NAMESPACE_URL, "thothii:demo:memory:memory:1")
|
|
)
|
|
|
|
|
|
def test_upsert_creates_collection_and_keyword_indexes_idempotently():
|
|
fake = FakeQdrantHttp()
|
|
store = _store(fake)
|
|
|
|
assert store.upsert("memory", [_write_record("memory:1", "memory")]) == 1
|
|
assert store.upsert("memory", [_write_record("memory:1", "memory")]) == 1
|
|
|
|
creates = [call for call in fake.calls if call[0] == "PUT" and call[1].endswith("/collections/workspace-semantic")]
|
|
assert len(creates) == 1
|
|
assert creates[0][2] == {"vectors": {"size": 1024, "distance": "Cosine"}}
|
|
assert fake.payload_indexes == {
|
|
"content_hash",
|
|
"document_id",
|
|
"kind",
|
|
"record_key",
|
|
"record_kind",
|
|
"vector_generation",
|
|
"workspace_id",
|
|
"workspace_revision",
|
|
}
|
|
|
|
|
|
def test_upsert_refuses_collection_dimension_or_distance_mismatch_without_recreating():
|
|
fake = FakeQdrantHttp(dimension=384, distance="Dot")
|
|
fake.collection = {"vectors": {"size": 384, "distance": "Dot"}}
|
|
store = _store(fake)
|
|
|
|
with pytest.raises(VectorStoreError, match="Qdrant collection configuration mismatch"):
|
|
store.upsert("memory", [_write_record("memory:1", "memory")])
|
|
|
|
creates = [call for call in fake.calls if call[0] == "PUT" and call[1].endswith("/collections/workspace-semantic")]
|
|
assert creates == []
|
|
|
|
|
|
def test_upsert_require_existing_refuses_missing_collection_without_creating():
|
|
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 fake.collection is None
|
|
creates = [
|
|
call
|
|
for call in fake.calls
|
|
if call[0] == "PUT" and call[1].endswith("/collections/workspace-semantic")
|
|
]
|
|
assert creates == []
|
|
|
|
|
|
def test_upsert_require_existing_refuses_incompatible_collection_without_mutating():
|
|
fake = FakeQdrantHttp(dimension=384, distance="Dot")
|
|
fake.collection = {"vectors": {"size": 384, "distance": "Dot"}}
|
|
store = _store(fake, collection_lifecycle="require_existing")
|
|
|
|
with pytest.raises(VectorStoreError, match="semantic_index_incompatible"):
|
|
store.upsert("memory", [_write_record("memory:1", "memory")])
|
|
|
|
assert fake.payload_indexes == set()
|
|
mutating = [call for call in fake.calls if call[0] == "PUT"]
|
|
assert mutating == []
|
|
|
|
|
|
def test_upsert_require_existing_writes_to_existing_compatible_collection():
|
|
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
|
|
|
|
create_or_index = [
|
|
call
|
|
for call in fake.calls
|
|
if call[0] == "PUT" and not call[1].endswith("/points?wait=true")
|
|
]
|
|
assert create_or_index == []
|
|
|
|
|
|
def test_upsert_require_existing_fails_if_collection_disappears_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",
|
|
}
|
|
fake.drop_collection_on_points = True
|
|
store = _store(fake, collection_lifecycle="require_existing")
|
|
|
|
with pytest.raises(VectorStoreError, match="HTTP 404"):
|
|
store.upsert("memory", [_write_record("memory:1", "memory")])
|
|
|
|
creates = [
|
|
call
|
|
for call in fake.calls
|
|
if call[0] == "PUT" and call[1].endswith("/collections/workspace-semantic")
|
|
]
|
|
assert creates == []
|
|
|
|
|
|
def test_health_fails_when_the_bound_collection_is_missing():
|
|
fake = FakeQdrantHttp()
|
|
|
|
health = _store(fake).health()
|
|
|
|
assert health.ok is False
|
|
assert health.read_reachable is False
|
|
assert health.write_reachable is False
|
|
assert "missing" in (health.detail or "").lower()
|
|
|
|
|
|
def test_health_fails_when_required_payload_indexes_are_missing_without_creating_them():
|
|
fake = FakeQdrantHttp()
|
|
fake.collection = {"vectors": {"size": 1024, "distance": "Cosine"}}
|
|
|
|
health = _store(fake).health()
|
|
|
|
assert health.ok is False
|
|
assert fake.payload_indexes == set()
|
|
|
|
|
|
def test_health_reports_when_evidence_bm25_is_unavailable_without_disabling_dense_callers():
|
|
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",
|
|
}
|
|
|
|
health = _store(fake).health()
|
|
|
|
assert health.ok is True
|
|
assert health.bm25_compatible is False
|
|
|
|
|
|
@pytest.mark.parametrize(
|
|
("record", "semantic_kind"),
|
|
[
|
|
(_write_record("schema_column:patients.id", "schema_column"), "schema"),
|
|
(
|
|
_write_record(
|
|
"demo:gen:11111111111111111111111111111111:chunk:1",
|
|
"evidence",
|
|
metadata={
|
|
"workspace_id": "demo",
|
|
"vector_generation": "gen:11111111111111111111111111111111",
|
|
"document_id": "doc:abc",
|
|
},
|
|
),
|
|
"evidence",
|
|
),
|
|
(_write_record("memory:1", "memory"), "memory"),
|
|
],
|
|
)
|
|
def test_upsert_serializes_qdrant_point_payloads(record, semantic_kind):
|
|
fake = FakeQdrantHttp()
|
|
store = _store(fake)
|
|
|
|
store.upsert("memory" if semantic_kind == "memory" else "evidence" if semantic_kind == "evidence" else "schema_records", [record])
|
|
|
|
point = next(iter(fake.points.values()))
|
|
expected_revision = "a" * 40 if semantic_kind in ("schema_table", "schema_column", "evidence") else None
|
|
assert point["id"] == point_id("demo", semantic_kind, record.record.id, expected_revision)
|
|
assert point["vector"] == record.embedding
|
|
assert point["payload"]["workspace_id"] == "demo"
|
|
assert point["payload"]["workspace_revision"] == "a" * 40
|
|
assert point["payload"]["kind"] == semantic_kind
|
|
assert point["payload"]["record_kind"] == record.record.kind
|
|
assert point["payload"]["record_key"] == record.record.id
|
|
assert point["payload"]["content_hash"] == record.content_hash
|
|
|
|
|
|
def test_evidence_upsert_sends_dense_and_server_side_italian_bm25():
|
|
fake = FakeQdrantHttp()
|
|
_ready_collection_with_bm25(fake)
|
|
store = _store(fake)
|
|
record = _write_record(
|
|
"demo:gen:11111111111111111111111111111111:chunk:1",
|
|
"evidence",
|
|
metadata={
|
|
"workspace_id": "demo",
|
|
"vector_generation": "gen:11111111111111111111111111111111",
|
|
"document_id": "doc:abc",
|
|
},
|
|
sparse_text="ricovero per cardiomiopatia dilatativa",
|
|
sparse_language="italian",
|
|
)
|
|
|
|
store.upsert("evidence", [record])
|
|
|
|
point = next(iter(fake.points.values()))
|
|
assert point["vector"] == {
|
|
"": record.embedding,
|
|
"bm25": {
|
|
"text": "ricovero per cardiomiopatia dilatativa",
|
|
"model": "qdrant/bm25",
|
|
"options": {"language": "italian"},
|
|
},
|
|
}
|
|
|
|
|
|
def test_evidence_search_uses_filtered_dense_and_bm25_prefetches_with_default_rrf():
|
|
fake = FakeQdrantHttp()
|
|
_ready_collection_with_bm25(fake)
|
|
store = _store(fake)
|
|
generation = "gen:" + "1" * 32
|
|
store.upsert("evidence", [
|
|
_write_record(
|
|
f"demo:{generation}:chunk:1",
|
|
"evidence",
|
|
metadata={"workspace_id": "demo", "vector_generation": generation, "document_id": "doc:abc"},
|
|
sparse_text="ricovero per cardiomiopatia dilatativa",
|
|
sparse_language="italian",
|
|
)
|
|
])
|
|
|
|
store.search(
|
|
["evidence"], [0.2] * 1024, limit=10, kinds=["evidence"],
|
|
query_text="cardiomiopatia", query_language="italian",
|
|
metadata_filter={"workspace_id": "demo", "vector_generation": generation, "document_ids": ["doc:abc"]},
|
|
)
|
|
|
|
query = next(call[2] for call in reversed(fake.calls) if call[1].endswith("/points/query"))
|
|
assert query["query"] == {"rrf": {}}
|
|
assert query["limit"] == 10
|
|
assert query["prefetch"] == [
|
|
{
|
|
"query": [0.2] * 1024,
|
|
"limit": 20,
|
|
"filter": {"must": [
|
|
{"key": "workspace_id", "match": {"value": "demo"}},
|
|
{"key": "workspace_revision", "match": {"value": "a" * 40}},
|
|
{"key": "kind", "match": {"any": ["evidence"]}},
|
|
{"key": "record_kind", "match": {"any": ["evidence"]}},
|
|
{"key": "vector_generation", "match": {"value": generation}},
|
|
{"key": "document_id", "match": {"any": ["doc:abc"]}},
|
|
]},
|
|
},
|
|
{
|
|
"query": {
|
|
"text": "cardiomiopatia", "model": "qdrant/bm25",
|
|
"options": {"language": "italian"},
|
|
},
|
|
"using": "bm25",
|
|
"limit": 20,
|
|
"filter": {"must": [
|
|
{"key": "workspace_id", "match": {"value": "demo"}},
|
|
{"key": "workspace_revision", "match": {"value": "a" * 40}},
|
|
{"key": "kind", "match": {"any": ["evidence"]}},
|
|
{"key": "record_kind", "match": {"any": ["evidence"]}},
|
|
{"key": "vector_generation", "match": {"value": generation}},
|
|
{"key": "document_id", "match": {"any": ["doc:abc"]}},
|
|
]},
|
|
},
|
|
]
|
|
|
|
|
|
@pytest.mark.parametrize("retrieval_mode, expected_query", [
|
|
("dense", None),
|
|
("bm25", {"text": "cardiomiopatia", "model": "qdrant/bm25", "options": {"language": "italian"}}),
|
|
])
|
|
def test_evidence_diagnostic_branch_searches_use_the_runtime_filter(retrieval_mode, expected_query):
|
|
fake = FakeQdrantHttp()
|
|
_ready_collection_with_bm25(fake)
|
|
store = _store(fake)
|
|
generation = "gen:" + "1" * 32
|
|
|
|
store.search(
|
|
["evidence"], [0.2] * 1024, limit=10, kinds=["evidence"],
|
|
query_text="cardiomiopatia", query_language="italian", retrieval_mode=retrieval_mode,
|
|
metadata_filter={"workspace_id": "demo", "vector_generation": generation, "document_ids": ["doc:abc"]},
|
|
)
|
|
|
|
query = next(call[2] for call in reversed(fake.calls) if call[1].endswith("/points/query"))
|
|
assert query["limit"] == 10
|
|
assert query["filter"]["must"][-2:] == [
|
|
{"key": "vector_generation", "match": {"value": generation}},
|
|
{"key": "document_id", "match": {"any": ["doc:abc"]}},
|
|
]
|
|
if retrieval_mode == "dense":
|
|
assert query["vector"] == [0.2] * 1024
|
|
else:
|
|
assert query["query"] == expected_query
|
|
assert query["using"] == "bm25"
|
|
|
|
|
|
def test_search_filters_by_workspace_and_allowed_record_kinds():
|
|
fake = FakeQdrantHttp()
|
|
store = _store(fake)
|
|
store.upsert("memory", [_write_record("memory:1", "memory")])
|
|
other = next(iter(fake.points.values())).copy()
|
|
other["id"] = point_id("other", "memory", "memory:2")
|
|
other["payload"] = {**other["payload"], "workspace_id": "other", "record_key": "memory:2"}
|
|
fake.points[other["id"]] = other
|
|
solved = next(iter(fake.points.values())).copy()
|
|
solved["id"] = point_id("demo", "memory", "solved:1")
|
|
solved["payload"] = {**solved["payload"], "record_key": "solved:1", "record_kind": "solved_question"}
|
|
fake.points[solved["id"]] = solved
|
|
|
|
hits = store.search(["memory"], [0.2] * 1024, limit=5, kinds=["memory"])
|
|
|
|
assert [hit.id for hit in hits] == ["memory:1"]
|
|
query_call = next(call for call in fake.calls if call[0] == "POST" and call[1].endswith("/points/query?wait=true") is False and call[1].endswith("/points/query"))
|
|
assert query_call[2]["filter"] == {
|
|
"must": [
|
|
{"key": "workspace_id", "match": {"value": "demo"}},
|
|
{"key": "kind", "match": {"any": ["memory"]}},
|
|
{"key": "record_kind", "match": {"any": ["memory"]}},
|
|
]
|
|
}
|
|
|
|
|
|
@pytest.mark.parametrize("kind,family", [("memory", "domain_clarification"),
|
|
("solved_question", "solved_question")])
|
|
def test_memory_hybrid_applies_identical_scope_to_both_prefetch_branches(kind, family):
|
|
from tht.memory.retrieval import RecallScope
|
|
|
|
fake = FakeQdrantHttp()
|
|
_ready_collection_with_bm25(fake)
|
|
store = _store(fake)
|
|
scope = RecallScope(database="dwh", schema_name="sales", table="orders", column="id",
|
|
scope="Sales", concepts=["grain"])
|
|
store.search(["memory"], [0.2] * 1024, limit=5, kinds=[kind],
|
|
query_text="Order grain", query_language="english",
|
|
metadata_filter=scope.vector_filter(family))
|
|
query = next(call[2] for call in reversed(fake.calls) if call[1].endswith("/points/query"))
|
|
dense, lexical = query["prefetch"]
|
|
assert dense["filter"] == lexical["filter"]
|
|
must = dense["filter"]["must"]
|
|
assert {"key": "workspace_id", "match": {"value": "demo"}} in must
|
|
assert {"key": "memory_family", "match": {"value": family}} in must
|
|
assert {"key": "memory_format", "match": {"value": 2}} in must
|
|
assert {"key": "memory_scope", "match": {"value": "Sales"}} in must
|
|
assert {"key": "memory_concepts", "match": {"value": "grain"}} in must
|
|
nested = must[-1]["should"][1]["nested"]
|
|
assert nested["key"] == "memory_dependencies"
|
|
assert nested["filter"]["must"] == [
|
|
{"key": "database", "match": {"value": "dwh"}},
|
|
{"key": "schema_name", "match": {"any": ["", "sales"]}},
|
|
{"key": "table", "match": {"any": ["", "orders"]}},
|
|
{"key": "column", "match": {"any": ["", "id"]}},
|
|
]
|
|
assert lexical["using"] == "bm25"
|
|
assert query["query"] == {"rrf": {}}
|
|
|
|
|
|
def test_evidence_search_refuses_dense_only_fallback():
|
|
store = _store(FakeQdrantHttp())
|
|
|
|
with pytest.raises(VectorStoreError, match="hybrid query text"):
|
|
store.search(["evidence"], [0.2] * 1024, limit=5, kinds=["evidence"])
|
|
|
|
|
|
def test_hybrid_evidence_search_refuses_a_collection_without_bm25():
|
|
fake = FakeQdrantHttp()
|
|
_ready_collection_with_bm25(fake)
|
|
fake.sparse_vectors = None
|
|
store = _store(fake)
|
|
|
|
with pytest.raises(VectorStoreError, match="BM25 collection configuration"):
|
|
store.search(
|
|
["evidence"], [0.2] * 1024, limit=5, kinds=["evidence"],
|
|
query_text="cardiomiopatia", query_language="italian",
|
|
)
|
|
|
|
|
|
def test_search_excludes_inconsistent_semantic_kind_in_bound_workspace():
|
|
fake = FakeQdrantHttp()
|
|
store = _store(fake)
|
|
store.upsert("memory", [_write_record("memory:1", "memory")])
|
|
contaminated = next(iter(fake.points.values())).copy()
|
|
contaminated["id"] = point_id("demo", "evidence", "memory:contaminated")
|
|
contaminated["payload"] = {
|
|
**contaminated["payload"],
|
|
"kind": "evidence",
|
|
"record_key": "memory:contaminated",
|
|
}
|
|
fake.points[contaminated["id"]] = contaminated
|
|
|
|
hits = store.search(["memory"], [0.2] * 1024, limit=5, kinds=["memory"])
|
|
|
|
assert [hit.id for hit in hits] == ["memory:1"]
|
|
|
|
|
|
def test_existing_hashes_health_and_exact_generation_inventory_and_delete():
|
|
fake = FakeQdrantHttp()
|
|
store = _store(fake)
|
|
generation = "gen:" + "1" * 32
|
|
keep = "gen:" + "2" * 32
|
|
store.upsert("evidence", [
|
|
_write_record(
|
|
f"demo:{generation}:chunk:1",
|
|
"evidence",
|
|
metadata={"workspace_id": "demo", "vector_generation": generation, "document_id": "doc:1"},
|
|
),
|
|
_write_record(
|
|
f"demo:{keep}:chunk:2",
|
|
"evidence",
|
|
metadata={"workspace_id": "demo", "vector_generation": keep, "document_id": "doc:2"},
|
|
),
|
|
])
|
|
|
|
assert store.existing_hashes("evidence", ["evidence"]) == {
|
|
f"demo:{generation}:chunk:1": "sha256:" + "a" * 64,
|
|
f"demo:{keep}:chunk:2": "sha256:" + "a" * 64,
|
|
}
|
|
assert store.list_evidence_generations("evidence", "demo") == [generation, keep]
|
|
assert store.delete_generation("evidence", generation, "demo") == 1
|
|
assert store.list_evidence_generations("evidence", "demo") == [keep]
|
|
|
|
health = store.health()
|
|
assert health.ok is True
|
|
assert health.read_reachable is True
|
|
assert health.write_reachable is True
|
|
assert health.observed_dimensions == (1024,)
|
|
assert health.dimension_compatible is True
|
|
|
|
|
|
def test_evidence_inventory_and_delete_ignore_inconsistent_semantic_kind():
|
|
fake = FakeQdrantHttp()
|
|
store = _store(fake)
|
|
generation = "gen:" + "1" * 32
|
|
contaminated_generation = "gen:" + "2" * 32
|
|
store.upsert("evidence", [
|
|
_write_record(
|
|
f"demo:{generation}:chunk:1",
|
|
"evidence",
|
|
metadata={
|
|
"workspace_id": "demo",
|
|
"vector_generation": generation,
|
|
"document_id": "doc:1",
|
|
},
|
|
),
|
|
])
|
|
contaminated_delete = next(iter(fake.points.values())).copy()
|
|
contaminated_delete["id"] = point_id("demo", "memory", "evidence:contaminated-delete")
|
|
contaminated_delete["payload"] = {
|
|
**contaminated_delete["payload"],
|
|
"kind": "memory",
|
|
"record_key": "evidence:contaminated-delete",
|
|
}
|
|
fake.points[contaminated_delete["id"]] = contaminated_delete
|
|
contaminated_list = next(iter(fake.points.values())).copy()
|
|
contaminated_list["id"] = point_id("demo", "memory", "evidence:contaminated-list")
|
|
contaminated_list["payload"] = {
|
|
**contaminated_list["payload"],
|
|
"kind": "memory",
|
|
"record_key": "evidence:contaminated-list",
|
|
"vector_generation": contaminated_generation,
|
|
}
|
|
fake.points[contaminated_list["id"]] = contaminated_list
|
|
|
|
assert store.list_evidence_generations("evidence", "demo") == [generation]
|
|
assert store.delete_generation("evidence", generation, "demo") == 1
|
|
assert contaminated_delete["id"] in fake.points
|
|
assert contaminated_list["id"] in fake.points
|
|
|
|
|
|
def test_metadata_search_rejects_a_workspace_id_different_from_the_bound_adapter():
|
|
fake = FakeQdrantHttp()
|
|
store = _store(fake)
|
|
generation = "gen:" + "1" * 32
|
|
store.upsert("evidence", [
|
|
_write_record(
|
|
f"demo:{generation}:chunk:1",
|
|
"evidence",
|
|
metadata={
|
|
"workspace_id": "demo", "vector_generation": generation,
|
|
"document_id": "doc:shared",
|
|
},
|
|
),
|
|
])
|
|
foreign = next(iter(fake.points.values())).copy()
|
|
foreign["id"] = point_id("other", "evidence", f"other:{generation}:chunk:1")
|
|
foreign["payload"] = {
|
|
**foreign["payload"],
|
|
"workspace_id": "other",
|
|
"record_key": f"other:{generation}:chunk:1",
|
|
"ref": "ref:foreign",
|
|
"title": "foreign",
|
|
"content": "foreign",
|
|
}
|
|
fake.points[foreign["id"]] = foreign
|
|
|
|
with pytest.raises(VectorStoreError, match="workspace namespace does not match"):
|
|
store.search(
|
|
["evidence"], [0.2] * 1024, limit=5, kinds=["evidence"],
|
|
metadata_filter={
|
|
"workspace_id": "other",
|
|
"vector_generation": generation,
|
|
"document_ids": ["doc:shared"],
|
|
},
|
|
)
|
|
|
|
|
|
def test_generation_inventory_rejects_a_workspace_id_different_from_the_bound_adapter():
|
|
fake = FakeQdrantHttp()
|
|
store = _store(fake)
|
|
|
|
with pytest.raises(VectorStoreError, match="workspace namespace does not match"):
|
|
store.list_evidence_generations("evidence", "other")
|
|
assert not any(call[1].endswith("/points/scroll") for call in fake.calls)
|
|
|
|
|
|
def test_generation_delete_cannot_mutate_foreign_workspace_or_non_evidence_points():
|
|
fake = FakeQdrantHttp()
|
|
store = _store(fake)
|
|
generation = "gen:" + "1" * 32
|
|
store.upsert("evidence", [
|
|
_write_record(
|
|
f"demo:{generation}:chunk:1",
|
|
"evidence",
|
|
metadata={
|
|
"workspace_id": "demo", "vector_generation": generation,
|
|
"document_id": "doc:demo",
|
|
},
|
|
),
|
|
])
|
|
demo = next(iter(fake.points.values()))
|
|
foreign = demo.copy()
|
|
foreign["id"] = point_id("other", "evidence", f"other:{generation}:chunk:1")
|
|
foreign["payload"] = {
|
|
**demo["payload"], "workspace_id": "other",
|
|
"record_key": f"other:{generation}:chunk:1",
|
|
}
|
|
fake.points[foreign["id"]] = foreign
|
|
memory = demo.copy()
|
|
memory["id"] = point_id("other", "memory", "memory:foreign")
|
|
memory["payload"] = {
|
|
**demo["payload"], "workspace_id": "other", "kind": "memory",
|
|
"record_kind": "memory", "record_key": "memory:foreign",
|
|
}
|
|
fake.points[memory["id"]] = memory
|
|
before = set(fake.points)
|
|
|
|
with pytest.raises(VectorStoreError, match="workspace namespace does not match"):
|
|
store.delete_generation("evidence", generation, "other")
|
|
|
|
assert set(fake.points) == before
|
|
assert not any(call[1].endswith("/points/delete?wait=true") for call in fake.calls)
|
|
|
|
|
|
def test_delete_kinds_is_workspace_scoped_and_preserves_other_semantic_kinds():
|
|
fake = FakeQdrantHttp()
|
|
store = _store(fake)
|
|
|
|
store.upsert("memory", [_write_record("memory:1", "memory")])
|
|
store.upsert("memory", [_write_record("solved:1", "solved_question")])
|
|
store.upsert("schema_records", [_write_record("schema_table:patients", "schema_table")])
|
|
|
|
other_workspace_memory = next(
|
|
point for point in fake.points.values() if point["payload"]["record_key"] == "memory:1"
|
|
).copy()
|
|
other_workspace_memory["id"] = point_id("other", "memory", "memory:other")
|
|
other_workspace_memory["payload"] = {
|
|
**other_workspace_memory["payload"],
|
|
"workspace_id": "other",
|
|
"record_key": "memory:other",
|
|
"title": "title:memory:other",
|
|
"content": "content:memory:other",
|
|
"ref": "ref:memory:other",
|
|
}
|
|
fake.points[other_workspace_memory["id"]] = other_workspace_memory
|
|
|
|
assert store.delete_kinds("memory", ["memory"]) == 1
|
|
|
|
delete_call = next(
|
|
call
|
|
for call in fake.calls
|
|
if call[0] == "POST" and call[1].endswith("/points/delete?wait=true")
|
|
)
|
|
assert delete_call[2]["filter"] == {
|
|
"must": [
|
|
{"key": "workspace_id", "match": {"value": "demo"}},
|
|
{"key": "kind", "match": {"any": ["memory"]}},
|
|
{"key": "record_kind", "match": {"any": ["memory"]}},
|
|
]
|
|
}
|
|
assert {
|
|
point["payload"]["record_key"]: point["payload"]["record_kind"]
|
|
for point in fake.points.values()
|
|
} == {
|
|
"solved:1": "solved_question",
|
|
"schema_table:patients": "schema_table",
|
|
"memory:other": "memory",
|
|
}
|
|
|
|
|
|
def test_existing_hashes_and_delete_kinds_ignore_inconsistent_semantic_kind():
|
|
fake = FakeQdrantHttp()
|
|
store = _store(fake)
|
|
store.upsert("memory", [_write_record("memory:1", "memory")])
|
|
contaminated = next(iter(fake.points.values())).copy()
|
|
contaminated["id"] = point_id("demo", "evidence", "memory:contaminated")
|
|
contaminated["payload"] = {
|
|
**contaminated["payload"],
|
|
"kind": "evidence",
|
|
"record_key": "memory:contaminated",
|
|
}
|
|
fake.points[contaminated["id"]] = contaminated
|
|
|
|
assert store.existing_hashes("memory", ["memory"]) == {
|
|
"memory:1": "sha256:" + "a" * 64,
|
|
}
|
|
assert store.delete_kinds("memory", ["memory"]) == 1
|
|
assert contaminated["id"] in fake.points
|
|
|
|
|
|
def test_sanitizes_timeout_and_malformed_responses():
|
|
fake = FakeQdrantHttp()
|
|
store = _store(fake)
|
|
fake.fail_request = requests.Timeout("dial tcp 10.0.0.9:6333: i/o timeout")
|
|
|
|
with pytest.raises(VectorStoreError, match="Qdrant request failed") as timeout:
|
|
store.search(["memory"], [0.2] * 1024, limit=1)
|
|
assert "10.0.0.9" not in str(timeout.value)
|
|
|
|
fake.fail_request = None
|
|
store.upsert("memory", [_write_record("memory:1", "memory")])
|
|
fake.malformed_query = True
|
|
with pytest.raises(VectorStoreError, match="Qdrant returned malformed query response"):
|
|
store.search(["memory"], [0.2] * 1024, limit=1)
|
|
|
|
fake.malformed_query = False
|
|
fake.malformed_scroll = True
|
|
with pytest.raises(VectorStoreError, match="Qdrant returned malformed scroll response"):
|
|
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"]
|