fix(livekit): process captured calls instantly by default

With the dynamic schedule a captured call took 10 to 20 minutes to become
memories, so a caller who rang back right away was not recalled. Capture
now defaults to dreaming="instant": on LiveKit Cloud a call's content was
recallable 19 seconds after it ended, and the callback answered from it.

Each capture write bills one extra operation. When the organization has
no balance for instant processing (HTTP 402), capture and remember fall
back to the dynamic schedule for the rest of the call instead of failing.
capture_dreaming="dynamic" keeps the old behaviour.
This commit is contained in:
Ishaan Gupta 2026-10-02 00:54:38 +05:30
parent 0e604c8384
commit b5d468b06e
5 changed files with 71 additions and 27 deletions

View file

@ -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

View file

@ -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-<session_id>`, 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

View file

@ -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": "<your user id>"}` 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.

View file

@ -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}")

View file

@ -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")