Skip to main content

reifydb_store_multi/store/
mod.rs

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