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