24#ifndef A11_STORES_CHUNK_STORE_WRITER_H_
25#define A11_STORES_CHUNK_STORE_WRITER_H_
32#include <absl/status/status.h>
33#include <absl/status/statusor.h>
114 static absl::StatusOr<std::shared_ptr<ChunkStoreWriter>>
Create(
170 std::optional<std::uint32_t>
seq = std::nullopt,
171 bool final =
false,
bool ensure_started =
true);
190 data::Chunk chunk, std::optional<std::uint32_t>
seq = std::nullopt,
194 [[nodiscard]] std::optional<absl::Status>
GetStatus()
const;
271 absl::Status
AttachStream(std::shared_ptr<net::WireStream> stream);
281 absl::Status
DetachStream(
const std::shared_ptr<net::WireStream>& stream);
284 [[nodiscard]] std::shared_ptr<ChunkStore>
store()
const;
294 : state_(std::move(state)) {}
296 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:100
void Flush()
Run the flush loop now, on the calling thread, if it is idle.
Definition chunk_store_writer.cc:745
a11::Task WaitForBufferToDrain()
Wait until the in-flight write buffer empties.
Definition chunk_store_writer.cc:958
absl::Status DetachStream(const std::shared_ptr< net::WireStream > &stream)
Stop mirroring fragments to a previously attached wire stream.
Definition chunk_store_writer.cc:997
absl::Status AttachStream(std::shared_ptr< net::WireStream > stream)
Mirror persisted fragments to an additional wire stream.
Definition chunk_store_writer.cc:981
a11::Task Cancel()
Stop the writer immediately, discarding any queued chunks.
Definition chunk_store_writer.cc:855
~ChunkStoreWriter()=default
bool IsWritable() const
Definition chunk_store_writer.cc:850
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:831
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:722
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:758
size_t queue_size() const
Definition chunk_store_writer.cc:1015
std::optional< absl::Status > GetAbortStatus() const
Definition chunk_store_writer.cc:845
void EnsureStarted()
Start the background flush loop if it is not already running.
Definition chunk_store_writer.cc:741
std::optional< absl::Status > GetStatus() const
Definition chunk_store_writer.cc:840
ChunkStoreWriterOptions options() const
Definition chunk_store_writer.cc:1011
std::shared_ptr< ChunkStore > store() const
Definition chunk_store_writer.cc:1007
a11::Task DrainAndClose()
Flush every queued chunk, then close the writer.
Definition chunk_store_writer.cc:888
a11::Task AbortWithStatus(absl::Status status)
Abort the writer with an error status.
Definition chunk_store_writer.cc:915
Completion values used by every asynchronous A11 operation.
std::uint32_t seq
Definition sqlite_chunk_store.cc:186
A unit of data: bytes plus optional descriptive metadata.
Definition types.h:185
The pair of awaitables returned when enqueuing a chunk, separating queue admission from backing-store...
Definition chunk_store_writer.h:82
a11::Task admitted
Resolves once the chunk is admitted into the bounded queue (backpressure).
Definition chunk_store_writer.h:84
a11::Future< std::uint32_t > confirmation
Resolves with the assigned sequence once the backing store accepts it.
Definition chunk_store_writer.h:86
Tunables controlling how a ChunkStoreWriter batches and buffers chunks.
Definition chunk_store_writer.h:54
absl::Status Validate() const
Validate that the options are internally consistent.
Definition chunk_store_writer.cc:58
std::uint64_t max_chunks_to_write_at_once
Maximum number of chunks flushed to the store per batch.
Definition chunk_store_writer.h:58
bool sticky_mimetype
Whether repeated contiguous chunk mimetypes are omitted when writing.
Definition chunk_store_writer.h:62
std::optional< std::uint64_t > num_chunks_to_buffer
Optional bound on the in-flight write buffer size.
Definition chunk_store_writer.h:60
std::uint32_t offset
Sequence number at which writing begins.
Definition chunk_store_writer.h:56
Definition chunk_store_writer.cc:77
A11's core wire value types: chunks, node fragments and messages.