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::{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}