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#[derive(Debug, Clone, PartialEq, Eq)]
20pub struct WalSegmentStamp {
21 pub path: PathBuf,
22 pub len: u64,
23 pub modified: Option<SystemTime>,
24}
25
26pub const QUARANTINE_DIR: &str = "quarantine";
29
30pub struct WriteAheadLog {
37 wal_dir: PathBuf,
39
40 current_file: Arc<RwLock<WALFile>>,
42
43 config: WALConfig,
45
46 stats: Arc<RwLock<WALStats>>,
48
49 sequence: Arc<RwLock<u64>>,
51
52 replication_tx: parking_lot::Mutex<Option<tokio::sync::broadcast::Sender<WALEntry>>>,
56
57 unreadable_segments: parking_lot::Mutex<HashSet<PathBuf>>,
59}
60
61#[derive(Debug, Clone)]
62pub struct WALConfig {
63 pub max_file_size: usize,
65
66 pub sync_on_write: bool,
68
69 pub max_wal_files: usize,
71
72 pub compress: bool,
74
75 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, 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#[derive(Debug, Clone, Serialize, Deserialize)]
106pub struct WALEntry {
107 pub sequence: u64,
109
110 pub wal_timestamp: DateTime<Utc>,
112
113 pub event: Event,
115
116 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 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
153struct 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 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 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 pub fn new(wal_dir: impl Into<PathBuf>, config: WALConfig) -> Result<Self> {
227 let wal_dir = wal_dir.into();
228
229 fs::create_dir_all(&wal_dir).map_err(|e| {
231 AllSourceError::StorageError(format!("Failed to create WAL directory: {e}"))
232 })?;
233
234 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 fn generate_wal_filename(dir: &Path, sequence: u64) -> PathBuf {
253 dir.join(format!("wal-{sequence:016x}.log"))
254 }
255
256 #[cfg_attr(feature = "hotpath", hotpath::measure)]
258 pub fn append(&self, event: Event) -> Result<u64> {
259 let mut seq = self.sequence.write();
261 *seq += 1;
262 let sequence = *seq;
263 drop(seq);
264
265 let entry = WALEntry::new(sequence, event);
267
268 let mut current = self.current_file.write();
270 let bytes_written = current.write_entry(&entry, self.config.sync_on_write)?;
271
272 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 if let Some(ref tx) = *self.replication_tx.lock() {
282 let _ = tx.send(entry);
283 }
284
285 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 #[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 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 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 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 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 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 pub fn segment_stamps(&self) -> Result<Vec<WalSegmentStamp>> {
431 let mut stamps = Vec::new();
432 for path in self.list_wal_files()? {
433 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 #[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 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 let mut seq = self.sequence.write();
525 *seq = max_sequence;
526 drop(seq);
527
528 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 #[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 #[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 pub fn fsync_interval_ms(&self) -> Option<u64> {
572 self.config.fsync_interval_ms
573 }
574
575 #[cfg_attr(feature = "hotpath", hotpath::measure)]
577 pub fn truncate(&self) -> Result<()> {
578 tracing::info!("๐งน Truncating WAL after checkpoint");
579
580 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 let new_file_path = Self::generate_wal_filename(&self.wal_dir, 0);
591 *current = WALFile::new(new_file_path)?;
592
593 let mut seq = self.sequence.write();
595 *seq = 0;
596
597 tracing::info!("โ
WAL truncated successfully");
598
599 Ok(())
600 }
601
602 pub fn stats(&self) -> WALStats {
604 (*self.stats.read()).clone()
605 }
606
607 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 pub fn current_sequence(&self) -> u64 {
632 *self.sequence.read()
633 }
634
635 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 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 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 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 for _ in 0..5 {
754 wal.append(create_test_event()).unwrap();
755 }
756
757 wal.flush().unwrap();
758
759 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 for _ in 0..3 {
773 wal.append(create_test_event()).unwrap();
774 }
775 wal.flush().unwrap();
776
777 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(); drop(f);
786
787 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, ..Default::default()
803 };
804
805 let wal = WriteAheadLog::new(temp_dir.path(), config).unwrap();
806
807 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 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 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, ..Default::default()
846 };
847 let wal = WriteAheadLog::new(temp_dir.path(), config).unwrap();
848
849 for _ in 0..5 {
851 wal.append(create_test_event()).unwrap();
852 }
853
854 wal.sync().unwrap();
856
857 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 for _ in 0..5 {
870 wal.append(create_test_event()).unwrap();
871 }
872
873 wal.truncate().unwrap();
875
876 assert_eq!(wal.current_sequence(), 0);
878
879 let recovered = wal.recover().unwrap();
881 assert_eq!(recovered.len(), 0);
882 }
883}