Nodes¶
An AsyncNode is A11's unit of streaming state:
an ordered, asynchronous sequence of chunks backed by a ChunkStore.
Nodes carry data between action ports and across network transports.
Using AsyncNode¶
Creating and writing¶
Create a node with create and write items
sequentially with put. A write has two
asynchronous stages: admission to the bounded writer, then acceptance by the
backing store. Await both when later work depends on durable acceptance:
import a11
node = a11.AsyncNode.create("events")
confirmation = await node.put({"event": "start"})
sequence = await confirmation
The first await applies local backpressure. The confirmation reports store acceptance and any local transport failure; it is not an acknowledgement from a remote reader. Code that will drain the node before shutdown can omit the second await.
finalize records the logical end
of the data and closes the writer:
await node.put({"event": "progress", "percent": 50})
await node.finalize({"event": "complete", "percent": 100})
finalize(value) makes value the final visible record. finalize() writes
an invisible final marker after previously admitted records. Pass wait=True
when the process must remain alive until finality and closure reach the store.
Use abort_with_status() when partial output must be reported as a failure.
Iterating streams and reading unary values¶
Read items one-by-one with async for or
next. For an action returning one
complete value, use
consume:
next() returns one independent value and None at the end of a successful
stream. consume() requires one logically complete value, including the final
marker, and rejects an incomplete unary result. Use next_chunk() or
next_fragment() when code needs media metadata, serialization tags, sequence
numbers, or routing fields.
The AsyncNode lifecycle defines finality, closure, reset, replay, and failure in detail.
Values and encoded chunks¶
put() serializes an application value through the node's registry.
put_chunk() writes bytes that already have their final representation, such
as a PNG or an HTTP body. The chunk metadata must state the media type because
readers use that metadata to interpret the bytes. The
data reference describes representations and type tags.
Pre-encoded binary assets use chunks directly. The reader uses
next_chunk() and an explicit size limit.
a11.nodes.async_node.AsyncNode
¶
AsyncNode(chunk_store: ChunkStore, node_map: NodeMap | None = None, *, serialization_registry: SerializationRegistry | None = None, reader_options: ChunkStoreReaderOptions | dict[str, Any] | None = None, writer_options: ChunkStoreWriterOptions | dict[str, Any] | None = None)
Build a node over chunk_store.
Prefer create, which constructs the store for you from a node
id. Pass reader_options / writer_options (as
ChunkStoreReaderOptions / ChunkStoreWriterOptions or plain dicts) to
tune buffering and ordering, and a custom serialization_registry to
control how Python objects map to chunks.
chunk_store
property
¶
chunk_store: ChunkStore
The underlying chunk store backing this node (see get_chunk_store).
reader_options
property
writable
¶
reader_options: ChunkStoreReaderOptions
Options controlling how this node reads from its chunk store, such as buffering and flow control.
serialization_registry
property
writable
¶
The registry used to (de)serialize Python objects for this node.
writer_options
property
writable
¶
writer_options: ChunkStoreWriterOptions
Options controlling how this node writes to its chunk store, such as buffering and flow control.
expect_types
¶
expect_types(**kwds) -> Iterator[AsyncNode]
Temporarily set the expected read types for the with block.
create
classmethod
¶
create(node_id: str, node_map: NodeMap | None = None, *, serialization_registry: SerializationRegistry | None = None, reader_options: ChunkStoreReaderOptions | dict[str, Any] | None = None, writer_options: ChunkStoreWriterOptions | dict[str, Any] | None = None, chunk_store_factory: Callable[[str], ChunkStore] = LocalChunkStore) -> AsyncNode
Create a standalone node identified by node_id.
chunk_store_factory builds the backing store from the id; it
defaults to an in-memory
LocalChunkStore,
so overriding it is how you place a node's data in a different backend.
Examples:
Create a stream used to deliver answer fragments:
abort_with_status
¶
abort_with_status(status: Status) -> Future[None]
Aborts the stream with the given error status and returns a future that resolves once the abort has propagated. Readers then observe the error rather than a normal end-of-stream.
attach_stream
¶
attach_stream(stream: WireStream) -> None
Attaches a wire stream so this node's chunks are mirrored over the network transport. The stream is kept alive for the node's lifetime.
cancel
¶
Cancels both the reader and the writer, tearing down all pending streaming operations on the node at once.
cancel_reader
¶
Cancels the node's reader, unblocking any pending next-chunk or next-fragment awaits on the read side of the stream.
cancel_writer
¶
Cancels the node's writer, unblocking any pending put or drain awaits on the producing side of the stream.
close
async
¶
Flush buffered writes and close the writer, marking nothing final.
The specialised half of finalize: closure without finality, for a
producer that cannot say which chunk was the last one -- a log, say --
but can say that no more are coming. Closing always drains, so this
resolves once the backing store is closed.
detach_stream
¶
detach_stream(stream: WireStream) -> None
Detaches a previously attached wire stream so the node stops mirroring its chunks over that transport.
finalize
async
¶
finalize(value: Any | None = None, seq: int | None = None, mimetype: str = '', wait: bool = False, close: bool = True) -> None
End the stream: mark the logical end, and close the writer.
The one call an ordinary producer needs. value is written as the
final fragment; passing nothing (or None) writes a null terminator
instead, which is the form to use once the last visible value has
already gone out with put. Unless close is cleared the writer is
closed as well, so readers waiting on data that can no longer arrive
are released and a peer's mirror of the node closes too.
It does not wait. The write and the close are carried out by the
writer's pump, which keeps running after this coroutine -- and the
enclosing action -- has returned, so a producer can finalise and walk
away. Nothing is swallowed: a failed write or close is logged and stays
visible through get_writer_status. Pass wait=True when the
producer must know the store accepted the end of the stream before
going on, or await wait_for_buffer_to_drain separately.
value may be a NodeFragment (whose
own seq is used), a Chunk, or any object the
node's serialization registry can encode.
Examples:
Finish a streamed answer whose last token is not known in advance:
Deliver a unary result, which is one chunk and one call:
get_chunk_store
¶
get_chunk_store() -> ChunkStore
Returns the underlying chunk store backing this node. The chunk store is the ordered storage boundary that the node's reader and writer stream through; reach for it when you need lower-level access than the async put/next API provides.
get_reader_options
¶
get_reader_options() -> ChunkStoreReaderOptions
Return a copy of the reader's current options.
get_reader_status
¶
get_reader_status() -> Status
Returns the current status of the node's reader. Check it to tell whether the read end of the stream is healthy, has completed, or has failed while streaming.
get_writer_abort_status
¶
get_writer_abort_status() -> Status | None
Returns the status the writer was aborted with, or None if the writer has not been aborted.
get_writer_options
¶
get_writer_options() -> ChunkStoreWriterOptions
Return a copy of the writer's current options.
get_writer_status
¶
get_writer_status() -> Status
Returns the current status of the node's writer. Check it to tell whether the producing end of the stream is healthy, has completed, or has failed while streaming.
is_writable
¶
Returns a future that resolves once it is known whether the node can currently accept writes. Await it before producing chunks to respect backpressure rather than blocking a busy stream.
iter_chunks
¶
Async-iterate raw chunks until the stream ends.
iter_fragments
¶
iter_fragments(timeout: Duration | None = None) -> AsyncIterator[NodeFragment]
Async-iterate raw fragments until the stream ends.
iter_values
¶
iter_values(obj_type: type[T] | None = None, timeout: Duration | None = None, mimetype_patterns: str | Sequence[str] = '') -> AsyncIterator[Any]
Async-iterate deserialized values, a batch of fragments per await.
The batched counterpart to async for value in node, which reads one
fragment per await and so pays an event-loop turn per value -- on a
selector loop that turn is a syscall, and it dominates everything else
about a value. This asks for ITER_BATCH fragments at a time, exactly
as iter_fragments does, and deserializes each.
Prefer this when the reader will iterate to the end. Keep
async for when something else may read the same node: fragments in
this iterator's batch have already left the reader, so abandoning it
part way through a batch abandons them -- the same hazard
iter_fragments carries, and the reason __anext__ was left reading
one at a time.
iter_with_deadline
¶
iter_with_deadline(deadline: Time)
Async-iterate deserialized values until deadline or end of
stream.
next
async
¶
next(obj_type: type[T] | None = None, timeout: Duration | None = None, mimetype_patterns: str | Sequence[str] = '') -> T | None
next_chunk
async
¶
Read the next raw chunk, or None at end of stream.
next_fragment
async
¶
next_fragment(timeout: Duration | None = None) -> NodeFragment | None
Read the next raw fragment, or None at end of stream.
next_fragments
¶
next_fragments(limit: SupportsInt, timeout: Duration | None = None) -> Future[list[NodeFragment | None]]
Returns a future resolving to a list of up to limit fragments, with a trailing None at end-of-stream. The batched counterpart to next_fragment, and the one to prefer when draining: every await costs an event-loop turn, so reading a hundred values one await at a time is a hundred turns. It returns whatever is already buffered and waits only when nothing is, so a live stream still yields each value as soon as it arrives.
next_object
async
¶
next_object(obj_type: type[T] | None = None, timeout: Duration | None = None, mimetype_patterns: str | Sequence[str] = '') -> T | None
Read and deserialize the next value, or None at end of stream.
put
async
¶
Write value and return its store-confirmation future.
value may be a NodeFragment, a
Chunk, or any Python object the node's
serialization registry can encode (mimetype selects the encoding).
Set final=True on the last data fragment so readers know where the
logical value ends. Finality does not close the writer: finalize is
the one call that does both. The returned asyncio.Future resolves to
the stored sequence number after the backing store accepts the
fragment. Attached WireStream sends are attempted or queued by the
writer but are not separately acknowledged.
Examples:
Add an intermediate token while a model response is produced:
put_chunk
async
¶
put_chunk(chunk: Chunk, seq: int | None = None, final: bool = False) -> Future[int]
Admit a native chunk and return its store-confirmation future.
Await this coroutine to respect the writer's bounded admission buffer, then await the returned future when the backing store must have accepted the fragment. Attached stream sends are attempted or queued as the writer processes the batch, but do not add a second delivery confirmation.
Awaiting the confirmation flushes the writer on this thread, so a store that accepts the write without waiting resolves it with no event-loop turn at all. A producer that awaits only admission does not flush: its write goes out on the writer's own pump, which is the point of the two being separate.
put_fragment
async
¶
put_fragment(fragment: NodeFragment) -> Future[int]
Enqueue a NodeFragment (carrying its seq/final).
reset_reader
¶
reset_reader(options: ChunkStoreReaderOptions | dict[str, Any] | None = None) -> AsyncNode
Rewind/reconfigure the reader (e.g. to re-read from an offset).
set_expected_types
¶
set_expected_types(mimetype_patterns: str | Sequence[str], obj_type: type | None) -> AsyncNode
Set the default MIME patterns and object type for reads.
Once set, next()/consume() and async for deserialize to
obj_type (matching mimetype_patterns) without repeating those
arguments on every call. Returns self for chaining.
set_reader_options
¶
set_reader_options(options: ChunkStoreReaderOptions | dict[str, Any]) -> AsyncNode
Replace the reader options and return self for chaining.
set_serialization_registry
¶
set_serialization_registry(registry: SerializationRegistry) -> AsyncNode
Set the serialization registry and return self for chaining.
set_writer_options
¶
set_writer_options(options: ChunkStoreWriterOptions | dict[str, Any]) -> AsyncNode
Replace the writer options and return self for chaining.
wait_for_buffer_to_drain
¶
Returns a future that resolves once the write buffer has drained. Await it to apply backpressure from a fast producer, letting readers catch up before you push more chunks.
consume
async
¶
consume(obj_type: type[T] | None = None, timeout: Duration | None = None, mimetype_patterns: str | Sequence[str] = '', allow_none: Literal[False] = False) -> T
consume(obj_type: type[T] | None = None, timeout: Duration | None = None, mimetype_patterns: str | Sequence[str] = '', allow_none: Literal[True] = True) -> T | None
consume(obj_type: type[T] | None = None, timeout: Duration | None = None, mimetype_patterns: str | Sequence[str] = '', allow_none: bool = False) -> T | None
consume(obj_type: type[T] | None = None, timeout: Duration | None = None, mimetype_patterns: str | Sequence[str] = '', allow_none: Literal[False] = False) -> T
Consume exactly one whole value and return it deserialized.
For a node that carries a single result, such as a unary action
output. Pass obj_type to deserialize to a specific type, or request
NodeFragment/Chunk to get the raw form.
Examples:
Read the unary customer input of an action handler:
consume_fragment
async
¶
consume_fragment(timeout: Duration | None = None, allow_none: Literal[False] = False) -> NodeFragment
consume_fragment(timeout: Duration | None = None, allow_none: Literal[True] = True) -> NodeFragment | None
consume_fragment(timeout: Duration | None = None, allow_none: bool = False) -> NodeFragment | None
consume_fragment(timeout: Duration | None = None, allow_none: Literal[False] = False) -> NodeFragment
Read exactly one whole value's fragment, enforcing the terminator.
Unlike next_fragment, this expects the node to hold exactly one
value, and raises if that shape is violated. Two spellings are
accepted: the value written as final, or the value followed by a null
final chunk. With allow_none a node that holds no value — closed
empty, or holding nothing but a null final — yields None instead of
raising. Requires an ordered reader.
AsyncNode and durable stream services¶
An AsyncNode backed by Redis or SQLite combines an ordered record stream with
an action port's lifecycle. Producers append chunks; readers can follow new
data or read stored data; finalization and failure are part of the same
contract. Media metadata and serialization tags remain attached to the data,
and changing the store does not change the action schema.
Modern stream services expose related facilities. For example, S2 provides managed, durable, ordered streams with append sessions, live tailing, and replay from retained positions. Its agent patterns include resumable token delivery, event sourcing, and coordination through a stream per run.
The scopes differ. S2 is a hosted stream storage API with service-specific
positions, access controls, reconnection, and scaling. AsyncNode is an A11
runtime primitive connected directly to action inputs and outputs. Its storage
is selectable: memory for local work, SQLite for embedded durability, or Redis
for readers and writers in independent processes. This is useful when the
application needs streamed action ports and already operates Redis, or when a
single-machine service can keep its state in SQLite.
NodeMap¶
NodeMap coordinates named streams shared by actions within a session:
node_map = a11.NodeMap()
input_node = node_map.get("user_input")
output_node = node_map.get("agent_response")
a11.nodes.async_node.NodeMap
¶
Creates a node map, optionally backed by a chunk-store factory callable invoked to construct the backing store for each new node.
discard
¶
Removes the node for the given id, optionally only if it matches the expected node.
get
¶
get(node_id: str) -> AsyncNode
Returns the node for the given id, creating it if it does not already exist.
get_if_exists
¶
get_if_exists(node_id: str) -> AsyncNode | None
Returns the node for the given id, or None if it does not exist.
ids
¶
The id of every node the map holds, sorted.
A snapshot: nodes are created on demand, so this is what was there when it was asked for.