Buffered (fire-and-forget) writes
|
Available on NAMS: No. The procedure on this page needs the Bolt backend. With the hosted NAMS backend its operations are unavailable: depending on the call, the client raises |
How to keep the agent’s user-facing response from blocking on Neo4j round-trips, using an opt-in fire-and-forget write API.
Default behavior is sync — every write awaits a Neo4j round-trip. In
production, that means each record_tool_call, add_step, custom
write, etc. delays the agent’s reply by the round-trip latency. With
write_mode="buffered", writes go through client.buffered.submit(…),
return after enqueueing when capacity is available, and a background drainer
attempts each write against Neo4j. Existing memory API calls still persist inline.
Prerequisites
-
A connected Bolt client and a Neo4j database on which you can run parameterized writes.
-
MemorySettingsconfigured as in Store, retrieve, and summarize messages. -
An application shutdown path that awaits
flush()orclose()and inspects background errors.
The snippets run inside an async function. generate_response below represents your application’s response function; it is not a library method.
Goal
The agent responds before persistence completes:
async def agent_turn(client, session_id, content):
# Value-returning APIs commit inline — the caller reads the result back,
# so they cannot be deferred.
message = await client.short_term.add_message(
session_id, "user", content, extraction_mode="skip",
)
# Derived writes that nothing in this turn reads back go on the buffer
# and return immediately.
await client.buffered.submit(
"""
MATCH (c:Conversation {session_id: $session})
SET c.turns_recorded = coalesce(c.turns_recorded, 0) + 1
""",
{"session": session_id},
)
# The agent's response is not blocked on the derived write.
return generate_response(message.content)
async with MemoryClient(settings) as client:
response = await agent_turn(client, "session-1", "Hello")
# Drain the queue at end of session / before shutdown.
await client.flush()
Steps
1. Enable buffered mode
# Start with the configured settings from the prerequisites.
settings.memory.write_mode = "buffered"
settings.memory.max_pending = 200 # back-pressure threshold
When max_pending writes are queued, submit() blocks until a worker
drains an item. This is back-pressure — preferable to dropping writes
silently or letting memory grow unbounded.
2. Submit fire-and-forget writes
await client.buffered.submit(
"""
MATCH (m:Message {id: $message_id})
MATCH (e:Entity {name: $name})
MERGE (m)-[:MENTIONS]->(e)
""",
{"message_id": message_id, "name": "Anthem"},
)
submit returns as soon as the work is queued. The drainer task pulls
from the queue and runs the actual execute_write against Neo4j.
In write_mode="sync" (default) submit is a thin passthrough — the
write awaits inline. This means tests that opt into buffered patterns
work in both modes without changes.
3. Drain the queue
client.flush() waits until every submitted job has finished its write attempt.
It does not raise captured write failures or retry them:
# At end of session, end of request handler, or before metrics flush.
await client.flush()
client.wait_for_pending() is an alias.
client.close() (and async with on its way out) automatically
drains the queue before disconnecting from Neo4j. A clean shutdown attempts
queued writes; check write_errors to determine whether any failed. The queue
is held in process memory and is not a durable retry log.
4. Inspect background errors
Errors during background writes are captured rather than raised into
the agent’s hot path. write_errors is cumulative for this client instance:
import logging
log = logging.getLogger(__name__)
for err in client.write_errors:
# BufferedWriteError carries the originating Cypher, its parameters, the
# exception and when it failed.
log.warning(
"Buffered write failed at %s: %s — %s",
err.when.isoformat(),
err.query.strip()[:40],
type(err.error).__name__,
)
5. Verify the persisted result
For the agent_turn example, record the counter before the turn, await the
turn and flush(), then read it back:
await client.flush()
if client.write_errors:
raise RuntimeError(f"{len(client.write_errors)} buffered writes failed")
rows = await client.query.cypher(
"MATCH (c:Conversation {session_id: $session}) "
"RETURN c.turns_recorded AS turns",
{"session": "session-1"},
)
print(rows)
The counter should increase by one for each successful turn in this example.
client.buffered.pending == 0 only measures queue depth; a job can still be
in flight before flush(), and a failed job also leaves the queue. Retain the
error details and decide whether each failed operation is safe to retry.
Tradeoffs
Fire-and-forget trades transactional guarantees for lower latency: a
queued-but-undrained write is lost if the process is killed before
flush() runs, and two related writes execute in submission order but are
not transactional with each other. See
Why the v0.2 operational primitives are opt-in for why buffered writes are
opt-in and what the other v0.2 primitives trade away for the same reasons.
See also
-
examples/buffered-writes/— a runnable example that times a 50-turn conversation insyncmode against the same conversation inbufferedmode, then dropsmax_pendingto 4 so the bounded queue’s back-pressure is visible in the output. Runs withllm=Noneand a local embedder; no API key.