A11 (C++ runtime)
Native C++ implementation of the A11 streaming action runtime
Loading...
Searching...
No Matches
chunk_store_writer.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
24#ifndef A11_STORES_CHUNK_STORE_WRITER_H_
25#define A11_STORES_CHUNK_STORE_WRITER_H_
26
27#include <cstddef>
28#include <cstdint>
29#include <memory>
30#include <optional>
31
32#include <absl/status/status.h>
33#include <absl/status/statusor.h>
34
36#include "a11/data/types.h"
38
39namespace a11::net {
40class WireStream;
41} // namespace a11::net
42
43namespace a11::stores {
44
56 std::uint32_t offset = 0;
58 std::uint64_t max_chunks_to_write_at_once = 8;
60 std::optional<std::uint64_t> num_chunks_to_buffer = {};
62 bool sticky_mimetype = false;
63
70 absl::Status Validate() const;
71};
72
88
101 public:
114 static absl::StatusOr<std::shared_ptr<ChunkStoreWriter>> Create(
115 std::shared_ptr<ChunkStore> store, ChunkStoreWriterOptions options = {});
116
117 ~ChunkStoreWriter() = default;
118
125 void EnsureStarted();
126
147 void Flush();
148
170 std::optional<std::uint32_t> seq = std::nullopt,
171 bool final = false, bool ensure_started = true);
172
190 data::Chunk chunk, std::optional<std::uint32_t> seq = std::nullopt,
191 bool final = false);
192
194 [[nodiscard]] std::optional<absl::Status> GetStatus() const;
196 [[nodiscard]] std::optional<absl::Status> GetAbortStatus() const;
199 [[nodiscard]] bool IsWritable() const;
200
208
228
240 a11::Task AbortWithStatus(absl::Status status);
241
255
271 absl::Status AttachStream(std::shared_ptr<net::WireStream> stream);
272
281 absl::Status DetachStream(const std::shared_ptr<net::WireStream>& stream);
282
284 [[nodiscard]] std::shared_ptr<ChunkStore> store() const;
286 [[nodiscard]] ChunkStoreWriterOptions options() const;
288 [[nodiscard]] size_t queue_size() const;
289
290 private:
291 struct State;
292
293 explicit ChunkStoreWriter(std::shared_ptr<State> state)
294 : state_(std::move(state)) {}
295
296 std::shared_ptr<State> state_;
297};
298
299} // namespace a11::stores
300
301#endif // A11_STORES_CHUNK_STORE_WRITER_H_
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
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.
Definition action.h:65
Definition node_map.h:45
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.