A11 (C++ runtime)
Native C++ implementation of the A11 streaming action runtime
Loading...
Searching...
No Matches
redis_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_REDIS_CHUNK_STORE_H_
23#define A11_STORES_REDIS_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 "redis/client.h"
40
41namespace a11::stores {
42
46 std::string key_prefix = "a11:";
47
49 size_t inline_data_threshold = 256 * 1024;
50
52 absl::Status Validate() const;
54 static absl::StatusOr<RedisChunkStoreOptions> FromEnvironment();
55
57 const RedisChunkStoreOptions&) = default;
58};
59
62 std::string metadata;
63 std::string stream;
64 std::string sequence_index;
65 std::string arrival_index;
66 std::string blobs;
67 std::string events;
68
70 [[nodiscard]] std::vector<std::string> ScriptKeys() const;
71
72 friend bool operator==(const RedisChunkStoreKeys&,
73 const RedisChunkStoreKeys&) = default;
74};
75
78 std::string id;
79 bool closed = false;
80 std::optional<absl::Status> status;
81 std::optional<std::uint32_t>
83 size_t size = 0;
84 std::uint64_t total_chunks_put = 0;
85 std::uint64_t next_cursor = 0;
86 std::optional<std::uint32_t> max_seq;
87 std::uint64_t revision = 0;
88};
89
107class RedisChunkStore final : public ChunkStore {
108 private:
109 struct ConstructorToken {};
110
111 public:
113 static absl::StatusOr<std::shared_ptr<RedisChunkStore>> Create(
114 std::string node_id, std::shared_ptr<redis::Client> client,
116
118 static absl::StatusOr<std::shared_ptr<RedisChunkStore>> Create(
119 std::string node_id, std::shared_ptr<redis::Client> client);
120
122 static absl::StatusOr<std::shared_ptr<RedisChunkStore>> Create(
123 std::string node_id);
124
125 ~RedisChunkStore() override = default;
126
128 using ChunkStore::Get;
130 using ChunkStore::Next;
131
133 absl::Time deadline) override;
135 std::uint64_t arrival_order, absl::Time deadline) override;
137 absl::Time deadline, size_t limit) override;
140 std::vector<data::NodeFragment> fragments) override;
141 a11::Future<data::NodeFragment> ClearData(std::uint32_t seq) override;
143 std::uint64_t arrival_order) override;
146 absl::Status status, bool return_status_if_already_closed) override;
147 a11::Future<size_t> Size() override;
148 absl::StatusOr<std::string> GetId() const override;
149
152
155
157 [[nodiscard]] std::shared_ptr<redis::Client> client() const {
158 return client_;
159 }
160
162 [[nodiscard]] const RedisChunkStoreOptions& options() const {
163 return options_;
164 }
165
167 [[nodiscard]] const RedisChunkStoreKeys& keys() const { return keys_; }
168
169 RedisChunkStore(ConstructorToken, std::string node_id,
170 std::shared_ptr<redis::Client> client,
172 : node_id_(std::move(node_id)),
173 client_(std::move(client)),
174 options_(std::move(options)),
175 keys_(std::move(keys)) {}
176
177 private:
178 enum class ReadKind { kSequence, kArrivalOrder };
179
180 a11::Future<data::NodeFragment> Read(ReadKind kind, std::uint64_t value,
181 absl::Time deadline);
182
183 const std::string node_id_;
184 const std::shared_ptr<redis::Client> client_;
185 const RedisChunkStoreOptions options_;
186 const RedisChunkStoreKeys keys_;
187};
188
189} // namespace a11::stores
190
191#endif // A11_STORES_REDIS_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
A persistent, multi-process ChunkStore backed by Redis Streams.
Definition redis_chunk_store.h:107
RedisChunkStore(ConstructorToken, std::string node_id, std::shared_ptr< redis::Client > client, RedisChunkStoreOptions options, RedisChunkStoreKeys keys)
Definition redis_chunk_store.h:169
a11::Future< std::uint32_t > GetSeqForArrivalOrder(std::uint64_t arrival_order) override
Translate an arrival order into the sequence number of that fragment.
Definition redis_chunk_store.cc:686
a11::Future< data::NodeFragment > Get(std::uint32_t seq, absl::Time deadline) override
Get the fragment stored at a sequence number.
Definition redis_chunk_store.cc:465
a11::Task Initialize()
Ensure that the metadata hash exists, without writing chunk data.
Definition redis_chunk_store.cc:798
std::shared_ptr< redis::Client > client() const
Return the shared Redis client used for commands and subscriptions.
Definition redis_chunk_store.h:157
const RedisChunkStoreKeys & keys() const
Return the sharding-safe key set owned by this node.
Definition redis_chunk_store.h:167
a11::Future< size_t > Size() override
Get the number of fragments currently in the store.
Definition redis_chunk_store.cc:776
a11::Future< std::vector< std::uint32_t > > PutMany(std::vector< data::NodeFragment > fragments) override
Append several fragments in one batch.
Definition redis_chunk_store.cc:571
a11::Future< RedisChunkStoreMetadata > GetMetadata()
Read all node-level metadata without iterating over stream entries.
Definition redis_chunk_store.cc:815
a11::Future< std::vector< std::optional< data::NodeFragment > > > Next()
Get the next logical-sequence fragment, waiting indefinitely.
Definition chunk_store.h:149
static absl::StatusOr< std::shared_ptr< RedisChunkStore > > Create(std::string node_id, std::shared_ptr< redis::Client > client, RedisChunkStoreOptions options)
Create a store with an injected client and explicit storage policy.
Definition redis_chunk_store.cc:346
a11::Future< std::uint32_t > Put(data::NodeFragment fragment) override
Append one fragment to the log.
Definition redis_chunk_store.cc:563
const RedisChunkStoreOptions & options() const
Return the validated storage policy captured at construction.
Definition redis_chunk_store.h:162
absl::StatusOr< std::string > GetId() const override
Get the store's node identifier.
Definition redis_chunk_store.cc:794
a11::Future< std::optional< std::uint32_t > > GetFinalSeq() override
Get the sequence number explicitly marked as the final fragment.
Definition redis_chunk_store.cc:708
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 redis_chunk_store.cc:470
~RedisChunkStore() override=default
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 redis_chunk_store.cc:727
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 redis_chunk_store.cc:668
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
The sharding-safe Redis keys owned by one node stream.
Definition redis_chunk_store.h:61
std::string events
Pub/Sub invalidation channel for waiting readers.
Definition redis_chunk_store.h:67
std::vector< std::string > ScriptKeys() const
Keys in the stable order expected by the store's Lua state machine.
Definition redis_chunk_store.cc:342
std::string blobs
Encoded chunks stored outside stream fields.
Definition redis_chunk_store.h:66
friend bool operator==(const RedisChunkStoreKeys &, const RedisChunkStoreKeys &)=default
std::string metadata
Node state hash.
Definition redis_chunk_store.h:62
std::string stream
Ordered chunk/control Redis Stream.
Definition redis_chunk_store.h:63
std::string arrival_index
Arrival-order-to-sequence hash.
Definition redis_chunk_store.h:65
std::string sequence_index
Sequence-to-stream-entry hash.
Definition redis_chunk_store.h:64
Node-level state read directly from the metadata hash, without chunks.
Definition redis_chunk_store.h:77
std::uint64_t revision
Monotonic state-change revision.
Definition redis_chunk_store.h:87
std::uint64_t total_chunks_put
Lifetime successful write count.
Definition redis_chunk_store.h:84
std::uint64_t next_cursor
Cursor used by Next() reads.
Definition redis_chunk_store.h:85
std::optional< absl::Status > status
Terminal status when closed.
Definition redis_chunk_store.h:80
std::string id
Node id persisted with the key set.
Definition redis_chunk_store.h:78
size_t size
Number of fragment slots currently indexed.
Definition redis_chunk_store.h:83
std::optional< std::uint32_t > final_seq
Logical final fragment, if marked.
Definition redis_chunk_store.h:82
bool closed
Whether writes have been sealed.
Definition redis_chunk_store.h:79
std::optional< std::uint32_t > max_seq
Greatest assigned sequence.
Definition redis_chunk_store.h:86
Storage layout policy for RedisChunkStore.
Definition redis_chunk_store.h:44
friend bool operator==(const RedisChunkStoreOptions &, const RedisChunkStoreOptions &)=default
static absl::StatusOr< RedisChunkStoreOptions > FromEnvironment()
Read storage policy from the A11_REDIS_CHUNK_STORE_* environment values.
Definition redis_chunk_store.cc:324
std::string key_prefix
Prefix before the per-node Redis Cluster hash tag.
Definition redis_chunk_store.h:46
size_t inline_data_threshold
Raw chunk bytes larger than this are moved to the separate blob hash.
Definition redis_chunk_store.h:49
absl::Status Validate() const
Validate the key prefix and inline-data threshold.
Definition redis_chunk_store.cc:313
A11's core wire value types: chunks, node fragments and messages.