A11 (C++ runtime)
Native C++ implementation of the A11 streaming action runtime
Loading...
Searching...
No Matches
a11::nodes::AsyncNode Class Reference

An asynchronous, ordered stream of chunks read from and written to A11. More...

#include <cpp/a11/nodes/async_node.h>

Inheritance diagram for a11::nodes::AsyncNode:
[legend]

Public Member Functions

 ~AsyncNode ()=default
 
 AsyncNode (const AsyncNode &)=delete
 
AsyncNode & operator= (const AsyncNode &)=delete
 
absl::StatusOr< std::string > GetId () const
 Return the node's stable identifier.
 
std::shared_ptr< stores::ChunkStore > GetChunkStore () const
 Return the underlying chunk store backing this node.
 
std::shared_ptr< data::SerializationRegistry > serialization_registry () const
 Return the registry used to (de)serialize typed values.
 
absl::Status SetSerializationRegistry (std::shared_ptr< data::SerializationRegistry > registry)
 Replace the serialization registry.
 
absl::StatusOr< std::shared_ptr< stores::ChunkStoreReader > > reader ()
 Return the node's chunk store reader (the consuming half).
 
absl::StatusOr< std::shared_ptr< stores::ChunkStoreWriter > > writer ()
 Return the node's chunk store writer (the producing half).
 
stores::ChunkStoreReaderOptions GetReaderOptions () const
 Return a copy of the reader's current options.
 
absl::Status SetReaderOptions (stores::ChunkStoreReaderOptions options)
 Replace the reader options.
 
absl::Status ResetReader (std::optional< stores::ChunkStoreReaderOptions > options=std::nullopt)
 Rewind/reconfigure the reader (e.g.
 
stores::ChunkStoreWriterOptions GetWriterOptions () const
 Return a copy of the writer's current options.
 
absl::Status SetWriterOptions (stores::ChunkStoreWriterOptions options)
 Replace the writer options.
 
absl::Status GetReaderStatus () const
 Return the current status of the node's reader.
 
absl::Status GetWriterStatus () const
 Return the current status of the node's writer.
 
std::optional< absl::Status > GetWriterAbortStatus () const
 Return the status the writer was aborted with.
 
a11::Future< bool > IsWritable ()
 Report whether the node can currently accept writes.
 
a11::Future< std::uint32_t > PutChunk (data::Chunk chunk, std::optional< std::uint32_t > seq=std::nullopt, bool final=false)
 Enqueue a raw chunk into the stream.
 
a11::Future< std::uint32_t > PutFragment (data::NodeFragment fragment)
 Enqueue a fragment carrying its own sequence and final flag.
 
template<typename T >
requires data::HasSerialTypeTag<T>
a11::Future< std::uint32_t > PutObject (T value, std::string_view mimetype, std::optional< std::uint32_t > seq=std::nullopt, bool final=false)
 Serialize and write a typed value to the stream.
 
template<typename T >
a11::Future< std::uint32_t > Put (const T &value, std::optional< std::uint32_t > seq=std::nullopt, bool final=false, std::string_view mimetype={})
 
a11::Task Finalize (FinalizeOptions options={})
 End the stream with an explicit null terminator (no value).
 
a11::Task Finalize (data::Chunk chunk, FinalizeOptions options={})
 End the stream with a raw chunk as its final data.
 
template<typename T >
a11::Task Finalize (const T &value, FinalizeOptions options={})
 Serialize value as the stream's final data, then end the stream.
 
a11::Future< std::vector< std::optional< data::NodeFragment > > > NextFragments (size_t limit, absl::Duration timeout=absl::InfiniteDuration())
 Read up to limit fragments in a single await.
 
a11::Future< std::optional< data::NodeFragment > > NextFragment (absl::Duration timeout=absl::InfiniteDuration())
 Read the next raw fragment.
 
a11::Future< std::optional< data::NodeFragment > > NextFragmentRaw (absl::Duration timeout=absl::InfiniteDuration())
 Read the next fragment without producing bytes for it.
 
a11::Future< std::optional< data::Chunk > > NextChunk (absl::Duration timeout=absl::InfiniteDuration())
 Read the next raw chunk.
 
template<typename T >
a11::Future< std::optional< T > > NextObject (absl::Duration timeout=absl::InfiniteDuration(), std::vector< std::string > mimetype_patterns={})
 Read and deserialize the next value.
 
a11::Task WaitForBufferToDrain ()
 Wait for the write buffer to drain.
 
a11::Task Close ()
 Flush all buffered chunks and close the writer, without finality.
 
a11::Task AbortWithStatus (absl::Status status)
 Fail the stream with an error status.
 
absl::Status AttachStream (std::shared_ptr< net::WireStream > stream)
 Tee this node's stored chunks onto a wire stream.
 
absl::Status DetachStream (const std::shared_ptr< net::WireStream > &stream)
 Stop mirroring chunks over a previously attached wire stream.
 
void CancelReader ()
 Cancel the reader, unblocking pending Next* awaits.
 
void CancelWriter ()
 Cancel the writer, unblocking pending Put/drain awaits.
 
void Cancel ()
 Cancel both reader and writer, tearing down all pending streaming operations on the node at once.
 

Static Public Member Functions

static absl::StatusOr< std::shared_ptr< AsyncNode > > Create (std::shared_ptr< stores::ChunkStore > store, std::shared_ptr< data::SerializationRegistry > serialization_registry=nullptr, stores::ChunkStoreReaderOptions reader_options={}, stores::ChunkStoreWriterOptions writer_options={})
 Create a node over an existing chunk store.
 

Detailed Description

An asynchronous, ordered stream of chunks read from and written to A11.

A node has two halves. The writer end admits values into the backing stores::ChunkStore in sequence; the reader end yields them back, optionally deserializing typed objects on the way out. Every Put* resolves once the backing store accepts the chunk. A stackless pump turn enqueues each batch on attached net::WireStreams. The store confirmation carries no mirror handover or remote-delivery acknowledgement. WaitForBufferToDrain() covers both stages and reports tee failures.

A producer ends a node with Finalize(): it marks the logical end of the data and, by default, closes the writer. Finality and closure remain two distinct facts – see the AsyncNode lifecycle guide – and Finalize() writes one and requests the other in the order readers expect. Close() is the rarer half on its own, for a producer that cannot say which chunk was last.

Instances are always heap-allocated and shared via Create; the class is non-copyable and derives from enable_shared_from_this.

Constructor & Destructor Documentation

◆ ~AsyncNode()

a11::nodes::AsyncNode::~AsyncNode ( )
default

◆ AsyncNode()

a11::nodes::AsyncNode::AsyncNode ( const AsyncNode &  )
delete

Member Function Documentation

◆ AbortWithStatus()

a11::Task a11::nodes::AsyncNode::AbortWithStatus ( absl::Status  status)

Fail the stream with an error status.

Parameters
statusThe error to abort with.
Returns
An awaitable that resolves once the abort has propagated, so consumers observe the error instead of a normal end-of-stream.

◆ AttachStream()

absl::Status a11::nodes::AsyncNode::AttachStream ( std::shared_ptr< net::WireStream >  stream)

Tee this node's stored chunks onto a wire stream.

WireStream::Send() confirms local transport admission, not receipt by the remote agent. A send failure stops later writes but cannot revoke the current batch's store confirmations. Closing the writer also tees a closure marker to every attached stream.

Parameters
streamThe transport to attach; kept alive for the node's lifetime.
Returns
OK, or an error status on failure.

◆ Cancel()

void a11::nodes::AsyncNode::Cancel ( )

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

◆ CancelReader()

void a11::nodes::AsyncNode::CancelReader ( )

Cancel the reader, unblocking pending Next* awaits.

◆ CancelWriter()

void a11::nodes::AsyncNode::CancelWriter ( )

Cancel the writer, unblocking pending Put/drain awaits.

◆ Close()

a11::Task a11::nodes::AsyncNode::Close ( )

Flush all buffered chunks and close the writer, without finality.

The specialised half of Finalize(): a storage-lifecycle operation and not an end-of-data marker. It appends no fragment and chooses no final sequence number, so an ordered reader learns only that nothing more can arrive – which is all a producer that cannot say which chunk was last, such as a log, is able to promise. Prefer Finalize() everywhere else.

Closing always drains: making a closure with nothing to mark asynchronous would buy nothing, so this awaitable resolves when the store is closed.

Attached streams learn of the closure: the writer follows the last teed batch with a closure marker, so a peer holding a mirror of this node closes its write half too and its readers reach a clean end.

Returns
An awaitable that resolves once every produced chunk has been flushed and the backing store is closed to further writes.

◆ Create()

absl::StatusOr< std::shared_ptr< AsyncNode > > a11::nodes::AsyncNode::Create ( std::shared_ptr< stores::ChunkStore >  store,
std::shared_ptr< data::SerializationRegistry >  serialization_registry = nullptr,
stores::ChunkStoreReaderOptions  reader_options = {},
stores::ChunkStoreWriterOptions  writer_options = {} 
)
static

Create a node over an existing chunk store.

Parameters
storeThe backing ordered buffer the reader and writer stream through.
serialization_registryRegistry used to (de)serialize typed values passed to Put/NextObject; defaults to none.
reader_optionsOptions controlling how the reader buffers and orders chunks.
writer_optionsOptions controlling how the writer buffers chunks.
Returns
The new node, or an error status on failure.

◆ DetachStream()

absl::Status a11::nodes::AsyncNode::DetachStream ( const std::shared_ptr< net::WireStream > &  stream)

Stop mirroring chunks over a previously attached wire stream.

Parameters
streamThe transport to detach.
Returns
OK, or an error status on failure.

◆ Finalize() [1/3]

template<typename T >
a11::Task a11::nodes::AsyncNode::Finalize ( const T &  value,
FinalizeOptions  options = {} 
)
inline

Serialize value as the stream's final data, then end the stream.

The unary case – a node that carries one result – and the streaming case where the producer knows which value is the last one and can save transmitting a separate terminator for it.

Template Parameters
TThe type of the value being written.
Parameters
valueThe value to encode and write as final.
optionsHow to end the stream.
Returns
An awaitable resolving as described by Finalize(FinalizeOptions), or a failed awaitable if serialization fails.

◆ Finalize() [2/3]

a11::Task a11::nodes::AsyncNode::Finalize ( data::Chunk  chunk,
FinalizeOptions  options = {} 
)

End the stream with a raw chunk as its final data.

Parameters
chunkThe last chunk of the stream.
optionsHow to end the stream; options.mimetype is ignored, since the chunk carries its own.
Returns
An awaitable resolving as described by Finalize(FinalizeOptions).

◆ Finalize() [3/3]

a11::Task a11::nodes::AsyncNode::Finalize ( FinalizeOptions  options = {})

End the stream with an explicit null terminator (no value).

The ordinary way to finish a node: it marks the logical end of the data so ordered readers stop immediately, and – unless options.close is cleared – closes the writer so the backing store admits nothing more.

With options.wait left false the returned awaitable is already resolved and both the write and the close proceed on the writer's pump, which is safe after the producing frame is gone. Nothing is silently dropped: a failed write or close is logged, and remains visible through GetWriterStatus(). Set options.wait when the producer must know the store accepted the end of the stream before continuing.

Parameters
optionsHow to end the stream; the defaults finalise and close without waiting.
Returns
An awaitable that resolves once the requested work is done, or immediately when not waiting.

◆ GetChunkStore()

std::shared_ptr< stores::ChunkStore > a11::nodes::AsyncNode::GetChunkStore ( ) const

Return the underlying chunk store backing this node.

Returns
The ordered storage boundary the reader and writer stream through; reach for it when you need lower-level access than the Put/Next API provides.

◆ GetId()

absl::StatusOr< std::string > a11::nodes::AsyncNode::GetId ( ) const

Return the node's stable identifier.

Returns
The id, or an error status if it cannot be resolved. Use it to correlate the node with the rest of an agent's state or to key it in a NodeMap.

◆ GetReaderOptions()

stores::ChunkStoreReaderOptions a11::nodes::AsyncNode::GetReaderOptions ( ) const

Return a copy of the reader's current options.

Returns
The reader options in effect.

◆ GetReaderStatus()

absl::Status a11::nodes::AsyncNode::GetReaderStatus ( ) const

Return the current status of the node's reader.

Returns
Whether the consuming end is healthy, has completed, or has failed while streaming.

◆ GetWriterAbortStatus()

std::optional< absl::Status > a11::nodes::AsyncNode::GetWriterAbortStatus ( ) const

Return the status the writer was aborted with.

Returns
The abort status, or nullopt if the writer has not been aborted. Use it to surface why a stream was cut short.

◆ GetWriterOptions()

stores::ChunkStoreWriterOptions a11::nodes::AsyncNode::GetWriterOptions ( ) const

Return a copy of the writer's current options.

Returns
The writer options in effect.

◆ GetWriterStatus()

absl::Status a11::nodes::AsyncNode::GetWriterStatus ( ) const

Return the current status of the node's writer.

Returns
Whether the producing end is healthy, has completed, or has failed while streaming.

◆ IsWritable()

a11::Future< bool > a11::nodes::AsyncNode::IsWritable ( )

Report whether the node can currently accept writes.

Returns
An awaitable that resolves once writability is known; await it before producing chunks to respect backpressure.

◆ NextChunk()

a11::Future< std::optional< data::Chunk > > a11::nodes::AsyncNode::NextChunk ( absl::Duration  timeout = absl::InfiniteDuration())

Read the next raw chunk.

Parameters
timeoutHow long to wait for a chunk.
Returns
An awaitable that resolves to the next chunk, or nullopt at end of stream.

◆ NextFragment()

a11::Future< std::optional< data::NodeFragment > > a11::nodes::AsyncNode::NextFragment ( absl::Duration  timeout = absl::InfiniteDuration())

Read the next raw fragment.

Parameters
timeoutHow long to wait for a fragment.
Returns
An awaitable that resolves to the next fragment, or nullopt at end of stream.

◆ NextFragmentRaw()

a11::Future< std::optional< data::NodeFragment > > a11::nodes::AsyncNode::NextFragmentRaw ( absl::Duration  timeout = absl::InfiniteDuration())

Read the next fragment without producing bytes for it.

Preserves an in-process object for NextObject<T>(); use NextFragment() when the caller needs materialized bytes.

◆ NextFragments()

a11::Future< std::vector< std::optional< data::NodeFragment > > > a11::nodes::AsyncNode::NextFragments ( size_t  limit,
absl::Duration  timeout = absl::InfiniteDuration() 
)

Read up to limit fragments in a single await.

The batched counterpart to NextFragment(), with the same end-of-stream marker: a trailing empty optional. It returns whatever is already buffered and waits only when nothing is, so it never trades latency for throughput on a live stream. See ChunkStoreReader::NextMany().

Parameters
limitMaximum number of fragments to return.
timeoutHow long to wait when no fragments are buffered.

◆ NextObject()

template<typename T >
a11::Future< std::optional< T > > a11::nodes::AsyncNode::NextObject ( absl::Duration  timeout = absl::InfiniteDuration(),
std::vector< std::string >  mimetype_patterns = {} 
)
inline

Read and deserialize the next value.

Template Parameters
TThe type to deserialize into.
Parameters
timeoutHow long to wait for a value.
mimetype_patternsOptional MIME patterns constraining which encodings are accepted.
Returns
An awaitable that resolves to the next value, or nullopt at end of stream.

◆ operator=()

AsyncNode & a11::nodes::AsyncNode::operator= ( const AsyncNode &  )
delete

◆ Put()

template<typename T >
a11::Future< std::uint32_t > a11::nodes::AsyncNode::Put ( const T &  value,
std::optional< std::uint32_t >  seq = std::nullopt,
bool  final = false,
std::string_view  mimetype = {} 
)
inline

◆ PutChunk()

a11::Future< std::uint32_t > a11::nodes::AsyncNode::PutChunk ( data::Chunk  chunk,
std::optional< std::uint32_t >  seq = std::nullopt,
bool  final = false 
)

Enqueue a raw chunk into the stream.

Parameters
chunkThe chunk to admit.
seqOptional explicit sequence number; assigned in order when omitted.
finalSet true on the last chunk to establish the logical final sequence. This does not close the writer.
Returns
An awaitable that resolves to the sequence number once the backing store accepts the chunk. Attached stream sends are attempted during the flush but do not acknowledge remote delivery.

◆ PutFragment()

a11::Future< std::uint32_t > a11::nodes::AsyncNode::PutFragment ( data::NodeFragment  fragment)

Enqueue a fragment carrying its own sequence and final flag.

Parameters
fragmentThe fragment to admit.
Returns
An awaitable that resolves to the stored sequence number.

◆ PutObject()

template<typename T >
requires data::HasSerialTypeTag<T>
a11::Future< std::uint32_t > a11::nodes::AsyncNode::PutObject ( T  value,
std::string_view  mimetype,
std::optional< std::uint32_t >  seq = std::nullopt,
bool  final = false 
)
inline

Serialize and write a typed value to the stream.

Write a typed value without encoding it.

The local fast path. Put() encodes on the way in and NextObject() decodes on the way out, and for a value that never leaves this process both are waste – decode especially, being the dearer half by an order of magnitude. This admits a chunk carrying the value itself; the bytes are produced only if something actually needs them, which is a peer, a persisting store, or a reader asking for a chunk rather than a value.

A reader calling NextObject<T>() with the same T gets a copy of the value, having encoded and decoded nothing.

Two obligations, and both are already true of anything put in a store:

  • The value must not change afterwards. A node replays its fragments to every reader, including late readers. Consumers receive copies and cannot mutate the stored value.
  • mimetype is stated rather than derived. Put() lets the registry choose it; here the caller says it, because a chunk whose mimetype differed from the one its bytes would have had would be filtered differently by | mime depending on whether anybody had asked for the bytes yet. Pass what Put() would have produced.
Template Parameters
TType of the value. Its serialisation tag identifies it to readers, so it must have one.

◆ reader()

absl::StatusOr< std::shared_ptr< stores::ChunkStoreReader > > a11::nodes::AsyncNode::reader ( )

Return the node's chunk store reader (the consuming half).

Returns
The reader, or an error status on failure.

◆ ResetReader()

absl::Status a11::nodes::AsyncNode::ResetReader ( std::optional< stores::ChunkStoreReaderOptions >  options = std::nullopt)

Rewind/reconfigure the reader (e.g.

to re-read from an offset).

Parameters
optionsOptional new reader options; the existing options are kept when omitted.
Returns
OK, or an error status on failure.

◆ serialization_registry()

std::shared_ptr< data::SerializationRegistry > a11::nodes::AsyncNode::serialization_registry ( ) const

Return the registry used to (de)serialize typed values.

Returns
The current serialization registry (may be null).

◆ SetReaderOptions()

absl::Status a11::nodes::AsyncNode::SetReaderOptions ( stores::ChunkStoreReaderOptions  options)

Replace the reader options.

Parameters
optionsThe new reader options.
Returns
OK, or an error status on failure.

◆ SetSerializationRegistry()

absl::Status a11::nodes::AsyncNode::SetSerializationRegistry ( std::shared_ptr< data::SerializationRegistry >  registry)

Replace the serialization registry.

Parameters
registryThe registry to install.
Returns
OK, or an error status on failure.

◆ SetWriterOptions()

absl::Status a11::nodes::AsyncNode::SetWriterOptions ( stores::ChunkStoreWriterOptions  options)

Replace the writer options.

Parameters
optionsThe new writer options.
Returns
OK, or an error status on failure.

◆ WaitForBufferToDrain()

a11::Task a11::nodes::AsyncNode::WaitForBufferToDrain ( )

Wait for the write buffer to drain.

Returns
An awaitable that resolves once the write buffer has drained; await it to let consumers catch up before pushing more chunks.

◆ writer()

absl::StatusOr< std::shared_ptr< stores::ChunkStoreWriter > > a11::nodes::AsyncNode::writer ( )

Return the node's chunk store writer (the producing half).

Returns
The writer, or an error status on failure.

The documentation for this class was generated from the following files: