Sessions and Services¶
A Session manages connection-scoped state. It
multiplexes action messages and node fragments over wire streams,
dispatches inbound actions, and tracks in-flight work.
A Service manages an action registry across
multiple sessions, allowing one instance to serve several listeners and
protocols.
Session boundary¶
A session owns the node namespace, action registry, active-action accounting,
transport buffers, and connection deadline for one peer relationship. Attaching
a stream in start or accept mode installs the session's callbacks and begins
its receive pump immediately.
One WireMessage can carry action control for several calls and fragments for
several nodes. This multiplexing lets concurrent action and data streams share
one physical connection. A session may also attach several wire streams, with
per-stream and session-wide limits bounding queued messages and bytes.
Clean shutdown is a two-sided operation. half_close() rejects new work and
half-closes active transports while already admitted actions and callbacks
finish. await session.done.wait() is the Python barrier for released stream
state. get_status() then reports the authoritative terminal status. The
Session lifecycle covers deadlines, aborts,
pull-style reception, and completion barriers.
Dynamic services over shared transports¶
A gRPC channel provides an efficient connection to methods normally described in an IDL and exposed through generated stubs. A11 uses a similar separation between logical calls and physical connections, but resolves action schemas from a live registry. Clients can discover and call the operations available in a particular session without generated service code.
The registry remains live for the connection, allowing clients to discover and call operations available to that session.
Using Session and Service¶
import asyncio
import a11
from a11.actions import ActionRegistry
from a11.service.serving import serving, websocket
registry = ActionRegistry()
@registry.action(name="greet")
async def greet(name: str) -> str:
return f"Hello, {name}!"
service = a11.Service(action_registry=registry)
async def run_server():
async with serving(service, websocket(port=8080)):
print("Server running on ws://localhost:8080")
await asyncio.Event().wait()
a11.service.session.Session
¶
Session(session_id: str = '', on_stream_message: Any | None = None, on_stream_done: Any | None = None, headers: Mapping[str, bytes] | None = None, options: SessionOptions | None = None, node_map: NodeMap | None = None, action_registry: ActionRegistry | None = None)
Create an A11 session that multiplexes wire streams and actions. Streams deliver messages asynchronously to the optional on_stream_message and on_stream_done callbacks, which may be coroutines. This is the top-level object an agent drives to exchange wire messages and run actions.
action_registry
property
writable
¶
action_registry: ActionRegistry | None
The ActionRegistry used to resolve action messages; assigning replaces it.
done
property
¶
An asyncio.Event-shaped view of full session completion.
Session.is_closed can become true as soon as shutdown starts. Await
this event (or wait_done) when streams and actions must all have
released their runtime state.
node_map
property
writable
¶
node_map: NodeMap
The NodeMap backing this session's node state; assigning replaces it.
actions
¶
actions() -> list[tuple[str, Action]]
Return the (action_id, action) pairs currently running in the session. Actions execute asynchronously, so this is a point-in-time snapshot of in-flight work.
add_done_callback
¶
Invoke callback(session) once this session fully completes.
The callback fires exactly once when the session finishes -- whether it
drained and closed cleanly, its deadline elapsed, or it was aborted
("dies") -- because every one of those paths resolves the completion
view exposed by done. If the session is already done, the callback
still runs on the next event-loop iteration.
A synchronous callback runs to completion; one returning an awaitable is awaited. This is the hook connection-scoped resources should register on so they are released when the session ends regardless of outcome (e.g. reaping the shells started within a session).
Returns the scheduled asyncio.Task. Must be called from within a
running event loop.
add_stream
¶
add_stream(stream: WireStream, mode: Any = 'start') -> Future[None]
Attach a wire stream and begin pumping its messages, returning an awaitable for the stream's lifetime. mode selects whether this side starts ("start") or accepts ("accept") the stream.
Examples:
Attach the client transport before exchanging messages:
await_all_actions
¶
await_all_actions(timeout: Duration | None = None) -> Future[None]
Return an awaitable that resolves once all in-flight actions have finished, or the optional timeout elapses. Await this to synchronize on the session's outstanding asynchronous work before proceeding.
cancel_action
¶
Request cancellation of the running action with the given id, raising if it is unknown. Cancellation is cooperative and completes asynchronously as the action unwinds.
cancel_all_actions
¶
Request cancellation of every action currently running in the session. Each action unwinds asynchronously; await await_all_actions to observe completion.
dispatch_action
¶
Dispatch an already-constructed Action to run within the session, returning an awaitable for its handling.
dispatch_action_message
¶
dispatch_action_message(action_message: ActionMessage, origin_stream: WireStream | None = None) -> Future[None]
Dispatch an action message, resolving it against the action registry and running the resulting action. Returns an awaitable that completes when the action has been handled; origin_stream attributes the message to a source stream.
dispatch_node_fragment
¶
dispatch_node_fragment(fragment: NodeFragment) -> Future[int]
Dispatch a node fragment into the session's NodeMap and return an awaitable resolving to the applied revision. Fragments are applied asynchronously in order, letting an agent stream incremental document updates.
dispatch_wire_message
¶
dispatch_wire_message(message: WireMessage, origin_stream: WireStream | None = None) -> Future[None]
Route a wire message through the session as though it arrived on a stream, returning an awaitable for its processing. origin_stream optionally records which stream the message is attributed to.
get_action
¶
get_action(action_id: str) -> Action
Look up a running action by its id, raising if none matches.
get_action_registry
¶
get_action_registry() -> ActionRegistry | None
Return the ActionRegistry used to resolve incoming action messages into runnable actions.
get_node_map
¶
get_node_map() -> NodeMap
Return the NodeMap backing this session's node state. Node fragments dispatched to the session are applied to this map as messages stream in.
get_status
¶
get_status() -> Status
Return the session's terminal status, indicating whether it completed successfully or was aborted.
get_stream
¶
get_stream(stream_id: str) -> WireStream
Look up an attached stream by its id, raising if no such stream exists. Because streams come and go over the session's lifetime, guard against a stream having been removed since you last observed it.
half_close
¶
is_closed
¶
Return whether the session has been closed and no longer accepts new streams or messages.
is_done
¶
Return whether the session has fully finished, including all streams and actions. Prefer awaiting done for asynchronous completion rather than polling this flag.
send
¶
send(message: WireMessage, stream_id: str = '') -> None
set_action_registry
¶
set_action_registry(registry: ActionRegistry | None) -> None
Replace the ActionRegistry used to resolve incoming action messages, raising on failure. Active actions are rebound for later nested-name resolution; configure it before dispatch to avoid mixing registry versions.
set_deadline
¶
set_deadline(deadline: Time | None = None) -> None
Set the absolute deadline after which the session is aborted; passing None clears it to no deadline. The session enforces this asynchronously as time passes.
set_node_map
¶
set_node_map(node_map: NodeMap) -> None
Replace the NodeMap backing this session's node state, raising on failure. Active actions are rebound, but existing fragments are not migrated; set it before traffic to avoid splitting state between maps.
streams
¶
streams() -> list[tuple[str, WireStream]]
Return the (stream_id, stream) pairs currently attached to the session. Streams are added and removed asynchronously as peers connect and disconnect, so treat the result as a snapshot taken at call time.
wait_done
¶
Return an awaitable that resolves when the session has fully finished. Await this to block until every stream and action has completed asynchronously.
SessionWithRecv¶
SessionWithRecv provides explicit
pull-based message reception for applications managing custom event loops or
multiplexed transports:
Pull reception replaces the default action and node dispatcher. Code that
still needs normal routing passes each received message to
dispatch_wire_message() after inspection or proxying.
a11.service.session.SessionWithRecv
¶
SessionWithRecv(session_id: str = '', headers: Mapping[str, bytes] | None = None, options: SessionOptions | None = None, node_map: NodeMap | None = None, action_registry: ActionRegistry | None = None)
Bases: Session
Create a session that buffers inbound messages for explicit pull-based reception instead of callbacks. receive and receive_with_stream_id await messages as they stream in.
receive
async
¶
Await the next inbound message, or None when the session ends.
Use this when one receive loop handles every attached stream. Choose
receive_with_stream_id when replies or diagnostics must retain their
transport identity. The optional absolute deadline limits only this
wait; it does not change the session deadline.
Examples:
Route messages from a session with one attached transport:
receive_with_stream_id
async
¶
Await (message, stream_id), or None after completion.
This is the pull-style counterpart to OnSessionStreamMessage and
is useful when an agent multiplexes several transports in one loop.
Examples:
Preserve the source while routing gateway traffic:
a11.service.session.SessionOptions
¶
SessionOptions(*, max_buffered_messages_total: SupportsInt | None = 256, max_buffered_messages_per_stream: SupportsInt | None = 32, max_concurrent_root_actions: SupportsInt | None = 32, max_concurrent_nested_actions: SupportsInt | None = 128, max_single_message_size: SupportsInt | None = 33554432, max_buffered_bytes_total: SupportsInt | None = 33554432, max_buffered_bytes_per_stream: SupportsInt | None = 4194304, no_stream_timeout: Duration | None = None, deadline: Time | None = None)
Construct session limits and timeouts; all parameters are keyword-only.
max_buffered_bytes_per_stream
property
writable
¶
Maximum bytes buffered per stream.
max_buffered_bytes_total
property
writable
¶
Maximum total bytes buffered across all streams.
max_buffered_messages_per_stream
property
writable
¶
Maximum number of messages buffered per stream.
max_buffered_messages_total
property
writable
¶
Maximum number of messages buffered across all streams.
max_concurrent_nested_actions
property
writable
¶
Maximum number of concurrently running nested actions.
max_concurrent_root_actions
property
writable
¶
Maximum number of concurrently running root actions.
max_single_message_size
property
writable
¶
Maximum size in bytes of a single wire message.
Configured telemetry records an a11.session span from session creation
through full completion. See observability.
Service¶
a11.service.service.Service
¶
Service(*, action_registry: ActionRegistry | None = None, on_connection: Any | None = None, options: ServiceOptions | None = None)
A service: an action registry plus the sessions serving it.
accept is shaped to be a transport's on-stream callback, so one service can be
bound to several listeners, or to none at all (hand it an in-process stream). The
optional on_connection(session, stream) coroutine runs once per connection,
after the session exists and before it starts pumping -- the only window in which
a connection can be specialised without racing its first message.
Examples:
Serve a gateway over WebSocket:
service = a11.Service(action_registry=registry, on_connection=prepare)
server = a11.net.WebSocketWireServer.create(service.accept, options)
action_registry
property
¶
action_registry: ActionRegistry | None
The template registry new connections are built from.
done
property
¶
An asyncio.Event-shaped view of the service being closed and empty.
Set once the service has stopped accepting and every session it was serving has finished.
accept
¶
accept(stream: WireStream) -> Future[None]
Serve an accepted stream, awaiting its session's whole lifetime.
aclose
async
¶
aclose(*, timeout: Duration | None = None) -> None
Stop accepting, then wait for what is in flight.
The graceful shutdown, in the order that makes it graceful: refusing new connections first means the set being waited on cannot grow.
add_stream_to_session
¶
add_stream_to_session(session_id: str, stream: WireStream, mode: Any = 'accept') -> None
Attach another transport to an existing session.
describe
¶
Describe this service's actions as an a11.actions/v1 JSON document. With a name, describes that one action or raises NOT_FOUND. query carries the same filters as the HTTP endpoint's query string.
drain
async
¶
drain(timeout: Duration | None = None) -> None
Await the completion of every session currently being served.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
timeout
|
Duration | None
|
How long to wait. |
None
|
Raises:
| Type | Description |
|---|---|
StatusException
|
|
expose_descriptors_on
¶
Point a listener's describe options at this service, so its transport answers GET /actions from the same describer the list_actions builtin uses.
get_session
¶
get_session(session_id: str) -> Session
The session with this id, raising NOT_FOUND if there is none.
get_session_for_stream
¶
get_session_for_stream(stream_id: str) -> Session
The session serving this stream.
serve
¶
serve(stream: WireStream, mode: Any = 'accept') -> Future[None]
Serve a stream in the given mode ("start" or "accept").
set_action_registry
¶
set_action_registry(action_registry: ActionRegistry | None) -> None
Replace the registry new connections are built from, without interrupting any stream.
start
¶
start(stream: WireStream) -> Future[None]
Serve a stream this side initiated, awaiting its whole lifetime.
start_stream_handler
¶
start_stream_handler(stream: WireStream, mode: Any = 'accept') -> Session
Begin serving a stream and return its session immediately.
a11.service.service.ServiceOptions
¶
ServiceOptions(*, session_options: SessionOptions | None = None, copy_registry_per_connection: bool = False, session_headers: Mapping[str, bytes] | None = None, drain_timeout: Duration | None = None)
Construct service options; all parameters are keyword-only.
copy_registry_per_connection
property
writable
¶
Give each connection its own copy of the registry. Leave false when the connection hook makes the copy itself.
drain_timeout
property
writable
¶
drain_timeout: Duration
How long draining waits for live sessions.
session_headers
property
writable
¶
Headers stamped on every session the service creates.
session_options
property
writable
¶
session_options: SessionOptions
Limits and timeouts for every session created.
Serving¶
serving binds a service to multiple transport
listeners, yields the active listeners, and performs ordered teardown on exit:
async with serving(service, websocket(ws_options), http_sse("0.0.0.0", 8012)):
await shutdown_event.wait()
a11.service.serving.serving
async
¶
serving(service: Service, *listeners: Listener, drain_timeout: Duration | None = None) -> AsyncIterator[list[Any]]
Bind listeners to service, yield them, then shut everything down.
Teardown is ordered: listeners are stopped first (in reverse order), so that no new connection can arrive while the service is draining the ones it already has.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
service
|
Service
|
The service to expose. |
required |
*listeners
|
Listener
|
Listener factories, e.g. |
()
|
drain_timeout
|
Duration | None
|
How long to wait for live sessions on the way out. |
None
|
Yields:
| Type | Description |
|---|---|
AsyncIterator[list[Any]]
|
The live listeners, in the order given. |
a11.service.serving.websocket
¶
A WebSocket listener bound to the service.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
options
|
WebSocketServerOptions
|
Where to listen, and how to frame accepted streams. |
required |
expose_descriptors
|
bool
|
Also answer |
True
|
a11.service.serving.http_sse
¶
http_sse(bind_address: str, port: int, options: HttpSseOptions | None = None, *, expose_descriptors: bool = True) -> Listener
An HTTP SSE listener bound to the service.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
bind_address
|
str
|
Local address to listen on. |
required |
port
|
int
|
Port to listen on; 0 requests an ephemeral one. |
required |
options
|
HttpSseOptions | None
|
Endpoint paths and transport tuning. |
None
|
expose_descriptors
|
bool
|
Also answer |
True
|
a11.service.serving.webrtc
¶
webrtc(signalling: SignallingTransport, configuration: WebRtcConfiguration | None = None, stream_options: WireStreamOptions | None = None) -> Listener
A WebRTC listener bound to the service.
Unlike the HTTP listeners this one does not open a port: peers find it
through signalling, and the server listens as whatever identity that
transport registered under. Pass a
WebSocketSignallingClient
connected to a signalling server to be reachable through it, or an endpoint
of an in-process
SignallingService to be reachable
within this process.
The server closes the transport when it is stopped, which serving does on
the way out along with every other listener.
CLI: a11 serve¶
The CLI command serves an action module over configured transports:
a11 serve mypkg.actions # REGISTRY over WebSocket
a11 serve mypkg.actions:TOOLS --ws --sse # a named registry, two endpoints
a11 serve ./examples/demo/main.py # a file, nothing installed
a11 serve mypkg.actions --webrtc \
--webrtc-signalling-server wss://a11.services/ice \
--webrtc-signalling-identity demoserver
a11 serve mypkg.app:SERVICE --ws --hosted demoserver # a Service, on the exchange
The symbol may be a Service as well as an ActionRegistry, which is how a
backend that specialises each connection -- a registry copy per caller, a
reverse-dispatch bridge bound to the session -- is served by this command rather
than by a loop of its own.
The
a11.demos.web_demos_server
module is a complete example. With --hosted, its actions are also available
through the exchange at a11.to/ui, without an inbound port.
a11.cli.commands.serve
¶
a11 serve: expose an Action registry from a module over one or more
transports.
a11 serve mypkg.actions # REGISTRY, WebSocket, defaults
a11 serve mypkg.actions:TOOLS --ws --sse # a named registry, two endpoints
a11 serve ./examples/demo/main.py # a file, nothing installed
The target names a module either way -- as an import path
(pkg.subpkg.module) or as a path to a .py file -- with an optional
:SYMBOL that defaults to REGISTRY. The module is loaded and the symbol
read; it has to be an
ActionRegistry. Writing one is the
annotated shape:
Or a whole Service¶
The symbol may also be a Service, which is how
a backend that needs more than a registry gets served by this command rather
than by a loop of its own. What "more" means in practice is an on_connection
hook: a registry copy per connection, a reverse-dispatch bridge onto it, and a
session where the caller's actions are registered. The
a11.demos.web_demos_server module uses these features. A Service arrives
with its listeners still unbound, so every transport and --hosted below works
on it unchanged.
One service, however many endpoints¶
Each --ws / --sse / --webrtc group adds a listener, and every
listener is bound to the same Service -- so an
action's state, its registry and its concurrency limits are shared no matter
which endpoint a caller arrived on. That is
serving's doing, which also stops the listeners
before draining the service, so nothing new arrives while it is finishing.
One endpoint per group: this is a command, not a load balancer.
Or as MCP tools¶
--mcp and --mcp-stdio publish the same registry to a host that speaks
Model Context Protocol rather than A11 -- Claude Desktop, an editor, another
agent runtime -- alongside whatever A11 transports are bound. Each action
becomes one tool, and --mcp-allow narrows which:
a11 serve mypkg.actions --mcp # Streamable HTTP on 8013
a11 serve mypkg.actions --mcp-stdio # launched by an MCP host
a11 serve mypkg.actions --ws --mcp --mcp-allow 'shell_.*'
--mcp-stdio gives the protocol this process's stdout, so the command prints
its summary to stderr instead and ends when the client hangs up. The MCP serving
guide documents the server API.
A file, or an import path¶
Which one was meant is read off the target: a .py suffix, a path separator,
a leading ./~, or a file that is simply there, and it is a path.
Otherwise it is imported the ordinary way. A file is loaded under its own stem
rather than as __main__, so a if __name__ == "__main__": block stays
asleep and the module's own entry point does not run; and its directory goes on
sys.path the way it would for python thatfile.py, so imports of its
siblings resolve.
HTTP protocol and TLS¶
--h11 (the default), --h2c and --h2 are mutually exclusive and apply to
every HTTP-based endpoint, as do --cert/--privkey. HTTP/1.1 is the default
for compatibility with RFC 6455 WebSocket clients, allowing browsers and
a11.client to use the same port.
SSE runs on HTTP/1.1 too, at the cost of a second connection: a connection
carries one request and the event stream has it, so the outbound direction gets
one of its own. See
HttpSseWireStream.
ServeError
¶
Bases: Exception
A configuration error that the CLI reports without a traceback.
Served
dataclass
¶
What a target resolved to, and where it came from.
Attributes:
| Name | Type | Description |
|---|---|---|
registry |
ActionRegistry | None
|
The actions to serve. Read off the service when the module gave one, so the count this command prints is the same either way. |
service |
Any | None
|
The module's own service, or |
module_path |
str
|
The module as it was named on the command line. |
symbol |
str
|
The symbol the above was read from. |
split_target
¶
Split MODULE[:SYMBOL] into the module part and the symbol.
Split from the right, and only where the tail is an identifier, so a
Windows path keeps its drive letter (C:\src\actions.py is all module)
and a dotted path without a symbol keeps all its dots.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
target
|
str
|
The command's positional argument. |
required |
Returns:
| Type | Description |
|---|---|
tuple[str, str]
|
The module part, and the symbol -- |
is_path_target
¶
Whether module names a file rather than an import path.
A .py suffix, a path separator, a leading . or ~, or a file that
is simply there. main.py is a file, not the py submodule of a
main package, because nobody has ever meant the latter.
resolve_served
¶
resolve_served(target: str) -> Served
Load MODULE[:SYMBOL] and return what it serves.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
target
|
str
|
An import path ( |
required |
Returns:
| Type | Description |
|---|---|
Served
|
The |
Raises:
| Type | Description |
|---|---|
ServeError
|
If the module cannot be loaded, has no such symbol, or the
symbol is neither an |
resolve_registry
¶
resolve_served, for a caller that wants the registry alone.
Raises:
| Type | Description |
|---|---|
ServeError
|
As |
configure
¶
Declare this command's flags on parser.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
parser
|
ArgumentParser
|
The parser to add to. |
required |
with_target
|
bool
|
Whether to add the positional target. A module that serves
itself through this machinery, such as
|
True
|
http2_options
¶
Transport options shared by every HTTP-based endpoint.
Each protocol is set explicitly rather than left at its default, because
Http2Options enables all three and "HTTP/1.1" has to mean only that --
otherwise a request that could be upgraded silently is, and the endpoint no
longer speaks what it was asked to.
HTTP/1.1 is the default for both transports. WebSocket needs it for RFC
6455; SSE runs over it by giving its outbound direction a connection of its
own, since an HTTP/1.1 connection carries one request and the event stream
has it. See
HttpSseWireStream.
Raises:
| Type | Description |
|---|---|
ServeError
|
For a combination that cannot be served: |
serve
async
¶
Import the registry, bind the listeners, run until signalled.
run
async
¶
Serve what args asks for, reporting a mistake rather than raising.
Public because it is the whole command: a module that declares its flags
with configure runs them through here, and so gets the same reporting and
the same exit codes as a11 serve itself.