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