diff --git a/apps/docs/integrations/livekit.mdx b/apps/docs/integrations/livekit.mdx index 7d9f5292..20964206 100644 --- a/apps/docs/integrations/livekit.mdx +++ b/apps/docs/integrations/livekit.mdx @@ -146,7 +146,7 @@ memory = SupermemoryLiveKit( search_threshold=0.1, recall_timeout=2.0, capture="always", - capture_dreaming="dynamic", + capture_dreaming="instant", ), ) ``` @@ -160,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. See below. | +| `capture_dreaming` | `"instant"` | How the call transcript becomes memories. See below. | ## Tools @@ -178,8 +178,8 @@ Recall waits at most `recall_timeout` seconds (default 2). If the profile call i 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. +- With the default `capture_dreaming="instant"`, each write is processed on its own and is usually recallable within a minute, so a caller who rings back is recalled from the last call. Each write bills one extra operation, and the call is written after every agent reply. If your organization has no balance for instant processing, the call is saved on the `dynamic` schedule instead. +- With `capture_dreaming="dynamic"`, related documents are processed together at no extra cost. 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. - 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 76af0a2e..6f5384b7 100644 --- a/packages/livekit-sdk-python/README.md +++ b/packages/livekit-sdk-python/README.md @@ -97,7 +97,7 @@ memory = SupermemoryLiveKit( search_threshold=0.1, recall_timeout=2.0, # seconds; a slow recall falls back to the preloaded profile capture="always", # "always" | "never" - capture_dreaming="dynamic", # see "When a call becomes recallable" below + capture_dreaming="instant", # see "When a call becomes recallable" below ), ) ``` @@ -112,7 +112,7 @@ One call is stored as a single document under custom id `lk-`, so a ### 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. +With the default `capture_dreaming="instant"`, a captured call is usually recallable within a minute. Each write bills one extra operation, and the call is written after every agent reply. If your organization has no balance for instant processing, the call falls back to the `dynamic` schedule. `capture_dreaming="dynamic"` costs nothing extra but took 10 to 20 minutes to become memories in our tests. 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 b2a7f8b6..6d01f2fa 100644 --- a/packages/livekit-sdk-python/examples/README.md +++ b/packages/livekit-sdk-python/examples/README.md @@ -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. 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. +Each call is stored as one document and is usually recallable within a minute, 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 35c85282..26f78e3e 100644 --- a/packages/livekit-sdk-python/src/supermemory_livekit/memory.py +++ b/packages/livekit-sdk-python/src/supermemory_livekit/memory.py @@ -44,7 +44,7 @@ class InputParams(BaseModel): mode: Literal["profile", "query", "full"] = "full" recall_timeout: float = Field(default=2.0, gt=0.0, le=8.0) capture: Literal["always", "never"] = "always" - capture_dreaming: Literal["dynamic", "instant"] = "dynamic" + capture_dreaming: Literal["dynamic", "instant"] = "instant" class SupermemoryLiveKit: @@ -85,6 +85,7 @@ class SupermemoryLiveKit: 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._instant_unavailable = False self._session: Any = None self._shutdown_registered = False @@ -275,23 +276,19 @@ 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(**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 + await asyncio.wait_for( + self._add( + "instant", + content=text, + container_tag=tag, + metadata={"source": "livekit", "kind": "explicit"}, + ), + 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: @@ -477,14 +474,28 @@ class SupermemoryLiveKit: f"{'User' if message['role'] == 'user' else 'Assistant'}: {message['content']}" for message in messages ] - await self._client.add( + await self._add( + self.params.capture_dreaming, content="\n".join(lines), container_tag=messages[0]["tag"], custom_id=messages[0]["custom_id"], metadata={"source": "livekit", "kind": "conversation"}, - dreaming=self.params.capture_dreaming, ) + async def _add(self, dreaming: str, **kwargs: Any) -> None: + """Add a document, on the default schedule when instant processing is not available.""" + if dreaming == "instant" and not self._instant_unavailable: + try: + await self._client.add(**kwargs, dreaming="instant") + return + except Exception as exc: + if getattr(exc, "status_code", None) != 402: + raise + # No balance for instant processing: use the default schedule for the rest of the call. + self._instant_unavailable = True + logger.warning("instant memory processing unavailable; using the default schedule") + await self._client.add(**kwargs) + def _custom_id(self) -> str: raw = self.session_id or self._generated_session return to_identifier(f"lk-{raw}") diff --git a/packages/livekit-sdk-python/tests/test_memory.py b/packages/livekit-sdk-python/tests/test_memory.py index 86d05a36..5e292547 100644 --- a/packages/livekit-sdk-python/tests/test_memory.py +++ b/packages/livekit-sdk-python/tests/test_memory.py @@ -361,7 +361,7 @@ class MemoryTests(unittest.TestCase): self.assertEqual(stored["custom_id"], to_identifier("lk-room 1")) self.assertNotIn("secret", stored["content"]) self.assertEqual(stored["metadata"]["source"], "livekit") - self.assertEqual(stored["dreaming"], "dynamic") + self.assertEqual(stored["dreaming"], "instant") def test_close_flushes_a_trailing_user_turn(self): client = FakeClient() @@ -467,6 +467,39 @@ class MemoryTests(unittest.TestCase): self.assertEqual(len(client.added), 1) self.assertNotIn("dreaming", client.added[0]) + def test_capture_falls_back_once_when_instant_is_not_available(self): + class NoBalance(Exception): + status_code = 402 + + client = FakeClient() + original = client.add + attempts = [] + + async def add(**kwargs): + attempts.append(kwargs.get("dreaming")) + if kwargs.get("dreaming") == "instant": + raise NoBalance("insufficient_balance") + return await original(**kwargs) + + client.add = add + plugin = memory(client, container_tag="user_1", session_id="room-5") + session = Session() + plugin.attach(session) + + async def run(): + for index, role in enumerate(["user", "assistant", "user", "assistant"]): + item = SimpleNamespace(id=f"m{index}", role=role, text_content=f"turn {index}") + session.emit("conversation_item_added", SimpleNamespace(item=item)) + await asyncio.sleep(0) + if plugin._flush_task: + await plugin._flush_task + await plugin.aclose() + + asyncio.run(run()) + + self.assertEqual(attempts, ["instant", None, None]) + self.assertEqual(len(client.added), 2) + def test_rebind_keeps_captured_turns_in_the_callers_scope(self): client = FakeClient() plugin = memory(client, session_id="room-1")