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