10#ifndef A11_STORES_CHUNK_STORE_WRITER_H_
11#define A11_STORES_CHUNK_STORE_WRITER_H_
18#include <absl/status/status.h>
19#include <absl/status/statusor.h>
100 static absl::StatusOr<std::shared_ptr<ChunkStoreWriter>>
Create(
134 std::optional<std::uint32_t> seq = std::nullopt,
154 data::Chunk chunk, std::optional<std::uint32_t> seq = std::nullopt,
223 absl::Status
AttachStream(std::shared_ptr<net::WireStream> stream);
233 absl::Status
DetachStream(
const std::shared_ptr<net::WireStream>& stream);
246 : state_(std::
move(state)) {}
248 std::shared_ptr<State> state_;
A11's pluggable storage interface for streamed node data: an ordered, appendable log of fragments key...
A buffered, backpressured write cursor over a ChunkStore.
Definition chunk_store_writer.h:86
a11::Task WaitForBufferToDrain()
Wait until the in-flight write buffer empties.
Definition chunk_store_writer.cc:791
absl::Status DetachStream(const std::shared_ptr< net::WireStream > &stream)
Stop mirroring fragments to a previously attached wire stream.
Definition chunk_store_writer.cc:833
absl::Status AttachStream(std::shared_ptr< net::WireStream > stream)
Mirror persisted fragments to an additional wire stream.
Definition chunk_store_writer.cc:811
a11::Task Cancel()
Stop the writer immediately, discarding any queued chunks.
Definition chunk_store_writer.cc:688
~ChunkStoreWriter()=default
bool IsWritable() const
Definition chunk_store_writer.cc:683
a11::Future< std::uint32_t > PutChunk(data::Chunk chunk, std::optional< std::uint32_t > seq=std::nullopt, bool final=false)
Write a chunk and await backing-store confirmation.
Definition chunk_store_writer.cc:668
static absl::StatusOr< std::shared_ptr< ChunkStoreWriter > > Create(std::shared_ptr< ChunkStore > store, ChunkStoreWriterOptions options={})
Create a writer over store.
Definition chunk_store_writer.cc:572
ChunkStoreWrite EnqueueChunk(data::Chunk chunk, std::optional< std::uint32_t > seq=std::nullopt, bool final=false, bool ensure_started=true)
Enqueue a chunk, exposing backpressure and confirmation separately.
Definition chunk_store_writer.cc:595
size_t queue_size() const
Definition chunk_store_writer.cc:851
std::optional< absl::Status > GetAbortStatus() const
Definition chunk_store_writer.cc:678
void EnsureStarted()
Start the background flush loop if it is not already running.
Definition chunk_store_writer.cc:591
std::optional< absl::Status > GetStatus() const
Definition chunk_store_writer.cc:673
ChunkStoreWriterOptions options() const
Definition chunk_store_writer.cc:847
std::shared_ptr< ChunkStore > store() const
Definition chunk_store_writer.cc:843
a11::Task DrainAndClose()
Flush every queued chunk, then close the writer.
Definition chunk_store_writer.cc:721
a11::Task AbortWithStatus(absl::Status status)
Abort the writer with an error status.
Definition chunk_store_writer.cc:748
Completion values used by every asynchronous A11 operation.
Future< T > SubmitWithCancellationHook(absl::AnyInvocable< absl::StatusOr< T >() && > work, std::function< void()> cancellation_hook, thread::TreeOptions tree_options)
Run work on A11's fiber pool with application-specific cancellation.
Definition executor.h:30
A unit of data: bytes plus optional descriptive metadata.
Definition types.h:95
The pair of awaitables returned when enqueuing a chunk, separating queue admission from backing-store...
Definition chunk_store_writer.h:68
a11::Task admitted
Resolves once the chunk is admitted into the bounded queue (backpressure).
Definition chunk_store_writer.h:70
a11::Future< std::uint32_t > confirmation
Resolves with the assigned sequence once the backing store accepts it.
Definition chunk_store_writer.h:72
Tunables controlling how a ChunkStoreWriter batches and buffers chunks.
Definition chunk_store_writer.h:40
absl::Status Validate() const
Validate that the options are internally consistent.
Definition chunk_store_writer.cc:45
std::uint64_t max_chunks_to_write_at_once
Maximum number of chunks flushed to the store per batch.
Definition chunk_store_writer.h:44
bool sticky_mimetype
Whether repeated contiguous chunk mimetypes are omitted when writing.
Definition chunk_store_writer.h:48
std::optional< std::uint64_t > num_chunks_to_buffer
Optional bound on the in-flight write buffer size.
Definition chunk_store_writer.h:46
std::uint32_t offset
Sequence number at which writing begins.
Definition chunk_store_writer.h:42
Definition chunk_store_writer.cc:64
A11's core wire value types: chunks, node fragments and messages.