/rust/registry/src/index.crates.io-1949cf8c6b5b557f/surrealmx-0.22.0/src/persistence.rs
Line | Count | Source |
1 | | // Copyright © SurrealDB Ltd |
2 | | // |
3 | | // Licensed under the Apache License, Version 2.0 (the "License"); |
4 | | // you may not use this file except in compliance with the License. |
5 | | // You may obtain a copy of the License at |
6 | | // |
7 | | // http://www.apache.org/licenses/LICENSE-2.0 |
8 | | // |
9 | | // Unless required by applicable law or agreed to in writing, software |
10 | | // distributed under the License is distributed on an "AS IS" BASIS, |
11 | | // WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
12 | | // See the License for the specific language governing permissions and |
13 | | // limitations under the License. |
14 | | |
15 | | //! This module stores the database persistence logic. |
16 | | |
17 | | #![cfg(not(target_arch = "wasm32"))] |
18 | | |
19 | | use crate::compression::CompressedReader; |
20 | | use crate::compression::CompressedWriter; |
21 | | use crate::compression::CompressionMode; |
22 | | use crate::err::PersistenceError; |
23 | | use crate::inner::Inner; |
24 | | use crate::version::Version; |
25 | | use crate::versions::Versions; |
26 | | use bincode::config; |
27 | | use bytes::Bytes; |
28 | | use crossbeam_deque::{Injector, Steal}; |
29 | | use parking_lot::RwLock; |
30 | | use std::collections::BTreeMap; |
31 | | use std::fs::{self, File, OpenOptions}; |
32 | | use std::io::Write; |
33 | | use std::io::{BufReader, BufWriter, Seek, SeekFrom}; |
34 | | use std::path::PathBuf; |
35 | | use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; |
36 | | use std::sync::{Arc, Mutex}; |
37 | | use std::thread::{self, JoinHandle}; |
38 | | use web_time::{Duration, Instant}; |
39 | | |
40 | | /// Represents a pending asynchronous append operation |
41 | | #[derive(Debug, Clone)] |
42 | | pub(crate) struct AsyncAppendOperation { |
43 | | pub version: u64, |
44 | | pub writeset: BTreeMap<Bytes, Option<Bytes>>, |
45 | | } |
46 | | |
47 | | /// Configuration for AOL (Append-Only Log) behavior |
48 | | #[derive(Debug, Default, Clone, Copy, PartialEq, Eq)] |
49 | | pub enum AolMode { |
50 | | /// Never use AOL |
51 | | #[default] |
52 | | Never, |
53 | | /// Write immediatelyto AOL on every commit |
54 | | SynchronousOnCommit, |
55 | | /// Write asynchronously to AOL on every commit |
56 | | AsynchronousAfterCommit, |
57 | | } |
58 | | |
59 | | /// Configuration for snapshot behavior |
60 | | #[derive(Debug, Default, Clone, Copy, PartialEq, Eq)] |
61 | | pub enum SnapshotMode { |
62 | | /// Never use snapshots |
63 | | #[default] |
64 | | Never, |
65 | | /// Periodically snapshot at the given interval |
66 | | Interval(Duration), |
67 | | } |
68 | | |
69 | | /// Configuration for fsync behavior |
70 | | #[derive(Debug, Default, Clone, Copy, PartialEq, Eq)] |
71 | | pub enum FsyncMode { |
72 | | /// Never call fsync (fastest, least durable) |
73 | | #[default] |
74 | | Never, |
75 | | /// Call fsync after every append operation (slowest, most durable) |
76 | | EveryAppend, |
77 | | /// Call fsync at most once per interval |
78 | | Interval(Duration), |
79 | | } |
80 | | |
81 | | /// Configuration options for persistence |
82 | | #[derive(Debug, Clone)] |
83 | | pub struct PersistenceOptions { |
84 | | /// Base path for persistence files |
85 | | pub base_path: PathBuf, |
86 | | /// AOL (append-only log) behavior mode |
87 | | pub aol_mode: AolMode, |
88 | | /// Snapshot behavior mode |
89 | | pub snapshot_mode: SnapshotMode, |
90 | | /// Configuration for fsync behavior |
91 | | pub fsync_mode: FsyncMode, |
92 | | /// Path to the append-only log file (relative to base path or absolute) |
93 | | pub aol_path: Option<PathBuf>, |
94 | | /// Path to the snapshot file (relative to base path or absolute) |
95 | | pub snapshot_path: Option<PathBuf>, |
96 | | /// Compression mode for snapshots |
97 | | pub compression_mode: CompressionMode, |
98 | | } |
99 | | |
100 | | impl Default for PersistenceOptions { |
101 | 0 | fn default() -> Self { |
102 | 0 | Self { |
103 | 0 | base_path: PathBuf::from("./data"), |
104 | 0 | aol_mode: AolMode::default(), |
105 | 0 | snapshot_mode: SnapshotMode::default(), |
106 | 0 | fsync_mode: FsyncMode::default(), |
107 | 0 | aol_path: None, |
108 | 0 | snapshot_path: None, |
109 | 0 | compression_mode: CompressionMode::default(), |
110 | 0 | } |
111 | 0 | } |
112 | | } |
113 | | |
114 | | impl PersistenceOptions { |
115 | | /// Create new persistence options with the given base path |
116 | 0 | pub fn new<P: Into<PathBuf>>(base_path: P) -> Self { |
117 | 0 | Self { |
118 | 0 | base_path: base_path.into(), |
119 | 0 | ..Self::default() |
120 | 0 | } |
121 | 0 | } Unexecuted instantiation: <surrealmx::persistence::PersistenceOptions>::new::<_> Unexecuted instantiation: <surrealmx::persistence::PersistenceOptions>::new::<&alloc::string::String> Unexecuted instantiation: <surrealmx::persistence::PersistenceOptions>::new::<&alloc::string::String> |
122 | | |
123 | | /// Set the base path for persistence files |
124 | 0 | pub fn with_base_path<P: Into<PathBuf>>(mut self, path: P) -> Self { |
125 | 0 | self.base_path = path.into(); |
126 | 0 | self |
127 | 0 | } |
128 | | |
129 | | /// Set the AOL (append-only log) behavior mode |
130 | 0 | pub fn with_aol_mode(mut self, mode: AolMode) -> Self { |
131 | 0 | self.aol_mode = mode; |
132 | 0 | self |
133 | 0 | } |
134 | | |
135 | | /// Set the snapshot behavior mode |
136 | 0 | pub fn with_snapshot_mode(mut self, mode: SnapshotMode) -> Self { |
137 | 0 | self.snapshot_mode = mode; |
138 | 0 | self |
139 | 0 | } |
140 | | |
141 | | /// Set the fsync mode |
142 | 0 | pub fn with_fsync_mode(mut self, mode: FsyncMode) -> Self { |
143 | 0 | self.fsync_mode = mode; |
144 | 0 | self |
145 | 0 | } |
146 | | |
147 | | /// Set a custom AOL file path |
148 | 0 | pub fn with_aol_path<P: Into<PathBuf>>(mut self, path: P) -> Self { |
149 | 0 | self.aol_path = Some(path.into()); |
150 | 0 | self |
151 | 0 | } |
152 | | |
153 | | /// Set a custom snapshot file path |
154 | 0 | pub fn with_snapshot_path<P: Into<PathBuf>>(mut self, path: P) -> Self { |
155 | 0 | self.snapshot_path = Some(path.into()); |
156 | 0 | self |
157 | 0 | } |
158 | | |
159 | | /// Set the compression mode for snapshots |
160 | 0 | pub fn with_compression(mut self, mode: CompressionMode) -> Self { |
161 | 0 | self.compression_mode = mode; |
162 | 0 | self |
163 | 0 | } |
164 | | } |
165 | | |
166 | | /// A persistence layer for storing and loading database state |
167 | | /// |
168 | | /// This struct handles the persistence of database state through: |
169 | | /// - Append-only log (AOL) for recording changes |
170 | | /// - Periodic snapshots for efficient recovery |
171 | | /// - Background worker for automatic snapshot creation |
172 | | #[derive(Clone)] |
173 | | pub struct Persistence { |
174 | | /// Reference to the inner database state |
175 | | pub(crate) inner: Arc<Inner>, |
176 | | /// File handle for the append-only log (None if AOL is disabled) |
177 | | pub(crate) aol: Option<Arc<Mutex<File>>>, |
178 | | /// Path to the append-only log file (None if AOL is disabled) |
179 | | pub(crate) aol_path: PathBuf, |
180 | | /// Path to the snapshot file |
181 | | pub(crate) snapshot_path: PathBuf, |
182 | | /// AOL (append-only log) behavior mode |
183 | | pub(crate) aol_mode: AolMode, |
184 | | /// Snapshot behavior mode |
185 | | pub(crate) snapshot_mode: SnapshotMode, |
186 | | /// Fsync configuration mode |
187 | | pub(crate) fsync_mode: FsyncMode, |
188 | | /// Compression mode for snapshots |
189 | | pub(crate) compression_mode: CompressionMode, |
190 | | /// Specifies whether background worker threads are enabled |
191 | | pub(crate) background_threads_enabled: Arc<AtomicBool>, |
192 | | /// Handle to the background fsync worker thread (for interval mode) |
193 | | pub(crate) fsync_handle: Arc<RwLock<Option<JoinHandle<()>>>>, |
194 | | /// Handle to the background snapshot worker thread |
195 | | pub(crate) snapshot_handle: Arc<RwLock<Option<JoinHandle<()>>>>, |
196 | | /// Handle to the background async append worker thread |
197 | | pub(crate) appender_handle: Arc<RwLock<Option<JoinHandle<()>>>>, |
198 | | /// Last fsync timestamp for interval mode |
199 | | pub(crate) last_fsync: Arc<Mutex<Instant>>, |
200 | | /// Counter for AOL appends since last fsync |
201 | | pub(crate) pending_syncs: Arc<AtomicU64>, |
202 | | /// Queue for asynchronous append operations |
203 | | pub(crate) async_append_injector: Arc<Injector<AsyncAppendOperation>>, |
204 | | } |
205 | | |
206 | | impl Persistence { |
207 | | /// Creates a new persistence layer with custom options |
208 | | /// |
209 | | /// # Arguments |
210 | | /// * `options` - Configuration options for persistence |
211 | | /// * `inner` - Reference to the database state |
212 | | /// |
213 | | /// # Returns |
214 | | /// * `Result<Self, PersistenceError>` - The created persistence layer or an |
215 | | /// error |
216 | 0 | pub(crate) fn new_with_options( |
217 | 0 | options: PersistenceOptions, |
218 | 0 | inner: Arc<Inner>, |
219 | 0 | ) -> Result<Self, PersistenceError> { |
220 | | // Get the base path from options |
221 | 0 | let base_path = &options.base_path; |
222 | | // Ensure the directory exists |
223 | 0 | fs::create_dir_all(base_path)?; |
224 | | // Determine the specified AOL file path |
225 | 0 | let aol_path = if let Some(path) = options.aol_path { |
226 | 0 | if path.is_absolute() { |
227 | 0 | path |
228 | | } else { |
229 | 0 | base_path.join(path) |
230 | | } |
231 | | } else { |
232 | 0 | base_path.join("aol.bin") |
233 | | }; |
234 | | // Determine the specified snapshot file path |
235 | 0 | let snapshot_path = if let Some(path) = options.snapshot_path { |
236 | 0 | if path.is_absolute() { |
237 | 0 | path |
238 | | } else { |
239 | 0 | base_path.join(path) |
240 | | } |
241 | | } else { |
242 | 0 | base_path.join("snapshot.bin") |
243 | | }; |
244 | | // Initialize AOL components if enabled |
245 | 0 | let aol = if !matches!(options.aol_mode, AolMode::Never) { |
246 | | // Ensure parent directories exist for AOL path |
247 | 0 | if let Some(parent) = aol_path.parent() { |
248 | 0 | fs::create_dir_all(parent)?; |
249 | 0 | } |
250 | | // Open the AOL file with append mode |
251 | 0 | let file = OpenOptions::new().create(true).append(true).read(true).open(&aol_path)?; |
252 | 0 | Some(Arc::new(Mutex::new(file))) |
253 | | } else { |
254 | 0 | None |
255 | | }; |
256 | | // Ensure parent directories exist for snapshot path |
257 | 0 | if let Some(parent) = snapshot_path.parent() { |
258 | 0 | fs::create_dir_all(parent)?; |
259 | 0 | } |
260 | | // Create the persistence instance |
261 | 0 | let this = Self { |
262 | 0 | inner, |
263 | 0 | aol, |
264 | 0 | aol_path, |
265 | 0 | snapshot_path, |
266 | 0 | aol_mode: options.aol_mode, |
267 | 0 | snapshot_mode: options.snapshot_mode, |
268 | 0 | fsync_mode: options.fsync_mode, |
269 | 0 | compression_mode: options.compression_mode, |
270 | 0 | background_threads_enabled: Arc::new(AtomicBool::new(true)), |
271 | 0 | fsync_handle: Arc::new(RwLock::new(None)), |
272 | 0 | snapshot_handle: Arc::new(RwLock::new(None)), |
273 | 0 | appender_handle: Arc::new(RwLock::new(None)), |
274 | 0 | last_fsync: Arc::new(Mutex::new(Instant::now())), |
275 | 0 | pending_syncs: Arc::new(AtomicU64::new(0)), |
276 | 0 | async_append_injector: Arc::new(Injector::new()), |
277 | 0 | }; |
278 | | // Load existing data from disk |
279 | 0 | this.load()?; |
280 | | // Start the background snapshot worker if snapshots are enabled |
281 | 0 | this.spawn_snapshot_worker(); |
282 | | // Start the fsync worker if needed (only when AOL is enabled) |
283 | 0 | this.spawn_fsync_worker(); |
284 | | // Start the async append worker if asynchronous mode is enabled |
285 | 0 | this.spawn_appender_worker(); |
286 | | // Return the persistence layer |
287 | 0 | Ok(this) |
288 | 0 | } |
289 | | |
290 | | /// Creates a new snapshot of the current database state |
291 | | /// |
292 | | /// This function: |
293 | | /// 1. Captures the current AOL file position as a cutoff point |
294 | | /// 2. Creates a new snapshot file atomically using a temporary file |
295 | | /// 3. Streams data to reduce memory usage |
296 | | /// 4. Truncates AOL only up to the cutoff position, preserving newer |
297 | | /// entries |
298 | | /// |
299 | | /// # Returns |
300 | | /// * `Result<(), PersistenceError>` - Success or an error |
301 | 0 | pub fn snapshot(&self) -> Result<(), PersistenceError> { |
302 | | // Create temporary file for atomic swap |
303 | 0 | let temp_path = self.snapshot_path.with_extension("tmp"); |
304 | | // Execute snapshot operation in closure for clean error handling |
305 | 0 | let result = (|| -> Result<(), PersistenceError> { |
306 | | // Create temporary file |
307 | 0 | let file = File::create(&temp_path)?; |
308 | | // Create compressed writer (handles buffering internally) |
309 | 0 | let mut writer = CompressedWriter::new(file, self.compression_mode)?; |
310 | | // Get the current position in the AOL file (if AOL is enabled) |
311 | 0 | let aol_cutoff_position = if let Some(ref aol) = self.aol { |
312 | 0 | aol.lock()?.metadata()?.len() |
313 | | } else { |
314 | 0 | 0 |
315 | | }; |
316 | | // Stream write each key-value pair to reduce memory usage |
317 | 0 | for entry in self.inner.datastore.iter() { |
318 | | // Get all versions for this key |
319 | 0 | let versions = entry.value().read().all_versions(); |
320 | | // Ensure that there are version entries |
321 | 0 | if !versions.is_empty() { |
322 | | // Serialize and write this single entry |
323 | 0 | bincode::serde::encode_into_std_write( |
324 | 0 | &(entry.key().clone(), versions), |
325 | 0 | &mut writer, |
326 | 0 | config::standard(), |
327 | 0 | )?; |
328 | 0 | } |
329 | | } |
330 | | // Flush the compressed writer |
331 | 0 | writer.flush()?; |
332 | | // Finish compression (finalizes LZ4 stream) |
333 | 0 | writer.finish()?; |
334 | | // Atomically rename temporary file to actual snapshot |
335 | 0 | fs::rename(&temp_path, &self.snapshot_path)?; |
336 | | // Sync the renamed file to disk for durability |
337 | | { |
338 | 0 | let final_file = File::open(&self.snapshot_path)?; |
339 | 0 | final_file.sync_all()?; |
340 | | } |
341 | | // Truncate AOL only up to the cutoff position |
342 | 0 | Self::truncate(&self.aol, aol_cutoff_position, &self.pending_syncs)?; |
343 | | // All ok |
344 | 0 | Ok(()) |
345 | | })(); |
346 | | // Clean up temporary file if operation failed |
347 | 0 | if result.is_err() { |
348 | 0 | // Ignore removal errors |
349 | 0 | let _ = fs::remove_file(&temp_path); |
350 | 0 | } |
351 | | // Return the operation result |
352 | 0 | result |
353 | 0 | } |
354 | | |
355 | | /// Loads the database state from disk |
356 | | /// |
357 | | /// This function: |
358 | | /// 1. Loads the latest snapshot if it exists |
359 | | /// 2. Applies any changes from the append-only log |
360 | 0 | fn load(&self) -> Result<(), PersistenceError> { |
361 | | // Check if snapshot file exists |
362 | 0 | if self.snapshot_path.exists() { |
363 | | // Read and deserialize the snapshot data |
364 | 0 | let file = File::open(&self.snapshot_path)?; |
365 | | // Get the metadata of the snapshot file |
366 | 0 | let metadata = file.metadata()?; |
367 | | // Check if the snapshot file is empty |
368 | 0 | if metadata.len() > 0 { |
369 | | // Create compressed reader that auto-detects compression mode |
370 | 0 | let mut reader = CompressedReader::new(file)?; |
371 | | // Initialize counters for tracking loaded entries |
372 | 0 | let mut count = 0; |
373 | | // Stream reading the snapshot to reduce memory usage |
374 | | loop { |
375 | | // Increment the counter |
376 | 0 | count += 1; |
377 | | // Trace the loading of the snapshot entry |
378 | 0 | tracing::trace!("Loading snapshot entry: {count}"); |
379 | | // Type alias for the entry |
380 | | type Entry = (Bytes, Vec<(u64, Option<Bytes>)>); |
381 | | // Attempt to decode the next entry, handling EOF gracefully |
382 | 0 | let result: Result<Entry, _> = |
383 | 0 | bincode::serde::decode_from_std_read(&mut reader, config::standard()); |
384 | | // Detech any end of file errors |
385 | 0 | match result { |
386 | 0 | Ok((k, versions)) => { |
387 | | // Ensure that there are version entries |
388 | 0 | if !versions.is_empty() { |
389 | | // Create a new versions entry |
390 | 0 | let mut entries = Versions::new(); |
391 | | // Add all of the version entries |
392 | 0 | for (version, value) in versions.into_iter() { |
393 | 0 | entries.push(Version { |
394 | 0 | version, |
395 | 0 | value, |
396 | 0 | }); |
397 | 0 | } |
398 | | // Insert the entry into the datastore |
399 | 0 | self.inner.datastore.insert(k, RwLock::new(entries)); |
400 | 0 | } |
401 | | } |
402 | 0 | Err(e) => match e { |
403 | | // Handle bincode decode errors that indicate EOF |
404 | | bincode::error::DecodeError::Io { |
405 | 0 | inner, |
406 | | .. |
407 | 0 | } if inner.kind() == std::io::ErrorKind::UnexpectedEof => { |
408 | 0 | break; |
409 | | } |
410 | 0 | e => return Err(PersistenceError::Deserialization(e)), |
411 | | }, |
412 | | } |
413 | | } |
414 | 0 | } |
415 | 0 | } |
416 | | // Check if append-only file exists |
417 | 0 | if self.aol_path.exists() { |
418 | | // Open and read the AOL file |
419 | 0 | let file = File::open(&self.aol_path)?; |
420 | | // Get the metadata of the append-only file |
421 | 0 | let metadata = file.metadata()?; |
422 | | // Check if the append-only file is empty |
423 | 0 | if metadata.len() > 0 { |
424 | | // Create buffered reader for efficient reading |
425 | 0 | let mut reader = BufReader::new(file); |
426 | | // Initialize counters for tracking loaded entries |
427 | 0 | let mut count = 0; |
428 | | // Read and apply each change from the AOL |
429 | | loop { |
430 | | // Increment the counter |
431 | 0 | count += 1; |
432 | | // Trace the loading of the append-only entry |
433 | 0 | tracing::trace!("Loading AOL entry: {count}"); |
434 | | // Type alias for the entry |
435 | | type Entry = (Bytes, u64, Option<Bytes>); |
436 | | // Explicitly type the result to help type inference |
437 | 0 | let result: Result<Entry, _> = |
438 | 0 | bincode::serde::decode_from_std_read(&mut reader, config::standard()); |
439 | | // Detech any end of file errors |
440 | 0 | match result { |
441 | 0 | Ok((k, version, val)) => { |
442 | | // Check if the key already exists |
443 | 0 | if let Some(entry) = self.inner.datastore.get(&k) { |
444 | 0 | // Update existing key with stored version |
445 | 0 | entry.value().write().push(Version { |
446 | 0 | version, |
447 | 0 | value: val, |
448 | 0 | }); |
449 | 0 | } else { |
450 | 0 | // Insert new key with stored version |
451 | 0 | self.inner.datastore.insert( |
452 | 0 | k.clone(), |
453 | 0 | RwLock::new(Versions::from(Version { |
454 | 0 | version, |
455 | 0 | value: val, |
456 | 0 | })), |
457 | 0 | ); |
458 | 0 | } |
459 | | } |
460 | 0 | Err(e) => match e { |
461 | | // Handle bincode decode errors that indicate EOF |
462 | | bincode::error::DecodeError::Io { |
463 | 0 | inner, |
464 | | .. |
465 | 0 | } if inner.kind() == std::io::ErrorKind::UnexpectedEof => { |
466 | 0 | break; |
467 | | } |
468 | 0 | e => return Err(PersistenceError::Deserialization(e)), |
469 | | }, |
470 | | } |
471 | | } |
472 | 0 | } |
473 | 0 | } |
474 | | // Return success |
475 | 0 | Ok(()) |
476 | 0 | } |
477 | | |
478 | | /// Truncate the AOL file up to the specified position, preserving any data |
479 | | /// after. |
480 | 0 | fn truncate( |
481 | 0 | aol: &Option<Arc<Mutex<File>>>, |
482 | 0 | position: u64, |
483 | 0 | pending_syncs: &Arc<AtomicU64>, |
484 | 0 | ) -> Result<(), PersistenceError> { |
485 | | // Check that we have a AOL file handle |
486 | 0 | if let Some(ref aol) = aol { |
487 | | // Lock the AOL file |
488 | 0 | let mut file = aol.lock()?; |
489 | | // Get the current file length |
490 | 0 | let file_len = file.metadata()?.len(); |
491 | | // Check if there is remaining data |
492 | 0 | if file_len > position { |
493 | | // Generate a unique name for the temporary file |
494 | 0 | let name = format!("aol_truncate_{}.tmp", std::process::id()); |
495 | | // Generate the path for the temporary file |
496 | 0 | let path = std::env::temp_dir().join(name); |
497 | | // Execute truncation in a closure for clean error handling |
498 | 0 | let result = (|| -> Result<(), PersistenceError> { |
499 | | // Create temporary file and copy remaining data |
500 | | { |
501 | 0 | file.seek(SeekFrom::Start(position))?; |
502 | | // Create the temporary file |
503 | 0 | let mut temp = File::create(&path)?; |
504 | | // Copy the remaining data to the temporary file |
505 | 0 | std::io::copy(&mut *file, &mut temp)?; |
506 | | // Sync the temporary file |
507 | 0 | temp.sync_all()?; |
508 | | } |
509 | | // Go to the beginning of the file |
510 | 0 | file.seek(SeekFrom::Start(0))?; |
511 | | // Truncate the AOL file |
512 | 0 | file.set_len(0)?; |
513 | | // Copy data from temporary file |
514 | | { |
515 | 0 | let mut temp = File::open(&path)?; |
516 | 0 | std::io::copy(&mut temp, &mut *file)?; |
517 | | } |
518 | | // Flush the file contents |
519 | 0 | file.flush()?; |
520 | | // All ok |
521 | 0 | Ok(()) |
522 | | })(); |
523 | | // Delete the temporary file |
524 | 0 | let _ = fs::remove_file(&path); |
525 | | // Return the result |
526 | 0 | result?; |
527 | | } else { |
528 | | // Truncate the AOL file |
529 | 0 | file.set_len(0)?; |
530 | | // Flush the file contents |
531 | 0 | file.flush()?; |
532 | | } |
533 | | // Reset pending syncs if we truncated to beginning |
534 | 0 | if position == 0 { |
535 | 0 | pending_syncs.store(0, Ordering::Release); |
536 | 0 | } |
537 | 0 | } |
538 | | // All ok |
539 | 0 | Ok(()) |
540 | 0 | } |
541 | | |
542 | | /// Spawns a background worker thread for periodic fsync |
543 | 0 | fn spawn_fsync_worker(&self) { |
544 | | // Check if AOL is enabled |
545 | 0 | if self.aol_mode == AolMode::Never { |
546 | 0 | return; |
547 | 0 | } |
548 | | // Get the specified fsync interval |
549 | 0 | let FsyncMode::Interval(interval) = self.fsync_mode else { |
550 | 0 | return; |
551 | | }; |
552 | | // Check if AOL is enabled |
553 | 0 | if let Some(ref aol) = self.aol { |
554 | | // Check if a background thread is already running |
555 | 0 | if self.fsync_handle.read().is_none() { |
556 | | // Clone necessary fields for the worker thread |
557 | 0 | let aol = aol.clone(); |
558 | 0 | let pending_syncs = self.pending_syncs.clone(); |
559 | 0 | let enabled = self.background_threads_enabled.clone(); |
560 | | // Spawn the background worker thread |
561 | 0 | let handle = thread::spawn(move || { |
562 | | // Check whether the persistence process is enabled |
563 | 0 | while enabled.load(Ordering::Acquire) { |
564 | | // Sleep for the configured interval |
565 | 0 | thread::park_timeout(interval); |
566 | | // Check shutdown flag again after waking |
567 | 0 | if !enabled.load(Ordering::Acquire) { |
568 | 0 | break; |
569 | 0 | } |
570 | | // Check if there are pending syncs |
571 | 0 | if pending_syncs.load(Ordering::Acquire) > 0 { |
572 | 0 | if let Ok(file) = aol.lock() { |
573 | 0 | if let Err(e) = file.sync_all() { |
574 | 0 | tracing::error!("Fsync worker error: {e}"); |
575 | 0 | } else { |
576 | 0 | pending_syncs.store(0, Ordering::Release); |
577 | 0 | } |
578 | 0 | } |
579 | 0 | } |
580 | | } |
581 | 0 | }); |
582 | | // Store and track the thread handle |
583 | 0 | *self.fsync_handle.write() = Some(handle); |
584 | 0 | } |
585 | 0 | } |
586 | 0 | } |
587 | | |
588 | | /// Spawns a background worker thread for periodic snapshots |
589 | | /// |
590 | | /// The worker thread: |
591 | | /// 1. Sleeps for the configured interval |
592 | | /// 2. Captures the current AOL file position |
593 | | /// 3. Creates a new snapshot |
594 | | /// 4. Truncates AOL up to the cutoff, preserving newer entries |
595 | 0 | fn spawn_snapshot_worker(&self) { |
596 | | // Check if snapshots are enabled |
597 | 0 | if self.snapshot_mode == SnapshotMode::Never { |
598 | 0 | return; |
599 | 0 | } |
600 | | // Only spawn if snapshot interval is configured |
601 | 0 | let SnapshotMode::Interval(interval) = self.snapshot_mode else { |
602 | 0 | return; |
603 | | }; |
604 | | // Check if a background thread is already running |
605 | 0 | if self.snapshot_handle.read().is_none() { |
606 | | // Clone necessary fields for the worker thread |
607 | 0 | let db = self.inner.clone(); |
608 | 0 | let aol = self.aol.clone(); |
609 | 0 | let snapshot_path = self.snapshot_path.clone(); |
610 | 0 | let pending_syncs = self.pending_syncs.clone(); |
611 | 0 | let enabled = self.background_threads_enabled.clone(); |
612 | 0 | let compression = self.compression_mode; |
613 | | // Spawn the background worker thread |
614 | 0 | let handle = thread::spawn(move || { |
615 | | // Check whether the persistence process is enabled |
616 | 0 | while enabled.load(Ordering::Acquire) { |
617 | | // Sleep for the configured interval |
618 | 0 | thread::park_timeout(interval); |
619 | | // Check shutdown flag again after waking |
620 | 0 | if !enabled.load(Ordering::Acquire) { |
621 | 0 | break; |
622 | 0 | } |
623 | | // Create temporary file for atomic swap |
624 | 0 | let temp_path = snapshot_path.with_extension("tmp"); |
625 | | // Ensure clean error handling in closure |
626 | 0 | let result = (|| -> Result<(), PersistenceError> { |
627 | | // Create temporary file |
628 | 0 | let file = File::create(&temp_path)?; |
629 | | // Create compressed writer (handles buffering internally) |
630 | 0 | let mut writer = CompressedWriter::new(file, compression)?; |
631 | | // Get the current position in the AOL file before snapshotting (if AOL |
632 | | // enabled) |
633 | 0 | let aol_cutoff_position = if let Some(ref aol) = aol { |
634 | 0 | aol.lock()?.metadata()?.len() |
635 | | } else { |
636 | 0 | 0 |
637 | | }; |
638 | | // Stream write each entry to reduce memory usage |
639 | 0 | for entry in db.datastore.iter() { |
640 | | // Get all versions for this key |
641 | 0 | let versions = entry.value().read().all_versions(); |
642 | | // Ensure that there are version entries |
643 | 0 | if !versions.is_empty() { |
644 | | // Serialize and write this single entry |
645 | 0 | bincode::serde::encode_into_std_write( |
646 | 0 | &(entry.key().clone(), versions), |
647 | 0 | &mut writer, |
648 | 0 | config::standard(), |
649 | 0 | )?; |
650 | 0 | } |
651 | | } |
652 | | // Flush the compressed writer |
653 | 0 | writer.flush()?; |
654 | | // Finish compression (finalizes LZ4 stream) |
655 | 0 | writer.finish()?; |
656 | | // Atomically rename temporary file |
657 | 0 | fs::rename(&temp_path, &snapshot_path)?; |
658 | | // Sync the renamed file to disk for durability |
659 | | { |
660 | 0 | let final_file = File::open(&snapshot_path)?; |
661 | 0 | final_file.sync_all()?; |
662 | | } |
663 | | // Truncate AOL to the cutoff position |
664 | 0 | Self::truncate(&aol, aol_cutoff_position, &pending_syncs)?; |
665 | | // All ok |
666 | 0 | Ok(()) |
667 | | })(); |
668 | | // Check if the snapshot operation failed |
669 | 0 | if let Err(e) = result { |
670 | | // Trace the snapshot worker error |
671 | 0 | tracing::error!("Snapshot worker error: {e}"); |
672 | | // Clean up temporary file if it exists |
673 | 0 | let _ = fs::remove_file(&temp_path); |
674 | 0 | } |
675 | | } |
676 | 0 | }); |
677 | | // Store the worker thread handle |
678 | 0 | *self.snapshot_handle.write() = Some(handle); |
679 | 0 | } |
680 | 0 | } |
681 | | |
682 | | /// Spawn the background worker thread for processing async append |
683 | | /// operations |
684 | 0 | fn spawn_appender_worker(&self) { |
685 | | // Check if asynchronous append mode is enabled |
686 | 0 | if self.aol_mode != AolMode::AsynchronousAfterCommit { |
687 | 0 | return; |
688 | 0 | } |
689 | | // Check if AOL is enabled |
690 | 0 | if let Some(ref aol) = self.aol { |
691 | | // Check if a background thread is already running |
692 | 0 | if self.appender_handle.read().is_none() { |
693 | | // Clone necessary fields for the worker thread |
694 | 0 | let injector = self.async_append_injector.clone(); |
695 | 0 | let aol = aol.clone(); |
696 | 0 | let fsync_mode = self.fsync_mode; |
697 | 0 | let enabled = self.background_threads_enabled.clone(); |
698 | 0 | let pending_syncs = self.pending_syncs.clone(); |
699 | 0 | let last_fsync = self.last_fsync.clone(); |
700 | | // Spawn the background worker thread |
701 | 0 | let handle = thread::spawn(move || { |
702 | | // Set the batch size and timeout |
703 | | const BATCH_SIZE: usize = 100; |
704 | | const TIMEOUT_MS: u64 = 10; |
705 | | // Initialize the batch vector |
706 | 0 | let mut batch = Vec::with_capacity(BATCH_SIZE); |
707 | | // Check whether the persistence process is enabled |
708 | 0 | while enabled.load(Ordering::Acquire) { |
709 | | // Check shutdown flag again after waking |
710 | 0 | if !enabled.load(Ordering::Acquire) { |
711 | 0 | break; |
712 | 0 | } |
713 | | // Clear the batch |
714 | 0 | batch.clear(); |
715 | | // Collect operations into a batch |
716 | | loop { |
717 | | // Check shutdown flag in the inner loop |
718 | 0 | if !enabled.load(Ordering::Acquire) { |
719 | 0 | break; |
720 | 0 | } |
721 | 0 | match injector.steal() { |
722 | | Steal::Retry => { |
723 | 0 | std::thread::yield_now(); |
724 | 0 | continue; |
725 | | } |
726 | 0 | Steal::Success(op) => { |
727 | 0 | batch.push(op); |
728 | 0 | if batch.len() == BATCH_SIZE { |
729 | 0 | break; |
730 | 0 | } |
731 | | } |
732 | | Steal::Empty => { |
733 | | // If we have items to append, break |
734 | 0 | if !batch.is_empty() { |
735 | 0 | break; |
736 | 0 | } |
737 | | // Park the thread to wait for work |
738 | 0 | thread::park_timeout(Duration::from_millis(TIMEOUT_MS)); |
739 | | } |
740 | | } |
741 | | } |
742 | | // Process the batch if we have operations |
743 | 0 | if !batch.is_empty() { |
744 | | // Ensure clean error handling in closure |
745 | 0 | let result = (|| -> Result<(), PersistenceError> { |
746 | | // Lock the AOL file for writing |
747 | 0 | if let Ok(mut file) = aol.lock() { |
748 | | // Create a new buffer for the AOL file |
749 | 0 | let mut writer = BufWriter::new(&mut *file); |
750 | | // Write all operations in the batch |
751 | 0 | for op in &batch { |
752 | 0 | for (k, v) in &op.writeset { |
753 | 0 | bincode::serde::encode_into_std_write( |
754 | 0 | (k, op.version, v), |
755 | 0 | &mut writer, |
756 | 0 | config::standard(), |
757 | 0 | )?; |
758 | | } |
759 | | } |
760 | | // Flush the buffer to the file on the operating system |
761 | 0 | writer.flush()?; |
762 | | // Drop the writer to release the mutable borrow |
763 | 0 | drop(writer); |
764 | | // Handle fsync based on mode |
765 | 0 | match fsync_mode { |
766 | | // Let the operating system handle syncing to disk |
767 | 0 | FsyncMode::Never => { |
768 | 0 | // No fsync, just increment pending counter |
769 | 0 | pending_syncs.fetch_add(1, Ordering::Release); |
770 | 0 | } |
771 | | // Sync immediately to diskafter every append |
772 | | FsyncMode::EveryAppend => { |
773 | | // Sync immediately |
774 | 0 | file.sync_all()?; |
775 | | } |
776 | | // Force sync to disk at a specified interval |
777 | 0 | FsyncMode::Interval(duration) => { |
778 | | // Check if we should sync based on time |
779 | 0 | let now = Instant::now(); |
780 | | // Check if we should sync based on time |
781 | 0 | let should_sync = { |
782 | | // Get the last fsync time |
783 | 0 | let mut last_fsync = last_fsync.lock()?; |
784 | | // Check if the last fsync time is greater than the |
785 | | // duration |
786 | 0 | if now.duration_since(*last_fsync) >= duration { |
787 | | // Update the last fsync time |
788 | 0 | *last_fsync = now; |
789 | 0 | true |
790 | | } else { |
791 | 0 | false |
792 | | } |
793 | | }; |
794 | | // Check if we should sync |
795 | 0 | if should_sync { |
796 | | // Force sync the AOL file to disk |
797 | 0 | file.sync_all()?; |
798 | | // Reset the pending syncs counter |
799 | 0 | pending_syncs.store(0, Ordering::Release); |
800 | 0 | } else { |
801 | 0 | // Increment the pending syncs counter |
802 | 0 | pending_syncs.fetch_add(1, Ordering::Release); |
803 | 0 | } |
804 | | } |
805 | | } |
806 | 0 | } |
807 | | // All ok |
808 | 0 | Ok(()) |
809 | | })(); |
810 | | // Check if the async append operation failed |
811 | 0 | if let Err(e) = result { |
812 | | // Trace the snapshot worker error |
813 | 0 | tracing::error!("Async append worker error: {e}"); |
814 | 0 | } |
815 | 0 | } |
816 | | } |
817 | 0 | }); |
818 | | // Store the thread handle |
819 | 0 | *self.appender_handle.write() = Some(handle); |
820 | 0 | } |
821 | 0 | } |
822 | 0 | } |
823 | | |
824 | | /// Appends a set of changes to the append-only log |
825 | | /// |
826 | | /// # Arguments |
827 | | /// * `version` - The version (timestamp) for these changes |
828 | | /// * `writeset` - Map of key-value changes to append |
829 | | /// |
830 | | /// # Returns |
831 | | /// * `Result<(), PersistenceError>` - Success or an error |
832 | 0 | pub(crate) fn append( |
833 | 0 | &self, |
834 | 0 | version: u64, |
835 | 0 | writeset: &BTreeMap<Bytes, Option<Bytes>>, |
836 | 0 | ) -> Result<(), PersistenceError> { |
837 | | // Skip AOL writing if AOL is disabled |
838 | 0 | if self.aol_mode == AolMode::Never { |
839 | 0 | return Ok(()); |
840 | 0 | } |
841 | | // AOL is enabled, proceed with append logic |
842 | 0 | if let Some(ref aol) = self.aol { |
843 | | // Handle asynchronous AOL mode by queuing the operation |
844 | 0 | if self.aol_mode == AolMode::AsynchronousAfterCommit { |
845 | | // Queue the append operation |
846 | 0 | self.async_append_injector.push(AsyncAppendOperation { |
847 | 0 | version, |
848 | 0 | writeset: writeset.clone(), |
849 | 0 | }); |
850 | | // Wake up the async append worker if available |
851 | 0 | if let Some(handle) = self.appender_handle.read().as_ref() { |
852 | 0 | handle.thread().unpark(); |
853 | 0 | } |
854 | 0 | } |
855 | 0 | if self.aol_mode == AolMode::SynchronousOnCommit { |
856 | | // Lock the AOL file for writing |
857 | 0 | let mut file = aol.lock()?; |
858 | | // Create a new buffer for the AOL file |
859 | 0 | let mut writer = BufWriter::new(&mut *file); |
860 | | // Serialize and write each change with version |
861 | 0 | for (k, v) in writeset { |
862 | 0 | bincode::serde::encode_into_std_write( |
863 | 0 | (k, version, v), |
864 | 0 | &mut writer, |
865 | 0 | config::standard(), |
866 | 0 | )?; |
867 | | } |
868 | | // Flush the buffer to the file on the operating system |
869 | 0 | writer.flush()?; |
870 | | // Drop the writer to release the mutable borrow |
871 | 0 | drop(writer); |
872 | | // Handle fsync based on mode |
873 | 0 | match self.fsync_mode { |
874 | | // Let the operating system handle syncing to disk |
875 | 0 | FsyncMode::Never => { |
876 | 0 | // No fsync, just increment pending counter |
877 | 0 | self.pending_syncs.fetch_add(1, Ordering::Release); |
878 | 0 | } |
879 | | // Sync immediately to diskafter every append |
880 | | FsyncMode::EveryAppend => { |
881 | | // Sync immediately |
882 | 0 | file.sync_all()?; |
883 | | } |
884 | | // Force sync to disk at a specified interval |
885 | 0 | FsyncMode::Interval(duration) => { |
886 | | // Check if we should sync based on time |
887 | 0 | let now = Instant::now(); |
888 | | // Check if we should sync based on time |
889 | 0 | let should_sync = { |
890 | | // Get the last fsync time |
891 | 0 | let mut last_fsync = self.last_fsync.lock()?; |
892 | | // Check if the last fsync time is greater than the duration |
893 | 0 | if now.duration_since(*last_fsync) >= duration { |
894 | | // Update the last fsync time |
895 | 0 | *last_fsync = now; |
896 | 0 | true |
897 | | } else { |
898 | 0 | false |
899 | | } |
900 | | }; |
901 | | // Check if we should sync |
902 | 0 | if should_sync { |
903 | | // Force sync the AOL file to disk |
904 | 0 | file.sync_all()?; |
905 | | // Reset the pending syncs counter |
906 | 0 | self.pending_syncs.store(0, Ordering::Release); |
907 | 0 | } else { |
908 | 0 | // Increment the pending syncs counter |
909 | 0 | self.pending_syncs.fetch_add(1, Ordering::Release); |
910 | 0 | } |
911 | | } |
912 | | } |
913 | 0 | } |
914 | 0 | } |
915 | | // All ok |
916 | 0 | Ok(()) |
917 | 0 | } |
918 | | } |
919 | | |
920 | | impl Drop for Persistence { |
921 | | /// Cleans up resources when the persistence layer is dropped |
922 | 0 | fn drop(&mut self) { |
923 | | // Signal shutdown to the worker threads |
924 | 0 | self.background_threads_enabled.store(false, Ordering::Release); |
925 | | // Stop the fsync worker if it exists |
926 | 0 | if let Some(handle) = self.fsync_handle.write().take() { |
927 | 0 | handle.thread().unpark(); |
928 | 0 | let _ = handle.join(); |
929 | 0 | } |
930 | | // Stop the snapshot worker if it exists |
931 | 0 | if let Some(handle) = self.snapshot_handle.write().take() { |
932 | 0 | handle.thread().unpark(); |
933 | 0 | let _ = handle.join(); |
934 | 0 | } |
935 | | // Stop the async append worker if it exists |
936 | 0 | if let Some(handle) = self.appender_handle.write().take() { |
937 | 0 | handle.thread().unpark(); |
938 | 0 | let _ = handle.join(); |
939 | 0 | } |
940 | | // Perform final fsync if there are pending syncs |
941 | 0 | if self.aol_mode != AolMode::Never && self.pending_syncs.load(Ordering::Acquire) > 0 { |
942 | | // Try to acquire lock on AOL file |
943 | 0 | if let Some(ref aol) = self.aol { |
944 | | // Lock the AOL file |
945 | 0 | if let Ok(file) = aol.lock() { |
946 | 0 | // Sync file contents to disk |
947 | 0 | let _ = file.sync_all(); |
948 | 0 | } |
949 | 0 | } |
950 | 0 | } |
951 | 0 | } |
952 | | } |