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