Skip to main content

reifydb_store_multi/store/
mod.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use 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}