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