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