Skip to main content

allsource_core/infrastructure/persistence/
wal.rs

1use crate::{
2    domain::entities::Event,
3    error::{AllSourceError, Result},
4};
5use chrono::{DateTime, Utc};
6use parking_lot::RwLock;
7use serde::{Deserialize, Serialize};
8use std::{
9    collections::HashSet,
10    fs::{self, File, OpenOptions},
11    io::{BufRead, BufReader, BufWriter, Read, Seek, SeekFrom, Write},
12    path::{Path, PathBuf},
13    sync::Arc,
14    time::SystemTime,
15};
16
17/// One WAL segment's identity as the filesystem reports it, so a reader can
18/// tell whether anything was appended, rotated or truncated without opening it.
19#[derive(Debug, Clone, PartialEq, Eq)]
20pub struct WalSegmentStamp {
21    pub path: PathBuf,
22    pub len: u64,
23    pub modified: Option<SystemTime>,
24}
25
26/// Subdirectory of the WAL dir that holds retired segments containing lines
27/// recovery could not read. Nothing lists, replays, or deletes it.
28pub const QUARANTINE_DIR: &str = "quarantine";
29
30/// Write-Ahead Log for durability and crash recovery.
31///
32/// Invariant: a segment is deleted only when every line in it is durable
33/// elsewhere. A segment holding a line recovery could not parse or verify is
34/// moved to [`QUARANTINE_DIR`] instead, so a format this binary cannot read is
35/// never destroyed by it (#287).
36pub struct WriteAheadLog {
37    /// Directory where WAL files are stored
38    wal_dir: PathBuf,
39
40    /// Current active WAL file
41    current_file: Arc<RwLock<WALFile>>,
42
43    /// Configuration
44    config: WALConfig,
45
46    /// Statistics
47    stats: Arc<RwLock<WALStats>>,
48
49    /// Current sequence number
50    sequence: Arc<RwLock<u64>>,
51
52    /// Optional broadcast channel for replication โ€” when set, each WAL append
53    /// publishes the entry so the WAL shipper can stream it to followers.
54    /// Wrapped in Mutex so it can be set at runtime (e.g. during follower โ†’ leader promotion).
55    replication_tx: parking_lot::Mutex<Option<tokio::sync::broadcast::Sender<WALEntry>>>,
56
57    /// Segments in which `recover()` met a line it could not read.
58    unreadable_segments: parking_lot::Mutex<HashSet<PathBuf>>,
59}
60
61#[derive(Debug, Clone)]
62pub struct WALConfig {
63    /// Maximum size of a single WAL file before rotation (in bytes)
64    pub max_file_size: usize,
65
66    /// Whether to sync to disk after each write (fsync)
67    pub sync_on_write: bool,
68
69    /// Maximum number of WAL files to keep
70    pub max_wal_files: usize,
71
72    /// Enable WAL compression
73    pub compress: bool,
74
75    /// Interval in milliseconds for background coalesced fsync.
76    /// When set, a background task calls `flush() + sync_all()` every N ms,
77    /// giving near-zero write latency with bounded data loss window.
78    /// When `Some`, `sync_on_write` is forced to `false` to prevent double-fsync.
79    pub fsync_interval_ms: Option<u64>,
80}
81
82impl Default for WALConfig {
83    fn default() -> Self {
84        Self {
85            max_file_size: 64 * 1024 * 1024, // 64 MB
86            sync_on_write: true,
87            max_wal_files: 10,
88            compress: false,
89            fsync_interval_ms: None,
90        }
91    }
92}
93
94#[derive(Debug, Clone, Default, Serialize)]
95pub struct WALStats {
96    pub total_entries: u64,
97    pub total_bytes_written: u64,
98    pub current_file_size: usize,
99    pub files_rotated: u64,
100    pub files_cleaned: u64,
101    pub recovery_count: u64,
102}
103
104/// WAL entry wrapping an event
105#[derive(Debug, Clone, Serialize, Deserialize)]
106pub struct WALEntry {
107    /// Sequence number for ordering
108    pub sequence: u64,
109
110    /// Timestamp when written to WAL
111    pub wal_timestamp: DateTime<Utc>,
112
113    /// The event being logged
114    pub event: Event,
115
116    /// Checksum for integrity verification
117    pub checksum: u32,
118}
119
120impl WALEntry {
121    pub fn new(sequence: u64, event: Event) -> Self {
122        let mut entry = Self {
123            sequence,
124            wal_timestamp: Utc::now(),
125            event,
126            checksum: 0,
127        };
128        entry.checksum = entry.calculate_checksum();
129        entry
130    }
131
132    fn calculate_checksum(&self) -> u32 {
133        // Simple CRC32 checksum
134        let data = format!("{}{}{}", self.sequence, self.wal_timestamp, self.event.id);
135        crc32fast::hash(data.as_bytes())
136    }
137
138    pub fn verify(&self) -> bool {
139        self.checksum == self.calculate_checksum()
140    }
141}
142
143fn ends_with_newline(file: &mut File) -> Result<bool> {
144    let io_err = |e: std::io::Error| {
145        AllSourceError::StorageError(format!("Failed to read WAL file tail: {e}"))
146    };
147    file.seek(SeekFrom::End(-1)).map_err(io_err)?;
148    let mut last = [0u8; 1];
149    file.read_exact(&mut last).map_err(io_err)?;
150    Ok(last[0] == b'\n')
151}
152
153/// Represents an active WAL file
154struct WALFile {
155    path: PathBuf,
156    writer: BufWriter<File>,
157    size: usize,
158    created_at: DateTime<Utc>,
159}
160
161impl WALFile {
162    fn new(path: PathBuf) -> Result<Self> {
163        let mut file = OpenOptions::new()
164            .create(true)
165            .read(true)
166            .append(true)
167            .open(&path)
168            .map_err(|e| AllSourceError::StorageError(format!("Failed to open WAL file: {e}")))?;
169
170        let mut size = file.metadata().map_or(0, |m| m.len() as usize);
171
172        // A crash can leave a torn final line with no newline. Appending
173        // straight after it would fuse the next entry onto the torn bytes and
174        // make that entry unreadable too.
175        if size > 0 && !ends_with_newline(&mut file)? {
176            file.write_all(b"\n").map_err(|e| {
177                AllSourceError::StorageError(format!("Failed to terminate torn WAL line: {e}"))
178            })?;
179            size += 1;
180        }
181
182        Ok(Self {
183            path,
184            writer: BufWriter::new(file),
185            size,
186            created_at: Utc::now(),
187        })
188    }
189
190    fn write_entry(&mut self, entry: &WALEntry, sync: bool) -> Result<usize> {
191        // Serialize entry as JSON line
192        let json = serde_json::to_string(entry)?;
193
194        let line = format!("{json}\n");
195        let bytes_written = line.len();
196
197        self.writer
198            .write_all(line.as_bytes())
199            .map_err(|e| AllSourceError::StorageError(format!("Failed to write to WAL: {e}")))?;
200
201        if sync {
202            self.writer
203                .flush()
204                .map_err(|e| AllSourceError::StorageError(format!("Failed to flush WAL: {e}")))?;
205
206            self.writer
207                .get_ref()
208                .sync_all()
209                .map_err(|e| AllSourceError::StorageError(format!("Failed to sync WAL: {e}")))?;
210        }
211
212        self.size += bytes_written;
213        Ok(bytes_written)
214    }
215
216    fn flush(&mut self) -> Result<()> {
217        self.writer
218            .flush()
219            .map_err(|e| AllSourceError::StorageError(format!("Failed to flush WAL: {e}")))?;
220        Ok(())
221    }
222}
223
224impl WriteAheadLog {
225    /// Create a new WAL
226    pub fn new(wal_dir: impl Into<PathBuf>, config: WALConfig) -> Result<Self> {
227        let wal_dir = wal_dir.into();
228
229        // Create WAL directory if it doesn't exist
230        fs::create_dir_all(&wal_dir).map_err(|e| {
231            AllSourceError::StorageError(format!("Failed to create WAL directory: {e}"))
232        })?;
233
234        // Create initial WAL file
235        let initial_file_path = Self::generate_wal_filename(&wal_dir, 0);
236        let current_file = WALFile::new(initial_file_path)?;
237
238        tracing::info!("โœ… WAL initialized at: {}", wal_dir.display());
239
240        Ok(Self {
241            wal_dir,
242            current_file: Arc::new(RwLock::new(current_file)),
243            config,
244            stats: Arc::new(RwLock::new(WALStats::default())),
245            sequence: Arc::new(RwLock::new(0)),
246            replication_tx: parking_lot::Mutex::new(None),
247            unreadable_segments: parking_lot::Mutex::new(HashSet::new()),
248        })
249    }
250
251    /// Generate a WAL filename based on sequence
252    fn generate_wal_filename(dir: &Path, sequence: u64) -> PathBuf {
253        dir.join(format!("wal-{sequence:016x}.log"))
254    }
255
256    /// Write an event to the WAL
257    #[cfg_attr(feature = "hotpath", hotpath::measure)]
258    pub fn append(&self, event: Event) -> Result<u64> {
259        // Get next sequence number
260        let mut seq = self.sequence.write();
261        *seq += 1;
262        let sequence = *seq;
263        drop(seq);
264
265        // Create WAL entry
266        let entry = WALEntry::new(sequence, event);
267
268        // Write to current file
269        let mut current = self.current_file.write();
270        let bytes_written = current.write_entry(&entry, self.config.sync_on_write)?;
271
272        // Update statistics
273        let mut stats = self.stats.write();
274        stats.total_entries += 1;
275        stats.total_bytes_written += bytes_written as u64;
276        stats.current_file_size = current.size;
277        drop(stats);
278
279        // Broadcast to replication channel (if enabled).
280        // Errors are ignored: no receivers simply means no followers are connected.
281        if let Some(ref tx) = *self.replication_tx.lock() {
282            let _ = tx.send(entry);
283        }
284
285        // Check if we need to rotate
286        let should_rotate = current.size >= self.config.max_file_size;
287        drop(current);
288
289        if should_rotate {
290            self.rotate()?;
291        }
292
293        tracing::trace!("WAL entry written: sequence={}", sequence);
294
295        Ok(sequence)
296    }
297
298    /// Rotate to a new WAL file
299    #[cfg_attr(feature = "hotpath", hotpath::measure)]
300    fn rotate(&self) -> Result<()> {
301        let seq = *self.sequence.read();
302        let new_file_path = Self::generate_wal_filename(&self.wal_dir, seq);
303
304        tracing::info!("๐Ÿ”„ Rotating WAL to new file: {:?}", new_file_path);
305
306        let new_file = WALFile::new(new_file_path)?;
307
308        let mut current = self.current_file.write();
309        current.flush()?;
310        *current = new_file;
311
312        let mut stats = self.stats.write();
313        stats.files_rotated += 1;
314        stats.current_file_size = 0;
315        drop(stats);
316
317        self.warn_if_segments_accumulate()?;
318
319        Ok(())
320    }
321
322    /// Rotation never deletes: a segment past `max_wal_files` may hold events
323    /// no checkpoint has flushed yet. Only a checkpoint retires segments.
324    fn warn_if_segments_accumulate(&self) -> Result<()> {
325        let count = self.list_wal_files()?.len();
326        if count > self.config.max_wal_files {
327            tracing::warn!(
328                "WAL holds {} segments (max_wal_files = {}); checkpoints are not keeping up, \
329                 so the WAL keeps growing rather than dropping unflushed events",
330                count,
331                self.config.max_wal_files
332            );
333        }
334        Ok(())
335    }
336
337    /// Delete a segment, or move it to quarantine if recovery could not read
338    /// every line in it.
339    fn retire_segment(&self, path: &Path) -> Result<()> {
340        if !self.unreadable_segments.lock().remove(path) {
341            fs::remove_file(path).map_err(|e| {
342                AllSourceError::StorageError(format!("Failed to remove WAL file: {e}"))
343            })?;
344            tracing::debug!("Removed WAL file: {:?}", path);
345            return Ok(());
346        }
347
348        let quarantine = self.wal_dir.join(QUARANTINE_DIR);
349        fs::create_dir_all(&quarantine).map_err(|e| {
350            AllSourceError::StorageError(format!("Failed to create WAL quarantine: {e}"))
351        })?;
352        let name = path.file_name().map_or_else(
353            || "wal-unknown.log".into(),
354            |n| n.to_string_lossy().into_owned(),
355        );
356        let target = quarantine.join(format!(
357            "{}-{name}",
358            Utc::now().format("%Y%m%dT%H%M%S%.6fZ")
359        ));
360        fs::rename(path, &target).map_err(|e| {
361            AllSourceError::StorageError(format!("Failed to quarantine WAL file: {e}"))
362        })?;
363        tracing::warn!(
364            "WAL segment {:?} held lines this binary could not read; kept at {:?} instead of \
365             deleting it",
366            path,
367            target
368        );
369        Ok(())
370    }
371
372    /// Start a new segment so every entry appended so far sits in a segment
373    /// older than the returned one. Pair with [`Self::remove_sealed`].
374    ///
375    /// The caller must stop appends to durable storage for the duration of
376    /// this call; otherwise an entry can land in the sealed segment after the
377    /// checkpoint's flush has already run.
378    pub fn seal(&self) -> Result<PathBuf> {
379        let seq = *self.sequence.read();
380        let new_path = Self::generate_wal_filename(&self.wal_dir, seq);
381
382        let mut current = self.current_file.write();
383        current.flush()?;
384        if current.path != new_path {
385            *current = WALFile::new(new_path.clone())?;
386            self.stats.write().current_file_size = current.size;
387        }
388        Ok(new_path)
389    }
390
391    /// Retire every segment older than `active`, the path [`Self::seal`]
392    /// returned. Segments at or after it are kept.
393    pub fn remove_sealed(&self, active: &Path) -> Result<()> {
394        let Some(active_name) = active.file_name() else {
395            return Ok(());
396        };
397        for path in self.list_wal_files()? {
398            if path.file_name().is_some_and(|name| name < active_name) {
399                self.retire_segment(&path)?;
400            }
401        }
402        Ok(())
403    }
404
405    /// List all WAL files in the directory
406    fn list_wal_files(&self) -> Result<Vec<PathBuf>> {
407        let entries = fs::read_dir(&self.wal_dir).map_err(|e| {
408            AllSourceError::StorageError(format!("Failed to read WAL directory: {e}"))
409        })?;
410
411        let mut wal_files = Vec::new();
412        for entry in entries {
413            let entry = entry.map_err(|e| {
414                AllSourceError::StorageError(format!("Failed to read directory entry: {e}"))
415            })?;
416
417            let path = entry.path();
418            if let Some(name) = path.file_name()
419                && name.to_string_lossy().starts_with("wal-")
420                && name.to_string_lossy().ends_with(".log")
421            {
422                wal_files.push(path);
423            }
424        }
425
426        Ok(wal_files)
427    }
428
429    /// Stat every WAL segment, sorted by path. Reads no segment contents.
430    pub fn segment_stamps(&self) -> Result<Vec<WalSegmentStamp>> {
431        let mut stamps = Vec::new();
432        for path in self.list_wal_files()? {
433            // A segment unlinked between the listing and the stat is simply gone.
434            let Ok(metadata) = fs::metadata(&path) else {
435                continue;
436            };
437            stamps.push(WalSegmentStamp {
438                path,
439                len: metadata.len(),
440                modified: metadata.modified().ok(),
441            });
442        }
443        stamps.sort_by(|a, b| a.path.cmp(&b.path));
444        Ok(stamps)
445    }
446
447    /// Recover events from WAL files
448    #[cfg_attr(feature = "hotpath", hotpath::measure)]
449    pub fn recover(&self) -> Result<Vec<Event>> {
450        tracing::info!("๐Ÿ”„ Starting WAL recovery...");
451
452        let mut wal_files = self.list_wal_files()?;
453        wal_files.sort();
454
455        let mut recovered_events = Vec::new();
456        let mut max_sequence = 0u64;
457        let mut corrupted_entries = 0;
458
459        for wal_file_path in &wal_files {
460            tracing::debug!("Reading WAL file: {:?}", wal_file_path);
461
462            let file = File::open(wal_file_path).map_err(|e| {
463                AllSourceError::StorageError(format!("Failed to open WAL file for recovery: {e}"))
464            })?;
465
466            let reader = BufReader::new(file);
467            let corrupted_before = corrupted_entries;
468
469            for (line_num, line) in reader.lines().enumerate() {
470                let line = match line {
471                    Ok(l) => l,
472                    Err(e) => {
473                        tracing::warn!(
474                            "I/O error reading WAL line at {:?}:{}: {}",
475                            wal_file_path,
476                            line_num + 1,
477                            e
478                        );
479                        corrupted_entries += 1;
480                        continue;
481                    }
482                };
483
484                if line.trim().is_empty() {
485                    continue;
486                }
487
488                match serde_json::from_str::<WALEntry>(&line) {
489                    Ok(entry) => {
490                        // Verify checksum
491                        if !entry.verify() {
492                            tracing::warn!(
493                                "Corrupted WAL entry at {:?}:{} (checksum mismatch)",
494                                wal_file_path,
495                                line_num + 1
496                            );
497                            corrupted_entries += 1;
498                            continue;
499                        }
500
501                        max_sequence = max_sequence.max(entry.sequence);
502                        recovered_events.push(entry.event);
503                    }
504                    Err(e) => {
505                        tracing::warn!(
506                            "Failed to parse WAL entry at {:?}:{}: {}",
507                            wal_file_path,
508                            line_num + 1,
509                            e
510                        );
511                        corrupted_entries += 1;
512                    }
513                }
514            }
515
516            if corrupted_entries > corrupted_before {
517                self.unreadable_segments
518                    .lock()
519                    .insert(wal_file_path.clone());
520            }
521        }
522
523        // Update sequence counter
524        let mut seq = self.sequence.write();
525        *seq = max_sequence;
526        drop(seq);
527
528        // Update stats
529        let mut stats = self.stats.write();
530        stats.recovery_count += 1;
531        drop(stats);
532
533        tracing::info!(
534            "โœ… WAL recovery complete: {} events recovered, {} corrupted entries",
535            recovered_events.len(),
536            corrupted_entries
537        );
538
539        Ok(recovered_events)
540    }
541
542    /// Manually flush the current WAL file
543    #[cfg_attr(feature = "hotpath", hotpath::measure)]
544    pub fn flush(&self) -> Result<()> {
545        let mut current = self.current_file.write();
546        current.flush()?;
547        Ok(())
548    }
549
550    /// Flush the BufWriter and fsync the current WAL file to disk.
551    ///
552    /// Called by the background interval-based fsync task. Acquires the write
553    /// lock, flushes buffered data, then issues `sync_all()` to ensure the
554    /// OS has persisted the data to durable storage.
555    #[cfg_attr(feature = "hotpath", hotpath::measure)]
556    pub fn sync(&self) -> Result<()> {
557        let mut current = self.current_file.write();
558        current
559            .writer
560            .flush()
561            .map_err(|e| AllSourceError::StorageError(format!("Failed to flush WAL: {e}")))?;
562        current
563            .writer
564            .get_ref()
565            .sync_all()
566            .map_err(|e| AllSourceError::StorageError(format!("Failed to sync WAL: {e}")))?;
567        Ok(())
568    }
569
570    /// Get the configured fsync interval (if any).
571    pub fn fsync_interval_ms(&self) -> Option<u64> {
572        self.config.fsync_interval_ms
573    }
574
575    /// Truncate WAL after successful checkpoint
576    #[cfg_attr(feature = "hotpath", hotpath::measure)]
577    pub fn truncate(&self) -> Result<()> {
578        tracing::info!("๐Ÿงน Truncating WAL after checkpoint");
579
580        // Close current file
581        let mut current = self.current_file.write();
582        current.flush()?;
583
584        let wal_files = self.list_wal_files()?;
585        for file_path in wal_files {
586            self.retire_segment(&file_path)?;
587        }
588
589        // Create new WAL file
590        let new_file_path = Self::generate_wal_filename(&self.wal_dir, 0);
591        *current = WALFile::new(new_file_path)?;
592
593        // Reset sequence
594        let mut seq = self.sequence.write();
595        *seq = 0;
596
597        tracing::info!("โœ… WAL truncated successfully");
598
599        Ok(())
600    }
601
602    /// Get WAL statistics
603    pub fn stats(&self) -> WALStats {
604        (*self.stats.read()).clone()
605    }
606
607    /// Sum the on-disk size of all WAL segment files and count them.
608    ///
609    /// Walks the WAL directory and `fs::metadata`s each `wal-*.log` segment.
610    /// Returns `(total_bytes, segment_count)`. Unlike `WALStats.total_bytes_written`
611    /// (a monotonic write counter that ignores rotation/cleanup), this reflects the
612    /// *current* on-disk footprint โ€” what we want to report as storage size and as
613    /// `allsource_wal_segments_total`.
614    ///
615    /// Cheap (one `statx` per segment; each checkpoint retires sealed segments,
616    /// so the count stays small), so it's safe to call on the checkpoint cadence. Errors reading the directory
617    /// surface as `Err`; a per-file metadata error skips that file rather than
618    /// aborting the whole tally.
619    pub fn on_disk_stats(&self) -> Result<(u64, usize)> {
620        let wal_files = self.list_wal_files()?;
621        let mut total_bytes = 0u64;
622        for path in &wal_files {
623            if let Ok(metadata) = fs::metadata(path) {
624                total_bytes += metadata.len();
625            }
626        }
627        Ok((total_bytes, wal_files.len()))
628    }
629
630    /// Get current sequence number
631    pub fn current_sequence(&self) -> u64 {
632        *self.sequence.read()
633    }
634
635    /// Get the oldest available WAL sequence number.
636    ///
637    /// Returns the first sequence found in the oldest WAL file, or `None`
638    /// if no WAL entries exist. Used by the replication catch-up protocol
639    /// to determine whether a follower can catch up from WAL alone.
640    pub fn oldest_sequence(&self) -> Option<u64> {
641        let Ok(mut wal_files) = self.list_wal_files() else {
642            return None;
643        };
644
645        if wal_files.is_empty() {
646            return None;
647        }
648
649        wal_files.sort();
650
651        // Read the first entry from the oldest WAL file
652        for wal_file_path in &wal_files {
653            let Ok(file) = File::open(wal_file_path) else {
654                continue;
655            };
656            let reader = BufReader::new(file);
657            for line in reader.lines() {
658                let Ok(line) = line else {
659                    continue;
660                };
661                if line.trim().is_empty() {
662                    continue;
663                }
664                if let Ok(entry) = serde_json::from_str::<WALEntry>(&line) {
665                    return Some(entry.sequence);
666                }
667            }
668        }
669
670        None
671    }
672
673    /// Attach a broadcast sender for WAL replication.
674    ///
675    /// When set, every `append()` call will publish the WAL entry to this
676    /// channel so the WAL shipper can stream it to connected followers.
677    pub fn set_replication_tx(&self, tx: tokio::sync::broadcast::Sender<WALEntry>) {
678        *self.replication_tx.lock() = Some(tx);
679    }
680}
681
682#[cfg(test)]
683mod tests {
684    use super::*;
685    use serde_json::json;
686    use tempfile::TempDir;
687    use uuid::Uuid;
688
689    fn create_test_event() -> Event {
690        Event::reconstruct_from_strings(
691            Uuid::new_v4(),
692            "test.event".to_string(),
693            "test-entity".to_string(),
694            "default".to_string(),
695            json!({"test": "data"}),
696            Utc::now(),
697            None,
698            1,
699        )
700    }
701
702    #[test]
703    fn test_wal_creation() {
704        let temp_dir = TempDir::new().unwrap();
705        let wal = WriteAheadLog::new(temp_dir.path(), WALConfig::default());
706        assert!(wal.is_ok());
707    }
708
709    #[test]
710    fn test_wal_append() {
711        let temp_dir = TempDir::new().unwrap();
712        let wal = WriteAheadLog::new(temp_dir.path(), WALConfig::default()).unwrap();
713
714        let event = create_test_event();
715        let seq = wal.append(event);
716        assert!(seq.is_ok());
717        assert_eq!(seq.unwrap(), 1);
718
719        let stats = wal.stats();
720        assert_eq!(stats.total_entries, 1);
721    }
722
723    #[test]
724    fn test_wal_on_disk_stats() {
725        let temp_dir = TempDir::new().unwrap();
726        let wal = WriteAheadLog::new(temp_dir.path(), WALConfig::default()).unwrap();
727
728        // Empty WAL dir (the active segment is created lazily on first append):
729        // zero or one segment, zero bytes.
730        let (bytes0, segs0) = wal.on_disk_stats().unwrap();
731        assert_eq!(bytes0, 0);
732        assert!(segs0 <= 1);
733
734        for _ in 0..5 {
735            wal.append(create_test_event()).unwrap();
736        }
737        wal.flush().unwrap();
738
739        let (bytes, segs) = wal.on_disk_stats().unwrap();
740        assert!(segs >= 1, "expected at least one WAL segment, got {segs}");
741        assert!(
742            bytes > 0,
743            "expected non-zero WAL bytes after appends, got {bytes}"
744        );
745    }
746
747    #[test]
748    fn test_wal_recovery() {
749        let temp_dir = TempDir::new().unwrap();
750        let wal = WriteAheadLog::new(temp_dir.path(), WALConfig::default()).unwrap();
751
752        // Write some events
753        for _ in 0..5 {
754            wal.append(create_test_event()).unwrap();
755        }
756
757        wal.flush().unwrap();
758
759        // Create new WAL instance (simulating restart)
760        let wal2 = WriteAheadLog::new(temp_dir.path(), WALConfig::default()).unwrap();
761        let recovered = wal2.recover().unwrap();
762
763        assert_eq!(recovered.len(), 5);
764    }
765
766    #[test]
767    fn test_wal_recovery_with_partial_write() {
768        let temp_dir = TempDir::new().unwrap();
769        let wal = WriteAheadLog::new(temp_dir.path(), WALConfig::default()).unwrap();
770
771        // Write 3 complete events
772        for _ in 0..3 {
773            wal.append(create_test_event()).unwrap();
774        }
775        wal.flush().unwrap();
776
777        // Simulate a partial write by appending malformed data to the WAL file
778        let wal_file_path = temp_dir.path().join("wal-0000000000000000.log");
779        use std::io::Write as _;
780        let mut f = std::fs::OpenOptions::new()
781            .append(true)
782            .open(&wal_file_path)
783            .unwrap();
784        f.write_all(b"{\"partial\": true, \"seq\"\n").unwrap(); // truncated JSON
785        drop(f);
786
787        // Recovery should succeed and return only the 3 valid events
788        let wal2 = WriteAheadLog::new(temp_dir.path(), WALConfig::default()).unwrap();
789        let recovered = wal2.recover().unwrap();
790        assert_eq!(
791            recovered.len(),
792            3,
793            "Should recover only the 3 valid events, not the partial one"
794        );
795    }
796
797    #[test]
798    fn test_wal_rotation() {
799        let temp_dir = TempDir::new().unwrap();
800        let config = WALConfig {
801            max_file_size: 1024, // Small size to trigger rotation
802            ..Default::default()
803        };
804
805        let wal = WriteAheadLog::new(temp_dir.path(), config).unwrap();
806
807        // Write enough events to trigger rotation
808        for _ in 0..50 {
809            wal.append(create_test_event()).unwrap();
810        }
811
812        let stats = wal.stats();
813        assert!(stats.files_rotated > 0);
814    }
815
816    #[test]
817    fn test_wal_entry_checksum() {
818        let event = create_test_event();
819        let entry = WALEntry::new(1, event);
820
821        assert!(entry.verify());
822
823        // Modify and verify it fails
824        let mut corrupted = entry.clone();
825        corrupted.checksum = 0;
826        assert!(!corrupted.verify());
827    }
828
829    #[test]
830    fn test_wal_fsync_interval_config() {
831        let config = WALConfig {
832            fsync_interval_ms: Some(100),
833            ..Default::default()
834        };
835        assert_eq!(config.fsync_interval_ms, Some(100));
836        // Default should have no interval
837        assert_eq!(WALConfig::default().fsync_interval_ms, None);
838    }
839
840    #[test]
841    fn test_wal_sync_method() {
842        let temp_dir = TempDir::new().unwrap();
843        let config = WALConfig {
844            sync_on_write: false, // Disable per-write sync to test explicit sync
845            ..Default::default()
846        };
847        let wal = WriteAheadLog::new(temp_dir.path(), config).unwrap();
848
849        // Write events without per-write fsync
850        for _ in 0..5 {
851            wal.append(create_test_event()).unwrap();
852        }
853
854        // Explicitly sync โ€” should flush + fsync without error
855        wal.sync().unwrap();
856
857        // Verify data survives by recovering from a new WAL instance
858        let wal2 = WriteAheadLog::new(temp_dir.path(), WALConfig::default()).unwrap();
859        let recovered = wal2.recover().unwrap();
860        assert_eq!(recovered.len(), 5);
861    }
862
863    #[test]
864    fn test_wal_truncate() {
865        let temp_dir = TempDir::new().unwrap();
866        let wal = WriteAheadLog::new(temp_dir.path(), WALConfig::default()).unwrap();
867
868        // Write events
869        for _ in 0..5 {
870            wal.append(create_test_event()).unwrap();
871        }
872
873        // Truncate
874        wal.truncate().unwrap();
875
876        // Verify sequence is reset
877        assert_eq!(wal.current_sequence(), 0);
878
879        // Verify recovery returns empty
880        let recovered = wal.recover().unwrap();
881        assert_eq!(recovered.len(), 0);
882    }
883}