reifydb_store_multi/store/
mod.rs1use std::{
5 ops::Deref,
6 sync::{Arc, OnceLock},
7};
8
9use reifydb_codec::key::encoded::EncodedKey;
10#[cfg(all(feature = "sqlite", not(target_arch = "wasm32")))]
11use reifydb_core::metrics::sample::MetricsSample;
12use reifydb_core::{
13 common::CommitVersion, event::EventBus, lifecycle::watermark::EvictionWatermark,
14 metrics::collect::MetricsCollector,
15};
16use reifydb_runtime::{actor::system::ActorSystem, context::clock::Clock, shutdown::Shutdown, sync::rwlock::RwLock};
17#[cfg(all(feature = "sqlite", not(target_arch = "wasm32")))]
18use reifydb_sqlite::SqliteTempPathGuard;
19use reifydb_value::util::cowvec::CowVec;
20use tracing::instrument;
21
22use crate::{
23 CommitBufferConfig,
24 config::MultiStoreConfig,
25 flush::{ObjectPersistence, engine::FlushEngine},
26 tier::{
27 commit::buffer::MultiCommitBufferTier,
28 persistent::MultiPersistentTier,
29 read::{MultiReadBufferTier, ReadBufferShardMetrics},
30 },
31};
32#[cfg(all(feature = "sqlite", not(target_arch = "wasm32")))]
33use crate::{config::PersistentConfig, tier::read::ReadBufferConfig};
34
35pub mod multi;
36pub mod router;
37
38use crate::Result;
39
40#[cfg(all(feature = "sqlite", not(target_arch = "wasm32")))]
41struct SqlitePageCacheCollector {
42 persistent: MultiPersistentTier,
43}
44
45#[cfg(all(feature = "sqlite", not(target_arch = "wasm32")))]
46impl MetricsCollector for SqlitePageCacheCollector {
47 fn collect(&self, out: &mut Vec<MetricsSample>) {
48 let metrics = self.persistent.page_cache_metrics();
49 out.push(MetricsSample::bytes("sqlite::multi", "page_cache_used_bytes", metrics.used));
50 out.push(MetricsSample::counter("sqlite::multi", "page_cache_hit_count", metrics.hits.as_u64()));
51 out.push(MetricsSample::counter("sqlite::multi", "page_cache_miss_count", metrics.misses.as_u64()));
52 out.push(MetricsSample::count(
53 "sqlite::multi",
54 "page_cache_sampled_connections",
55 metrics.connections_sampled.as_u64(),
56 ));
57 }
58}
59
60#[derive(Clone)]
61pub struct StandardMultiStore(Arc<StandardMultiStoreInner>);
62
63pub struct StandardMultiStoreInner {
64 pub(crate) commit: MultiCommitBufferTier,
65 pub(crate) persistent: Option<MultiPersistentTier>,
66 pub(crate) read: Option<MultiReadBufferTier>,
67
68 #[allow(dead_code)]
69 pub(crate) flush_engine: Option<Arc<FlushEngine>>,
70 #[allow(dead_code)]
71 pub(crate) row_settings_provider: Arc<OnceLock<Arc<dyn ObjectPersistence>>>,
72 #[allow(dead_code)]
73 pub(crate) eviction_watermark: Arc<RwLock<Option<Arc<dyn EvictionWatermark>>>>,
74
75 pub(crate) event_bus: EventBus,
76}
77
78impl StandardMultiStore {
79 #[instrument(name = "store::multi::new", level = "debug", skip(config), fields(
80 has_persistent = config.persistent.is_some(),
81 ))]
82 pub fn new(config: MultiStoreConfig) -> Result<Self> {
83 let commit = config.commit.storage;
84
85 let row_settings_provider: Arc<OnceLock<Arc<dyn ObjectPersistence>>> = Arc::new(OnceLock::new());
86
87 let eviction_watermark: Arc<RwLock<Option<Arc<dyn EvictionWatermark>>>> = Arc::new(RwLock::new(None));
88
89 let read =
90 config.persistent.is_some().then(|| config.read.and_then(MultiReadBufferTier::new)).flatten();
91
92 #[cfg(all(feature = "sqlite", not(target_arch = "wasm32")))]
93 let (persistent, flush_engine) = {
94 let persistent_config = config.persistent.clone();
95 let persistent = persistent_config.as_ref().map(|c| c.storage.clone());
96 let flush_engine = match (persistent.as_ref(), persistent_config.as_ref()) {
97 (Some(persistent_storage), Some(_)) => Some(Arc::new(FlushEngine::new(
98 commit.clone(),
99 persistent_storage.clone(),
100 row_settings_provider.clone(),
101 eviction_watermark.clone(),
102 read.clone(),
103 config.clock.clone(),
104 config.event_bus.clone(),
105 ))),
106 _ => None,
107 };
108 (persistent, flush_engine)
109 };
110
111 #[cfg(not(all(feature = "sqlite", not(target_arch = "wasm32"))))]
112 let (persistent, flush_engine): (Option<MultiPersistentTier>, Option<Arc<FlushEngine>>) = {
113 let _ = config.persistent;
114 (None, None)
115 };
116
117 let read = persistent.as_ref().and(read);
118
119 Ok(Self(Arc::new(StandardMultiStoreInner {
120 commit,
121 persistent,
122 read,
123 flush_engine,
124 row_settings_provider,
125 eviction_watermark,
126 event_bus: config.event_bus,
127 })))
128 }
129
130 pub fn flush_engine(&self) -> Option<Arc<FlushEngine>> {
131 self.flush_engine.clone()
132 }
133
134 pub fn insert_read_key(&self, key: EncodedKey, version: CommitVersion, value: Option<CowVec<u8>>) {
135 if let Some(read) = &self.read {
136 read.insert(key, version, value);
137 }
138 }
139
140 pub fn invalidate_read_key(&self, key: &EncodedKey) {
141 if let Some(read) = &self.read {
142 read.invalidate(key);
143 }
144 }
145
146 pub fn clear_read(&self) {
147 if let Some(read) = &self.read {
148 read.clear();
149 }
150 }
151
152 pub fn set_row_settings_provider(&self, provider: Arc<dyn ObjectPersistence>) {
153 let _ = self.row_settings_provider.set(provider);
154 }
155
156 pub fn set_eviction_watermark(&self, watermark: Arc<dyn EvictionWatermark>) {
157 *self.eviction_watermark.write() = Some(watermark);
158 }
159
160 pub fn clear_eviction_watermark(&self) {
161 *self.eviction_watermark.write() = None;
162 }
163
164 pub fn commit(&self) -> &MultiCommitBufferTier {
165 &self.commit
166 }
167
168 pub fn event_bus(&self) -> &EventBus {
169 &self.event_bus
170 }
171
172 pub fn metrics_collectors(&self) -> Vec<Arc<dyn MetricsCollector>> {
173 let mut collectors: Vec<Arc<dyn MetricsCollector>> = Vec::new();
174 if let Some(read) = &self.read {
175 collectors.push(Arc::new(read.clone()));
176 }
177 collectors.push(Arc::new(self.commit.clone()));
178 #[cfg(all(feature = "sqlite", not(target_arch = "wasm32")))]
179 if let Some(persistent) = &self.persistent {
180 collectors.push(Arc::new(SqlitePageCacheCollector {
181 persistent: persistent.clone(),
182 }));
183 }
184 collectors
185 }
186
187 pub fn persistent(&self) -> Option<&MultiPersistentTier> {
188 self.persistent.as_ref()
189 }
190
191 pub fn read_buffer_shard_metrics(&self) -> Vec<ReadBufferShardMetrics> {
192 self.read.as_ref().map(|read| read.shard_metrics()).unwrap_or_default()
193 }
194
195 pub fn flush_pending_blocking(&self) {
196 #[cfg(all(feature = "sqlite", not(target_arch = "wasm32")))]
197 if let Some(engine) = self.flush_engine.as_ref() {
198 self.event_bus.wait_for_completion();
199 engine.flush_pending();
200 }
201 }
202
203 pub fn flush_all_blocking(&self) {
204 #[cfg(all(feature = "sqlite", not(target_arch = "wasm32")))]
205 if let Some(engine) = self.flush_engine.as_ref() {
206 self.event_bus.wait_for_completion();
207 engine.flush_all();
208 }
209 }
210}
211
212impl Deref for StandardMultiStore {
213 type Target = StandardMultiStoreInner;
214
215 fn deref(&self) -> &Self::Target {
216 &self.0
217 }
218}
219
220impl Shutdown for StandardMultiStore {
221 fn shutdown(&self) {
222 if let Some(persistent) = self.persistent.as_ref() {
223 persistent.shutdown();
224 }
225 }
226}
227
228impl StandardMultiStore {
229 pub fn testing_memory() -> Self {
230 let clock = Clock::testing();
231 let actor_system = ActorSystem::testing(clock.clone());
232 let spawner = actor_system.spawner();
233 let event_bus = EventBus::new(&spawner);
234 Self::new(MultiStoreConfig {
235 commit: CommitBufferConfig {
236 storage: MultiCommitBufferTier::memory(),
237 },
238 persistent: None,
239 read: None,
240 retention: Default::default(),
241 merge_config: Default::default(),
242 event_bus,
243 spawner,
244 clock,
245 })
246 .unwrap()
247 }
248
249 pub fn testing_memory_with_eventbus(event_bus: EventBus) -> Self {
250 let clock = Clock::testing();
251 let actor_system = ActorSystem::testing(clock.clone());
252 let spawner = actor_system.spawner();
253 Self::new(MultiStoreConfig {
254 commit: CommitBufferConfig {
255 storage: MultiCommitBufferTier::memory(),
256 },
257 persistent: None,
258 read: None,
259 retention: Default::default(),
260 merge_config: Default::default(),
261 event_bus,
262 spawner,
263 clock,
264 })
265 .unwrap()
266 }
267
268 #[cfg(all(feature = "sqlite", not(target_arch = "wasm32")))]
269 pub fn testing_memory_with_persistent_sqlite() -> (Self, SqliteTempPathGuard) {
270 Self::testing_memory_with_persistent_sqlite_read(ReadBufferConfig::default())
271 }
272
273 #[cfg(all(feature = "sqlite", not(target_arch = "wasm32")))]
274 pub fn testing_memory_with_persistent_sqlite_read(read: ReadBufferConfig) -> (Self, SqliteTempPathGuard) {
275 let clock = Clock::testing();
276 let actor_system = ActorSystem::testing(clock.clone());
277 let spawner = actor_system.spawner();
278 let event_bus = EventBus::new(&spawner);
279 let (persistent, guard) = PersistentConfig::sqlite_in_memory();
280 let store = Self::new(MultiStoreConfig {
281 commit: CommitBufferConfig {
282 storage: MultiCommitBufferTier::memory(),
283 },
284 persistent: Some(persistent),
285 read: Some(read),
286 retention: Default::default(),
287 merge_config: Default::default(),
288 event_bus,
289 spawner,
290 clock,
291 })
292 .unwrap();
293 (store, guard)
294 }
295
296 #[cfg(all(feature = "sqlite", not(target_arch = "wasm32")))]
297 pub fn testing_memory_with_persistent_sqlite_with_eventbus(event_bus: EventBus) -> (Self, SqliteTempPathGuard) {
298 let clock = Clock::testing();
299 let actor_system = ActorSystem::testing(clock.clone());
300 let spawner = actor_system.spawner();
301 let (persistent, guard) = PersistentConfig::sqlite_in_memory();
302 let store = Self::new(MultiStoreConfig {
303 commit: CommitBufferConfig {
304 storage: MultiCommitBufferTier::memory(),
305 },
306 persistent: Some(persistent),
307 read: Some(ReadBufferConfig::default()),
308 retention: Default::default(),
309 merge_config: Default::default(),
310 event_bus,
311 spawner,
312 clock,
313 })
314 .unwrap();
315 (store, guard)
316 }
317}