8#ifndef A11_STORES_REDIS_CHUNK_STORE_H_
9#define A11_STORES_REDIS_CHUNK_STORE_H_
18#include <absl/status/status.h>
19#include <absl/status/statusor.h>
20#include <absl/time/time.h>
25#include "redis/client.h"
67 std::optional<std::uint32_t>
95 struct ConstructorToken {};
99 static absl::StatusOr<std::shared_ptr<RedisChunkStore>>
Create(
100 std::string node_id, std::shared_ptr<redis::Client>
client,
104 static absl::StatusOr<std::shared_ptr<RedisChunkStore>>
Create(
105 std::string node_id, std::shared_ptr<redis::Client>
client);
108 static absl::StatusOr<std::shared_ptr<RedisChunkStore>>
Create(
109 std::string node_id);
119 absl::Time deadline)
override;
123 absl::Time deadline,
size_t limit)
override;
126 std::vector<data::NodeFragment>
fragments)
override;
134 absl::StatusOr<std::string>
GetId()
const override;
156 std::shared_ptr<redis::Client>
client,
158 : node_id_(std::
move(node_id)),
164 enum class ReadKind { kSequence, kArrivalOrder };
167 absl::Time deadline);
169 const std::string node_id_;
170 const std::shared_ptr<redis::Client> client_;
171 const RedisChunkStoreOptions options_;
172 const RedisChunkStoreKeys keys_;
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:110
Abstract, pluggable backing store for the data of a node: an ordered, appendable log of fragments.
Definition chunk_store.h:53
a11::Future< data::NodeFragment > GetByArrivalOrder(std::uint64_t arrival_order)
Get a fragment by arrival order, waiting indefinitely.
Definition chunk_store.h:102
a11::Future< data::NodeFragment > Get(std::uint32_t seq)
Get the fragment at a sequence number, waiting indefinitely.
Definition chunk_store.h:68
a11::Future< std::vector< std::optional< data::NodeFragment > > > Next()
Get the next logical-sequence fragment, waiting indefinitely.
Definition chunk_store.h:135
a11::Future< absl::Status > CloseWritesWithStatus(absl::Status status)
Seal the store against further writes with a terminal status.
Definition chunk_store.h:250
A persistent, multi-process ChunkStore backed by Redis Streams.
Definition redis_chunk_store.h:93
RedisChunkStore(ConstructorToken, std::string node_id, std::shared_ptr< redis::Client > client, RedisChunkStoreOptions options, RedisChunkStoreKeys keys)
Definition redis_chunk_store.h:155
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:757
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:512
a11::Task Initialize()
Ensure that the metadata hash exists, without writing chunk data.
Definition redis_chunk_store.cc:890
std::shared_ptr< redis::Client > client() const
Return the shared Redis client used for commands and subscriptions.
Definition redis_chunk_store.h:143
const RedisChunkStoreKeys & keys() const
Return the sharding-safe key set owned by this node.
Definition redis_chunk_store.h:153
a11::Future< size_t > Size() override
Get the number of fragments currently in the store.
Definition redis_chunk_store.cc:861
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:629
a11::Future< RedisChunkStoreMetadata > GetMetadata()
Read all node-level metadata without iterating over stream entries.
Definition redis_chunk_store.cc:912
a11::Future< std::vector< std::optional< data::NodeFragment > > > Next()
Get the next logical-sequence fragment, waiting indefinitely.
Definition chunk_store.h:135
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:376
a11::Future< std::uint32_t > Put(data::NodeFragment fragment) override
Append a single fragment to the log.
Definition redis_chunk_store.cc:606
const RedisChunkStoreOptions & options() const
Return the validated storage policy captured at construction.
Definition redis_chunk_store.h:148
absl::StatusOr< std::string > GetId() const override
Get the store's node identifier.
Definition redis_chunk_store.cc:886
a11::Future< std::optional< std::uint32_t > > GetFinalSeq() override
Get the sequence number explicitly marked as the final fragment.
Definition redis_chunk_store.cc:786
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:517
~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:812
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:736
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
std::vector< std::optional< data::NodeFragment > > fragments
Definition redis_chunk_store.cc:235
One piece of a node's stream: an inline chunk or a node reference.
Definition types.h:167
The sharding-safe Redis keys owned by one node stream.
Definition redis_chunk_store.h:47
std::string events
Pub/Sub invalidation channel for waiting readers.
Definition redis_chunk_store.h:53
std::vector< std::string > ScriptKeys() const
Keys in the stable order expected by the store's Lua state machine.
Definition redis_chunk_store.cc:372
std::string blobs
Encoded chunks stored outside stream fields.
Definition redis_chunk_store.h:52
friend bool operator==(const RedisChunkStoreKeys &, const RedisChunkStoreKeys &)=default
std::string metadata
Node state hash.
Definition redis_chunk_store.h:48
std::string stream
Ordered chunk/control Redis Stream.
Definition redis_chunk_store.h:49
std::string arrival_index
Arrival-order-to-sequence hash.
Definition redis_chunk_store.h:51
std::string sequence_index
Sequence-to-stream-entry hash.
Definition redis_chunk_store.h:50
Storage layout policy for RedisChunkStore.
Definition redis_chunk_store.h:30
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:352
std::string key_prefix
Prefix before the per-node Redis Cluster hash tag.
Definition redis_chunk_store.h:32
size_t inline_data_threshold
Raw chunk bytes larger than this are moved to the separate blob hash.
Definition redis_chunk_store.h:35
absl::Status Validate() const
Validate the key prefix and inline-data threshold.
Definition redis_chunk_store.cc:341
A11's core wire value types: chunks, node fragments and messages.