|
A11 (C++ runtime)
Native C++ implementation of the A11 streaming action runtime
|
An asynchronous, ordered stream of chunks read from and written to A11. More...
#include <cpp/a11/nodes/async_node.h>
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. | |
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.
|
default |
|
delete |
| a11::Task a11::nodes::AsyncNode::AbortWithStatus | ( | absl::Status | status | ) |
Fail the stream with an error status.
| status | The error to abort with. |
| 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.
| stream | The transport to attach; kept alive for the node's lifetime. |
| void a11::nodes::AsyncNode::Cancel | ( | ) |
Cancel both reader and writer, tearing down all pending streaming operations on the node at once.
| void a11::nodes::AsyncNode::CancelReader | ( | ) |
Cancel the reader, unblocking pending Next* awaits.
| void a11::nodes::AsyncNode::CancelWriter | ( | ) |
Cancel the writer, unblocking pending Put/drain awaits.
| 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.
|
static |
Create a node over an existing chunk store.
| store | The backing ordered buffer the reader and writer stream through. |
| serialization_registry | Registry used to (de)serialize typed values passed to Put/NextObject; defaults to none. |
| reader_options | Options controlling how the reader buffers and orders chunks. |
| writer_options | Options controlling how the writer buffers chunks. |
| absl::Status a11::nodes::AsyncNode::DetachStream | ( | const std::shared_ptr< net::WireStream > & | stream | ) |
Stop mirroring chunks over a previously attached wire stream.
| stream | The transport to detach. |
|
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.
| T | The type of the value being written. |
| value | The value to encode and write as final. |
| options | How to end the stream. |
| a11::Task a11::nodes::AsyncNode::Finalize | ( | data::Chunk | chunk, |
| FinalizeOptions | options = {} |
||
| ) |
End the stream with a raw chunk as its final data.
| chunk | The last chunk of the stream. |
| options | How to end the stream; options.mimetype is ignored, since the chunk carries its own. |
| 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.
| options | How to end the stream; the defaults finalise and close without waiting. |
| std::shared_ptr< stores::ChunkStore > a11::nodes::AsyncNode::GetChunkStore | ( | ) | const |
Return the underlying chunk store backing this node.
| absl::StatusOr< std::string > a11::nodes::AsyncNode::GetId | ( | ) | const |
Return the node's stable identifier.
| stores::ChunkStoreReaderOptions a11::nodes::AsyncNode::GetReaderOptions | ( | ) | const |
Return a copy of the reader's current options.
| absl::Status a11::nodes::AsyncNode::GetReaderStatus | ( | ) | const |
Return the current status of the node's reader.
| std::optional< absl::Status > a11::nodes::AsyncNode::GetWriterAbortStatus | ( | ) | const |
Return the status the writer was aborted with.
| stores::ChunkStoreWriterOptions a11::nodes::AsyncNode::GetWriterOptions | ( | ) | const |
Return a copy of the writer's current options.
| absl::Status a11::nodes::AsyncNode::GetWriterStatus | ( | ) | const |
Return the current status of the node's writer.
| a11::Future< bool > a11::nodes::AsyncNode::IsWritable | ( | ) |
Report whether the node can currently accept writes.
| a11::Future< std::optional< data::Chunk > > a11::nodes::AsyncNode::NextChunk | ( | absl::Duration | timeout = absl::InfiniteDuration() | ) |
Read the next raw chunk.
| timeout | How long to wait for a chunk. |
| a11::Future< std::optional< data::NodeFragment > > a11::nodes::AsyncNode::NextFragment | ( | absl::Duration | timeout = absl::InfiniteDuration() | ) |
Read the next raw fragment.
| timeout | How long to wait for a fragment. |
| 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.
| 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().
| limit | Maximum number of fragments to return. |
| timeout | How long to wait when no fragments are buffered. |
|
inline |
Read and deserialize the next value.
| T | The type to deserialize into. |
| timeout | How long to wait for a value. |
| mimetype_patterns | Optional MIME patterns constraining which encodings are accepted. |
|
inline |
| 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.
| chunk | The chunk to admit. |
| seq | Optional explicit sequence number; assigned in order when omitted. |
| final | Set true on the last chunk to establish the logical final sequence. This does not close the writer. |
| a11::Future< std::uint32_t > a11::nodes::AsyncNode::PutFragment | ( | data::NodeFragment | fragment | ) |
Enqueue a fragment carrying its own sequence and final flag.
| fragment | The fragment to admit. |
|
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:
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.| T | Type of the value. Its serialisation tag identifies it to readers, so it must have one. |
| absl::StatusOr< std::shared_ptr< stores::ChunkStoreReader > > a11::nodes::AsyncNode::reader | ( | ) |
Return the node's chunk store reader (the consuming half).
| 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).
| options | Optional new reader options; the existing options are kept when omitted. |
| std::shared_ptr< data::SerializationRegistry > a11::nodes::AsyncNode::serialization_registry | ( | ) | const |
Return the registry used to (de)serialize typed values.
| absl::Status a11::nodes::AsyncNode::SetReaderOptions | ( | stores::ChunkStoreReaderOptions | options | ) |
Replace the reader options.
| options | The new reader options. |
| absl::Status a11::nodes::AsyncNode::SetSerializationRegistry | ( | std::shared_ptr< data::SerializationRegistry > | registry | ) |
Replace the serialization registry.
| registry | The registry to install. |
| absl::Status a11::nodes::AsyncNode::SetWriterOptions | ( | stores::ChunkStoreWriterOptions | options | ) |
Replace the writer options.
| options | The new writer options. |
| a11::Task a11::nodes::AsyncNode::WaitForBufferToDrain | ( | ) |
Wait for the write buffer to drain.
| absl::StatusOr< std::shared_ptr< stores::ChunkStoreWriter > > a11::nodes::AsyncNode::writer | ( | ) |
Return the node's chunk store writer (the producing half).