35#ifndef A11_NODES_ASYNC_NODE_H_
36#define A11_NODES_ASYNC_NODE_H_
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>
61#include "thread/boost_primitives.h"
84 std::optional<std::uint32_t>
seq = {};
111class AsyncNode :
public std::enable_shared_from_this<AsyncNode> {
126 static absl::StatusOr<std::shared_ptr<AsyncNode>>
Create(
127 std::shared_ptr<stores::ChunkStore> store,
144 absl::StatusOr<std::string>
GetId()
const;
152 [[nodiscard]] std::shared_ptr<stores::ChunkStore>
GetChunkStore()
const;
158 [[nodiscard]] std::shared_ptr<data::SerializationRegistry>
167 std::shared_ptr<data::SerializationRegistry> registry);
173 absl::StatusOr<std::shared_ptr<stores::ChunkStoreReader>>
reader();
179 absl::StatusOr<std::shared_ptr<stores::ChunkStoreWriter>>
writer();
201 std::optional<stores::ChunkStoreReaderOptions> options = std::nullopt);
256 data::Chunk chunk, std::optional<std::uint32_t>
seq = std::nullopt,
294 template <
typename T>
297 std::optional<std::uint32_t>
seq = std::nullopt,
bool final =
false) {
298 std::shared_ptr<data::SerializationRegistry> registry;
300 thread::MutexLock lock(&mu_);
301 registry = serialization_registry_;
304 data::MakeChunkObject<T>(std::move(
value), data::SerialTypeTag<T>(),
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;
315 thread::MutexLock lock(&mu_);
316 registry = serialization_registry_;
318 absl::StatusOr<data::Chunk> chunk = registry->ToChunk<T>(
value,
mimetype);
320 return a11::FailedFuture<std::uint32_t>(chunk.status());
368 template <
typename T>
370 std::shared_ptr<data::SerializationRegistry> registry;
372 thread::MutexLock lock(&mu_);
373 registry = serialization_registry_;
375 absl::StatusOr<data::Chunk> chunk =
376 registry->ToChunk<T>(
value, options.mimetype);
380 return Finalize(std::move(*chunk), options);
394 size_t limit, absl::Duration
timeout = absl::InfiniteDuration());
403 absl::Duration
timeout = absl::InfiniteDuration());
412 absl::Duration
timeout = absl::InfiniteDuration());
421 absl::Duration
timeout = absl::InfiniteDuration());
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>> {
444 absl::StatusOr<std::optional<data::NodeFragment>> fragment =
445 self->NextFragmentRaw(
timeout).Await();
446 if (!fragment.ok()) {
447 return fragment.status();
449 if (!fragment->has_value()) {
452 absl::StatusOr<const data::Chunk*> chunk = (*fragment)->GetChunk();
454 return chunk.status();
459 if ((*chunk)->HasObject()) {
463 if constexpr (data::HasSerialTypeTag<T>) {
464 if (std::optional<T> taken =
465 data::TryTakeObject<T>(**chunk, data::SerialTypeTag<T>());
473 data::Chunk materialised = **chunk;
474 ABSL_RETURN_IF_ERROR(materialised.Materialize());
475 std::shared_ptr<data::SerializationRegistry> registry;
477 thread::MutexLock lock(&self->mu_);
478 registry = self->serialization_registry_;
480 absl::StatusOr<T>
value =
481 registry->FromChunk<T>(materialised, mimetype_patterns);
483 return value.status();
485 return std::optional<T>(std::move(*
value));
487 if ((*chunk)->IsNull()) {
488 if ((*fragment)->continued) {
489 return absl::FailedPreconditionError(
490 "A null stream marker must be final");
494 std::shared_ptr<data::SerializationRegistry> registry;
496 thread::MutexLock lock(&self->mu_);
497 registry = self->serialization_registry_;
499 absl::StatusOr<T>
value =
500 registry->FromChunk<T>(**chunk, mimetype_patterns);
502 return value.status();
504 return std::optional<T>(std::move(*
value));
554 absl::Status
AttachStream(std::shared_ptr<net::WireStream> stream);
561 absl::Status
DetachStream(
const std::shared_ptr<net::WireStream>& stream);
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) {}
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_);
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
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.