A11 (C++ runtime)
Native C++ implementation of the A11 streaming action runtime
Loading...
Searching...
No Matches
async_node.h
Go to the documentation of this file.
1/*
2 * Copyright 2026 The A11 Authors
3 *
4 * Licensed under the Apache License, Version 2.0 (the "License");
5 * you may not use this file except in compliance with the License.
6 * You may obtain a copy of the License at
7 *
8 * http://www.apache.org/licenses/LICENSE-2.0
9 *
10 * Unless required by applicable law or agreed to in writing, software
11 * distributed under the License is distributed on an "AS IS" BASIS,
12 * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
13 * See the License for the specific language governing permissions and
14 * limitations under the License.
15 */
16
35#ifndef A11_NODES_ASYNC_NODE_H_
36#define A11_NODES_ASYNC_NODE_H_
37
38#include <cstddef>
39#include <cstdint>
40#include <memory>
41#include <optional>
42#include <string>
43#include <string_view>
44#include <utility>
45#include <vector>
46
47#include <absl/base/thread_annotations.h>
48#include <absl/status/status.h>
49#include <absl/status/status_macros.h>
50#include <absl/status/statusor.h>
51#include <absl/time/time.h>
52
57#include "a11/data/types.h"
61#include "thread/boost_primitives.h"
62
63namespace a11::net {
64class WireStream;
65} // namespace a11::net
66
67namespace a11::nodes {
68
78 bool wait = false;
81 bool close = true;
84 std::optional<std::uint32_t> seq = {};
87 std::string_view mimetype = {};
88};
89
111class AsyncNode : public std::enable_shared_from_this<AsyncNode> {
112 public:
126 static absl::StatusOr<std::shared_ptr<AsyncNode>> Create(
127 std::shared_ptr<stores::ChunkStore> store,
128 std::shared_ptr<data::SerializationRegistry> serialization_registry =
129 nullptr,
130 stores::ChunkStoreReaderOptions reader_options = {},
131 stores::ChunkStoreWriterOptions writer_options = {});
132
133 ~AsyncNode() = default;
134
135 AsyncNode(const AsyncNode&) = delete;
136 AsyncNode& operator=(const AsyncNode&) = delete;
137
144 absl::StatusOr<std::string> GetId() const;
145
152 [[nodiscard]] std::shared_ptr<stores::ChunkStore> GetChunkStore() const;
153
158 [[nodiscard]] std::shared_ptr<data::SerializationRegistry>
160
166 absl::Status SetSerializationRegistry(
167 std::shared_ptr<data::SerializationRegistry> registry);
168
173 absl::StatusOr<std::shared_ptr<stores::ChunkStoreReader>> reader();
174
179 absl::StatusOr<std::shared_ptr<stores::ChunkStoreWriter>> writer();
180
186
193
200 absl::Status ResetReader(
201 std::optional<stores::ChunkStoreReaderOptions> options = std::nullopt);
202
208
215
221 [[nodiscard]] absl::Status GetReaderStatus() const;
222
228 [[nodiscard]] absl::Status GetWriterStatus() const;
229
235 [[nodiscard]] std::optional<absl::Status> GetWriterAbortStatus() const;
236
243
256 data::Chunk chunk, std::optional<std::uint32_t> seq = std::nullopt,
257 bool final = false);
258
265
294 template <typename T>
296 T value, std::string_view mimetype,
297 std::optional<std::uint32_t> seq = std::nullopt, bool final = false) {
298 std::shared_ptr<data::SerializationRegistry> registry;
299 {
300 thread::MutexLock lock(&mu_);
301 registry = serialization_registry_;
302 }
303 return PutChunk(
304 data::MakeChunkObject<T>(std::move(value), data::SerialTypeTag<T>(),
305 std::string(mimetype), registry),
306 seq, final);
307 }
308
309 template <typename T>
311 const T& value, std::optional<std::uint32_t> seq = std::nullopt,
312 bool final = false, std::string_view mimetype = {}) {
313 std::shared_ptr<data::SerializationRegistry> registry;
314 {
315 thread::MutexLock lock(&mu_);
316 registry = serialization_registry_;
317 }
318 absl::StatusOr<data::Chunk> chunk = registry->ToChunk<T>(value, mimetype);
319 if (!chunk.ok()) {
320 return a11::FailedFuture<std::uint32_t>(chunk.status());
321 }
322 return PutChunk(std::move(*chunk), seq, final);
323 }
324
344 a11::Task Finalize(FinalizeOptions options = {});
345
353 a11::Task Finalize(data::Chunk chunk, FinalizeOptions options = {});
354
368 template <typename T>
369 a11::Task Finalize(const T& value, FinalizeOptions options = {}) {
370 std::shared_ptr<data::SerializationRegistry> registry;
371 {
372 thread::MutexLock lock(&mu_);
373 registry = serialization_registry_;
374 }
375 absl::StatusOr<data::Chunk> chunk =
376 registry->ToChunk<T>(value, options.mimetype);
377 if (!chunk.ok()) {
378 return a11::FailedTask(chunk.status());
379 }
380 return Finalize(std::move(*chunk), options);
381 }
382
394 size_t limit, absl::Duration timeout = absl::InfiniteDuration());
395
403 absl::Duration timeout = absl::InfiniteDuration());
404
412 absl::Duration timeout = absl::InfiniteDuration());
413
421 absl::Duration timeout = absl::InfiniteDuration());
422
432 template <typename T>
434 absl::Duration timeout = absl::InfiniteDuration(),
435 std::vector<std::string> mimetype_patterns = {}) {
436 std::shared_ptr<AsyncNode> self = shared_from_this();
437 return a11::Submit<std::optional<T>>(
438 [self = std::move(self), timeout,
439 mimetype_patterns = std::move(
440 mimetype_patterns)]() mutable -> absl::StatusOr<std::optional<T>> {
441 // Raw, so a chunk carrying a value still is one when it gets here.
442 // NextFragment() would have materialised it, which is right for every
443 // caller that wants bytes and is exactly what this one does not.
444 absl::StatusOr<std::optional<data::NodeFragment>> fragment =
445 self->NextFragmentRaw(timeout).Await();
446 if (!fragment.ok()) {
447 return fragment.status();
448 }
449 if (!fragment->has_value()) {
450 return std::nullopt;
451 }
452 absl::StatusOr<const data::Chunk*> chunk = (*fragment)->GetChunk();
453 if (!chunk.ok()) {
454 return chunk.status();
455 }
456 // The fast path: the producer put a T and this reader wants a T, so
457 // nothing was encoded and nothing is decoded. Guarded by the tag, not
458 // by type identity -- see a11::data::ChunkObject.
459 if ((*chunk)->HasObject()) {
460 // `if constexpr`, because NextObject is instantiated for types with
461 // no tag at all -- a JSON-native value, a bare string -- and those
462 // simply have no fast path to take.
463 if constexpr (data::HasSerialTypeTag<T>) {
464 if (std::optional<T> taken =
465 data::TryTakeObject<T>(**chunk, data::SerialTypeTag<T>());
466 taken.has_value()) {
467 return taken;
468 }
469 }
470 // Carrying something else. Fall through to the bytes, which means
471 // producing them now: a reader asking for a different type than the
472 // producer wrote is exactly when the wire format earns its keep.
473 data::Chunk materialised = **chunk;
474 ABSL_RETURN_IF_ERROR(materialised.Materialize());
475 std::shared_ptr<data::SerializationRegistry> registry;
476 {
477 thread::MutexLock lock(&self->mu_);
478 registry = self->serialization_registry_;
479 }
480 absl::StatusOr<T> value =
481 registry->FromChunk<T>(materialised, mimetype_patterns);
482 if (!value.ok()) {
483 return value.status();
484 }
485 return std::optional<T>(std::move(*value));
486 }
487 if ((*chunk)->IsNull()) {
488 if ((*fragment)->continued) {
489 return absl::FailedPreconditionError(
490 "A null stream marker must be final");
491 }
492 return std::nullopt;
493 }
494 std::shared_ptr<data::SerializationRegistry> registry;
495 {
496 thread::MutexLock lock(&self->mu_);
497 registry = self->serialization_registry_;
498 }
499 absl::StatusOr<T> value =
500 registry->FromChunk<T>(**chunk, mimetype_patterns);
501 if (!value.ok()) {
502 return value.status();
503 }
504 return std::optional<T>(std::move(*value));
505 });
506 }
507
514
534
541 a11::Task AbortWithStatus(absl::Status status);
542
554 absl::Status AttachStream(std::shared_ptr<net::WireStream> stream);
555
561 absl::Status DetachStream(const std::shared_ptr<net::WireStream>& stream);
562
566 void CancelReader();
567
571 void CancelWriter();
572
577 void Cancel();
578
579 private:
580 AsyncNode(std::shared_ptr<stores::ChunkStore> store,
581 std::shared_ptr<data::SerializationRegistry> registry,
582 stores::ChunkStoreReaderOptions reader_options,
583 stores::ChunkStoreWriterOptions writer_options)
584 : store_(std::move(store)),
585 serialization_registry_(std::move(registry)),
586 reader_options_(reader_options),
587 writer_options_(writer_options) {}
588
589 const std::shared_ptr<stores::ChunkStore> store_;
590 mutable thread::Mutex mu_;
591 std::shared_ptr<data::SerializationRegistry> serialization_registry_
592 ABSL_GUARDED_BY(mu_);
593 stores::ChunkStoreReaderOptions reader_options_ ABSL_GUARDED_BY(mu_);
594 stores::ChunkStoreWriterOptions writer_options_ ABSL_GUARDED_BY(mu_);
595 std::shared_ptr<stores::ChunkStoreReader> reader_ ABSL_GUARDED_BY(mu_);
596 std::shared_ptr<stores::ChunkStoreWriter> writer_ ABSL_GUARDED_BY(mu_);
597};
598
599} // namespace a11::nodes
600
601#endif // A11_NODES_ASYNC_NODE_H_
A11's pluggable storage interface for streamed node data: an ordered, appendable log of fragments key...
A read cursor over a ChunkStore that pulls fragments out in order (or by arrival),...
A write cursor over a ChunkStore that admits chunks into the store in sequence, applying backpressure...
Shared handle to one asynchronous result.
Definition future.h:126
An asynchronous, ordered stream of chunks read from and written to A11.
Definition async_node.h:111
absl::Status SetSerializationRegistry(std::shared_ptr< data::SerializationRegistry > registry)
Replace the serialization registry.
Definition async_node.cc:110
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.
Definition async_node.h:295
std::shared_ptr< stores::ChunkStore > GetChunkStore() const
Return the underlying chunk store backing this node.
Definition async_node.cc:100
a11::Task WaitForBufferToDrain()
Wait for the write buffer to drain.
Definition async_node.cc:379
a11::Task Finalize(FinalizeOptions options={})
End the stream with an explicit null terminator (no value).
Definition async_node.cc:272
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.
Definition async_node.cc:253
stores::ChunkStoreWriterOptions GetWriterOptions() const
Return a copy of the writer's current options.
Definition async_node.cc:182
a11::Task Finalize(const T &value, FinalizeOptions options={})
Serialize value as the stream's final data, then end the stream.
Definition async_node.h:369
a11::Task Close()
Flush all buffered chunks and close the writer, without finality.
Definition async_node.cc:388
stores::ChunkStoreReaderOptions GetReaderOptions() const
Return a copy of the reader's current options.
Definition async_node.cc:146
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.
Definition async_node.cc:67
a11::Task AbortWithStatus(absl::Status status)
Fail the stream with an error status.
Definition async_node.cc:394
absl::Status DetachStream(const std::shared_ptr< net::WireStream > &stream)
Stop mirroring chunks over a previously attached wire stream.
Definition async_node.cc:408
AsyncNode(const AsyncNode &)=delete
a11::Future< std::optional< data::Chunk > > NextChunk(absl::Duration timeout=absl::InfiniteDuration())
Read the next raw chunk.
Definition async_node.cc:356
absl::Status SetReaderOptions(stores::ChunkStoreReaderOptions options)
Replace the reader options.
Definition async_node.cc:151
absl::Status SetWriterOptions(stores::ChunkStoreWriterOptions options)
Replace the writer options.
Definition async_node.cc:187
a11::Future< std::uint32_t > Put(const T &value, std::optional< std::uint32_t > seq=std::nullopt, bool final=false, std::string_view mimetype={})
Definition async_node.h:310
void CancelWriter()
Cancel the writer, unblocking pending Put/drain awaits.
Definition async_node.cc:428
absl::Status GetWriterStatus() const
Return the current status of the node's writer.
Definition async_node.cc:208
a11::Future< std::optional< data::NodeFragment > > NextFragment(absl::Duration timeout=absl::InfiniteDuration())
Read the next raw fragment.
Definition async_node.cc:334
absl::Status GetReaderStatus() const
Return the current status of the node's reader.
Definition async_node.cc:199
absl::Status ResetReader(std::optional< stores::ChunkStoreReaderOptions > options=std::nullopt)
Rewind/reconfigure the reader (e.g.
Definition async_node.cc:163
AsyncNode & operator=(const AsyncNode &)=delete
std::optional< absl::Status > GetWriterAbortStatus() const
Return the status the writer was aborted with.
Definition async_node.cc:221
absl::StatusOr< std::shared_ptr< stores::ChunkStoreWriter > > writer()
Return the node's chunk store writer (the producing half).
Definition async_node.cc:133
void CancelReader()
Cancel the reader, unblocking pending Next* awaits.
Definition async_node.cc:417
void Cancel()
Cancel both reader and writer, tearing down all pending streaming operations on the node at once.
Definition async_node.cc:439
absl::Status AttachStream(std::shared_ptr< net::WireStream > stream)
Tee this node's stored chunks onto a wire stream.
Definition async_node.cc:400
a11::Future< bool > IsWritable()
Report whether the node can currently accept writes.
Definition async_node.cc:230
std::shared_ptr< data::SerializationRegistry > serialization_registry() const
Return the registry used to (de)serialize typed values.
Definition async_node.cc:104
a11::Future< std::optional< T > > NextObject(absl::Duration timeout=absl::InfiniteDuration(), std::vector< std::string > mimetype_patterns={})
Read and deserialize the next value.
Definition async_node.h:433
absl::StatusOr< std::shared_ptr< stores::ChunkStoreReader > > reader()
Return the node's chunk store reader (the consuming half).
Definition async_node.cc:120
a11::Future< std::optional< data::NodeFragment > > NextFragmentRaw(absl::Duration timeout=absl::InfiniteDuration())
Read the next fragment without producing bytes for it.
Definition async_node.cc:325
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.
Definition async_node.cc:316
absl::StatusOr< std::string > GetId() const
Return the node's stable identifier.
Definition async_node.cc:96
a11::Future< std::uint32_t > PutFragment(data::NodeFragment fragment)
Enqueue a fragment carrying its own sequence and final flag.
Definition async_node.cc:263
Whether T publishes a serialisation tag.
Definition serializable.h:171
std::string value
Definition discover.cc:114
Completion values used by every asynchronous A11 operation.
std::optional< absl::Duration > timeout
Definition main.cc:144
Definition action.h:65
Definition action.h:69
Task FailedTask(absl::Status status)
Return an already-failed Task.
Definition future.h:419
ADL-based customization points that make a custom C++ type serializable through a11::data::Serializat...
Type-and-mimetype indexed serialization of values to/from Chunks.
std::uint32_t seq
Definition sqlite_chunk_store.cc:186
std::string mimetype
Definition sqlite_chunk_store.cc:191
A unit of data: bytes plus optional descriptive metadata.
Definition types.h:185
One piece of a node's stream: an inline chunk or a node reference.
Definition types.h:335
How an AsyncNode::Finalize() call ends the stream.
Definition async_node.h:75
std::string_view mimetype
MIME type selecting the encoding of a typed value.
Definition async_node.h:87
bool close
Whether to close the writer after the final chunk.
Definition async_node.h:81
std::optional< std::uint32_t > seq
Explicit sequence number for the final chunk; assigned in order when omitted.
Definition async_node.h:84
bool wait
Whether to resolve only once the store confirmed the write and, if close is set, closed.
Definition async_node.h:78
Tunables controlling how a ChunkStoreReader delivers and buffers fragments.
Definition chunk_store_reader.h:51
Tunables controlling how a ChunkStoreWriter batches and buffers chunks.
Definition chunk_store_writer.h:54
A11's core wire value types: chunks, node fragments and messages.