A11 (C++ runtime)
Native C++ implementation of the A11 streaming action runtime
Loading...
Searching...
No Matches
sqlite_chunk_store.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
22#ifndef A11_STORES_SQLITE_CHUNK_STORE_H_
23#define A11_STORES_SQLITE_CHUNK_STORE_H_
24
25#include <cstddef>
26#include <cstdint>
27#include <memory>
28#include <optional>
29#include <string>
30#include <vector>
31
32#include <absl/status/status.h>
33#include <absl/status/statusor.h>
34#include <absl/time/time.h>
35
37#include "a11/data/types.h"
39#include "a11/stores/internal/sqlite_database.h"
40
41namespace a11::stores {
42
51 size_t inline_data_threshold = 128 * 1024;
52
57 std::string owner_id;
58
60 internal::SqliteSynchronous synchronous =
61 internal::SqliteSynchronous::kNormal;
62
70 absl::Duration cross_process_poll_interval = absl::ZeroDuration();
71
75 absl::Duration blob_grace_period = absl::Hours(1);
76
78 absl::Status Validate() const;
80 static absl::StatusOr<SQLiteChunkStoreOptions> FromEnvironment();
81
83 const SQLiteChunkStoreOptions&) = default;
84};
85
88 std::string id;
89 std::string owner_id;
90 bool closed = false;
91 std::optional<absl::Status> status;
92 std::optional<std::uint32_t> final_seq;
93 size_t size = 0;
94 std::uint64_t total_chunks_put = 0;
95 std::uint64_t next_cursor = 0;
96 std::uint64_t data_bytes = 0;
97 std::optional<std::uint32_t> max_seq;
98 std::uint64_t revision = 0;
99 absl::Time created_at = absl::InfinitePast();
100 absl::Time updated_at = absl::InfinitePast();
101};
102
104
132class SQLiteChunkStore final : public ChunkStore {
133 private:
134 struct ConstructorToken {};
135
136 public:
145 static absl::StatusOr<std::shared_ptr<SQLiteChunkStore>> Create(
146 std::string node_id);
147
158 static absl::StatusOr<std::shared_ptr<SQLiteChunkStore>> Create(
159 std::string node_id, std::string root);
160
173 static absl::StatusOr<std::shared_ptr<SQLiteChunkStore>> Create(
174 std::string node_id, const std::string& root,
176
177 ~SQLiteChunkStore() override = default;
178
180 using ChunkStore::Get;
182 using ChunkStore::Next;
183
185 absl::Time deadline) override;
187 std::uint64_t arrival_order, absl::Time deadline) override;
189 absl::Time deadline, size_t limit) override;
192 std::vector<data::NodeFragment> fragments) override;
193 a11::Future<data::NodeFragment> ClearData(std::uint32_t seq) override;
195 std::uint64_t arrival_order) override;
198 absl::Status status, bool return_status_if_already_closed) override;
199 a11::Future<size_t> Size() override;
200 absl::StatusOr<std::string> GetId() const override;
201
204
219
222
224 [[nodiscard]] const SQLiteChunkStoreOptions& options() const {
225 return options_;
226 }
227
229 [[nodiscard]] std::string root() const;
230
231 SQLiteChunkStore(ConstructorToken, std::string node_id,
232 std::shared_ptr<internal::SqliteDatabase> database,
234 : node_id_(std::move(node_id)),
235 database_(std::move(database)),
236 options_(std::move(options)) {}
237
238 private:
240
241 enum class ReadKind { kSequence, kArrivalOrder };
242
243 a11::Future<data::NodeFragment> Read(ReadKind kind, std::uint64_t value,
244 absl::Time deadline);
245
246 const std::string node_id_;
247 const std::shared_ptr<internal::SqliteDatabase> database_;
248 const SQLiteChunkStoreOptions options_;
249};
250
260 private:
261 struct ConstructorToken {};
262
263 public:
274 static absl::StatusOr<std::shared_ptr<SQLiteChunkStoreFactory>> Create(
275 const std::string& root, SQLiteChunkStoreOptions options);
276
278 static absl::StatusOr<std::shared_ptr<SQLiteChunkStoreFactory>> Create(
279 const std::string& root);
280
284 static absl::StatusOr<std::shared_ptr<SQLiteChunkStoreFactory>> Create();
285
295 static std::string DefaultRoot();
296
305 absl::StatusOr<std::shared_ptr<SQLiteChunkStore>> Open(std::string node_id);
306
308 [[nodiscard]] std::string root() const;
309
311 [[nodiscard]] const SQLiteChunkStoreOptions& options() const {
312 return options_;
313 }
314
317
318 SQLiteChunkStoreFactory(ConstructorToken,
319 std::shared_ptr<internal::SqliteDatabase> database,
321 : database_(std::move(database)), options_(std::move(options)) {}
322
323 private:
324 const std::shared_ptr<internal::SqliteDatabase> database_;
325 const SQLiteChunkStoreOptions options_;
326};
327
328} // namespace a11::stores
329
330#endif // A11_STORES_SQLITE_CHUNK_STORE_H_
A11's pluggable storage interface for streamed node data: an ordered, appendable log of fragments key...
Shared handle to one asynchronous result.
Definition future.h:126
Abstract, pluggable backing store for the data of a node: an ordered, appendable log of fragments.
Definition chunk_store.h:67
a11::Future< data::NodeFragment > GetByArrivalOrder(std::uint64_t arrival_order)
Get a fragment by arrival order, waiting indefinitely.
Definition chunk_store.h:116
a11::Future< data::NodeFragment > Get(std::uint32_t seq)
Get the fragment at a sequence number, waiting indefinitely.
Definition chunk_store.h:82
a11::Future< std::vector< std::optional< data::NodeFragment > > > Next()
Get the next logical-sequence fragment, waiting indefinitely.
Definition chunk_store.h:149
a11::Future< absl::Status > CloseWritesWithStatus(absl::Status status)
Seal the store against further writes with a terminal status.
Definition chunk_store.h:268
Creates SQLiteChunkStores rooted at one directory.
Definition sqlite_chunk_store.h:259
a11::Future< size_t > SweepOrphanBlobs()
Delete unreferenced blob files older than the configured grace period.
Definition sqlite_chunk_store.cc:588
const SQLiteChunkStoreOptions & options() const
The validated storage policy applied to every store created here.
Definition sqlite_chunk_store.h:311
std::string root() const
The root this factory creates stores under.
Definition sqlite_chunk_store.cc:584
static std::string DefaultRoot()
The process-wide default storage root.
Definition sqlite_chunk_store.cc:529
SQLiteChunkStoreFactory(ConstructorToken, std::shared_ptr< internal::SqliteDatabase > database, SQLiteChunkStoreOptions options)
Definition sqlite_chunk_store.h:318
absl::StatusOr< std::shared_ptr< SQLiteChunkStore > > Open(std::string node_id)
Open a store for node_id under this factory's root.
Definition sqlite_chunk_store.cc:576
static absl::StatusOr< std::shared_ptr< SQLiteChunkStoreFactory > > Create()
Create a factory at the default root with the environment/default policy.
Definition sqlite_chunk_store.cc:572
A persistent ChunkStore backed by one SQLite database and a blob directory.
Definition sqlite_chunk_store.h:132
absl::StatusOr< std::string > GetId() const override
Get the store's node identifier.
Definition sqlite_chunk_store.cc:620
a11::Future< std::uint32_t > GetSeqForArrivalOrder(std::uint64_t arrival_order) override
Translate an arrival order into the sequence number of that fragment.
Definition sqlite_chunk_store.cc:1239
a11::Future< size_t > Size() override
Get the number of fragments currently in the store.
Definition sqlite_chunk_store.cc:1291
a11::Future< data::NodeFragment > Get(std::uint32_t seq, absl::Time deadline) override
Get the fragment stored at a sequence number.
Definition sqlite_chunk_store.cc:631
a11::Future< size_t > SweepOrphanBlobs()
Delete unreferenced blob files older than the configured grace period.
Definition sqlite_chunk_store.cc:1400
a11::Future< SQLiteChunkStoreMetadata > GetMetadata()
Read all node-level state in one row read, without listing fragments.
Definition sqlite_chunk_store.cc:1307
static absl::StatusOr< std::shared_ptr< SQLiteChunkStore > > Create(std::string node_id)
Create a store for node_id under the process-default root.
Definition sqlite_chunk_store.cc:615
a11::Future< std::vector< std::uint32_t > > PutMany(std::vector< data::NodeFragment > fragments) override
Append several fragments in one batch.
Definition sqlite_chunk_store.cc:869
a11::Future< absl::Status > CloseWritesWithStatus(absl::Status status, bool return_status_if_already_closed) override
Seal the store against further writes with a terminal status.
Definition sqlite_chunk_store.cc:1172
~SQLiteChunkStore() override=default
a11::Future< data::NodeFragment > ClearData(std::uint32_t seq) override
Erase the payload of the fragment at a sequence number while keeping its slot.
Definition sqlite_chunk_store.cc:1085
a11::Future< std::vector< std::optional< data::NodeFragment > > > Next()
Get the next logical-sequence fragment, waiting indefinitely.
Definition chunk_store.h:149
SQLiteChunkStore(ConstructorToken, std::string node_id, std::shared_ptr< internal::SqliteDatabase > database, SQLiteChunkStoreOptions options)
Definition sqlite_chunk_store.h:231
a11::Future< std::uint32_t > Put(data::NodeFragment fragment) override
Append one fragment to the log.
Definition sqlite_chunk_store.cc:861
const SQLiteChunkStoreOptions & options() const
The validated storage policy captured at construction.
Definition sqlite_chunk_store.h:224
std::string root() const
The storage root this store reads and writes under.
Definition sqlite_chunk_store.cc:624
a11::Future< std::optional< std::uint32_t > > GetFinalSeq() override
Get the sequence number explicitly marked as the final fragment.
Definition sqlite_chunk_store.cc:1274
a11::Future< data::NodeFragment > GetByArrivalOrder(std::uint64_t arrival_order, absl::Time deadline) override
Get the fragment identified by the order in which it arrived, rather than by its sequence number.
Definition sqlite_chunk_store.cc:636
a11::Future< std::vector< data::NodeFragment > > FindReferrers(size_t limit)
Find fragments elsewhere in the database whose NodeRef points at this node, newest sequence last.
Definition sqlite_chunk_store.cc:1336
std::string value
Definition discover.cc:114
TokenKind kind
Definition format.cc:913
Completion values used by every asynchronous A11 operation.
Definition node_map.h:45
std::vector< std::optional< data::NodeFragment > > fragments
Definition redis_chunk_store.cc:240
std::uint32_t seq
Definition sqlite_chunk_store.cc:186
std::uint64_t arrival_order
Definition sqlite_chunk_store.cc:187
One piece of a node's stream: an inline chunk or a node reference.
Definition types.h:335
Node-level state read from the node row, without touching any fragment.
Definition sqlite_chunk_store.h:87
std::uint64_t total_chunks_put
Lifetime accepted write count.
Definition sqlite_chunk_store.h:94
absl::Time updated_at
Most recent mutation.
Definition sqlite_chunk_store.h:100
absl::Time created_at
First accepted write.
Definition sqlite_chunk_store.h:99
std::optional< absl::Status > status
Terminal status when closed.
Definition sqlite_chunk_store.h:91
std::uint64_t data_bytes
Cached total of stored payload bytes.
Definition sqlite_chunk_store.h:96
std::uint64_t next_cursor
Next seq the shared Next() will want.
Definition sqlite_chunk_store.h:95
std::string owner_id
Recorded owner, possibly empty.
Definition sqlite_chunk_store.h:89
std::uint64_t revision
Monotonic state-change revision.
Definition sqlite_chunk_store.h:98
std::optional< std::uint32_t > max_seq
Greatest assigned sequence.
Definition sqlite_chunk_store.h:97
bool closed
Whether writes have been sealed.
Definition sqlite_chunk_store.h:90
size_t size
Fragment slots, tombstones included.
Definition sqlite_chunk_store.h:93
std::string id
Node id, as stored.
Definition sqlite_chunk_store.h:88
std::optional< std::uint32_t > final_seq
Logical final fragment, if any.
Definition sqlite_chunk_store.h:92
Storage layout policy for SQLiteChunkStore.
Definition sqlite_chunk_store.h:44
std::string owner_id
Recorded on the node row as its owner.
Definition sqlite_chunk_store.h:57
friend bool operator==(const SQLiteChunkStoreOptions &, const SQLiteChunkStoreOptions &)=default
absl::Status Validate() const
Validate the threshold, owner id and durations.
Definition sqlite_chunk_store.cc:458
absl::Duration blob_grace_period
How long an unreferenced blob must survive before a sweep may remove it.
Definition sqlite_chunk_store.h:75
size_t inline_data_threshold
Payloads strictly larger than this move out of the row into a blob file.
Definition sqlite_chunk_store.h:51
absl::Duration cross_process_poll_interval
How often to look for commits made by other processes, or zero to disable.
Definition sqlite_chunk_store.h:70
static absl::StatusOr< SQLiteChunkStoreOptions > FromEnvironment()
Read storage policy from the A11_SQLITE_CHUNK_STORE_* environment values.
Definition sqlite_chunk_store.cc:478
internal::SqliteSynchronous synchronous
Durability level applied with PRAGMA synchronous.
Definition sqlite_chunk_store.h:60
A11's core wire value types: chunks, node fragments and messages.