1use std::cmp::Ordering;
2use std::io::{Error, ErrorKind, Result};
3use std::num::NonZeroUsize;
4use std::path::{Path, PathBuf};
5use std::str::FromStr;
6use std::{fmt, mem, ops, result, thread};
7
8use fs_err as fs;
9use fs_err::File;
10use log::{debug, info, trace};
11pub use segment::{Entry, Segment};
12
13use crate::wal::segment_creator::SegmentCreatorV2;
14
15mod mmap_view_sync;
16mod segment;
17mod segment_creator;
18pub mod test_utils;
19
20#[cfg(test)]
21mod test_segment_recovery;
22
23#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
24pub struct WalOptions {
25 pub segment_capacity: usize,
27
28 pub segment_queue_len: usize,
31
32 pub retain_closed: NonZeroUsize,
34}
35
36impl Default for WalOptions {
37 fn default() -> WalOptions {
38 WalOptions {
39 segment_capacity: 32 * 1024 * 1024,
40 segment_queue_len: 0,
41 retain_closed: NonZeroUsize::new(1).unwrap(),
42 }
43 }
44}
45
46#[derive(Debug)]
48struct OpenSegment {
49 pub id: u64,
50 pub segment: Segment,
51}
52
53#[derive(Debug)]
55struct ClosedSegment {
56 pub start_index: u64,
57 pub segment: Segment,
58}
59
60enum WalSegment {
61 Open(OpenSegment),
62 Closed(ClosedSegment),
63}
64
65pub struct Wal {
76 open_segment: OpenSegment,
78 closed_segments: Vec<ClosedSegment>,
79 creator: SegmentCreatorV2,
80
81 retain_closed: NonZeroUsize,
83
84 #[expect(dead_code)]
87 dir: File,
88
89 path: PathBuf,
91
92 flush: Option<thread::JoinHandle<Result<()>>>,
95}
96
97impl Wal {
98 pub fn open<P>(path: P) -> Result<Wal>
99 where
100 P: AsRef<Path>,
101 {
102 Wal::with_options(path, &WalOptions::default())
103 }
104
105 pub fn generate_empty_wal_starting_at_index(
106 path: impl Into<PathBuf>,
107 options: &WalOptions,
108 index: u64,
109 ) -> Result<()> {
110 let open_id = 0;
111 let mut path_buf = path.into();
112 path_buf.push(format!("open-{open_id}"));
113 let segment = OpenSegment {
114 id: index + 1,
115 segment: Segment::create(&path_buf, options.segment_capacity)?,
116 };
117
118 let mut close_segment = close_segment(segment, index + 1)?;
119
120 close_segment.segment.flush()
121 }
122
123 pub fn with_options<P>(path: P, options: &WalOptions) -> Result<Wal>
124 where
125 P: AsRef<Path>,
126 {
127 debug!("Wal {{ path: {:?} }}: opening", path.as_ref());
128
129 #[cfg(not(target_os = "windows"))]
130 let path = path.as_ref().to_path_buf();
131 #[cfg(not(target_os = "windows"))]
132 let dir = File::open(&path)?;
133
134 #[cfg(target_os = "windows")]
140 let mut path = path.as_ref().to_path_buf();
141 #[cfg(target_os = "windows")]
142 let dir = {
143 path.push(".wal");
144 let dir = File::options()
145 .create(true)
146 .read(true)
147 .write(true)
148 .truncate(true)
149 .open(&path)?;
150 path.pop();
151 dir
152 };
153
154 fs4::FileExt::try_lock(dir.file())?;
164
165 let mut open_segments: Vec<OpenSegment> = Vec::new();
167 let mut closed_segments: Vec<ClosedSegment> = Vec::new();
168
169 for entry in fs::read_dir(&path)? {
170 match open_dir_entry(entry?)? {
171 Some(WalSegment::Open(open_segment)) => open_segments.push(open_segment),
172 Some(WalSegment::Closed(closed_segment)) => closed_segments.push(closed_segment),
173 None => {}
174 }
175 }
176
177 closed_segments.sort_by_key(|s| s.start_index);
179 let mut next_start_index = closed_segments
180 .first()
181 .map_or(0, |segment| segment.start_index);
182 for &ClosedSegment {
183 start_index,
184 ref segment,
185 } in &closed_segments
186 {
187 match start_index.cmp(&next_start_index) {
188 Ordering::Less => {
189 return Err(Error::new(
190 ErrorKind::InvalidData,
191 format!(
192 "overlapping segments: segment at {start_index} overlaps with already-covered range up to {next_start_index}"
193 ),
194 ));
195 }
196 Ordering::Equal => {
197 next_start_index = start_index + segment.len() as u64;
198 }
199 Ordering::Greater => {
200 return Err(Error::new(
201 ErrorKind::InvalidData,
202 format!(
203 "missing segment(s) containing wal entries {next_start_index} to {start_index}"
204 ),
205 ));
206 }
207 }
208 }
209
210 open_segments.sort_by_key(|s| s.id);
212
213 let mut open_segment: Option<OpenSegment> = None;
215 let mut unused_segments: Vec<OpenSegment> = Vec::new();
217
218 for segment in open_segments {
219 if !segment.segment.is_empty() {
220 let stranded_segment = open_segment.take();
226 open_segment = Some(segment);
227 if let Some(segment) = stranded_segment {
228 let closed_segment = close_segment(segment, next_start_index)?;
229 next_start_index += closed_segment.segment.len() as u64;
230 closed_segments.push(closed_segment);
231 }
232 } else if open_segment.is_none() {
233 open_segment = Some(segment);
234 } else {
235 unused_segments.push(segment);
236 }
237 }
238
239 let mut creator = SegmentCreatorV2::new(
240 &path,
241 open_segment.as_ref(),
242 unused_segments,
243 options.segment_capacity,
244 options.segment_queue_len,
245 );
246
247 let open_segment = match open_segment {
248 Some(segment) => segment,
249 None => creator.next()?,
250 };
251
252 let wal = Wal {
253 open_segment,
254 closed_segments,
255 retain_closed: options.retain_closed,
256 creator,
257 dir,
258 path,
259 flush: None,
260 };
261 info!("{wal:?}: opened");
262 Ok(wal)
263 }
264
265 fn retire_open_segment(&mut self) -> Result<()> {
266 trace!("{self:?}: retiring open segment");
267 let mut segment = self.creator.next()?;
268 mem::swap(&mut self.open_segment, &mut segment);
269
270 if let Some(flush) = self.flush.take() {
271 flush
272 .join()
273 .map_err(|err| Error::other(format!("wal flush thread panicked: {err:?}")))??;
274 };
275
276 self.flush = Some(segment.segment.flush_async());
277
278 let start_index = self.open_segment_start_index();
279
280 if let Some(empty_segment) = self
282 .closed_segments
283 .pop_if(|last_closed| last_closed.segment.is_empty())
284 {
285 empty_segment.segment.delete()?;
286 }
287
288 self.closed_segments
289 .push(close_segment(segment, start_index)?);
290 debug!("{self:?}: open segment retired. start_index: {start_index}");
291 Ok(())
292 }
293
294 pub fn append<T>(&mut self, entry: &T) -> Result<u64>
295 where
296 T: ops::Deref<Target = [u8]>,
297 {
298 trace!("{:?}: appending entry of length {}", self, entry.len());
299 if !self.open_segment.segment.sufficient_capacity(entry.len()) {
300 if !self.open_segment.segment.is_empty() {
301 self.retire_open_segment()?;
302 }
303 self.open_segment.segment.ensure_capacity(entry.len())?;
304 }
305
306 Ok(self.open_segment_start_index()
307 + self.open_segment.segment.append(entry).unwrap() as u64)
308 }
309
310 pub fn flush_open_segment(&mut self) -> Result<()> {
311 trace!("{self:?}: flushing open segments");
312 self.open_segment.segment.flush()?;
313 Ok(())
314 }
315
316 pub fn flush_open_segment_async(&mut self) -> thread::JoinHandle<Result<()>> {
317 trace!("{self:?}: flushing open segments");
318 self.open_segment.segment.flush_async()
319 }
320
321 pub fn entry(&self, index: u64) -> Option<Entry> {
323 let open_start_index = self.open_segment_start_index();
324 if index >= open_start_index {
325 return self
326 .open_segment
327 .segment
328 .entry((index - open_start_index) as usize);
329 }
330
331 match self.find_closed_segment(index) {
332 Ok(segment_index) => {
333 let segment = &self.closed_segments[segment_index];
334 segment
335 .segment
336 .entry((index - segment.start_index) as usize)
337 }
338 Err(i) => {
339 assert_eq!(0, i);
341 None
342 }
343 }
344 }
345
346 pub fn truncate(&mut self, from: u64) -> Result<()> {
352 trace!("{self:?}: truncate from entry {from}");
353 let open_start_index = self.open_segment_start_index();
354 if from >= open_start_index {
355 self.open_segment
356 .segment
357 .truncate((from - open_start_index) as usize);
358 } else {
359 self.open_segment.segment.truncate(0);
361
362 match self.find_closed_segment(from) {
363 Ok(index) => {
364 if from == self.closed_segments[index].start_index {
365 for segment in self.closed_segments.drain(index..) {
366 segment.segment.delete()?;
368 }
369 } else {
370 {
371 let segment = &mut self.closed_segments[index];
372 segment
373 .segment
374 .truncate((from - segment.start_index) as usize);
375 segment.segment.flush()?;
377 }
378 if index + 1 < self.closed_segments.len() {
379 for segment in self.closed_segments.drain(index + 1..) {
380 segment.segment.delete()?;
382 }
383 }
384 }
385 }
386 Err(index) => {
387 assert!(
389 from <= self
390 .closed_segments
391 .get(index)
392 .map_or(0, |segment| segment.start_index)
393 );
394 for segment in self.closed_segments.drain(..) {
395 segment.segment.delete()?;
397 }
398 }
399 }
400 }
401 Ok(())
402 }
403
404 pub fn prefix_truncate(&mut self, until: u64) -> Result<()> {
411 trace!("{self:?}: prefix_truncate until entry {until}");
412
413 if until
415 <= self
416 .closed_segments
417 .first()
418 .map_or(0, |segment| segment.start_index)
419 {
420 return Ok(());
421 }
422
423 let retain_start_index = self
426 .closed_segments
427 .len()
428 .saturating_sub(self.retain_closed.get());
429
430 if until >= self.open_segment_start_index() {
432 for segment in self.closed_segments.drain(..retain_start_index) {
433 segment.segment.delete()?
434 }
435 return Ok(());
436 }
437
438 let index = self.find_closed_segment(until).unwrap();
440 let truncate_until_index = index.min(retain_start_index);
441 trace!("{self:?}: prefix truncating until closed segment {truncate_until_index}");
442 for segment in self.closed_segments.drain(..truncate_until_index) {
443 segment.segment.delete()?
444 }
445 Ok(())
446 }
447
448 fn open_segment_start_index(&self) -> u64 {
450 self.closed_segments
451 .last()
452 .map_or(0, |segment: &ClosedSegment| {
453 segment.start_index + segment.segment.len() as u64
454 })
455 }
456
457 fn find_closed_segment(&self, index: u64) -> result::Result<usize, usize> {
458 self.closed_segments.binary_search_by(|segment| {
459 if index < segment.start_index {
460 Ordering::Greater
461 } else if index >= segment.start_index + segment.segment.len() as u64 {
462 Ordering::Less
463 } else {
464 Ordering::Equal
465 }
466 })
467 }
468
469 pub fn path(&self) -> &Path {
470 &self.path
471 }
472
473 pub fn num_segments(&self) -> usize {
474 self.closed_segments.len() + 1
475 }
476
477 pub fn num_entries(&self) -> u64 {
478 self.open_segment_start_index()
479 - self
480 .closed_segments
481 .first()
482 .map_or(0, |segment| segment.start_index)
483 + self.open_segment.segment.len() as u64
484 }
485
486 pub fn first_index(&self) -> u64 {
488 self.closed_segments
489 .first()
490 .map_or(0, |segment| segment.start_index)
491 }
492
493 pub fn last_index(&self) -> u64 {
495 let num_entries = self.num_entries();
496 self.first_index() + num_entries.saturating_sub(1)
497 }
498
499 pub fn clear(&mut self) -> Result<()> {
501 self.truncate(self.first_index())
502 }
503
504 pub fn set_retention(&mut self, retain_closed: usize) {
508 self.retain_closed = NonZeroUsize::new(retain_closed.max(1)).unwrap();
509 }
510}
511
512impl fmt::Debug for Wal {
513 fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
514 let Self {
515 open_segment,
516 closed_segments,
517 path,
518 creator: _,
519 retain_closed: _,
520 dir: _,
521 flush: _,
522 } = self;
523 let start_index = closed_segments
524 .first()
525 .map_or(0, |segment| segment.start_index);
526 let end_index = self.open_segment_start_index() + open_segment.segment.len() as u64;
527 f.debug_struct("Wal")
528 .field("path", path)
529 .field("segment-count", &(closed_segments.len() + 1))
530 .field("entries", &format_args!("[{start_index}, {end_index})"))
531 .finish_non_exhaustive()
532 }
533}
534
535fn close_segment(mut segment: OpenSegment, start_index: u64) -> Result<ClosedSegment> {
536 let new_path = segment
537 .segment
538 .path()
539 .with_file_name(format!("closed-{start_index}"));
540 segment.segment.rename(new_path)?;
541 segment.segment.close();
542 Ok(ClosedSegment {
543 start_index,
544 segment: segment.segment,
545 })
546}
547
548fn open_dir_entry(entry: fs::DirEntry) -> Result<Option<WalSegment>> {
549 let metadata = entry.metadata()?;
550
551 let error = || {
552 Error::new(
553 ErrorKind::InvalidData,
554 format!("unexpected entry in wal directory: {:?}", entry.path()),
555 )
556 };
557
558 if !metadata.is_file() {
559 return Ok(None); }
561
562 let filename = entry.file_name().into_string().map_err(|_| error())?;
563 match filename.split_once('-') {
564 Some(("tmp", _)) => {
565 fs::remove_file(entry.path())?;
567 Ok(None)
568 }
569 Some(("open", id)) => {
570 let id = u64::from_str(id).map_err(|_| error())?;
571 let segment = Segment::open(entry.path())?;
572 Ok(Some(WalSegment::Open(OpenSegment { id, segment })))
573 }
574 Some(("closed", start)) => {
575 let start = u64::from_str(start).map_err(|_| error())?;
576 let segment = Segment::open(entry.path())?;
577 Ok(Some(WalSegment::Closed(ClosedSegment {
578 start_index: start,
579 segment,
580 })))
581 }
582 _ => Ok(None), }
584}
585
586#[cfg(test)]
587mod test {
588 use std::io::{ErrorKind, Write};
589 use std::num::{NonZeroU8, NonZeroUsize};
590
591 use fs_err as fs;
592 use log::trace;
593 use quickcheck::{QuickCheck, TestResult};
594 use tempfile::Builder;
595
596 use super::{Wal, WalOptions};
597 use crate::wal::test_utils::EntryGenerator;
598
599 #[cfg(target_os = "windows")]
601 const QC_TESTS: u64 = 3;
602
603 #[cfg(not(target_os = "windows"))]
604 const QC_TESTS: u64 = 50;
605
606 fn init_logger() {
607 let _ = env_logger::builder().is_test(true).try_init();
608 }
609
610 #[test]
611 fn test_generate_empty_wal() {
612 init_logger();
613 let dir = Builder::new().prefix("wal").tempdir().unwrap();
614 let options = WalOptions {
615 segment_capacity: 80,
616 segment_queue_len: 3,
617 retain_closed: NonZeroUsize::new(1).unwrap(),
618 };
619
620 let init_offset = 10;
622 Wal::generate_empty_wal_starting_at_index(dir.path(), &options, init_offset).unwrap();
623
624 let mut wal = Wal::with_options(dir.path(), &options).unwrap();
625
626 let first_index = wal.first_index();
627 let last_index = wal.last_index();
628 let num_entries = wal.num_entries();
629
630 assert!(first_index <= last_index);
631 assert_eq!(num_entries, 0);
632
633 let next_entry: Vec<u8> = vec![1, 2, 3];
634 let op = wal.append(&next_entry).unwrap();
635
636 assert!(op > init_offset);
637
638 let first_index = wal.first_index();
639 let last_index = wal.last_index();
640 let num_entries = wal.num_entries();
641
642 assert!(first_index <= last_index);
643 assert_eq!(num_entries, 1);
644
645 wal.append(&next_entry).unwrap();
646 wal.append(&next_entry).unwrap();
647
648 let first_index = wal.first_index();
649 let last_index = wal.last_index();
650 let num_entries = wal.num_entries();
651
652 assert!(first_index <= last_index);
653 assert_eq!(num_entries, 3);
654 }
655
656 #[test]
657 fn test_create_empty_wal_with_initial_id() {
658 init_logger();
659 let dir = Builder::new().prefix("wal").tempdir().unwrap();
660 let options = WalOptions {
661 segment_capacity: 80,
662 segment_queue_len: 3,
663 retain_closed: NonZeroUsize::new(1).unwrap(),
664 };
665
666 let init_offset = 10;
668 Wal::generate_empty_wal_starting_at_index(dir.path(), &options, init_offset).unwrap();
669
670 let mut wal = Wal::with_options(dir.path(), &options).unwrap();
671
672 let last_index = wal.last_index();
673
674 assert!(last_index > init_offset);
675
676 assert_eq!(wal.num_entries(), 0);
677
678 let next_entry: Vec<u8> = vec![1, 2, 3];
679
680 wal.append(&next_entry).unwrap();
681
682 let last_index = wal.last_index();
683 assert_eq!(last_index, init_offset + 1);
684
685 assert_eq!(wal.num_entries(), 1);
686
687 let entry_count = 50;
688
689 let entries = EntryGenerator::new().take(entry_count).collect::<Vec<_>>();
690
691 for entry in &entries {
692 wal.append(entry).unwrap();
693 }
694
695 let last_index = wal.last_index();
696 assert_eq!(last_index, init_offset + 1 + entry_count as u64);
697
698 assert_eq!(wal.num_entries(), 1 + entry_count as u64);
699
700 {
702 let entry_index = init_offset + 1;
703 let entry = wal.entry(entry_index).unwrap();
704 assert_eq!(next_entry[..], entry[..]);
705
706 let entry_index = init_offset + 1 + entry_count as u64;
707 let entry = wal.entry(entry_index).unwrap();
708 assert_eq!(entries[entry_count - 1][..], entry[..]);
709
710 let entry_index = init_offset + 1 + 10;
711 let entry = wal.entry(entry_index).unwrap();
712 assert_eq!(entries[9][..], entry[..]);
713 }
714
715 wal.prefix_truncate(init_offset).unwrap();
716
717 assert_eq!(wal.num_entries(), entry_count as u64 + 1);
718
719 wal.prefix_truncate(init_offset + 20).unwrap();
720
721 assert!(wal.num_entries() < entry_count as u64 + 1);
722
723 let truncate_index = init_offset + 30;
724 wal.truncate(truncate_index).unwrap();
725
726 let last_index = wal.last_index();
727 assert_eq!(last_index, truncate_index - 1);
728 }
729
730 #[test]
732 fn check_wal() {
733 init_logger();
734 fn wal(entry_count: u8) -> TestResult {
735 let dir = Builder::new().prefix("wal").tempdir().unwrap();
736 let mut wal = Wal::with_options(
737 dir.path(),
738 &WalOptions {
739 segment_capacity: 80,
740 segment_queue_len: 3,
741 retain_closed: NonZeroUsize::new(1).unwrap(),
742 },
743 )
744 .unwrap();
745 let entries = EntryGenerator::new()
746 .take(entry_count as usize)
747 .collect::<Vec<_>>();
748
749 for entry in &entries {
750 wal.append(entry).unwrap();
751 }
752
753 for (index, expected) in entries.iter().enumerate() {
754 match wal.entry(index as u64) {
755 Some(ref entry) if entry[..] != expected[..] => return TestResult::failed(),
756 None => return TestResult::failed(),
757 _ => (),
758 }
759 }
760 TestResult::passed()
761 }
762
763 QuickCheck::new()
764 .tests(QC_TESTS)
765 .quickcheck(wal as fn(u8) -> TestResult);
766 }
767
768 #[test]
769 fn check_last_index() {
770 init_logger();
771 fn check(entry_count: u8) -> TestResult {
772 let dir = Builder::new().prefix("wal").tempdir().unwrap();
773 let mut wal = Wal::with_options(
774 dir.path(),
775 &WalOptions {
776 segment_capacity: 80,
777 segment_queue_len: 3,
778 retain_closed: NonZeroUsize::new(1).unwrap(),
779 },
780 )
781 .unwrap();
782 let entries = EntryGenerator::new()
783 .take(entry_count as usize)
784 .collect::<Vec<_>>();
785
786 for entry in &entries {
787 wal.append(entry).unwrap();
788 }
789 if entries.is_empty() {
790 assert_eq!(wal.last_index(), 0);
791 } else {
792 assert_eq!(wal.last_index(), entries.len() as u64 - 1);
793 }
794
795 let last_index = wal.last_index();
796 if wal.entry(last_index).is_none() && wal.num_entries() != 0 {
797 return TestResult::failed();
798 }
799 if wal.entry(last_index + 1).is_some() {
800 return TestResult::failed();
801 }
802 TestResult::passed()
803 }
804
805 QuickCheck::new()
806 .tests(QC_TESTS)
807 .quickcheck(check as fn(u8) -> TestResult)
808 }
809
810 #[test]
811 fn check_clear() {
812 init_logger();
813 fn check(entry_count: u8) -> TestResult {
814 let dir = Builder::new().prefix("wal").tempdir().unwrap();
815 let mut wal = Wal::with_options(
816 dir.path(),
817 &WalOptions {
818 segment_capacity: 80,
819 segment_queue_len: 3,
820 retain_closed: NonZeroUsize::new(1).unwrap(),
821 },
822 )
823 .unwrap();
824 let entries = EntryGenerator::new()
825 .take(entry_count as usize)
826 .collect::<Vec<_>>();
827
828 for entry in &entries {
829 wal.append(entry).unwrap();
830 }
831 wal.clear().unwrap();
832 TestResult::from_bool(wal.num_entries() == 0)
833 }
834
835 QuickCheck::new()
836 .tests(QC_TESTS)
837 .quickcheck(check as fn(u8) -> TestResult)
838 }
839
840 #[test]
842 fn check_reopen() {
843 init_logger();
844 fn wal(entry_count: u8) -> TestResult {
845 let entries = EntryGenerator::new()
846 .take(entry_count as usize)
847 .collect::<Vec<_>>();
848 let dir = Builder::new().prefix("wal").tempdir().unwrap();
849 {
850 let mut wal = Wal::with_options(
851 dir.path(),
852 &WalOptions {
853 segment_capacity: 80,
854 segment_queue_len: 3,
855 retain_closed: NonZeroUsize::new(1).unwrap(),
856 },
857 )
858 .unwrap();
859 for entry in &entries {
860 let _ = wal.append(entry);
861 }
862 }
863
864 {
865 let mut file = fs::OpenOptions::new()
867 .read(true)
868 .write(true)
869 .create(true)
870 .truncate(true)
871 .open(dir.path().join("tmp-open-123"))
872 .unwrap();
873
874 let _ = file.write(b"123").unwrap();
875 }
876
877 let wal = Wal::with_options(
878 dir.path(),
879 &WalOptions {
880 segment_capacity: 80,
881 segment_queue_len: 3,
882 retain_closed: NonZeroUsize::new(1).unwrap(),
883 },
884 )
885 .unwrap();
886 for (index, expected) in entries.iter().enumerate() {
888 match wal.entry(index as u64) {
889 Some(ref entry) if entry[..] != expected[..] => return TestResult::failed(),
890 None => return TestResult::failed(),
891 _ => (),
892 }
893 }
894 TestResult::passed()
895 }
896
897 QuickCheck::new()
898 .tests(QC_TESTS)
899 .quickcheck(wal as fn(u8) -> TestResult);
900 }
901
902 #[test]
903 fn check_truncate() {
904 init_logger();
905 fn truncate(entry_count: u8, truncate: u8) -> TestResult {
906 if truncate > entry_count {
907 return TestResult::discard();
908 }
909 let dir = Builder::new().prefix("wal").tempdir().unwrap();
910 let mut wal = Wal::with_options(
911 dir.path(),
912 &WalOptions {
913 segment_capacity: 80,
914 segment_queue_len: 3,
915 retain_closed: NonZeroUsize::new(1).unwrap(),
916 },
917 )
918 .unwrap();
919 let entries = EntryGenerator::new()
920 .take(entry_count as usize)
921 .collect::<Vec<_>>();
922
923 for entry in &entries {
924 if let Err(error) = wal.append(entry) {
925 return TestResult::error(error.to_string());
926 }
927 }
928
929 wal.truncate(u64::from(truncate)).unwrap();
930
931 for (index, expected) in entries.iter().take(truncate as usize).enumerate() {
932 match wal.entry(index as u64) {
933 Some(ref entry) if entry[..] != expected[..] => return TestResult::failed(),
934 None => return TestResult::failed(),
935 _ => (),
936 }
937 }
938
939 TestResult::from_bool(wal.entry(u64::from(truncate)).is_none())
940 }
941
942 QuickCheck::new()
943 .tests(QC_TESTS)
944 .quickcheck(truncate as fn(u8, u8) -> TestResult);
945 }
946
947 #[test]
948 fn check_prefix_truncate() {
949 init_logger();
950 fn prefix_truncate(entry_count: u8, until: u8, retain_closed: NonZeroU8) -> TestResult {
951 trace!(
952 "prefix truncate; entry_count: {entry_count}, until: {until}, retain_closed: {retain_closed}",
953 );
954 if until > entry_count {
955 return TestResult::discard();
956 }
957 let dir = Builder::new().prefix("wal").tempdir().unwrap();
958 let mut wal = Wal::with_options(
959 dir.path(),
960 &WalOptions {
961 segment_capacity: 80,
962 segment_queue_len: 3,
963 retain_closed: NonZeroUsize::from(retain_closed),
964 },
965 )
966 .unwrap();
967 let entries = EntryGenerator::new()
968 .take(entry_count as usize)
969 .collect::<Vec<_>>();
970
971 let mut has_ever_reached_max = false;
972 let retain_closed = retain_closed.get() as usize;
973
974 for entry in &entries {
975 wal.append(entry).unwrap();
976 if wal.closed_segments.len() >= retain_closed {
977 has_ever_reached_max = true;
978 }
979 }
980
981 wal.prefix_truncate(u64::from(until)).unwrap();
982
983 let retained = if has_ever_reached_max {
984 if until < entry_count {
986 wal.closed_segments.len() >= retain_closed
988 } else {
989 wal.closed_segments.len() == retain_closed
990 }
991 } else {
992 wal.closed_segments.len() < retain_closed
993 };
994
995 let num_entries = wal.num_entries() as u8;
996 TestResult::from_bool(
997 num_entries <= entry_count && num_entries >= entry_count - until && retained,
998 )
999 }
1000 QuickCheck::new()
1001 .tests(QC_TESTS)
1002 .quickcheck(prefix_truncate as fn(u8, u8, NonZeroU8) -> TestResult);
1003 }
1004
1005 #[test]
1006 fn test_append() {
1007 init_logger();
1008 let dir = Builder::new().prefix("wal").tempdir().unwrap();
1009 let mut wal = Wal::open(dir.path()).unwrap();
1010
1011 let entry: &[u8] = &[42u8; 4096];
1012 for _ in 1..10 {
1013 wal.append(&entry).unwrap();
1014 }
1015 }
1016
1017 #[test]
1018 fn test_truncate() {
1019 init_logger();
1020 let dir = Builder::new().prefix("wal").tempdir().unwrap();
1021 let mut wal = Wal::with_options(
1023 dir.path(),
1024 &WalOptions {
1025 segment_capacity: 4096,
1026 segment_queue_len: 3,
1027 retain_closed: NonZeroUsize::new(1).unwrap(),
1028 },
1029 )
1030 .unwrap();
1031
1032 let entry: [u8; 2000] = [42u8; 2000];
1033
1034 for truncate_index in 0..10 {
1035 assert!(wal.entry(0).is_none());
1036 for i in 0..10 {
1037 assert_eq!(i, wal.append(&&entry[..]).unwrap());
1038 }
1039
1040 wal.truncate(truncate_index).unwrap();
1041
1042 assert!(wal.entry(truncate_index).is_none());
1043
1044 if truncate_index > 0 {
1045 assert!(wal.entry(truncate_index - 1).is_some());
1046 }
1047 wal.truncate(0).unwrap();
1048 }
1049 }
1050
1051 fn run_test_with_retain_closed(retain_closed: usize) {
1052 init_logger();
1053 let num_entries = 10;
1054 let dir = Builder::new().prefix("wal").tempdir().unwrap();
1055 let entry: [u8; 2000] = [42u8; 2000];
1056 let mut wal = Wal::with_options(
1057 dir.path(),
1058 &WalOptions {
1059 segment_capacity: 4096,
1060 segment_queue_len: 3,
1061 retain_closed: NonZeroUsize::new(retain_closed).unwrap(),
1062 },
1063 )
1064 .unwrap();
1065
1066 let entries_per_segment = 2; for truncate_index in 0..num_entries {
1069 assert!(wal.entry(0).is_none());
1070 for i in 0..num_entries {
1071 assert_eq!(i, wal.append(&&entry[..]).unwrap());
1072 } let initial_closed_segments = wal.closed_segments.len();
1075
1076 wal.prefix_truncate(truncate_index).unwrap();
1078
1079 let segments_that_can_be_truncated = truncate_index / entries_per_segment;
1081 let expected_closed_segments = (initial_closed_segments
1082 - segments_that_can_be_truncated as usize)
1083 .max(retain_closed);
1084 let expected_trimmed_until =
1085 (initial_closed_segments - expected_closed_segments) as u64 * entries_per_segment;
1086
1087 assert!(wal.entry(truncate_index).is_some()); for i in 0..expected_trimmed_until {
1090 assert!(wal.entry(i).is_none());
1092 }
1093
1094 for i in expected_trimmed_until..num_entries {
1095 assert!(wal.entry(i).is_some());
1097 }
1098
1099 assert_eq!(wal.closed_segments.len(), expected_closed_segments);
1100 wal.truncate(0).unwrap(); }
1102 }
1103
1104 #[test]
1105 fn test_prefix_truncate_parametric() {
1106 run_test_with_retain_closed(1);
1107 run_test_with_retain_closed(2);
1108 run_test_with_retain_closed(3);
1109 }
1110
1111 #[test]
1112 fn test_truncate_flush() {
1113 init_logger();
1114 let dir = Builder::new().prefix("wal").tempdir().unwrap();
1115 let mut wal = Wal::with_options(
1117 dir.path(),
1118 &WalOptions {
1119 segment_capacity: 4096,
1120 segment_queue_len: 3,
1121 retain_closed: NonZeroUsize::new(1).unwrap(),
1122 },
1123 )
1124 .unwrap();
1125
1126 let entry: [u8; 2000] = [42u8; 2000];
1127 assert!(wal.entry(0).is_none());
1129
1130 for i in 0..10 {
1132 assert_eq!(i, wal.append(&&entry[..]).unwrap());
1133 }
1134
1135 assert_eq!(wal.num_entries(), 10);
1137 assert_eq!(wal.first_index(), 0);
1138 assert_eq!(wal.last_index(), 9);
1139 assert_eq!(wal.closed_segments.len(), 4); assert_eq!(wal.closed_segments[0].segment.len(), 2);
1141 assert_eq!(wal.closed_segments[1].segment.len(), 2);
1142 assert_eq!(wal.closed_segments[2].segment.len(), 2);
1143 assert_eq!(wal.closed_segments[3].segment.len(), 2);
1144 assert_eq!(wal.open_segment.segment.len(), 2); wal.flush_open_segment().unwrap();
1148
1149 assert_eq!(wal.num_entries(), 10);
1151 assert_eq!(wal.first_index(), 0);
1152 assert_eq!(wal.last_index(), 9);
1153 assert_eq!(wal.closed_segments.len(), 4); assert_eq!(wal.closed_segments[0].segment.len(), 2);
1155 assert_eq!(wal.closed_segments[1].segment.len(), 2);
1156 assert_eq!(wal.closed_segments[2].segment.len(), 2);
1157 assert_eq!(wal.closed_segments[3].segment.len(), 2);
1158 assert_eq!(wal.open_segment.segment.len(), 2); wal.truncate(9).unwrap();
1161
1162 assert_eq!(wal.open_segment.segment.len(), 1); wal.truncate(5).unwrap();
1166
1167 for i in 5..10 {
1169 assert!(wal.entry(i).is_none());
1170 }
1171
1172 wal.flush_open_segment().unwrap();
1174
1175 assert_eq!(wal.num_entries(), 5); assert_eq!(wal.first_index(), 0);
1177 assert_eq!(wal.last_index(), 4);
1178 assert_eq!(wal.closed_segments.len(), 3); assert_eq!(wal.closed_segments[0].segment.len(), 2);
1180 assert_eq!(wal.closed_segments[1].segment.len(), 2);
1181 assert_eq!(wal.closed_segments[2].segment.len(), 1);
1182 assert_eq!(wal.open_segment.segment.len(), 0); for i in 0..5 {
1186 assert_eq!(i + 5, wal.append(&&entry[..]).unwrap());
1187 }
1188
1189 assert_eq!(wal.num_entries(), 10);
1191 assert_eq!(wal.first_index(), 0);
1192 assert_eq!(wal.last_index(), 9);
1193 assert_eq!(wal.closed_segments.len(), 5);
1194 assert_eq!(wal.closed_segments[0].segment.len(), 2); assert_eq!(wal.closed_segments[1].segment.len(), 2); assert_eq!(wal.closed_segments[2].segment.len(), 1); assert_eq!(wal.closed_segments[3].segment.len(), 2); assert_eq!(wal.closed_segments[4].segment.len(), 2); assert_eq!(wal.open_segment.segment.len(), 1); eprintln!("wal: {wal:?}");
1202 eprintln!("wal open: {:?}", wal.open_segment);
1203 eprintln!("wal closed: {:?}", wal.closed_segments);
1204
1205 drop(wal);
1207 let wal = Wal::open(dir.path()).unwrap();
1208 assert_eq!(wal.num_entries(), 10);
1209 assert_eq!(wal.first_index(), 0);
1210 assert_eq!(wal.last_index(), 9);
1211 assert_eq!(wal.closed_segments.len(), 5);
1212 assert_eq!(wal.closed_segments[0].segment.len(), 2);
1213 assert_eq!(wal.closed_segments[1].segment.len(), 2);
1214 assert_eq!(wal.closed_segments[2].segment.len(), 1); assert_eq!(wal.closed_segments[3].segment.len(), 2);
1216 assert_eq!(wal.closed_segments[4].segment.len(), 2);
1217 assert_eq!(wal.open_segment.segment.len(), 1);
1218 }
1219
1220 #[test]
1222 fn test_exclusive_lock() {
1223 init_logger();
1224 let dir = Builder::new().prefix("wal").tempdir().unwrap();
1225 let wal = Wal::open(dir.path()).unwrap();
1226 assert_eq!(
1227 ErrorKind::WouldBlock,
1228 Wal::open(dir.path()).unwrap_err().kind()
1229 );
1230 drop(wal);
1231 assert!(Wal::open(dir.path()).is_ok());
1232 }
1233
1234 #[test]
1235 fn test_record_id_preserving() {
1236 init_logger();
1237 let entry_count = 55;
1238 let dir = Builder::new().prefix("wal").tempdir().unwrap();
1239 let options = WalOptions {
1240 segment_capacity: 512,
1241 segment_queue_len: 3,
1242 retain_closed: NonZeroUsize::new(1).unwrap(),
1243 };
1244
1245 let mut wal = Wal::with_options(dir.path(), &options).unwrap();
1246 let entries = vec![vec![0; 32]; entry_count];
1248
1249 for entry in &entries {
1250 wal.append(entry).unwrap();
1251 }
1252 let closed_segments = wal.closed_segments.len();
1253 let start_index = wal.open_segment_start_index();
1254
1255 wal.prefix_truncate(25).unwrap();
1256 let half_trunk_closed_segments = wal.closed_segments.len();
1257 let half_trunk_start_index = wal.open_segment_start_index();
1258
1259 wal.prefix_truncate((entry_count - 2) as u64).unwrap();
1260 let full_trunk_closed_segments = wal.closed_segments.len();
1261 let full_trunk_start_index = wal.open_segment_start_index();
1262
1263 assert!(closed_segments > half_trunk_closed_segments);
1264 assert!(half_trunk_closed_segments > full_trunk_closed_segments);
1265
1266 assert_eq!(start_index, half_trunk_start_index);
1267 assert_eq!(start_index, full_trunk_start_index);
1268 }
1269
1270 #[test]
1271 fn test_offset_after_open() {
1272 init_logger();
1273 let entry_count = 55;
1274 let dir = Builder::new().prefix("wal").tempdir().unwrap();
1275 let options = WalOptions {
1276 segment_capacity: 512,
1277 segment_queue_len: 3,
1278 retain_closed: NonZeroUsize::new(1).unwrap(),
1279 };
1280 let start_index;
1281 {
1282 let mut wal = Wal::with_options(dir.path(), &options).unwrap();
1283 let entries = EntryGenerator::new().take(entry_count).collect::<Vec<_>>();
1284
1285 for entry in &entries {
1286 wal.append(entry).unwrap();
1287 }
1288 start_index = wal.open_segment_start_index();
1289 wal.prefix_truncate(25).unwrap();
1290 assert_eq!(start_index, wal.open_segment_start_index());
1291 }
1292 {
1293 let wal2 = Wal::with_options(dir.path(), &options).unwrap();
1294 assert_eq!(start_index, wal2.open_segment_start_index());
1295 }
1296 }
1297}