Files
ratatoskr/tests/test_provider_memory.py
T
vh cd12951aca feat(provider): memory plane — SQLite+sqlite-vec store + dev shell
The second plane of ratatoskr's Tier-3 Bifrost consumer: a durable memory
store Worldtree writes agent memory chunks into (upsert_many) and recalls
by vector similarity (search), with point reads + deletes. Implements
bifrost's own MemoryDataStore Protocol; conformance is #195 parity vs
InMemoryMemoryStore through the real dispatch_memory_call.

Store (memory_store.py): open_memory_store, describe_store, upsert_many
(replay/conflict idempotency, optimistic locking, injection rule, atomic
batch), search (cosine over sqlite-vec vec0, scope isolation INV-005,
over-fetch-then-filter so top_k counts in-scope), get/get_many,
delete_many, build_memory_provider_app. Dev shell (serve_memory.py):
ratatoskr-memory-provider entrypoint, port 8391.

TDD + heid-code-review (panel Groa/Hulda/Regin, zero true drift). Adopted
fixups: scope_filter dict guard, top_k<=0 -> [], stronger scope-isolation
+ delete-hit-search + handshake-POST tests. Partial-map optimistic-lock
semantics pinned against the reference via a new expected_revisions
parity test.

26 memory + 4 serve tests; #195 parity (upsert/search/expected_revisions)
green; ruff clean. Deps: +sqlite-vec.
2026-06-15 21:39:42 -07:00

415 lines
17 KiB
Python

"""Tests for the Tier-3 Bifrost memory provider (ratatoskr.provider.memory_store).
Contract: docs/contracts/bifrost_memory_provider.contract.md (v1.1)
Vertical tracer-first: fresh_db -> basic_upsert (round-trip) -> replay -> conflict
-> optimistic_lock -> injection_rule -> search/scope_isolation -> get/get_many ->
delete_many -> build_memory_provider_app -> #195 parity vs InMemoryMemoryStore.
"""
from __future__ import annotations
import types
import pytest
from bifrost.memory import IdempotencyConflict, InvalidArguments, RevisionMismatch
from ratatoskr.provider.memory_store import (
build_memory_provider_app,
open_memory_store,
)
EMBEDDING_DIM = 8
def _ctx(sub: str = "sub-1"):
# Mirrors bifrost reference _ctx_actor: actor = job_id | jwt_sub | session_id.
return types.SimpleNamespace(jwt_sub=sub)
def _vec(*head: float) -> list[float]:
v = list(head) + [0.0] * EMBEDDING_DIM
return v[:EMBEDDING_DIM]
def _chunk(cid: str = "c1", *, embedding=None, scope=None, **extra) -> dict:
rec = {
"id": cid,
"embedding": embedding if embedding is not None else _vec(1.0),
"scope": scope if scope is not None else {"end_user": "u1"},
"origin": "worldtree",
"distillate": {"summary": f"distillate-{cid}"},
"content": f"content-{cid}",
}
rec.update(extra)
return rec
def _row_count(store, table: str) -> int:
return store._conn.execute(f"SELECT COUNT(*) FROM {table}").fetchone()[0]
# --- open_memory_store ---
def test_fresh_db_advertises_v1_caps_and_schema():
store = open_memory_store(":memory:", embedding_dim=EMBEDDING_DIM)
caps = store.describe_store()
assert caps["relational_edges_supported"] is False
assert caps["optimistic_locking_supported"] is True
assert caps["atomic_supersede_supported"] is False
assert caps["transaction_supported"] is False
assert caps["filterable_metadata_fields"] == []
# tables + vec index queryable
store._conn.execute("SELECT * FROM memory_chunks")
store._conn.execute("SELECT * FROM memory_idempotency")
store._conn.execute("SELECT * FROM memory_vec")
def test_reopen_existing_file_is_idempotent(tmp_path):
db = str(tmp_path / "memory.db")
open_memory_store(db, embedding_dim=EMBEDDING_DIM) # first open creates schema
store = open_memory_store(db, embedding_dim=EMBEDDING_DIM) # reopen: IF NOT EXISTS no-op
assert isinstance(store.describe_store(), dict)
store._conn.execute("SELECT * FROM memory_chunks")
store._conn.execute("SELECT * FROM memory_vec")
# --- upsert_many + get (tracer round-trip) ---
async def test_basic_upsert_round_trips_verbatim_with_revision():
store = open_memory_store(":memory:", embedding_dim=EMBEDDING_DIM)
c1 = _chunk("c1", embedding=_vec(1.0))
c2 = _chunk("c2", embedding=_vec(0.0, 1.0))
result = await store.upsert_many([c1, c2], idempotency_key="k1", ctx=_ctx())
assert result == {"upserted": 2, "replayed": False}
# INV-001: each chunk round-trips verbatim, with a revision key attached (first insert -> 1)
assert await store.get("c1") == {**c1, "revision": 1}
assert await store.get("c2") == {**c2, "revision": 1}
async def test_replay_same_key_same_payload_no_rewrite():
# INV-002: same idempotency_key + same digest -> replay (no second write, revision frozen)
store = open_memory_store(":memory:", embedding_dim=EMBEDDING_DIM)
c1 = _chunk("c1")
assert await store.upsert_many([c1], idempotency_key="k1", ctx=_ctx()) == {
"upserted": 1,
"replayed": False,
}
assert await store.upsert_many([c1], idempotency_key="k1", ctx=_ctx()) == {
"upserted": 1,
"replayed": True,
}
assert (await store.get("c1"))["revision"] == 1 # replay did not re-write / re-increment
assert _row_count(store, "memory_chunks") == 1
async def test_conflict_same_key_different_payload_raises_and_keeps_first():
# INV-002: same key, different digest -> IdempotencyConflict; the first batch is intact
store = open_memory_store(":memory:", embedding_dim=EMBEDDING_DIM)
first = _chunk("c1", content="first")
await store.upsert_many([first], idempotency_key="k1", ctx=_ctx())
with pytest.raises(IdempotencyConflict):
await store.upsert_many(
[_chunk("c1", content="second")], idempotency_key="k1", ctx=_ctx()
)
assert await store.get("c1") == {**first, "revision": 1} # untouched
async def test_optimistic_lock_stale_expected_revision_raises_nothing_written():
# INV-003: a stale expected_revisions entry rolls back the whole batch
store = open_memory_store(":memory:", embedding_dim=EMBEDDING_DIM)
c1 = _chunk("c1", content="v1")
await store.upsert_many([c1], idempotency_key="k1", ctx=_ctx()) # revision 1
with pytest.raises(RevisionMismatch):
await store.upsert_many(
[_chunk("c1", content="v2")],
idempotency_key="k2", # distinct key: not replay/conflict
ctx=_ctx(),
expected_revisions={"c1": 5}, # stale: stored revision is 1
)
assert await store.get("c1") == {**c1, "revision": 1} # nothing written
assert _row_count(store, "memory_chunks") == 1
async def test_optimistic_lock_match_upserts_and_increments_revision():
# INV-003: a matching expected_revisions writes and increments (1 -> 2)
store = open_memory_store(":memory:", embedding_dim=EMBEDDING_DIM)
await store.upsert_many([_chunk("c1", content="v1")], idempotency_key="k1", ctx=_ctx())
v2 = _chunk("c1", content="v2")
assert await store.upsert_many(
[v2], idempotency_key="k2", ctx=_ctx(), expected_revisions={"c1": 1}
) == {"upserted": 1, "replayed": False}
assert await store.get("c1") == {**v2, "revision": 2} # re-upsert increments
assert _row_count(store, "memory_chunks") == 1
async def test_injection_rule_injected_without_source_raises_no_write():
# INV-007: origin == injected_context requires injection_source
store = open_memory_store(":memory:", embedding_dim=EMBEDDING_DIM)
bad = _chunk("c1", origin="injected_context") # no injection_source
with pytest.raises(InvalidArguments):
await store.upsert_many([bad], idempotency_key="k1", ctx=_ctx())
assert _row_count(store, "memory_chunks") == 0
async def test_injection_rule_non_injected_with_source_raises_no_write():
# INV-007: a non-injected record carrying injection_source is rejected
store = open_memory_store(":memory:", embedding_dim=EMBEDDING_DIM)
bad = _chunk("c1", origin="worldtree", injection_source="elsewhere")
with pytest.raises(InvalidArguments):
await store.upsert_many([bad], idempotency_key="k1", ctx=_ctx())
assert _row_count(store, "memory_chunks") == 0
# --- search ---
async def test_basic_search_ranks_by_cosine_with_recalled_view():
store = open_memory_store(":memory:", embedding_dim=EMBEDDING_DIM)
scope = {"end_user": "u1"}
c1 = _chunk("c1", embedding=_vec(1.0, 0.0), scope=scope)
c3 = _chunk("c3", embedding=_vec(0.9, 0.1), scope=scope)
await store.upsert_many(
[c1, _chunk("c2", embedding=_vec(0.0, 1.0), scope=scope), c3],
idempotency_key="k1",
ctx=_ctx(),
)
results = await store.search(_vec(1.0, 0.0), top_k=2, scope_filter=scope)
assert [r["chunk_id"] for r in results] == ["c1", "c3"] # nearest to [1,0] by cosine
top = results[0]
assert top["chunk"] == c1 # verbatim chunk, no revision attached
assert top["recalled_view"] == {"summary": "distillate-c1"} # = chunk["distillate"]
assert top["revision"] == 1
assert isinstance(top["score"], float)
async def test_scope_isolation_excludes_other_scope_even_if_closer():
# INV-005: an out-of-scope chunk that scores HIGHER must not leak; only in-scope returned
store = open_memory_store(":memory:", embedding_dim=EMBEDDING_DIM)
await store.upsert_many(
[
_chunk("u2-near", embedding=_vec(1.0, 0.0), scope={"end_user": "u2"}), # closest
_chunk("u1-far", embedding=_vec(0.0, 1.0), scope={"end_user": "u1"}), # in-scope, far
],
idempotency_key="k1",
ctx=_ctx(),
)
results = await store.search(_vec(1.0, 0.0), top_k=2, scope_filter={"end_user": "u1"})
assert [r["chunk_id"] for r in results] == ["u1-far"] # u2-near excluded despite ranking first
async def test_search_empty_store_returns_empty():
store = open_memory_store(":memory:", embedding_dim=EMBEDDING_DIM)
assert await store.search(_vec(1.0), top_k=5) == []
async def test_search_non_empty_metadata_filter_rejected():
# PRE-002: v1 advertises no filterable metadata fields
store = open_memory_store(":memory:", embedding_dim=EMBEDDING_DIM)
with pytest.raises(InvalidArguments):
await store.search(_vec(1.0), top_k=5, metadata_filter={"x": 1})
async def test_search_wrong_vector_dim_rejected():
# PRE-001: vector length must equal the pinned embedding_dim
store = open_memory_store(":memory:", embedding_dim=EMBEDDING_DIM)
with pytest.raises(InvalidArguments):
await store.search([1.0, 0.0], top_k=5)
async def test_search_non_dict_scope_filter_rejected():
# search STEP 1: scope_filter must be a flat {axis: value} dict
store = open_memory_store(":memory:", embedding_dim=EMBEDDING_DIM)
with pytest.raises(InvalidArguments):
await store.search(_vec(1.0), top_k=5, scope_filter="u1")
async def test_search_top_k_zero_returns_empty():
# POST-001: at most top_k — zero means zero
store = open_memory_store(":memory:", embedding_dim=EMBEDDING_DIM)
await store.upsert_many([_chunk("c1")], idempotency_key="k1", ctx=_ctx())
assert await store.search(_vec(1.0), top_k=0, scope_filter={"end_user": "u1"}) == []
async def test_scope_isolation_fills_top_k_from_in_scope_past_higher_out_of_scope():
# INV-005: top_k counts IN-SCOPE hits. An out-of-scope chunk ranking #1 is skipped,
# and top_k is still filled from the in-scope set when enough in-scope chunks exist.
store = open_memory_store(":memory:", embedding_dim=EMBEDDING_DIM)
await store.upsert_many(
[
_chunk("u2-nearest", embedding=_vec(1.0, 0.0), scope={"end_user": "u2"}), # ranks #1
_chunk("u1-near", embedding=_vec(0.95, 0.05), scope={"end_user": "u1"}),
_chunk("u1-mid", embedding=_vec(0.8, 0.2), scope={"end_user": "u1"}),
_chunk("u1-far", embedding=_vec(0.0, 1.0), scope={"end_user": "u1"}),
],
idempotency_key="k1",
ctx=_ctx(),
)
results = await store.search(_vec(1.0, 0.0), top_k=2, scope_filter={"end_user": "u1"})
# exactly top_k in-scope (the 2 nearest u1 chunks); the higher-ranked u2 chunk is excluded
assert [r["chunk_id"] for r in results] == ["u1-near", "u1-mid"]
# --- get / get_many ---
async def test_get_absent_returns_none():
store = open_memory_store(":memory:", embedding_dim=EMBEDDING_DIM)
assert await store.get("nope") is None
async def test_get_many_returns_found_records_only():
store = open_memory_store(":memory:", embedding_dim=EMBEDDING_DIM)
c1 = _chunk("c1")
await store.upsert_many([c1], idempotency_key="k1", ctx=_ctx())
assert await store.get_many(["c1", "absent"]) == [{**c1, "revision": 1}]
# --- delete_many ---
async def test_delete_hit_removes_chunk_and_vec_row():
store = open_memory_store(":memory:", embedding_dim=EMBEDDING_DIM)
await store.upsert_many([_chunk("c1"), _chunk("c2")], idempotency_key="k1", ctx=_ctx())
assert await store.delete_many(["c1"]) == {"deleted": 1}
assert await store.get("c1") is None
assert _row_count(store, "memory_chunks") == 1
assert _row_count(store, "memory_vec") == 1 # c1's vec row gone too (no orphan)
# delete_hit: search no longer surfaces it (vec/chunk coupling held)
hits = await store.search(_vec(1.0), top_k=5, scope_filter={"end_user": "u1"})
assert all(r["chunk_id"] != "c1" for r in hits)
async def test_delete_absent_counts_zero():
store = open_memory_store(":memory:", embedding_dim=EMBEDDING_DIM)
assert await store.delete_many(["nope"]) == {"deleted": 0}
# --- build_memory_provider_app ---
def test_build_app_exposes_handshake_and_memory_routes():
store = open_memory_store(":memory:", embedding_dim=EMBEDDING_DIM)
app = build_memory_provider_app(store, heimdall_key=b"secret-key")
routes = {getattr(r, "path", None): r for r in app.routes}
assert "/bifrost/handshake" in routes
assert "/bifrost/memory-call" in routes
assert "POST" in routes["/bifrost/memory-call"].methods
assert "POST" in routes["/bifrost/handshake"].methods # both routes are POST (incl. POST)
def test_build_app_rejects_empty_key():
store = open_memory_store(":memory:", embedding_dim=EMBEDDING_DIM)
with pytest.raises(ValueError):
build_memory_provider_app(store, heimdall_key=b"")
# --- #195 conformance: parity vs the reference store through the real engine ---
def _dispatch_ctx(*scopes: str, session_id: str = "actor-1"):
return types.SimpleNamespace(
scope=list(scopes), session_id=session_id, jwt_sub=session_id, job_id=None
)
def _ref_record(chunk_id: str, vector: list[float], *, end_user: str = "u1") -> dict:
# Mirrors bifrost's reference `record` helper so the envelope validates.
return {
"id": chunk_id,
"embedding": vector,
"distillate": {"text": chunk_id},
"metadata": {"worldtree.appraisal_confidence": 0.8},
"scope": {"end_user": end_user, "tenant": "t1"},
"origin": "worldtree",
"source_role": "assistant",
"trust_tier": "tier-3",
"provenance": {"trace": chunk_id},
}
async def test_parity_upsert_many_vs_reference_through_dispatch():
from bifrost.consumer.testing import InMemoryMemoryStore
from bifrost.memory import dispatch_memory_call
ref = InMemoryMemoryStore()
mine = open_memory_store(":memory:", embedding_dim=2)
wctx = _dispatch_ctx("memory:write")
env = {
"operation": "upsert_many",
"args": {"records": [_ref_record("a", [1.0, 0.0]), _ref_record("b", [0.0, 1.0])]},
"idempotency_key": "k1",
}
# happy persist + replay: wire bodies must agree
assert await dispatch_memory_call(env, wctx, ref) == await dispatch_memory_call(env, wctx, mine)
assert await dispatch_memory_call(env, wctx, ref) == await dispatch_memory_call(env, wctx, mine)
async def test_parity_search_ranked_ids_vs_reference_through_dispatch():
from bifrost.consumer.testing import InMemoryMemoryStore
from bifrost.memory import dispatch_memory_call
ref = InMemoryMemoryStore()
mine = open_memory_store(":memory:", embedding_dim=2)
wctx = _dispatch_ctx("memory:write")
rctx = _dispatch_ctx("memory:read")
up = {
"operation": "upsert_many",
"args": {
"records": [
_ref_record("a", [1.0, 0.0]),
_ref_record("b", [0.0, 1.0]),
_ref_record("c", [0.9, 0.1]),
]
},
"idempotency_key": "k1",
}
await dispatch_memory_call(up, wctx, ref)
await dispatch_memory_call(up, wctx, mine)
search_env = {
"operation": "search",
"args": {"vector": [1.0, 0.0], "top_k": 2, "scope_filter": {"end_user": "u1"}},
}
rstatus, rbody = await dispatch_memory_call(search_env, rctx, ref)
mstatus, mbody = await dispatch_memory_call(search_env, rctx, mine)
assert rstatus == mstatus == 200
# #195: same ranked chunk_ids and the same per-result shape (scores may differ in the
# last float digit between vec0's cosine and the reference's Python cosine).
assert [r["chunk_id"] for r in rbody["results"]] == [r["chunk_id"] for r in mbody["results"]]
assert set(rbody["results"][0]) == set(mbody["results"][0])
async def test_parity_expected_revisions_vs_reference_through_dispatch():
# #195: pins the partial-map optimistic-lock semantics against the reference
# (does an expected_revisions map that omits some batch records lock only the
# listed ones?). Resolves the contract's ambiguous "each record's stored revision".
from bifrost.consumer.testing import InMemoryMemoryStore
from bifrost.memory import dispatch_memory_call
ref = InMemoryMemoryStore()
mine = open_memory_store(":memory:", embedding_dim=2)
wctx = _dispatch_ctx("memory:write")
seed = {
"operation": "upsert_many",
"args": {"records": [_ref_record("a", [1.0, 0.0]), _ref_record("b", [0.0, 1.0])]},
"idempotency_key": "seed",
}
assert await dispatch_memory_call(seed, wctx, ref) == await dispatch_memory_call(seed, wctx, mine)
# partial map: only "a" is locked (revision 1); "b" is omitted from expected_revisions
partial = {
"operation": "upsert_many",
"args": {
"records": [_ref_record("a", [1.0, 0.0]), _ref_record("b", [0.0, 1.0])],
"expected_revisions": {"a": 1},
},
"idempotency_key": "partial",
}
assert await dispatch_memory_call(partial, wctx, ref) == await dispatch_memory_call(
partial, wctx, mine
)
# stale lock: both map to the same RevisionMismatch wire error
stale = {
"operation": "upsert_many",
"args": {"records": [_ref_record("a", [1.0, 0.0])], "expected_revisions": {"a": 99}},
"idempotency_key": "stale",
}
assert await dispatch_memory_call(stale, wctx, ref) == await dispatch_memory_call(
stale, wctx, mine
)