Coverage Report

Created: 2026-08-14 08:14

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
/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
}