Skip to content

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.

deadline property

deadline: Time

The absolute time after which the session will be aborted.

done property

done: _DoneEvent

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.

id property

id: str

The session's unique identifier string.

node_map property writable

node_map: NodeMap

The NodeMap backing this session's node state; assigning replaces it.

abort

abort(status: Status) -> None

Abort the session immediately with the given error status, cancelling streams and actions.

Examples:

Propagate an authentication failure to the peer:

session.abort(Status(
    code=StatusCode.PERMISSION_DENIED,
    message=str(error),
))

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

add_done_callback(callback: Callable[['Session'], Any | Awaitable[Any]]) -> Task

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:

stream_lifetime = session.add_stream(websocket_stream)

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

cancel_action(action_id: str) -> None

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

cancel_all_actions() -> None

Request cancellation of every action currently running in the session. Each action unwinds asynchronously; await await_all_actions to observe completion.

dispatch_action

dispatch_action(action: Any) -> Future[None]

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_id

get_id() -> str

Return the session's unique identifier string.

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

half_close() -> None

Signal that this side will send no more messages, allowing the session to drain and finish once peers do the same. Remaining inbound messages continue to be processed asynchronously.

Examples:

Finish an exchange after sending the last message:

session.half_close()
await session.done.wait()

is_closed

is_closed() -> bool

Return whether the session has been closed and no longer accepts new streams or messages.

is_done

is_done() -> bool

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

Enqueue a wire message for delivery on the named stream (or the default stream), raising on failure. Delivery happens asynchronously as the stream drains.

Examples:

Route a response through a particular attached transport:

session.send(response, stream_id=websocket_stream.get_id())

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

wait_done() -> Future[None]

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

receive(deadline=None)

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:

while message := await session.receive():
    await route_message(message)

receive_with_stream_id async

receive_with_stream_id(deadline=None)

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:

while item := await session.receive_with_stream_id():
    message, stream_id = item
    await route_message(message, source=stream_id)

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.

deadline property writable

deadline: Time

Absolute time after which the session is aborted.

max_buffered_bytes_per_stream property writable

max_buffered_bytes_per_stream: int

Maximum bytes buffered per stream.

max_buffered_bytes_total property writable

max_buffered_bytes_total: int

Maximum total bytes buffered across all streams.

max_buffered_messages_per_stream property writable

max_buffered_messages_per_stream: int

Maximum number of messages buffered per stream.

max_buffered_messages_total property writable

max_buffered_messages_total: int

Maximum number of messages buffered across all streams.

max_concurrent_nested_actions property writable

max_concurrent_nested_actions: int

Maximum number of concurrently running nested actions.

max_concurrent_root_actions property writable

max_concurrent_root_actions: int

Maximum number of concurrently running root actions.

max_single_message_size property writable

max_single_message_size: int

Maximum size in bytes of a single wire message.

no_stream_timeout property writable

no_stream_timeout: Duration

How long the session waits with no active stream before finishing.

validate

validate() -> None

Validate the option values, raising on invalid configuration.

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)

accepting property

accepting: bool

Whether new connections are still admitted.

action_registry property

action_registry: ActionRegistry | None

The template registry new connections are built from.

done property

done: _ServiceDoneEvent

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.

session_count property

session_count: int

How many sessions are being served.

abort

abort(status: Status) -> None

Stop accepting and abort every live session.

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(name: str = '', query: str = '') -> str

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 waits indefinitely.

None

Raises:

Type Description
StatusException

DEADLINE_EXCEEDED when sessions remain after timeout; they are left running, so follow with abort if they must go.

expose_descriptors_on

expose_descriptors_on(options: DescribeEndpointOptions) -> None

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

session_ids

session_ids() -> list[str]

The ids of the sessions currently being served.

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.

stop_accepting

stop_accepting() -> None

Refuse new connections, leaving live ones alone.

wait_done

wait_done() -> Future[None]

Await the service being closed and empty.

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

copy_registry_per_connection: bool

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

session_headers: dict[str, bytes]

Headers stamped on every session the service creates.

session_options property writable

session_options: SessionOptions

Limits and timeouts for every session created.

validate

validate() -> None

Validate the options, raising on error.

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. websocket. None is valid -- a service with no listener still serves streams handed to it directly.

()
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

websocket(options: WebSocketServerOptions, *, expose_descriptors: bool = True) -> Listener

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 GET /actions on this port. A WebSocket client can ask __list_actions__ over the stream it already has; this is for whoever has the port number and no A11 client -- a curl, a health check, a person. Pass False to leave the port answering nothing but the upgrade.

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 GET /actions. On by default: an HTTP service is a thing people read with a browser, and the document is the same one __list_actions__ returns because the same describer produces both.

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:

REGISTRY = ActionRegistry()

@REGISTRY.action
async def summarise(document: str) -> str:
    ...

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

Served(registry: ActionRegistry | None, service: Any | None, module_path: str, symbol: str)

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 None to build the plain one.

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_target(target: str) -> tuple[str, str]

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 -- DEFAULT_SYMBOL if none was given.

is_path_target

is_path_target(module: str) -> bool

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 (pkg.subpkg.module) or a path to a .py file, optionally with :SYMBOL; the symbol defaults to REGISTRY. See is_path_target for how the two are told apart.

required

Returns:

Type Description
Served

The Served the symbol named -- a registry, or a whole service.

Raises:

Type Description
ServeError

If the module cannot be loaded, has no such symbol, or the symbol is neither an ActionRegistry nor a Service.

resolve_registry

resolve_registry(target: str) -> tuple[ActionRegistry, str, str]

resolve_served, for a caller that wants the registry alone.

Raises:

Type Description
ServeError

As resolve_served, and for a target that names a service rather than a registry.

configure

configure(parser: ArgumentParser, *, with_target: bool = True) -> None

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 a11.demos.web_demos_server, knows its own target and takes every transport flag below. It requests the flags without the positional argument and supplies target itself, keeping the transport flags consistent across entry points.

True

http2_options

http2_options(args: Namespace) -> Http2Options

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: --h2 without a certificate, --h2c with one, or half a TLS identity.

serve async

serve(args: Namespace) -> int

Import the registry, bind the listeners, run until signalled.

run async

run(args: Namespace) -> int

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.