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::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}