Transports¶
A WireStream is A11's transport abstraction:
a bidirectional, message channel connecting two peers. Concrete implementations
support in-process channels, WebSocket, HTTP SSE, and WebRTC data channels.
Transport contract¶
Each endpoint starts once in one role: start() for the initiator or
accept() for the responder. The message callback receives None when the
peer half-closes its write direction. The local endpoint can continue sending
until it also half-closes.
send() validates and admits a WireMessage to bounded transport state. A
successful return does not mean that the peer has received it. Half-close
after the final send, then drain the outgoing direction when application
shutdown depends on delivery through the local transport buffers. Full
completion occurs after both directions finish or an error aborts the stream.
WireStream does not provide global message ordering. Ordered action data uses
sequenced NodeFragment records, which the receiving node reconstructs. A
session can also attach several streams whose callbacks advance independently.
See the WireStream lifecycle for the complete
shutdown and failure contract.
Transport selection¶
| Implementation | Boundary |
|---|---|
| In-process | Two endpoints in one process, including tests and internal bridges |
| WebSocket | Long-lived client/server connections over HTTP/1.1 or HTTP/2 |
| HTTP SSE | Browser-compatible HTTP request and event-stream routes |
| WebRTC | Browser or native peers connected through signalling and ICE |
All implementations carry the same WireMessage format. Sessions and actions
therefore retain their protocol when deployment changes transport.
WireStream example¶
WireStreams connect endpoints, deliver serialized frames, and handle independent bidirectional shutdowns:
import a11
client_stream, server_stream = a11.create_in_process_wire_stream_pair()
async def on_message(msg):
if msg is None:
print("Peer closed write direction")
return
print("Received:", msg)
async def on_done(status):
print("Stream finished:", status)
await server_stream.accept(on_message, on_done)
await client_stream.start(on_message, on_done)
a11.net.wire_stream.WireStream
¶
Construct the abstract WireStream base. Subclass this in Python to implement a custom asynchronous, bidirectional transport for an agent; the abstract operations (send, start/accept, get_status, get_trailers, ...) are dispatched to your overrides.
deadline
property
¶
deadline: Time
The stream's current absolute deadline, after which it is automatically aborted.
accept
¶
accept(on_message: Callable[[WireMessage | None], Any], on_done: Callable[[], Any]) -> Future[None]
Begin driving the stream as the responding (server) side, delivering each inbound message to the asynchronous on_message callback and end-of-stream to on_done. Use this instead of start() when this endpoint is answering an incoming agent connection. Returns an awaitable that resolves when acceptance completes; use on_done as the terminal barrier.
drain_outgoing_messages
¶
get_id
¶
Return the stream's stable identifier, which also seeds its tracing trace id.
get_impl
¶
Return an opaque native handle to the underlying implementation, or None. Intended for advanced interop, not normal agent code.
get_status
¶
get_status() -> Status
Return the stream's terminal status once it has finished, or OK while it is still active. Inspect this after the stream completes to learn whether the agent exchange succeeded or failed.
get_trailers
¶
Return the trailers (final metadata) the peer sent at half-close, or None if none were received. Read this after the stream ends to recover end-of-turn metadata from the agent exchange.
half_close
¶
send
¶
send(message: WireMessage) -> None
set_deadline
¶
set_deadline(deadline: Time | None = None) -> None
Set an absolute wall-clock deadline after which the stream is automatically aborted; pass None to clear it.
start
¶
start(on_message: Callable[[WireMessage | None], Any], on_done: Callable[[], Any]) -> Future[None]
a11.net.wire_stream.WireStreamOptions
¶
WireStreamOptions(max_buffered_incoming_messages: SupportsInt | None = 100, max_single_message_size: SupportsInt | None = 33554432, max_buffered_incoming_bytes: SupportsInt | None = 33554432, message_timeout_millis: Any | None = None, deadline: Time | None = None)
Construct wire-stream options controlling buffering and timeouts for an agent stream. All arguments are keyword-friendly and validated on construction.
deadline
property
writable
¶
deadline: Time
Absolute wall-clock deadline after which the stream is aborted.
max_buffered_incoming_bytes
property
writable
¶
Maximum total bytes of buffered inbound messages before backpressure is applied.
max_buffered_incoming_messages
property
writable
¶
Maximum number of inbound messages buffered before backpressure is applied.
max_single_message_size
property
writable
¶
Maximum size, in bytes, of a single wire message.
message_timeout
property
writable
¶
message_timeout: Duration
Per-message inactivity timeout as a duration.
In-process¶
a11.net.in_process_wire_stream.InProcessWireStream
¶
Bases: WireStream
create_pair
staticmethod
¶
create_pair(options: WireStreamOptions | None = None, first_options: WireStreamOptions | None = None, second_options: WireStreamOptions | None = None) -> tuple[InProcessWireStream, InProcessWireStream]
Create a connected pair of in-process wire streams that talk to each other directly in memory, with no network involved. One endpoint drives start() while the other drives accept(). Pass shared options, or per-endpoint first_options/second_options, to tune buffering and timeouts.
wait
¶
Await until this in-process stream has fully finished. Block on this to know a local agent exchange has completed before tearing the pair down.
a11.net.in_process_wire_stream.create_in_process_wire_stream_pair
¶
create_in_process_wire_stream_pair(options: WireStreamOptions | None = None, *, first_options: WireStreamOptions | None = None, second_options: WireStreamOptions | None = None) -> tuple[InProcessWireStream, InProcessWireStream]
WebSocket¶
stream = a11.WebSocketWireStream.connect("ws://127.0.0.1:8080/ws")
server = a11.WebSocketWireServer.create(accept_callback, port=8080)
a11.net.websocket_wire_stream.WebSocketWireStream
¶
Bases: WireStream
request_headers
property
¶
The headers the accepted request carried, as (name, value) pairs, or empty for a client stream. A per-connection credential arrives here, which is what lets a server authenticate a stream rather than a port.
request_path
property
¶
The path this stream was accepted on, query string included, or empty for a client stream. On a server accepting under WebSocketServerOptions.path_prefix this is the only place the rest of the path survives, and so the only way one port can serve more than one thing.
connect
staticmethod
¶
connect(url: str, options: WireStreamOptions = ..., websocket_options: WebSocketClientOptions = ...) -> WebSocketWireStream
Open a client WebSocket connection to url and return a WireStream over it. This is the standard way for an agent to dial out to a remote A11 endpoint; the returned stream is then driven asynchronously via start()/send(). Tune transport buffering with options and the handshake (headers, framing, HTTP/2, TLS) with websocket_options.
a11.net.websocket_wire_stream.WebSocketWireServer
¶
port
property
¶
The actual TCP port the server is listening on, resolved even when an ephemeral port (0) was requested.
create
staticmethod
¶
create(on_stream: Any, options: WebSocketServerOptions = ...) -> WebSocketWireServer
Start a WebSocket server that accepts incoming A11 connections, invoking the asynchronous on_stream callback with a fresh WireStream for each accepted client. This is the server-side entry point for hosting an agent: each callback runs concurrently and typically drives accept() on its stream. Configure the listen address, port, path and TLS via options.
get_impl
¶
Return an opaque native handle to the underlying implementation, or None. Intended for advanced interop.
stop
¶
Stop the server and close the listening socket, releasing the bound port. Call this to shut the agent host down cleanly; it blocks until shutdown completes.
HTTP SSE¶
a11.net.http_sse_wire_stream.HttpSseWireStream
¶
Bases: WireStream
request_path
property
¶
The path a server stream was accepted on, query string included, or empty for a client stream. On a server accepting under HttpSseOptions.connect_endpoint_prefix this is the only place the rest of the path survives.
get_http_request_headers
¶
Return the HTTP headers carried on the underlying SSE request.
get_http_response_headers
¶
Return the HTTP response headers negotiated for the SSE connection, or None if they have not arrived yet. Because the connection is established asynchronously, prefer awaiting wait_for_http_headers() before relying on this value.
set_http_request_headers
¶
Set the HTTP headers to send on the underlying SSE request. Call this before the stream connects to attach auth or routing metadata that your agent's transport needs.
set_http_response_headers
¶
Set the HTTP headers to send on the SSE response. Used on the server side to attach transport metadata before the streaming response is flushed to the client.
wait_for_http_headers
¶
Await the exchange of HTTP headers for the SSE connection. Because SSE wire streams connect asynchronously, await this future before reading response headers or assuming the stream is live.
a11.net.http_sse_wire_stream.HttpSseServer
¶
create
staticmethod
¶
create(bind_address: str = '127.0.0.1', port: SupportsInt = 0, on_connect: Any | None = None, options: HttpSseOptions = ...) -> HttpSseServer
Create and start an SSE server that accepts A11 wire streams, invoking the optional async on_connect callback for each client.
wait_for_stream
¶
Await the next incoming SSE wire stream from a connecting client.
WebRTC¶
a11.net.webrtc_wire_stream.WebRtcWireStream
¶
Bases: WireStream
data_channel
property
¶
Opaque capsule around the underlying libdatachannel DataChannel. Exposed for advanced interop and diagnostics; agent code normally reads and writes through the WireStream API rather than touching this directly.
peer_connection
property
¶
Opaque capsule around the underlying libdatachannel PeerConnection, for inspecting ICE/connection state during debugging.
signalling_endpoint
property
¶
Signalling transport this stream negotiated over, the channel that carried the asynchronous SDP/ICE handshake.
create_client
staticmethod
¶
create_client(identity: str, peer_identity: str, signalling: SignallingService, configuration: WebRtcConfiguration = ..., options: WireStreamOptions = ...) -> WebRtcWireStream
create_client(peer_identity: str, signalling: SignallingTransport, configuration: WebRtcConfiguration = ..., options: WireStreamOptions = ...) -> WebRtcWireStream
create_client(identity: str, peer_identity: str, signalling: SignallingService, configuration: WebRtcConfiguration = ..., options: WireStreamOptions = ...) -> WebRtcWireStream
Open a WebRTC data-channel wire stream to a named peer over a shared in-process signalling service. It performs the ICE/SDP handshake and resolves to a WireStream carrying A11-framed messages, fragmenting large payloads transparently.
a11.net.webrtc_wire_stream.WebRtcWireServer
¶
signalling_endpoint
property
¶
Signalling endpoint the server negotiates over.
create
staticmethod
¶
create(identity: str, signalling: SignallingService, on_stream: Any, configuration: WebRtcConfiguration = ..., stream_options: WireStreamOptions = ...) -> WebRtcWireServer
create(signalling: SignallingTransport, on_stream: Any, configuration: WebRtcConfiguration = ..., stream_options: WireStreamOptions = ...) -> WebRtcWireServer
create(identity: str, signalling: SignallingService, on_stream: Any, configuration: WebRtcConfiguration = ..., stream_options: WireStreamOptions = ...) -> WebRtcWireServer
Create a WebRTC server that accepts peer connections and invokes the async on_stream callback with each new WebRtcWireStream, negotiating over a signalling service shared within this process.
Configured telemetry records one a11.wire_stream lifecycle span per endpoint
with its stream ID, terminal status, and send-size events. See
observability.
Signalling¶
Signalling provides out-of-band coordination for WebRTC peer connections:
a11.net.signalling.WebSocketSignallingServer
¶
create
staticmethod
¶
create(service: SignallingService, options: WebSocketSignallingServerOptions = ...) -> WebSocketSignallingServer
Create a WebSocket signalling server that fronts the given in-process signalling service.
disconnect
¶
Close one identity's connection, if this server holds it. The counterpart to admission: whatever authorised a registration can be withdrawn, and the socket has to go with it rather than surviving until its next message. Raises NOT_FOUND when this server is not holding that identity.
get_impl
¶
Opaque capsule around the native implementation, for interop.
a11.net.signalling.WebSocketSignallingClient
¶
Bases: SignallingTransport
connect
staticmethod
¶
connect(url: str, identity: str, on_message: Any | None = None, options: WebSocketSignallingClientOptions = ...) -> Future[WebSocketSignallingClient]
Asynchronously connect to a WebSocket signalling server, resolving to a client once registered under the given identity.
get_impl
¶
Opaque capsule around the native implementation, for interop.
a11.net.signalling.SignallingService
¶
Create a new in-process signalling service.
connect
¶
Register an identity and its async inbound-message callback, returning a signalling endpoint.
contains
¶
Return whether the given identity is currently connected.
deliver
¶
Deliver a message to a locally connected recipient, as though it had been routed from an endpoint of this service. This is the ingress half of a federated signalling fabric: pair it with WebSocketSignallingServerOptions.on_unroutable, which is the egress half, to make several servers behave as one. Raises NOT_FOUND when the recipient is not connected here.
HTTP/2 Primitives¶
a11.net.http2.Http2Client
¶
multiplexed
property
¶
Whether this connection can carry several exchanges at once: true for HTTP/2, false for HTTP/1.1, which A11 limits to one request per connection.
connect
staticmethod
¶
connect(host: str, port: SupportsInt, options: Http2Options = ...) -> Future[Http2Client]
Asynchronously connect to an HTTP/2 server, returning a future that resolves to the connected client.
extended_connect
¶
extended_connect(protocol: str, path: str, headers: Iterable[tuple[str, str]] | None = None, scheme: str = '') -> Http2DuplexStream
Open an extended CONNECT duplex stream for bidirectional data.
get_impl
¶
Return an opaque capsule wrapping the native client handle.
request
¶
request(method: str, path: str, headers: Iterable[tuple[str, str]] | None = None, body: Any = b'', scheme: str = '') -> Future[HttpResponse]
Send a request and await the full buffered response.
request_stream
¶
request_stream(method: str, path: str, headers: Iterable[tuple[str, str]] | None = None, body: Any = b'', scheme: str = '') -> Http2ResponseStream
Open a request and return a pull-oriented response stream for reading the response body incrementally.
request_streaming_body
¶
request_streaming_body(method: str, path: str, headers: Iterable[tuple[str, str]] | None = None, scheme: str = '') -> Http2DuplexStream
Open a request whose body is written incrementally afterwards, for an upload of unknown or unbounded length. Returns a duplex stream: write() sends more of the body, finish() ends it, and the response is read from the same handle. Do not set content-length.
a11.net.http2.Http2Server
¶
create
staticmethod
¶
create(bind_address: str = '127.0.0.1', port: SupportsInt = 0, handler: Any | None = None, options: Http2Options = ...) -> Http2Server
Create and start an HTTP/2 server bound to the given address and port, dispatching each request to the async handler.
get_impl
¶
Return an opaque capsule wrapping the native server handle.