From 0e604c83843073925733f79c809363ed54bdfa43 Mon Sep 17 00:00:00 2001 From: Ishaan Gupta Date: Fri, 2 Oct 2026 00:35:38 +0530 Subject: [PATCH] fix(livekit): keep memory scoped to the caller and bound capture writes From review of 5b2a435e: - Captured turns now keep the container tag and document id they were spoken under. Rebinding the instance used to flush earlier turns into the new caller's scope. Turns captured before any bind still go to the first caller bound. - Recall strips earlier injected memory before it runs, so a failed or empty recall can no longer leave another caller's memory in context. For the same caller, a slow or failed recall falls back to the profile loaded by preload. - The recall cache is keyed on the user message, not its text, so a later turn with the same words recalls again. - Capture writes time out after 10s and keep their turns for retry, including on cancellation. Large calls are written in chunks of at most 100k characters, and the buffer is capped. - remember falls back to the default processing schedule when the organization has no balance for instant processing (HTTP 402). - Docs state the measured delays: about a minute for remember, 10 to 20 minutes for captured calls on the default dynamic schedule. --- apps/docs/integrations/livekit.mdx | 16 ++- packages/livekit-sdk-python/README.md | 10 +- .../livekit-sdk-python/examples/README.md | 4 +- .../src/supermemory_livekit/memory.py | 135 +++++++++++++----- .../livekit-sdk-python/tests/test_memory.py | 133 ++++++++++++++++- 5 files changed, 255 insertions(+), 43 deletions(-) diff --git a/apps/docs/integrations/livekit.mdx b/apps/docs/integrations/livekit.mdx index 5b090084..7d9f5292 100644 --- a/apps/docs/integrations/livekit.mdx +++ b/apps/docs/integrations/livekit.mdx @@ -36,6 +36,8 @@ Set the attribute in the access token your backend issues. Do not grant the call One call is stored as a single document. Pass the LiveKit room name as `session_id`. The document id is `lk-`, so a later update with the same id appends to that call instead of creating another. +Use one `SupermemoryLiveKit` per call. If you do rebind it, turns already captured stay with the caller who said them, and memory from the earlier caller is never shown to the new one. + ## Quick start Create one `SupermemoryLiveKit` inside the job. A worker handles many calls, and each call needs its own scope. @@ -149,7 +151,7 @@ memory = SupermemoryLiveKit( ) ``` -Recall waits at most `recall_timeout` seconds (default 2). If the profile call is slower than that, the turn proceeds without memory. Retrieved text is inserted immediately before the user message and is not stored back as something the agent said. +Recall waits at most `recall_timeout` seconds (default 2). If the profile call is slower than that, or fails, the turn falls back to the profile loaded by `preload` at the start of the call, or proceeds without memory if there was none. In our tests most recalls took 0.4 to 1.4 seconds, with occasional calls near 3 seconds. Raise `recall_timeout` if you would rather wait than miss memory on a slow turn. Retrieved text is inserted immediately before the user message and is not stored back as something the agent said. | Parameter | Default | Description | | --- | --- | --- | @@ -158,7 +160,7 @@ Recall waits at most `recall_timeout` seconds (default 2). If the profile call i | `mode` | `"full"` | `profile`, `query`, or `full` | | `recall_timeout` | `2.0` | Seconds to wait before skipping recall | | `capture` | `"always"` | `never` disables storing the call | -| `capture_dreaming` | `"dynamic"` | How the call transcript becomes memories. `dynamic` batches related documents and can take several minutes. `instant` is usually ready within a minute and bills one extra operation per document. Until then, `search_memories` still finds the raw call text. | +| `capture_dreaming` | `"dynamic"` | How the call transcript becomes memories. See below. | ## Tools @@ -170,7 +172,15 @@ Recall waits at most `recall_timeout` seconds (default 2). If the profile call i | `remember` | Save one explicit fact, preference, or correction | | `forget` | Forget one fact by id from `search_memories`, or by exact text | -`remember` stores a standalone fact and processes it right away, so the next call can recall it. It is not appended to the call transcript. Automatic capture is what records the conversation. +`remember` stores a standalone fact and processes it immediately (`dreaming="instant"`), so it is usually recallable within a minute. If your organization has no balance for instant processing, it is saved on the default schedule instead. It is not appended to the call transcript. Automatic capture is what records the conversation. + +### When a call becomes recallable + +Captured turns are written after each agent reply, but they become memories on Supermemory's processing schedule: + +- With the default `capture_dreaming="dynamic"`, related documents are processed together. In our tests this took 10 to 20 minutes, so a caller who rings back right away is not recalled from the last call yet. +- With `capture_dreaming="instant"`, each write is processed on its own and is usually recallable within a minute. Each write bills one extra operation, and the call is written after every agent reply. +- In both modes, `search_memories` finds the raw call text as soon as it is stored, so the model can still look up a recent call. ## Self-hosting diff --git a/packages/livekit-sdk-python/README.md b/packages/livekit-sdk-python/README.md index c6d918e1..76af0a2e 100644 --- a/packages/livekit-sdk-python/README.md +++ b/packages/livekit-sdk-python/README.md @@ -95,9 +95,9 @@ memory = SupermemoryLiveKit( mode="full", # "profile" | "query" | "full" search_limit=10, search_threshold=0.1, - recall_timeout=2.0, # seconds; a slow recall is skipped + recall_timeout=2.0, # seconds; a slow recall falls back to the preloaded profile capture="always", # "always" | "never" - capture_dreaming="dynamic", # "dynamic" | "instant" (ready within a minute, extra operation) + capture_dreaming="dynamic", # see "When a call becomes recallable" below ), ) ``` @@ -108,7 +108,11 @@ memory = SupermemoryLiveKit( | `query` | No | Yes | You only need memories related to this turn | | `full` | Yes | Yes | Default | -One call is stored as a single document under custom id `lk-`, so a reconnect with the same session id updates that document instead of creating another. Explicit `remember` calls are separate facts, processed right away so the next call can recall them, and are not tied to the call document. +One call is stored as a single document under custom id `lk-`, so a reconnect with the same session id updates that document instead of creating another. Explicit `remember` calls are separate facts, processed immediately (usually recallable within a minute), and are not tied to the call document. + +### When a call becomes recallable + +With the default `capture_dreaming="dynamic"`, a captured call took 10 to 20 minutes to become memories in our tests, so a caller who rings back right away is not recalled from the last call yet. `capture_dreaming="instant"` is usually ready within a minute but bills one extra operation per write, and the call is written after every agent reply. In both modes `search_memories` finds the raw call text as soon as it is stored. ## Links diff --git a/packages/livekit-sdk-python/examples/README.md b/packages/livekit-sdk-python/examples/README.md index 85cfea14..b2a7f8b6 100644 --- a/packages/livekit-sdk-python/examples/README.md +++ b/packages/livekit-sdk-python/examples/README.md @@ -27,7 +27,7 @@ In the browser: run `python voice_agent.py dev`, open the [Agents Playground](ht 1. Say "My name is Priya, and please remember I'm vegetarian." 2. Hang up and start a new call. -3. The agent greets you by name. Ask "What should I order for dinner?" +3. The agent welcomes you back. Ask "What should I order for dinner?" Memory is scoped per caller: @@ -35,4 +35,4 @@ Memory is scoped per caller: - In rooms, the agent uses the participant attribute `supermemory_container_tag`, or else the participant identity. The Playground gives each session a new identity, so set `SUPERMEMORY_CONTAINER_TAG` in `.env` to keep one caller across Playground calls. - In production, dispatch the agent from your backend with `{"container_tag": ""}` as job metadata, and set `LIVEKIT_AGENT_NAME`. -Each call is stored as one document, and facts the caller asks it to remember are saved right away. See the [integration docs](https://supermemory.ai/docs/integrations/livekit) for configuration. +Each call is stored as one document. Facts the caller asks it to remember are usually recallable within a minute; the rest of the call takes longer, see [when a call becomes recallable](https://supermemory.ai/docs/integrations/livekit#when-a-call-becomes-recallable). See the [integration docs](https://supermemory.ai/docs/integrations/livekit) for configuration. diff --git a/packages/livekit-sdk-python/src/supermemory_livekit/memory.py b/packages/livekit-sdk-python/src/supermemory_livekit/memory.py index edbae245..35c85282 100644 --- a/packages/livekit-sdk-python/src/supermemory_livekit/memory.py +++ b/packages/livekit-sdk-python/src/supermemory_livekit/memory.py @@ -27,6 +27,14 @@ logger = logging.getLogger("supermemory_livekit") _UNAVAILABLE = "I couldn't reach memory just now." _ATTRIBUTE = "supermemory_container_tag" +_STORE_TIMEOUT = 10.0 +_MAX_DOCUMENT_CHARS = 100_000 +_MAX_BUFFERED = 2_000 + + +def _log_flush_failure(task: "asyncio.Task[None]") -> None: + if not task.cancelled() and task.exception() is not None: + logger.warning("memory capture failed; will retry", exc_info=task.exception()) class InputParams(BaseModel): @@ -71,10 +79,12 @@ class SupermemoryLiveKit: self._generated_session = f"session-{uuid4().hex[:12]}" self._client = client if client is not None else self._build_client(base_url) self._seen: set[str] = set() - self._buffer: list[dict[str, str]] = [] + # Each captured turn keeps the scope it was spoken in, so a rebind cannot move it. + self._buffer: list[dict[str, Optional[str]]] = [] self._lock = asyncio.Lock() self._flush_task: Optional[asyncio.Task[None]] = None self._recall: Optional[tuple[tuple[Optional[str], str], asyncio.Future[Optional[str]]]] = None + self._preloaded: Optional[tuple[str, str]] = None self._session: Any = None self._shutdown_registered = False @@ -120,8 +130,12 @@ class SupermemoryLiveKit: self.container_tag = scoped if session_id: self.session_id = str(session_id) - if self._buffer and self.container_tag: - self._schedule_flush() + if self.container_tag: + custom_id = self._custom_id() + for message in self._buffer: + if message["tag"] is None: + message["tag"], message["custom_id"] = self.container_tag, custom_id + self._schedule_flush() def tools(self) -> list[Any]: from .tools import build_tools @@ -130,9 +144,11 @@ class SupermemoryLiveKit: async def preload(self, chat_ctx: Any) -> bool: """Inject the caller profile into the initial chat context. Returns whether anything was added.""" + tag = self.container_tag text = await self._recall_text(query=None) - if not text: + if not text or tag is None or tag != self.container_tag: return False + self._preloaded = (tag, text) self._inject(chat_ctx, text, created_at=None) return True @@ -146,6 +162,8 @@ class SupermemoryLiveKit: user = self._last_user_item(chat_ctx) if user is not None: await self._recall_into(chat_ctx, user) + elif not self._preload_matches(): + self._strip_injected(chat_ctx) async def on_user_turn_completed(self, turn_ctx: Any, new_message: Any) -> None: """Recall from the turn hook, for realtime models that skip ``llm_node``. @@ -159,24 +177,40 @@ class SupermemoryLiveKit: query = message_text(user) if not query: return + tag = self.container_tag + # Drop earlier injections first, so a failed recall leaves no memory rather than stale + # or foreign memory. The caller's own call-start profile is the fallback. + self._strip_injected(chat_ctx) + text: Optional[str] = None + try: + text = await self._recall_once(user, query) + except Exception: + logger.warning("memory recall failed", exc_info=True) + if tag != self.container_tag: + return + if not text and self._preload_matches(): + text = self._preloaded[1] if self._preloaded else None + if not text: + return + created_at = getattr(user, "created_at", None) + before = created_at - 0.001 if isinstance(created_at, (int, float)) else None try: - text = await self._recall_once(query) - if not text: - return - created_at = getattr(user, "created_at", None) - before = created_at - 0.001 if isinstance(created_at, (int, float)) else None - self._strip_injected(chat_ctx) self._inject(chat_ctx, text, created_at=before) except Exception: logger.warning("memory inject failed", exc_info=True) - async def _recall_once(self, query: str) -> Optional[str]: - key = (self.container_tag, query) + async def _recall_once(self, user: Any, query: str) -> Optional[str]: + # Keyed on the user message, so retries and tool follow-ups of one turn share a recall + # while a later turn with the same words recalls again. + key = (self.container_tag, str(getattr(user, "id", None) or query)) if self._recall is None or self._recall[0] != key: self._recall = (key, asyncio.ensure_future(self._recall_text(query=query))) # A cancelled preemptive generation must not cancel a recall the next attempt reuses. return await asyncio.shield(self._recall[1]) + def _preload_matches(self) -> bool: + return self._preloaded is not None and self._preloaded[0] == self.container_tag + def attach( self, session: Any, @@ -241,19 +275,23 @@ class SupermemoryLiveKit: return "Memory is not scoped to a caller yet." if not text: return "Nothing to remember." + add = { + "content": text, + "container_tag": tag, + "metadata": {"source": "livekit", "kind": "explicit"}, + } try: - await asyncio.wait_for( - self._client.add( - content=text, - container_tag=tag, - metadata={"source": "livekit", "kind": "explicit"}, - dreaming="instant", - ), - timeout=4.0, - ) - except Exception: - logger.warning("memory remember failed", exc_info=True) - return _UNAVAILABLE + await asyncio.wait_for(self._client.add(**add, dreaming="instant"), timeout=4.0) + except Exception as exc: + if getattr(exc, "status_code", None) != 402: + logger.warning("memory remember failed", exc_info=True) + return _UNAVAILABLE + # No balance for instant processing: save it on the default schedule instead. + try: + await asyncio.wait_for(self._client.add(**add), timeout=4.0) + except Exception: + logger.warning("memory remember failed", exc_info=True) + return _UNAVAILABLE return "Saved." async def forget(self, *, memory_id: str = "", memory_text: str = "") -> str: @@ -373,7 +411,18 @@ class SupermemoryLiveKit: if item_id in self._seen: return self._seen.add(item_id) - self._buffer.append({"role": role, "content": text}) + tag = self.container_tag + self._buffer.append( + { + "role": role, + "content": text, + "tag": tag, + "custom_id": self._custom_id() if tag else None, + } + ) + if len(self._buffer) > _MAX_BUFFERED: + del self._buffer[: len(self._buffer) - _MAX_BUFFERED] + logger.warning("memory capture buffer full; dropped the oldest turns") if role == "assistant": self._schedule_flush() @@ -381,7 +430,7 @@ class SupermemoryLiveKit: self._schedule_flush() def _schedule_flush(self) -> None: - if self.params.capture != "always" or not self._buffer or not self.container_tag: + if self.params.capture != "always" or not self._buffer or self._buffer[0]["tag"] is None: return try: loop = asyncio.get_running_loop() @@ -390,30 +439,48 @@ class SupermemoryLiveKit: if self._flush_task is not None and not self._flush_task.done(): return self._flush_task = loop.create_task(self._flush()) + self._flush_task.add_done_callback(_log_flush_failure) async def _flush(self) -> None: while True: async with self._lock: - if self.params.capture != "always" or not self._buffer or not self.container_tag: + if self.params.capture != "always": + return + batch = self._take_batch() + if not batch: return - batch = self._buffer - self._buffer = [] try: - await self._store(batch) - except Exception: + await asyncio.wait_for(self._store(batch), timeout=_STORE_TIMEOUT) + except BaseException: async with self._lock: self._buffer = batch + self._buffer raise - async def _store(self, messages: list[dict[str, str]]) -> None: + def _take_batch(self) -> list[dict[str, Optional[str]]]: + """Take the leading turns that share one scope, capped in size.""" + if not self._buffer or self._buffer[0]["tag"] is None: + return [] + scope = (self._buffer[0]["tag"], self._buffer[0]["custom_id"]) + count, size = 0, 0 + for message in self._buffer: + length = len(message["content"] or "") + if (message["tag"], message["custom_id"]) != scope: + break + if count and size + length > _MAX_DOCUMENT_CHARS: + break + count, size = count + 1, size + length + batch, self._buffer = self._buffer[:count], self._buffer[count:] + return batch + + async def _store(self, messages: list[dict[str, Optional[str]]]) -> None: lines = [ f"{'User' if message['role'] == 'user' else 'Assistant'}: {message['content']}" for message in messages ] await self._client.add( content="\n".join(lines), - container_tag=self.container_tag, - custom_id=self._custom_id(), + container_tag=messages[0]["tag"], + custom_id=messages[0]["custom_id"], metadata={"source": "livekit", "kind": "conversation"}, dreaming=self.params.capture_dreaming, ) diff --git a/packages/livekit-sdk-python/tests/test_memory.py b/packages/livekit-sdk-python/tests/test_memory.py index 237ec291..86d05a36 100644 --- a/packages/livekit-sdk-python/tests/test_memory.py +++ b/packages/livekit-sdk-python/tests/test_memory.py @@ -7,6 +7,7 @@ import unittest from types import SimpleNamespace from supermemory_livekit import ConfigurationError, InputParams, SupermemoryLiveKit +from supermemory_livekit import memory as memory_module from supermemory_livekit.identifiers import to_identifier from supermemory_livekit.utils import format_tool_results, is_injected_memory, wrap_memory @@ -213,7 +214,7 @@ class MemoryTests(unittest.TestCase): first, second = ChatCtx(), ChatCtx() for ctx in (first, second): ctx.add_message(role="user", content="hi, first time calling", created_at=5.0) - # A preemptive attempt is cancelled while the retry reuses its recall. + # A cancelled attempt must not cancel the recall a follow-up of the same turn reuses. preemptive = asyncio.ensure_future(plugin.enrich(first)) await asyncio.sleep(0.01) preemptive.cancel() @@ -224,6 +225,56 @@ class MemoryTests(unittest.TestCase): self.assertEqual(len(client.profile.calls), 1) + def test_same_words_in_a_later_turn_recall_again(self): + client = FakeClient(FakeProfile(static=["Name is Ada"])) + plugin = memory(client, container_tag="user_1") + ctx = ChatCtx() + + ctx.add_message(role="user", content="what's my name?", created_at=1.0) + asyncio.run(plugin.enrich(ctx)) + ctx.add_message(role="user", content="what's my name?", created_at=2.0) + asyncio.run(plugin.enrich(ctx)) + + self.assertEqual(len(client.profile.calls), 2) + + def test_failed_recall_never_shows_another_callers_memory(self): + class ScopedProfile: + def __init__(self): + self.calls = [] + + async def __call__(self, **kwargs): + self.calls.append(kwargs) + if kwargs["container_tag"] == "bob": + raise RuntimeError("down") + return SimpleNamespace( + profile=SimpleNamespace(static=["Alice's door code is 4471"], dynamic=[]), + search_results=SimpleNamespace(results=[]), + ) + + plugin = memory(FakeClient(ScopedProfile()), container_tag="alice") + ctx = ChatCtx() + asyncio.run(plugin.preload(ctx)) + plugin.bind(container_tag="bob") + ctx.add_message(role="user", content="what's my door code?", created_at=50.0) + + asyncio.run(plugin.enrich(ctx)) + + self.assertFalse(any("4471" in item.content for item in ctx.items)) + + def test_slow_recall_falls_back_to_the_callers_call_start_profile(self): + profile = FakeProfile(static=["Name is Ada"]) + plugin = memory(FakeClient(profile), container_tag="user_1", params=InputParams(recall_timeout=0.01)) + ctx = ChatCtx() + asyncio.run(plugin.preload(ctx)) + profile.delay = 0.05 + ctx.add_message(role="user", content="what's my name?", created_at=50.0) + + asyncio.run(plugin.enrich(ctx)) + + injected = [item for item in ctx.items if is_injected_memory(item.content)] + self.assertEqual(len(injected), 1) + self.assertIn("Name is Ada", injected[0].content) + def test_new_user_message_recalls_again(self): client = FakeClient(FakeProfile(static=["Name is Ada"])) plugin = memory(client, container_tag="user_1") @@ -397,6 +448,86 @@ class MemoryTests(unittest.TestCase): self.assertEqual(client.memories.calls[0]["container_tag"], "user_1") self.assertIn("memory id", missing) + def test_remember_falls_back_when_instant_is_not_available(self): + class NoBalance(Exception): + status_code = 402 + + client = FakeClient() + original = client.add + + async def add(**kwargs): + if kwargs.get("dreaming") == "instant": + raise NoBalance("insufficient_balance") + return await original(**kwargs) + + client.add = add + plugin = memory(client, container_tag="user_1") + + self.assertEqual(asyncio.run(plugin.remember("Likes tea")), "Saved.") + self.assertEqual(len(client.added), 1) + self.assertNotIn("dreaming", client.added[0]) + + def test_rebind_keeps_captured_turns_in_the_callers_scope(self): + client = FakeClient() + plugin = memory(client, session_id="room-1") + session = Session() + plugin.attach(session) + + def say(item_id, role, text): + item = SimpleNamespace(id=item_id, role=role, text_content=text) + session.emit("conversation_item_added", SimpleNamespace(item=item)) + + say("u0", "user", "hello before we know who you are") + plugin.bind(container_tag="alice") + say("u1", "user", "Alice secret: door code 4471") + plugin.bind(container_tag="bob") + say("u2", "user", "Bob here") + asyncio.run(plugin.aclose()) + + stored = {call["container_tag"]: call["content"] for call in client.added} + self.assertEqual(set(stored), {"alice", "bob"}) + self.assertIn("4471", stored["alice"]) + self.assertIn("hello before we know who you are", stored["alice"]) + self.assertNotIn("4471", stored["bob"]) + + def test_hung_store_times_out_and_keeps_the_turns(self): + client = FakeClient() + + async def hang(**kwargs): + await asyncio.sleep(60) + + client.add = hang + plugin = memory(client, container_tag="user_1", session_id="room-3") + session = Session() + plugin.attach(session) + session.emit( + "conversation_item_added", + SimpleNamespace(item=SimpleNamespace(id="u1", role="user", text_content="keep me")), + ) + + original = memory_module._STORE_TIMEOUT + memory_module._STORE_TIMEOUT = 0.01 + try: + asyncio.run(asyncio.wait_for(plugin.aclose(), timeout=2)) + finally: + memory_module._STORE_TIMEOUT = original + + self.assertEqual([m["content"] for m in plugin._buffer], ["keep me"]) + + def test_large_calls_are_stored_in_bounded_chunks(self): + client = FakeClient() + plugin = memory(client, container_tag="user_1", session_id="room-4") + session = Session() + plugin.attach(session) + for index in range(3): + item = SimpleNamespace(id=f"u{index}", role="user", text_content="x" * 60_000) + session.emit("conversation_item_added", SimpleNamespace(item=item)) + + asyncio.run(plugin.aclose()) + + self.assertEqual(len(client.added), 3) + self.assertEqual({call["custom_id"] for call in client.added}, {to_identifier("lk-room-4")}) + def test_unscoped_tools_do_not_call_the_api(self): client = FakeClient() plugin = memory(client)