reifydb_store_multi/store/
mod.rs1use 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}