1use std::{
5 ops::Deref,
6 sync::{
7 Arc, OnceLock,
8 atomic::{AtomicU64, Ordering},
9 },
10};
11
12use reifydb_codec::key::encoded::EncodedKey;
13#[cfg(all(feature = "sqlite", not(target_arch = "wasm32")))]
14use reifydb_core::metrics::sample::MetricsSample;
15use reifydb_core::{
16 common::CommitVersion,
17 event::EventBus,
18 interface::store::{EntryKind, storage_key},
19 lifecycle::watermark::EvictionWatermark,
20 metrics::collect::MetricsCollector,
21};
22#[cfg(all(feature = "sqlite", not(target_arch = "wasm32")))]
23use reifydb_filter::{actor::FilterActor, config::FilterConfig};
24use reifydb_filter::{actor::FilterMessage, adaptive::FilterMetrics};
25use reifydb_runtime::{
26 actor::{mailbox::ActorRef, system::ActorSystem},
27 context::clock::Clock,
28 shutdown::Shutdown,
29 sync::rwlock::RwLock,
30};
31#[cfg(all(feature = "sqlite", not(target_arch = "wasm32")))]
32use reifydb_sqlite::SqliteTempPathGuard;
33use reifydb_store::metrics::PageCacheMetrics;
34use reifydb_store_commit::store::{CommitStore, MultiCommitMetrics};
35use reifydb_value::{count::Count, util::cowvec::CowVec};
36use tracing::instrument;
37
38#[cfg(all(feature = "sqlite", not(target_arch = "wasm32")))]
39use crate::filter::source::MultiCurrentKeySource;
40use crate::{
41 CommitStoreConfig,
42 config::MultiStoreConfig,
43 flush::ObjectPersistence,
44 tier::{
45 persistent::MultiPersistentTier,
46 point::{MultiPointShardMetrics, MultiPointTier},
47 range::{MultiRangeShardMetrics, MultiRangeTier},
48 },
49};
50#[cfg(all(feature = "sqlite", not(target_arch = "wasm32")))]
51use crate::{
52 config::PersistentConfig,
53 flush::engine::FlushEngine,
54 tier::{point::MultiPointConfig, range::MultiRangeConfig},
55};
56
57pub mod multi;
58pub mod router;
59
60use crate::Result;
61
62#[cfg(all(feature = "sqlite", not(target_arch = "wasm32")))]
63struct SqlitePageCacheCollector {
64 persistent: MultiPersistentTier,
65}
66
67#[cfg(all(feature = "sqlite", not(target_arch = "wasm32")))]
68impl MetricsCollector for SqlitePageCacheCollector {
69 fn collect(&self, out: &mut Vec<MetricsSample>) {
70 let metrics = self.persistent.page_cache_metrics();
71 out.push(MetricsSample::bytes("sqlite::multi", "page_cache_used_bytes", metrics.used));
72 out.push(MetricsSample::counter("sqlite::multi", "page_cache_hit_count", metrics.hits.as_u64()));
73 out.push(MetricsSample::counter("sqlite::multi", "page_cache_miss_count", metrics.misses.as_u64()));
74 out.push(MetricsSample::count(
75 "sqlite::multi",
76 "page_cache_sampled_connections",
77 metrics.connections_sampled.as_u64(),
78 ));
79 }
80}
81
82#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
83pub struct MultiPersistentProbeMetrics {
84 pub persistent_probes: Count,
85 pub persistent_absent: Count,
86}
87
88#[derive(Clone)]
89pub struct StandardMultiStore(Arc<StandardMultiStoreInner>);
90
91pub struct StandardMultiStoreInner {
92 pub(crate) commit: CommitStore,
93 pub(crate) persistent: Option<MultiPersistentTier>,
94 pub(crate) point: Option<MultiPointTier>,
95 pub(crate) range: Option<MultiRangeTier>,
96
97 #[cfg(all(feature = "sqlite", not(target_arch = "wasm32")))]
98 pub(crate) flush_engine: Option<Arc<FlushEngine>>,
99 #[allow(dead_code)]
100 pub(crate) row_settings_provider: Arc<OnceLock<Arc<dyn ObjectPersistence>>>,
101 #[allow(dead_code)]
102 pub(crate) eviction_watermark: Arc<RwLock<Option<Arc<dyn EvictionWatermark>>>>,
103
104 pub(crate) event_bus: EventBus,
105
106 pub(crate) persistent_probes: AtomicU64,
107 pub(crate) persistent_absent: AtomicU64,
108
109 pub(crate) filter: Option<ActorRef<FilterMessage>>,
110}
111
112impl StandardMultiStore {
113 #[instrument(name = "store::multi::new", level = "debug", skip(config), fields(
114 has_persistent = config.persistent.is_some(),
115 ))]
116 pub fn new(config: MultiStoreConfig) -> Result<Self> {
117 let commit = config.commit.storage;
118
119 let row_settings_provider: Arc<OnceLock<Arc<dyn ObjectPersistence>>> = Arc::new(OnceLock::new());
120
121 let eviction_watermark: Arc<RwLock<Option<Arc<dyn EvictionWatermark>>>> = Arc::new(RwLock::new(None));
122
123 let point = config.persistent.is_some().then(|| config.point.and_then(MultiPointTier::new)).flatten();
124
125 let range = config.persistent.is_some().then(|| config.range.and_then(MultiRangeTier::new)).flatten();
126
127 #[cfg(all(feature = "sqlite", not(target_arch = "wasm32")))]
128 let (persistent, flush_engine, filter) = {
129 let persistent_config = config.persistent.clone();
130 let persistent = persistent_config.as_ref().map(|c| c.storage.clone());
131 let filter = persistent.as_ref().map(|tier| {
132 let storage = tier.sqlite_storage().clone();
133 let actor = FilterActor::spawn(&config.spawner);
134 actor.send(FilterMessage::Register {
135 filter: storage.filter().handle(),
136 source: Box::new(MultiCurrentKeySource::new(storage)),
137 config: FilterConfig::default(),
138 })
139 .expect("multi current filter source could not be registered");
140 actor
141 });
142 let flush_engine = match (persistent.as_ref(), persistent_config.as_ref()) {
143 (Some(persistent_storage), Some(_)) => Some(Arc::new(
144 FlushEngine::new(
145 commit.clone(),
146 persistent_storage.clone(),
147 row_settings_provider.clone(),
148 eviction_watermark.clone(),
149 config.clock.clone(),
150 config.event_bus.clone(),
151 )
152 .with_point(point.clone())
153 .with_range(range.clone()),
154 )),
155 _ => None,
156 };
157 (persistent, flush_engine, filter)
158 };
159
160 #[cfg(not(all(feature = "sqlite", not(target_arch = "wasm32"))))]
161 let (persistent, filter): (Option<MultiPersistentTier>, Option<ActorRef<FilterMessage>>) = {
162 let _ = config.persistent;
163 (None, None)
164 };
165
166 let point = persistent.as_ref().and(point);
167 let range = persistent.as_ref().and(range);
168
169 Ok(Self(Arc::new(StandardMultiStoreInner {
170 commit,
171 persistent,
172 point,
173 range,
174 #[cfg(all(feature = "sqlite", not(target_arch = "wasm32")))]
175 flush_engine,
176 row_settings_provider,
177 eviction_watermark,
178 event_bus: config.event_bus,
179 persistent_probes: AtomicU64::new(0),
180 persistent_absent: AtomicU64::new(0),
181 filter,
182 })))
183 }
184
185 #[cfg(all(feature = "sqlite", not(target_arch = "wasm32")))]
186 pub fn flush_engine(&self) -> Option<Arc<FlushEngine>> {
187 self.flush_engine.clone()
188 }
189
190 pub fn insert_read_key(
191 &self,
192 table: EntryKind,
193 key: EncodedKey,
194 version: CommitVersion,
195 value: Option<CowVec<u8>>,
196 ) {
197 if let Some(range) = &self.range {
198 range.insert(table, key.clone(), version, value.clone());
199 }
200 if let Some(point) = &self.point {
201 let storage_key = storage_key(&key).1;
202 point.insert(table, storage_key, key, version, value);
203 }
204 }
205
206 pub fn invalidate_read_key(&self, table: EntryKind, key: &EncodedKey) {
207 if let Some(range) = &self.range {
208 range.invalidate(table, key);
209 }
210 if let Some(point) = &self.point {
211 point.invalidate(table, storage_key(key).1, key);
212 }
213 }
214
215 pub fn clear_read(&self) {
216 if let Some(range) = &self.range {
217 range.clear();
218 }
219 if let Some(point) = &self.point {
220 point.clear();
221 }
222 }
223
224 pub fn set_row_settings_provider(&self, provider: Arc<dyn ObjectPersistence>) {
225 let _ = self.row_settings_provider.set(provider);
226 }
227
228 pub fn set_eviction_watermark(&self, watermark: Arc<dyn EvictionWatermark>) {
229 *self.eviction_watermark.write() = Some(watermark);
230 }
231
232 pub fn clear_eviction_watermark(&self) {
233 *self.eviction_watermark.write() = None;
234 }
235
236 pub fn commit(&self) -> &CommitStore {
237 &self.commit
238 }
239
240 pub fn event_bus(&self) -> &EventBus {
241 &self.event_bus
242 }
243
244 pub fn metrics_collectors(&self) -> Vec<Arc<dyn MetricsCollector>> {
245 let mut collectors: Vec<Arc<dyn MetricsCollector>> = Vec::new();
246 collectors.push(Arc::new(self.commit.clone()));
247 if let Some(point) = &self.point {
248 collectors.push(Arc::new(point.clone()));
249 }
250 if let Some(range) = &self.range {
251 collectors.push(Arc::new(range.clone()));
252 }
253 #[cfg(all(feature = "sqlite", not(target_arch = "wasm32")))]
254 if let Some(persistent) = &self.persistent {
255 collectors.push(Arc::new(SqlitePageCacheCollector {
256 persistent: persistent.clone(),
257 }));
258 }
259 collectors
260 }
261
262 pub fn persistent(&self) -> Option<&MultiPersistentTier> {
263 self.persistent.as_ref()
264 }
265
266 pub fn point_shard_metrics(&self) -> Vec<MultiPointShardMetrics> {
267 self.point.as_ref().map(|point| point.shard_metrics()).unwrap_or_default()
268 }
269
270 pub fn range_shard_metrics(&self) -> Vec<MultiRangeShardMetrics> {
271 self.range.as_ref().map(|range| range.full_shard_metrics()).unwrap_or_default()
272 }
273
274 pub fn commit_metrics(&self) -> MultiCommitMetrics {
275 self.commit.metrics()
276 }
277
278 pub fn persistent_page_cache_metrics(&self) -> Option<PageCacheMetrics> {
279 self.persistent.as_ref().map(MultiPersistentTier::page_cache_metrics)
280 }
281
282 pub fn persistent_filter_metrics(&self) -> Option<FilterMetrics> {
283 self.persistent.as_ref().map(|tier| tier.filter().metrics())
284 }
285
286 pub fn persistent_probe_metrics(&self) -> Option<MultiPersistentProbeMetrics> {
287 self.persistent.as_ref().map(|_| MultiPersistentProbeMetrics {
288 persistent_probes: Count::new(self.persistent_probes.load(Ordering::Relaxed)),
289 persistent_absent: Count::new(self.persistent_absent.load(Ordering::Relaxed)),
290 })
291 }
292
293 pub fn flush_pending_blocking(&self) {
294 #[cfg(all(feature = "sqlite", not(target_arch = "wasm32")))]
295 if let Some(engine) = self.flush_engine.as_ref() {
296 self.event_bus.wait_for_completion();
297 engine.flush_pending();
298 }
299 }
300
301 pub fn flush_all_blocking(&self) {
302 #[cfg(all(feature = "sqlite", not(target_arch = "wasm32")))]
303 if let Some(engine) = self.flush_engine.as_ref() {
304 self.event_bus.wait_for_completion();
305 engine.flush_all();
306 }
307 }
308}
309
310impl Deref for StandardMultiStore {
311 type Target = StandardMultiStoreInner;
312
313 fn deref(&self) -> &Self::Target {
314 &self.0
315 }
316}
317
318impl Shutdown for StandardMultiStore {
319 fn shutdown(&self) {
320 if let Some(filter) = self.filter.as_ref() {
321 let _ = filter.send(FilterMessage::Shutdown);
322 }
323 if let Some(persistent) = self.persistent.as_ref() {
324 persistent.shutdown();
325 }
326 }
327}
328
329impl StandardMultiStore {
330 pub fn testing_memory() -> Self {
331 let clock = Clock::testing();
332 let actor_system = ActorSystem::testing(clock.clone());
333 let spawner = actor_system.spawner();
334 let event_bus = EventBus::new(&spawner);
335 Self::new(MultiStoreConfig {
336 commit: CommitStoreConfig {
337 storage: CommitStore::new(),
338 },
339 persistent: None,
340 point: None,
341 range: None,
342 retention: Default::default(),
343 merge_config: Default::default(),
344 event_bus,
345 spawner,
346 clock,
347 })
348 .unwrap()
349 }
350
351 pub fn testing_memory_with_eventbus(event_bus: EventBus) -> Self {
352 let clock = Clock::testing();
353 let actor_system = ActorSystem::testing(clock.clone());
354 let spawner = actor_system.spawner();
355 Self::new(MultiStoreConfig {
356 commit: CommitStoreConfig {
357 storage: CommitStore::new(),
358 },
359 persistent: None,
360 point: None,
361 range: None,
362 retention: Default::default(),
363 merge_config: Default::default(),
364 event_bus,
365 spawner,
366 clock,
367 })
368 .unwrap()
369 }
370
371 #[cfg(all(feature = "sqlite", not(target_arch = "wasm32")))]
372 pub fn testing_memory_with_persistent_sqlite() -> (Self, SqliteTempPathGuard) {
373 Self::testing_memory_with_persistent_sqlite_tiers(
374 MultiPointConfig::testing(),
375 MultiRangeConfig::testing(),
376 )
377 }
378
379 #[cfg(all(feature = "sqlite", not(target_arch = "wasm32")))]
380 pub fn testing_memory_with_persistent_sqlite_tiers(
381 point: MultiPointConfig,
382 range: MultiRangeConfig,
383 ) -> (Self, SqliteTempPathGuard) {
384 let clock = Clock::testing();
385 let actor_system = ActorSystem::testing(clock.clone());
386 let spawner = actor_system.spawner();
387 let event_bus = EventBus::new(&spawner);
388 let (persistent, guard) = PersistentConfig::sqlite_in_memory();
389 let store = Self::new(MultiStoreConfig {
390 commit: CommitStoreConfig {
391 storage: CommitStore::new(),
392 },
393 persistent: Some(persistent),
394 point: Some(point),
395 range: Some(range),
396 retention: Default::default(),
397 merge_config: Default::default(),
398 event_bus,
399 spawner,
400 clock,
401 })
402 .unwrap();
403 (store, guard)
404 }
405
406 #[cfg(all(feature = "sqlite", not(target_arch = "wasm32")))]
407 pub fn testing_memory_with_persistent_sqlite_with_eventbus(event_bus: EventBus) -> (Self, SqliteTempPathGuard) {
408 let clock = Clock::testing();
409 let actor_system = ActorSystem::testing(clock.clone());
410 let spawner = actor_system.spawner();
411 let (persistent, guard) = PersistentConfig::sqlite_in_memory();
412 let store = Self::new(MultiStoreConfig {
413 commit: CommitStoreConfig {
414 storage: CommitStore::new(),
415 },
416 persistent: Some(persistent),
417 point: Some(MultiPointConfig::testing()),
418 range: Some(MultiRangeConfig::testing()),
419 retention: Default::default(),
420 merge_config: Default::default(),
421 event_bus,
422 spawner,
423 clock,
424 })
425 .unwrap();
426 (store, guard)
427 }
428}