Serve Many Sessions¶
The isolated workbench is correct for one user. A server that shares its
context and sandbox across users is not.
Alice’s next model request can contain Bob’s messages. Their commands can edit the same workspace. If two turns finish concurrently, their messages can also enter one history in the wrong order.
Outcome¶
Build a small in-process session registry. Each stable session ID selects one context, sandbox, agent copy, cleanup stack, and turn lock.
Fast Track¶
Resolve each authenticated request to a stable server-issued session ID.
Bundle one context, sandbox, agent copy, and turn lock per session.
Protect registry creation with one lock and each turn with another.
Own session resources through an
AsyncExitStack.Close every session before closing shared infrastructure.
Hands-on delta¶
1. Bundle session-owned state¶
One object keeps the resources that must never cross a session boundary:
type ContextFactory = Callable[[str], ContextStore]
type SandboxFactory = Callable[[str], DockerSandbox]
@dataclass(slots=True)
class CloudSession:
agent: Agent
context: ContextStore
resources: AsyncExitStack
turn_lock: asyncio.Lock = field(default_factory=asyncio.Lock)
2. Use one lock for each scope¶
The registry lock prevents two cold requests from creating duplicate resources for one ID. A per-session lock serializes history mutation without blocking other sessions.
The lock covers the registry and not the work. Opening a session starts a container, which takes seconds; held under the registry lock, one cold session makes every other session’s first turn wait behind it. The registry therefore holds the task that opens the session. The lock is released as soon as that task exists, and each caller awaits it outside. A task that fails is removed, so the next request for that ID tries again instead of being handed the same failure.
@staticmethod
def _failed(opening: asyncio.Task[CloudSession]) -> bool:
"""Whether this open finished without producing a session, which is when to forget it.
Cancelled counts. Kept, a cancelled open hands `CancelledError` to every later turn for
that ID for the life of the harness.
"""
if not opening.done():
return False
return opening.cancelled() or opening.exception() is not None
async def _session(self, session_id: str) -> CloudSession:
async with self._registry_lock:
if self._closed:
raise RuntimeError("harness is closed")
opening = self._sessions.get(session_id)
if opening is not None and self._failed(opening):
# A session that would not open is not a session. Kept, its failure is handed to
# every later turn for that ID; dropped here, the next request tries again.
opening = None
if opening is None:
opening = asyncio.create_task(self._open(session_id))
self._sessions[session_id] = opening
# Started under the lock and awaited outside it: held while the container starts, one
# cold session made every other first turn wait. `shield` keeps the open running when
# this caller gives up, because it is already creating a container.
return await asyncio.shield(opening)
async def _open(self, session_id: str) -> CloudSession:
resources = AsyncExitStack()
try:
context = self._context_factory(session_id)
resources.push_async_callback(context.close)
sandbox = await resources.enter_async_context(
self._sandbox_factory(session_id)
)
return CloudSession(
agent=self._prototype_agent.copy(
tools=[
*self._prototype_agent.tools,
*(bound_text_tool(tool) for tool in sandbox.tools),
],
),
context=context,
resources=resources,
)
except BaseException:
await resources.aclose()
raise
Creation and turn execution use different locks. Closing the stream releases the turn lock after completion, failure, or client cancellation:
async def stream_turn(
self,
session_id: str,
prompt: str,
) -> AsyncIterator[StreamEvent]:
session = await self._session(session_id)
async with session.turn_lock:
stream = session.agent.run_stream(prompt, session.context)
try:
async for event in stream:
yield event
finally:
await stream.aclose()
async def close(self) -> None:
async with self._registry_lock:
self._closed = True
opening = list(self._sessions.values())
self._sessions.clear()
# A session still opening owns a container already. Wait for it before closing, or the
# resources it is in the middle of taking outlive the harness that asked for them.
opened = await asyncio.gather(*opening, return_exceptions=True)
await asyncio.gather(
*(
session.resources.aclose()
for session in opened
if isinstance(session, CloudSession)
)
)
The prototype carries shared configuration and non-execution tools. It must not carry local execution tools or another session’s Docker tools.
Agent.copy() creates a distinct Agent and replaces its tool list. The
prototype and transport remain shared, while Docker tools point only to that
session’s sandbox through Tool.context.
AsyncExitStack owns resources created during a cold start. If sandbox startup
fails or the task is cancelled, the partial context and Docker client still
close. stream_turn() also closes AgentStream when a client disconnects.
Production note
The example holds the registry lock during cold startup. Use keyed single-flight creation when serialized container startup becomes a measured bottleneck.
Stop accepting requests and drain active turns before closing the harness.
3. Reuse one stable ID¶
Issue IDs on the server, for example with uuid4().hex. Do not place an
arbitrary client string directly in a Docker name or database key.
Constructing this factory is daemon-free. Docker is contacted when the harness enters the returned sandbox.
def docker_for_session(session_id: str) -> DockerSandbox:
return DockerSandbox(
image="python:3.12-slim",
name=session_id,
remove=False,
memory="512m",
cpus="1.0",
network=False,
read_only=True,
cap_drop=["ALL"],
ulimits={"nofile": (256, 256), "nproc": 128},
tmpfs={"/tmp": "size=64m,mode=1777"},
named_volumes={
"/workspace": f"{session_id}-workspace",
},
volumes_remove=False,
workdir="/workspace",
)
name=session_id reattaches to an existing named container and starts it when
necessary. The named volume preserves /workspace if that container was
deleted and must be created again.
Existing containers keep their original sandbox policy. Version that policy, and replace stale containers before reuse. Keep fixed dependencies in the image or under the persistent workspace.
4. Recover the conversation¶
Open one SQLite connection during application startup. Replace the fixed
SESSION_ID from Persist the Conversation
with a factory. Preserve compaction by wrapping each session store with its own
AutoCompactStore:
async def serve_application(
database_path: Path,
prototype_agent: Agent,
transport,
sandbox_factory: SandboxFactory,
serve,
) -> None:
async with AsyncExitStack() as resources:
connection = await connect(database_path)
resources.push_async_callback(connection.close)
def context_for_session(session_id: str) -> ContextStore:
base_context = SQLiteContextStore(
connection,
session_id=session_id,
project="public-repository",
)
return AutoCompactStore(
base_context,
transport,
keep_recent=6,
threshold=0.75,
)
harness = CloudHarness(
prototype_agent=prototype_agent,
context_factory=context_for_session,
sandbox_factory=sandbox_factory,
)
resources.push_async_callback(harness.close)
await serve(harness)
The application owns the shared connection. Closing one session closes its context wrapper, not the connection.
See Context and Messages for store behavior, compaction, and custom backends.
Try It¶
This check does not require Docker or a model API. It proves that same-session turns serialize, different sessions overlap, resources stay separate, and SQLite recovers one stable ID.
Run uv run python examples/tutorial/serve_many_sessions.py from the repository
root.
async def verify_sqlite_recovery(database: Path) -> None:
session_id = "4f0f29b8a5d94f688306c231d86aa531"
first_connection = await connect(database)
try:
first = SQLiteContextStore(
first_connection,
session_id=session_id,
project="public-repository",
)
await first.append(
Message(
role="user",
content=[TextBlock(text="Review README.md")],
)
)
finally:
await first_connection.close()
second_connection = await connect(database)
try:
recovered = SQLiteContextStore(
second_connection,
session_id=session_id,
project="public-repository",
)
history = await recovered.get_history()
assert len(history) == 1
assert history[0].role == "user"
assert history[0].content == [TextBlock(text="Review README.md")]
finally:
await second_connection.close()
After a service restart, the registry starts empty. The first request recreates
the in-memory CloudSession; SQLite reloads its messages, and Docker reattaches
or remounts its workspace using the same ID.
If the database survives but both the container and volume are missing, history recovers into an empty workspace. Report that condition instead of claiming full recovery.
Done when¶
[ ] Every request resolves to an authenticated, stable server-issued ID.
[ ] Same-session turns never overlap.
[ ] Different sessions can stream concurrently.
[ ] Contexts, agent copies, and sandboxes stay session-local.
[ ]
AsyncExitStackcloses partial and completed resources.[ ] SQLite recovers history by the same stable ID.
Next failure¶
The next integration request creates another boundary. External tool connections have their own names, permissions, and lifetimes.
Continue with Extend the Tool Set without Changing the Loop.