Skip to content

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:

async for event in node:
    print(event)

result = await unary_node.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).

id property

id: str

The node's stable identifier (see get_id).

reader property

The node's ChunkStoreReader.

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

serialization_registry: SerializationRegistry

The registry used to (de)serialize Python objects for this node.

writer property

The node's ChunkStoreWriter.

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:

answer = AsyncNode.create("answer-tokens")

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

cancel() -> None

Cancels both the reader and the writer, tearing down all pending streaming operations on the node at once.

cancel_reader

cancel_reader() -> None

Cancels the node's reader, unblocking any pending next-chunk or next-fragment awaits on the read side of the stream.

cancel_writer

cancel_writer() -> None

Cancels the node's writer, unblocking any pending put or drain awaits on the producing side of the stream.

close async

close() -> None

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:

async for token in model:
    await answer.put(token)
await answer.finalize()

Deliver a unary result, which is one chunk and one call:

await action["customer"].finalize(customer)

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_id

get_id() -> str

Returns the node's stable identifier. Raises if the id cannot be resolved.

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

is_writable() -> Future[bool]

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

iter_chunks(timeout: Duration | None = None) -> AsyncIterator[Chunk]

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

Alias for next_object: the next deserialized value or None.

Examples:

Process a live audit stream one event at a time:

events.set_expected_types("application/json", AuditEvent)
while (event := await events.next()) is not None:
    await audit_index.store(event)

next_chunk async

next_chunk(timeout: Duration | None = None) -> Chunk | None

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

put(value: Any, seq: int | None = None, final: bool = False, mimetype: str = '') -> Future[int]

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:

await answer.put("The shipment ")

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

wait_for_buffer_to_drain() -> Future[None]

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:

customer = await action["customer"].consume(obj_type=Customer)

consume_chunk async

consume_chunk(timeout: Duration | None = None, allow_none: Literal[False] = False) -> Chunk
consume_chunk(timeout: Duration | None = None, allow_none: Literal[True] = True) -> Chunk | None
consume_chunk(timeout: Duration | None = None, allow_none: bool = False) -> Chunk | None
consume_chunk(timeout: Duration | None = None, allow_none: Literal[False] = False) -> Chunk

Consume exactly one whole value and return its raw chunk.

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

NodeMap(chunk_store_factory: Any | None = None)

Creates a node map, optionally backed by a chunk-store factory callable invoked to construct the backing store for each new node.

contains

contains(node_id: str) -> bool

Returns whether a node with the given id exists.

discard

discard(node_id: str, expected: AsyncNode | None = None) -> AsyncNode | None

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

ids() -> list[str]

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.

size

size() -> int

Number of nodes in the map.