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	mem,
6	ops::Deref,
7	sync::{Arc, OnceLock},
8};
9
10use reifydb_codec::key::encoded::EncodedKey;
11use reifydb_core::{common::CommitVersion, event::EventBus};
12use reifydb_runtime::{
13	actor::{mailbox::ActorRef, system::ActorSystem},
14	context::clock::Clock,
15	pool::{PoolConfig, Pools},
16	shutdown::Shutdown,
17	sync::{rwlock::RwLock, waiter::WaiterHandle},
18};
19#[cfg(all(feature = "sqlite", not(target_arch = "wasm32")))]
20use reifydb_sqlite::SqliteTempPathGuard;
21use reifydb_value::{util::cowvec::CowVec, value::duration::Duration};
22use tracing::instrument;
23
24use crate::{
25	CommitBufferConfig,
26	config::MultiStoreConfig,
27	flush::{ShapePersistence, actor::FlushMessage},
28	gc::EvictionWatermark,
29	tier::{
30		commit::buffer::MultiCommitBufferTier,
31		persistent::MultiPersistentTier,
32		read::{MultiReadBufferTier, ReadBufferConfig},
33	},
34};
35#[cfg(all(feature = "sqlite", not(target_arch = "wasm32")))]
36use crate::{config::PersistentConfig, flush::actor::FlushActor};
37
38pub mod drop;
39pub mod multi;
40pub mod pending;
41pub mod router;
42pub mod worker;
43
44use pending::PendingDrops;
45use reifydb_core::actors::drop::DropMessage;
46use worker::{DropActor, DropWorkerConfig};
47
48use crate::Result;
49
50#[derive(Clone)]
51pub struct StandardMultiStore(Arc<StandardMultiStoreInner>);
52
53pub struct StandardMultiStoreInner {
54	pub(crate) commit: Option<MultiCommitBufferTier>,
55	pub(crate) persistent: Option<MultiPersistentTier>,
56	pub(crate) read: Option<MultiReadBufferTier>,
57	pub(crate) pending_drops: PendingDrops,
58	pub(crate) drop_actor: Option<ActorRef<DropMessage>>,
59
60	#[allow(dead_code)]
61	pub(crate) flush_actor: Option<ActorRef<FlushMessage>>,
62	#[allow(dead_code)]
63	pub(crate) row_settings_provider: Arc<OnceLock<Arc<dyn ShapePersistence>>>,
64	#[allow(dead_code)]
65	pub(crate) eviction_watermark: Arc<RwLock<Option<Arc<dyn EvictionWatermark>>>>,
66
67	pub(crate) event_bus: EventBus,
68}
69
70impl StandardMultiStore {
71	#[instrument(name = "store::multi::new", level = "debug", skip(config), fields(
72		has_commit = config.commit.is_some(),
73		has_persistent = config.persistent.is_some(),
74	))]
75	pub fn new(config: MultiStoreConfig) -> Result<Self> {
76		let commit = config.commit.map(|c| c.storage);
77
78		let spawner = config.spawner.clone();
79
80		let row_settings_provider: Arc<OnceLock<Arc<dyn ShapePersistence>>> = Arc::new(OnceLock::new());
81
82		let eviction_watermark: Arc<RwLock<Option<Arc<dyn EvictionWatermark>>>> = Arc::new(RwLock::new(None));
83
84		let read = match (commit.as_ref(), config.persistent.is_some()) {
85			(Some(_), true) => Some(MultiReadBufferTier::new(ReadBufferConfig::default())),
86			_ => None,
87		};
88
89		#[cfg(all(feature = "sqlite", not(target_arch = "wasm32")))]
90		let (persistent, flush_actor) = {
91			let persistent_config = config.persistent.clone();
92			let persistent = persistent_config.as_ref().map(|c| c.storage.clone());
93			let flush_actor = match (commit.as_ref(), persistent.as_ref(), persistent_config.as_ref()) {
94				(Some(buf), Some(persistent_storage), Some(persistent_cfg)) => {
95					let actor_ref = FlushActor::spawn(
96						&spawner,
97						buf.clone(),
98						persistent_storage.clone(),
99						persistent_cfg.flush_interval,
100						row_settings_provider.clone(),
101						eviction_watermark.clone(),
102						read.clone(),
103					);
104					Some(actor_ref)
105				}
106				_ => None,
107			};
108			(persistent, flush_actor)
109		};
110
111		#[cfg(not(all(feature = "sqlite", not(target_arch = "wasm32"))))]
112		let (persistent, flush_actor): (Option<MultiPersistentTier>, Option<ActorRef<FlushMessage>>) = {
113			let _ = config.persistent;
114			(None, None)
115		};
116
117		let read = persistent.as_ref().and(read);
118
119		let pending_drops = PendingDrops::default();
120
121		let drop_actor = commit.as_ref().map(|storage| {
122			let drop_config = DropWorkerConfig::default();
123			DropActor::spawn(
124				&spawner,
125				drop_config,
126				storage.clone(),
127				config.event_bus.clone(),
128				config.clock,
129				persistent.clone(),
130				read.clone(),
131				pending_drops.clone(),
132			)
133		});
134
135		Ok(Self(Arc::new(StandardMultiStoreInner {
136			commit,
137			persistent,
138			read,
139			pending_drops,
140			drop_actor,
141			flush_actor,
142			row_settings_provider,
143			eviction_watermark,
144			event_bus: config.event_bus,
145		})))
146	}
147
148	pub fn configure_read_buffer_capacity(&self, capacity: usize) {
149		if let Some(read) = &self.read {
150			read.set_capacity(capacity);
151		}
152	}
153
154	pub fn configure_read_buffer(&self, resident_pages: usize, page_size_rows: u64) {
155		if let Some(read) = &self.read {
156			read.reconfigure(resident_pages, page_size_rows);
157		}
158	}
159
160	pub fn configure_flush_interval(&self, interval: Duration) {
161		if let Some(actor) = &self.flush_actor {
162			let _ = actor.send(FlushMessage::SetInterval(interval));
163		}
164	}
165
166	pub fn configure_wal_autocheckpoint(&self, frames: u32) {
167		if let Some(persistent) = &self.persistent {
168			persistent.set_checkpoint_threshold(frames);
169		}
170	}
171
172	pub fn insert_read_key(&self, key: EncodedKey, version: CommitVersion, value: Option<CowVec<u8>>) {
173		if let Some(read) = &self.read {
174			read.insert(key, version, value);
175		}
176	}
177
178	pub fn invalidate_read_key(&self, key: &EncodedKey) {
179		if let Some(read) = &self.read {
180			read.invalidate(key);
181		}
182	}
183
184	pub fn purge_pending_drops(&self) {
185		self.pending_drops.purge(self.persistent.as_ref(), self.read.as_ref());
186	}
187
188	pub fn remove_dropped_read_key(&self, key: &EncodedKey) {
189		if let Some(read) = &self.read {
190			read.remove_dropped(key);
191		}
192	}
193
194	pub fn clear_read(&self) {
195		if let Some(read) = &self.read {
196			read.clear();
197		}
198	}
199
200	pub fn set_row_settings_provider(&self, provider: Arc<dyn ShapePersistence>) {
201		let _ = self.row_settings_provider.set(provider);
202	}
203
204	pub fn set_eviction_watermark(&self, watermark: Arc<dyn EvictionWatermark>) {
205		*self.eviction_watermark.write() = Some(watermark);
206	}
207
208	pub fn clear_eviction_watermark(&self) {
209		*self.eviction_watermark.write() = None;
210	}
211
212	pub fn commit(&self) -> Option<&MultiCommitBufferTier> {
213		self.commit.as_ref()
214	}
215
216	pub fn persistent(&self) -> Option<&MultiPersistentTier> {
217		self.persistent.as_ref()
218	}
219
220	pub fn flush_pending_blocking(&self) {
221		let Some(actor_ref) = self.flush_actor.as_ref() else {
222			return;
223		};
224
225		self.event_bus.wait_for_completion();
226
227		let waiter = Arc::new(WaiterHandle::new());
228		let waiter_for_msg = Arc::clone(&waiter);
229		if actor_ref
230			.send_blocking(FlushMessage::FlushPending {
231				waiter: waiter_for_msg,
232			})
233			.is_err()
234		{
235			return;
236		}
237
238		waiter.wait_timeout(Duration::from_seconds(60).unwrap());
239	}
240
241	pub fn flush_all_blocking(&self) {
242		let Some(actor_ref) = self.flush_actor.as_ref() else {
243			return;
244		};
245
246		self.event_bus.wait_for_completion();
247
248		let waiter = Arc::new(WaiterHandle::new());
249		let waiter_for_msg = Arc::clone(&waiter);
250		if actor_ref
251			.send_blocking(FlushMessage::FlushAll {
252				waiter: waiter_for_msg,
253			})
254			.is_err()
255		{
256			return;
257		}
258
259		waiter.wait_timeout(Duration::from_seconds(60).unwrap());
260	}
261}
262
263impl Deref for StandardMultiStore {
264	type Target = StandardMultiStoreInner;
265
266	fn deref(&self) -> &Self::Target {
267		&self.0
268	}
269}
270
271impl Shutdown for StandardMultiStore {
272	fn shutdown(&self) {
273		if let Some(persistent) = self.persistent.as_ref() {
274			persistent.shutdown();
275		}
276	}
277}
278
279impl StandardMultiStore {
280	pub fn testing_memory() -> Self {
281		let pools = Pools::new(PoolConfig::sync_only());
282		let clock = Clock::testing();
283		let actor_system = ActorSystem::new(pools, clock.clone());
284		let spawner = actor_system.spawner();
285		let event_bus = EventBus::new(&spawner);
286		mem::forget(actor_system);
287		Self::new(MultiStoreConfig {
288			commit: Some(CommitBufferConfig {
289				storage: MultiCommitBufferTier::memory(),
290			}),
291			persistent: None,
292			retention: Default::default(),
293			merge_config: Default::default(),
294			event_bus,
295			spawner,
296			clock,
297		})
298		.unwrap()
299	}
300
301	pub fn testing_memory_with_eventbus(event_bus: EventBus) -> Self {
302		let pools = Pools::new(PoolConfig::sync_only());
303		let clock = Clock::testing();
304		let actor_system = ActorSystem::new(pools, clock.clone());
305		let spawner = actor_system.spawner();
306		mem::forget(actor_system);
307		Self::new(MultiStoreConfig {
308			commit: Some(CommitBufferConfig {
309				storage: MultiCommitBufferTier::memory(),
310			}),
311			persistent: None,
312			retention: Default::default(),
313			merge_config: Default::default(),
314			event_bus,
315			spawner,
316			clock,
317		})
318		.unwrap()
319	}
320
321	#[cfg(all(feature = "sqlite", not(target_arch = "wasm32")))]
322	pub fn testing_memory_with_persistent_sqlite() -> (Self, SqliteTempPathGuard) {
323		let pools = Pools::new(PoolConfig::default());
324		let clock = Clock::testing();
325		let actor_system = ActorSystem::new(pools, clock.clone());
326		let spawner = actor_system.spawner();
327		let event_bus = EventBus::new(&spawner);
328		mem::forget(actor_system);
329		let (persistent, guard) = PersistentConfig::sqlite_in_memory();
330		let store = Self::new(MultiStoreConfig {
331			commit: Some(CommitBufferConfig {
332				storage: MultiCommitBufferTier::memory(),
333			}),
334			persistent: Some(persistent),
335			retention: Default::default(),
336			merge_config: Default::default(),
337			event_bus,
338			spawner,
339			clock,
340		})
341		.unwrap();
342		(store, guard)
343	}
344
345	#[cfg(all(feature = "sqlite", not(target_arch = "wasm32")))]
346	pub fn testing_memory_with_persistent_sqlite_with_eventbus(event_bus: EventBus) -> (Self, SqliteTempPathGuard) {
347		let pools = Pools::new(PoolConfig::default());
348		let clock = Clock::testing();
349		let actor_system = ActorSystem::new(pools, clock.clone());
350		let spawner = actor_system.spawner();
351		mem::forget(actor_system);
352		let (persistent, guard) = PersistentConfig::sqlite_in_memory();
353		let store = Self::new(MultiStoreConfig {
354			commit: Some(CommitBufferConfig {
355				storage: MultiCommitBufferTier::memory(),
356			}),
357			persistent: Some(persistent),
358			retention: Default::default(),
359			merge_config: Default::default(),
360			event_bus,
361			spawner,
362			clock,
363		})
364		.unwrap();
365		(store, guard)
366	}
367
368	#[cfg(all(feature = "sqlite", not(target_arch = "wasm32")))]
369	pub fn testing_persistent_sqlite_only() -> (Self, SqliteTempPathGuard) {
370		let pools = Pools::new(PoolConfig::default());
371		let clock = Clock::testing();
372		let actor_system = ActorSystem::new(pools, clock.clone());
373		let spawner = actor_system.spawner();
374		let event_bus = EventBus::new(&spawner);
375		mem::forget(actor_system);
376		let (persistent, guard) = PersistentConfig::sqlite_in_memory();
377		let store =
378			Self::new(MultiStoreConfig::sqlite_unbuffered(persistent, spawner, clock, event_bus)).unwrap();
379		(store, guard)
380	}
381}