feat(#18): composite Bifrost endpoint — build_combined_app (Deliverable 1)
One ASGI app fronting BOTH the memory.* and affect.* planes (:8392), so a single bound Worldtree session both remembers AND shows live PAD. Closes #18 end-to-end (D2 PAD read-endpoint shipped v0.17.14; D1 was bifrost-blocked, now unparked by bifrost 0.10.0's public build_combined_app + FR-1 resolved — zero Worldtree change). - provider/combined.py: build_combined_provider_app wraps bifrost.consumer.build_combined_app over both stores + mounts the shared affect read route. Advertises both caps by store presence; per-route call-time isolation is bifrost's (INV-013). - affect_store.py: extract add_affect_read_route shared helper (the D2 INV-007 promise — composite + standalone mount the SAME read route over the same affect.db, INV-011). - opfeed.py: plane='combined' derives the OpEvent plane per request path (memory-call->memory, affect-call->affect, handshake->combined; INV-012). - serve_combined.py + ratatoskr-combined-provider console script on :8392 (additive — standalone :8390/:8391 untouched, INV-014). - contract: 18.contract.md § Deliverable 1 (INV-009..INV-014); D1 un-deferred. Latent bug fixed (exposed by the contract-mandated memory `search` dispatch test running through TestClient = a worker thread): open_memory_store lacked check_same_thread=False — the SAME sqlite thread-safety bug already fixed in the affect store (D2). The composite serves the memory plane over HTTP, so a memory-call on uvicorn's threadpool would trip it. Fix: check_same_thread=False + PRAGMA busy_timeout=5000 (memory contract Concurrency note). heid-code-review panel (Groa/Hulda/Regin): ZERO drift findings; the implementation matches INV-009..INV-014 at function-block level. Folded the genuine test-fidelity fix (memory leg describe_store -> search per the contract TEST) + added the PRE-001/PRE-002 guard tests. Suite 486 -> 502 green.
This commit is contained in:
@@ -0,0 +1,274 @@
|
||||
"""Tests for the combined Bifrost provider (ratatoskr.provider.combined) — issue #18
|
||||
Deliverable 1.
|
||||
|
||||
ONE app fronting BOTH planes (memory.* + affect.*) + the shared affect read route.
|
||||
Mirrors bifrost's tests/consumer/test_build_combined_app.py shapes (handshake +
|
||||
dispatch) and ratatoskr's op-feed test style (mint_dispatch_jwt, RecordingSink), so
|
||||
the envelopes and JWTs are the real wire shapes, not hand-mocked guesses ("test
|
||||
against the shipped lib").
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import base64
|
||||
import hashlib
|
||||
import hmac
|
||||
import json
|
||||
import time
|
||||
|
||||
import httpx
|
||||
import pytest
|
||||
from bifrost.core.dispatch_jwt import mint_dispatch_jwt
|
||||
from starlette.testclient import TestClient
|
||||
|
||||
from ratatoskr.provider.affect_store import open_affect_store
|
||||
from ratatoskr.provider.combined import build_combined_provider_app
|
||||
from ratatoskr.provider.memory_store import open_memory_store
|
||||
from ratatoskr.provider.opfeed import instrument_provider_app
|
||||
|
||||
_KEY = b"deterministic-test-heimdall-key-32-bytes!"
|
||||
_CONSUMER = "ratatoskr"
|
||||
_DIM = 8
|
||||
|
||||
|
||||
def _combined_app():
|
||||
memory_store = open_memory_store(":memory:", embedding_dim=_DIM)
|
||||
affect_store = open_affect_store(":memory:")
|
||||
app = build_combined_provider_app(
|
||||
memory_store, affect_store, heimdall_key=_KEY, consumer_id=_CONSUMER
|
||||
)
|
||||
return app, memory_store, affect_store
|
||||
|
||||
|
||||
def _dispatch_headers(*scopes: str, session_id: str = "sess-1") -> dict:
|
||||
token = mint_dispatch_jwt(
|
||||
session_id=session_id,
|
||||
consumer_id=_CONSUMER,
|
||||
issuer="worldtree",
|
||||
scope=list(scopes),
|
||||
secret_or_key=_KEY,
|
||||
algorithm="HS256",
|
||||
)
|
||||
return {"Authorization": f"Bearer {token}"}
|
||||
|
||||
|
||||
def _b64url(data: bytes) -> str:
|
||||
return base64.urlsafe_b64encode(data).rstrip(b"=").decode("ascii")
|
||||
|
||||
|
||||
def _handshake_jwt(session_id: str = "sess-1") -> str:
|
||||
"""Replicate bifrost's consumer conftest jwt_factory (HS256 handshake JWT)."""
|
||||
header = {"alg": "HS256", "typ": "JWT"}
|
||||
now = time.time()
|
||||
payload = {
|
||||
"session_id": session_id,
|
||||
"consumer_id": _CONSUMER,
|
||||
"issued_at": now,
|
||||
"expires_at": now + 3600,
|
||||
}
|
||||
h = _b64url(json.dumps(header, separators=(",", ":")).encode())
|
||||
p = _b64url(json.dumps(payload, separators=(",", ":")).encode())
|
||||
sig = hmac.new(_KEY, f"{h}.{p}".encode("ascii"), hashlib.sha256).digest()
|
||||
return f"{h}.{p}.{_b64url(sig)}"
|
||||
|
||||
|
||||
def _handshake_body(session_id: str = "sess-1") -> dict:
|
||||
return {
|
||||
"bifrost_version": "0.4.0",
|
||||
"mcp_version": "0.4.0",
|
||||
"session_id": session_id,
|
||||
"consumer_id": _CONSUMER,
|
||||
"auth": {"scheme": "Bearer", "token": _handshake_jwt(session_id)},
|
||||
"capabilities": ["memory", "affect"],
|
||||
}
|
||||
|
||||
|
||||
def _snapshot(agent: str = "ratatoskr:sindra", user: str = "vuong") -> dict:
|
||||
return {
|
||||
"agent_id": agent,
|
||||
"end_user_id": user,
|
||||
"pad": {"pleasure": 0.5, "arousal": 0.2, "dominance": -0.1},
|
||||
"valence": [{"entity_id": "e1", "regard": 0.7, "familiarity": 0.3}],
|
||||
"emitted_at": "2026-06-14T12:00:00Z",
|
||||
}
|
||||
|
||||
|
||||
def _emit_envelope(snap: dict) -> dict:
|
||||
return {
|
||||
"operation": "affect.emit",
|
||||
"idempotency_key": "sess-1:1:affect",
|
||||
"idempotency_class": "short-retry",
|
||||
"args": snap,
|
||||
}
|
||||
|
||||
|
||||
# --- build_combined_provider_app ---
|
||||
|
||||
def test_builds_both_planes_and_read_route():
|
||||
"""builds_both_planes [tracer]: the composite exposes handshake + memory-call +
|
||||
affect-call + the non-bifrost /affect/state read route (INV-011)."""
|
||||
app, _m, _a = _combined_app()
|
||||
paths = {getattr(r, "path", None) for r in app.routes}
|
||||
assert "/bifrost/handshake" in paths
|
||||
assert "/bifrost/memory-call" in paths
|
||||
assert "/bifrost/affect-call" in paths
|
||||
assert "/affect/state/{agent_id}" in paths
|
||||
|
||||
|
||||
def test_handshake_grants_both_caps():
|
||||
"""handshake_grants_both [scenario]: a handshake requesting [memory, affect] is
|
||||
granted BOTH by store PRESENCE (INV-010) — my wiring doesn't break it."""
|
||||
app, _m, _a = _combined_app()
|
||||
resp = TestClient(app).post("/bifrost/handshake", json=_handshake_body())
|
||||
assert resp.status_code == 200
|
||||
granted = resp.json()["capabilities_granted"]
|
||||
assert "memory" in granted
|
||||
assert "affect" in granted
|
||||
|
||||
|
||||
def test_memory_and_affect_dispatch_through_one_app():
|
||||
"""memory_and_affect_dispatch [scenario]: a memory SEARCH AND an affect emit each
|
||||
round-trip through the SINGLE combined app (INV-013; contract TEST + Acceptance §2
|
||||
name a memory `search`)."""
|
||||
app, _m, _a = _combined_app()
|
||||
client = TestClient(app)
|
||||
|
||||
mem = client.post(
|
||||
"/bifrost/memory-call",
|
||||
json={
|
||||
"operation": "search",
|
||||
"args": {"vector": [0.0] * _DIM, "top_k": 1, "scope_all": {}},
|
||||
},
|
||||
headers=_dispatch_headers("memory:read"),
|
||||
)
|
||||
assert mem.status_code == 200
|
||||
assert mem.json()["success"] is True
|
||||
|
||||
aff = client.post(
|
||||
"/bifrost/affect-call",
|
||||
json=_emit_envelope(_snapshot()),
|
||||
headers=_dispatch_headers("affect:write"),
|
||||
)
|
||||
assert aff.status_code == 200
|
||||
assert aff.json()["success"] is True
|
||||
assert aff.json()["stored"] is True
|
||||
|
||||
|
||||
def test_affect_read_route_on_composite_colon_id():
|
||||
"""affect_read_route_on_composite [happy]: after an emit, GET /affect/state for a
|
||||
colon-id agent returns the snapshot verbatim from the SAME store (INV-011 / INV-008)."""
|
||||
app, _m, _a = _combined_app()
|
||||
client = TestClient(app)
|
||||
snap = _snapshot()
|
||||
client.post(
|
||||
"/bifrost/affect-call",
|
||||
json=_emit_envelope(snap),
|
||||
headers=_dispatch_headers("affect:write"),
|
||||
)
|
||||
r = client.get("/affect/state/ratatoskr:sindra", params={"end_user_id": "vuong"})
|
||||
assert r.status_code == 200
|
||||
assert r.json() == snap
|
||||
|
||||
|
||||
def test_missing_affect_store_raises():
|
||||
"""missing_affect_store [adversarial]: affect_store=None → ValueError (INV-009)."""
|
||||
memory_store = open_memory_store(":memory:", embedding_dim=_DIM)
|
||||
with pytest.raises(ValueError):
|
||||
build_combined_provider_app(memory_store, None, heimdall_key=_KEY)
|
||||
|
||||
|
||||
def test_missing_memory_store_raises():
|
||||
"""INV-009 (other half): memory_store=None → ValueError (bifrost build_combined_app)."""
|
||||
affect_store = open_affect_store(":memory:")
|
||||
with pytest.raises(ValueError):
|
||||
build_combined_provider_app(None, affect_store, heimdall_key=_KEY)
|
||||
|
||||
|
||||
def test_empty_heimdall_key_raises():
|
||||
"""PRE-002: empty heimdall_key → ValueError (combined-level guard)."""
|
||||
memory_store = open_memory_store(":memory:", embedding_dim=_DIM)
|
||||
affect_store = open_affect_store(":memory:")
|
||||
with pytest.raises(ValueError):
|
||||
build_combined_provider_app(memory_store, affect_store, heimdall_key=b"")
|
||||
|
||||
|
||||
def test_non_advertising_affect_store_raises():
|
||||
"""PRE-001 / INV-010: affect_store with affect_supported=False → ValueError."""
|
||||
memory_store = open_memory_store(":memory:", embedding_dim=_DIM)
|
||||
affect_store = open_affect_store(":memory:")
|
||||
affect_store.affect_supported = False
|
||||
with pytest.raises(ValueError):
|
||||
build_combined_provider_app(memory_store, affect_store, heimdall_key=_KEY)
|
||||
|
||||
|
||||
# --- op-feed plane='combined' (per-path derivation, INV-012) ---
|
||||
|
||||
class _RecordingSink:
|
||||
def __init__(self) -> None:
|
||||
self.events: list = []
|
||||
|
||||
def emit(self, event) -> None:
|
||||
self.events.append(event)
|
||||
|
||||
|
||||
async def _post(app, path: str, body: dict, headers: dict | None = None) -> httpx.Response:
|
||||
transport = httpx.ASGITransport(app=app)
|
||||
async with httpx.AsyncClient(transport=transport, base_url="http://provider") as client:
|
||||
return await client.post(path, json=body, headers=headers or {})
|
||||
|
||||
|
||||
async def test_opfeed_combined_memory_call_stamps_memory():
|
||||
sink = _RecordingSink()
|
||||
app, _m, _a = _combined_app()
|
||||
wrapped = instrument_provider_app(app, plane="combined", sink=sink)
|
||||
resp = await _post(
|
||||
wrapped,
|
||||
"/bifrost/memory-call",
|
||||
{"operation": "search", "args": {"vector": [0.0] * _DIM, "top_k": 1, "scope_all": {}}},
|
||||
_dispatch_headers("memory:read"),
|
||||
)
|
||||
assert resp.status_code == 200
|
||||
assert len(sink.events) == 1
|
||||
assert sink.events[0].plane == "memory" # derived from path (INV-012)
|
||||
assert sink.events[0].op == "search"
|
||||
|
||||
|
||||
async def test_opfeed_combined_affect_call_stamps_affect():
|
||||
sink = _RecordingSink()
|
||||
app, _m, _a = _combined_app()
|
||||
wrapped = instrument_provider_app(app, plane="combined", sink=sink)
|
||||
resp = await _post(
|
||||
wrapped,
|
||||
"/bifrost/affect-call",
|
||||
_emit_envelope(_snapshot()),
|
||||
_dispatch_headers("affect:write"),
|
||||
)
|
||||
assert resp.status_code == 200
|
||||
assert len(sink.events) == 1
|
||||
assert sink.events[0].plane == "affect" # derived from path (INV-012)
|
||||
assert sink.events[0].op == "emit" # affect. prefix stripped
|
||||
|
||||
|
||||
async def test_opfeed_combined_handshake_stamps_combined():
|
||||
"""handshake isn't plane-specific → stamp plane='combined' (INV-012). A bad-version
|
||||
handshake is cleanly rejected but still emits exactly one OpEvent."""
|
||||
sink = _RecordingSink()
|
||||
app, _m, _a = _combined_app()
|
||||
wrapped = instrument_provider_app(app, plane="combined", sink=sink)
|
||||
resp = await _post(
|
||||
wrapped, "/bifrost/handshake", {"bifrost_version": "99.0.0", "mcp_version": "0.4.0"}
|
||||
)
|
||||
assert resp.status_code != 200 # major-version mismatch, cleanly rejected
|
||||
assert len(sink.events) == 1
|
||||
assert sink.events[0].plane == "combined"
|
||||
assert sink.events[0].op == "handshake"
|
||||
|
||||
|
||||
async def test_opfeed_combined_read_route_emits_no_event():
|
||||
"""INV-012/INV-004: the non-bifrost read route is outside _BIFROST_PATHS → NO OpEvent."""
|
||||
sink = _RecordingSink()
|
||||
app, _m, _a = _combined_app()
|
||||
wrapped = instrument_provider_app(app, plane="combined", sink=sink)
|
||||
transport = httpx.ASGITransport(app=wrapped)
|
||||
async with httpx.AsyncClient(transport=transport, base_url="http://provider") as client:
|
||||
await client.get("/affect/state/ratatoskr:sindra", params={"end_user_id": "vuong"})
|
||||
assert sink.events == []
|
||||
Reference in New Issue
Block a user