Skip to content

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

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.

abort

abort(status: Status) -> None

Terminate the stream immediately with an error status, discarding buffered messages and propagating failure to the peer and pending receivers.

Examples:

End an exchange when its upstream disappears:

stream.abort(Status(
    code=StatusCode.UNAVAILABLE,
    message="upstream connection was lost",
))

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

drain_outgoing_messages() -> Future[None]

Await until every queued outbound message has been handed to the transport. Call half_close first so buffered output is not dropped.

Examples:

Use the transport delivery barrier during orderly shutdown:

stream.half_close()
await stream.drain_outgoing_messages()

get_id

get_id() -> str

Return the stream's stable identifier, which also seeds its tracing trace id.

get_impl

get_impl() -> CapsuleType | None

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

get_trailers() -> dict[str, bytes] | None

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

half_close(trailers: Mapping[str, bytes] | None = None) -> None

Signal that this endpoint has finished sending, optionally attaching trailers. The stream stays open for inbound messages.

Examples:

End the local half and wait until queued messages reach the transport:

stream.half_close()
await stream.drain_outgoing_messages()

send

send(message: WireMessage) -> None

Queue a message for asynchronous delivery to the peer. This call is non-blocking: the message enters the ordered outbound queue and the transport applies backpressure.

Examples:

Admit a request before closing the local sending side:

stream.send(request_message)

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]

Begin driving the stream as the initiating side, delivering inbound messages to on_message and completion to on_done. Callbacks are awaited as data arrives.

Examples:

Start a client transport with application callbacks:

await stream.start(on_message, on_transport_done)

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

max_buffered_incoming_bytes: int

Maximum total bytes of buffered inbound messages before backpressure is applied.

max_buffered_incoming_messages property writable

max_buffered_incoming_messages: int

Maximum number of inbound messages buffered before backpressure is applied.

max_single_message_size property writable

max_single_message_size: int

Maximum size, in bytes, of a single wire message.

message_timeout property writable

message_timeout: Duration

Per-message inactivity timeout as a duration.

message_timeout_millis property writable

message_timeout_millis: Duration

Per-message inactivity timeout expressed in milliseconds.

validate

validate() -> None

Validate the options, raising on invalid configuration.

In-process

a11.net.in_process_wire_stream.InProcessWireStream

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

wait() -> Future[None]

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

WebSocketWireStream()

Bases: WireStream

request_headers property

request_headers: list[tuple[str, str]]

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

request_path: str

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

port: int

The actual TCP port the server is listening on, resolved even when an ephemeral port (0) was requested.

running property

running: bool

Whether the server is currently accepting connections.

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

get_impl() -> CapsuleType | None

Return an opaque native handle to the underlying implementation, or None. Intended for advanced interop.

stop

stop() -> None

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

HttpSseWireStream()

Bases: WireStream

request_path property

request_path: str

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

get_http_request_headers() -> list[tuple[str, str]]

Return the HTTP headers carried on the underlying SSE request.

get_http_response_headers

get_http_response_headers() -> list[tuple[str, str]] | None

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_http_request_headers(headers: Iterable[tuple[str, str]] | None) -> None

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_http_response_headers(headers: Iterable[tuple[str, str]] | None) -> None

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

wait_for_http_headers() -> Future[None]

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

http2_server property

http2_server: Http2Server

The underlying HTTP/2 server.

port property

port: int

The port the server is listening on.

running property

running: bool

Whether the server is currently running.

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.

stop

stop() -> None

Stop the server and release its resources.

wait_for_stream

wait_for_stream() -> Future[HttpSseServerWireStream]

Await the next incoming SSE wire stream from a connecting client.

WebRTC

a11.net.webrtc_wire_stream.WebRtcWireStream

WebRtcWireStream()

Bases: WireStream

data_channel property

data_channel: CapsuleType | None

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

peer_connection: CapsuleType | None

Opaque capsule around the underlying libdatachannel PeerConnection, for inspecting ICE/connection state during debugging.

signalling_endpoint property

signalling_endpoint: SignallingTransport

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

identity property

identity: str

Local identity this server listens as.

pending_peer_count property

pending_peer_count: int

Number of peers still completing negotiation.

running property

running: bool

Whether the server is currently running.

signalling_endpoint property

signalling_endpoint: SignallingTransport

Signalling endpoint the server negotiates over.

stop

stop() -> None

Stop the server and stop accepting new peer connections.

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

port property

port: int

Port the server is listening on.

running property

running: bool

Whether the server is currently running.

service property

Signalling service this server fronts.

create staticmethod

create(service: SignallingService, options: WebSocketSignallingServerOptions = ...) -> WebSocketSignallingServer

Create a WebSocket signalling server that fronts the given in-process signalling service.

disconnect

disconnect(identity: str) -> None

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

get_impl() -> CapsuleType | None

Opaque capsule around the native implementation, for interop.

stop

stop() -> None

Stop the server and close all client connections.

a11.net.signalling.WebSocketSignallingClient

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

get_impl() -> CapsuleType | None

Opaque capsule around the native implementation, for interop.

a11.net.signalling.SignallingService

SignallingService()

Create a new in-process signalling service.

create staticmethod

create() -> SignallingService

Create a new in-process signalling service.

connect

connect(identity: str, on_message: Any) -> SignallingEndpoint

Register an identity and its async inbound-message callback, returning a signalling endpoint.

contains

contains(identity: str) -> bool

Return whether the given identity is currently connected.

deliver

deliver(message: SignallingMessage) -> None

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.

identities

identities() -> list[str]

Return the list of currently connected identities.

stop

stop() -> None

Stop the service and disconnect all endpoints.

HTTP/2 Primitives

a11.net.http2.Http2Client

connected property

connected: bool

Whether the client is currently connected.

host property

host: str

The host the client is connected to.

multiplexed property

multiplexed: bool

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.

port property

port: int

The port the client is connected to.

secure property

secure: bool

Whether the connection is using TLS.

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.

close

close() -> None

Close the client connection.

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

get_impl() -> CapsuleType | None

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

bind_address property

bind_address: str

The address the server is bound to.

port property

port: int

The port the server is listening on.

running property

running: bool

Whether the server is currently running.

secure property

secure: bool

Whether the server is using TLS.

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

get_impl() -> CapsuleType | None

Return an opaque capsule wrapping the native server handle.

stop

stop() -> None

Stop the server and release its resources.