Skip to main content

qdrant_edge/wal/
mod.rs

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    /// The segment capacity. Defaults to 32MiB.
27    pub segment_capacity: usize,
28
29    /// The number of segments to create ahead of time, so that appends never
30    /// need to wait on creating a new segment.
31    pub segment_queue_len: usize,
32
33    /// The number of "closed-*" wal files to retain. Defaults to 1.
34    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/// An open segment and its ID.
48#[derive(Debug)]
49struct OpenSegment {
50    pub id: u64,
51    pub segment: Segment,
52}
53
54/// A closed segment, and the associated start and stop indices.
55#[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
66/// A write ahead log.
67///
68/// ### Logging
69///
70/// Wal operations are logged. Metadata operations (open) are logged at `info`
71/// level. Segment operations (create, close, delete) are logged at `debug`
72/// level. Flush operations are logged at `debug` level. Entry operations
73/// (append, truncate) are logged at `trace` level. Long-running or multi-step
74/// operations will log a message at a lower level when beginning, and a final
75/// completion message.
76pub struct Wal {
77    /// The segment currently being appended to.
78    open_segment: OpenSegment,
79    closed_segments: Vec<ClosedSegment>,
80    creator: SegmentCreatorV2,
81
82    /// The number of closed segments to retain.
83    retain_closed: NonZeroUsize,
84
85    /// The directory which contains the write ahead log. Used to hold an open
86    /// file lock for the lifetime of the log.
87    #[allow(dead_code)]
88    dir: File,
89
90    /// The directory path.
91    path: PathBuf,
92
93    /// Tracks the flush status of recently closed segments between user calls
94    /// to `Wal::flush`.
95    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        // Windows workaround. Directories cannot be exclusively held so we create a proxy file
136        // inside the tmp directory which is used for locking. This is done because:
137        // - A Windows directory is not a file unlike in Linux, so we cannot open it with
138        //   `File::open` nor lock it with `try_lock`
139        // - We want this to be auto-deleted together with the `TempDir`
140        #[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        // Holds open segments in the directory.
158        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        // Validate the closed segments. They must be non-overlapping, and contiguous.
170        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                    // TODO: figure out what to do here.
182                    // Current thinking is the previous segment should be truncated.
183                    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        // Validate the open segments.
200        open_segments.sort_by_key(|s| s.id);
201
202        // The latest open segment, may already have segments.
203        let mut open_segment: Option<OpenSegment> = None;
204        // Unused open segments.
205        let mut unused_segments: Vec<OpenSegment> = Vec::new();
206
207        for segment in open_segments {
208            if !segment.segment.is_empty() {
209                // This segment has already been written to. If a previous open
210                // segment has also already been written to, we close it out and
211                // replace it with this new one. This may happen because when a
212                // segment is closed it is renamed, but the directory is not
213                // sync'd, so the operation is not guaranteed to be durable.
214                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 there is an empty closed segment, remove it before adding the new one.
270        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    /// Retrieve the entry with the provided index from the log.
311    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                // Sanity check that the missing index is less than the start of the log.
329                assert_eq!(0, i);
330                None
331            }
332        }
333    }
334
335    /// Truncates entries in the log beginning with `from`.
336    ///
337    /// Entries can be immediately appended to the log once this method returns,
338    /// but the truncated entries are not guaranteed to be removed until the
339    /// wal is flushed.
340    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            // Truncate the open segment completely.
349            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                            // TODO: this should be async
356                            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                            // flushing closed segment after truncation
365                            segment.segment.flush()?;
366                        }
367                        if index + 1 < self.closed_segments.len() {
368                            for segment in self.closed_segments.drain(index + 1..) {
369                                // TODO: this should be async
370                                segment.segment.delete()?;
371                            }
372                        }
373                    }
374                }
375                Err(index) => {
376                    // The truncate index is before the first entry of the wal
377                    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                        // TODO: this should be async
385                        segment.segment.delete()?;
386                    }
387                }
388            }
389        }
390        Ok(())
391    }
392
393    /// Possibly removes entries from the beginning of the log before the given index.
394    ///
395    /// After calling this method, the `first_index` will be between the current
396    /// `first_index` (inclusive), and `until` (exclusive).
397    ///
398    /// This always keeps at least one closed segment.
399    pub fn prefix_truncate(&mut self, until: u64) -> Result<()> {
400        trace!("{self:?}: prefix_truncate until entry {until}");
401
402        // Return early if everything up to `until` has already been truncated
403        if until
404            <= self
405                .closed_segments
406                .first()
407                .map_or(0, |segment| segment.start_index)
408        {
409            return Ok(());
410        }
411
412        // Retain closed segments starting from this index.
413        // If `retain_closed` is 1 (default), only the last closed segment will be preserved.
414        let retain_start_index = self
415            .closed_segments
416            .len()
417            .saturating_sub(self.retain_closed.get());
418
419        // If `until` goes into or above our open segment, delete till preserved closed segment index
420        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        // Delete all closed segments before the one `until` is in
428        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    /// Returns the start index of the open segment.
438    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    /// The index of the first entry.
476    pub fn first_index(&self) -> u64 {
477        self.closed_segments
478            .first()
479            .map_or(0, |segment| segment.start_index)
480    }
481
482    /// The index of the last entry
483    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    /// Remove all entries
489    pub fn clear(&mut self) -> Result<()> {
490        self.truncate(self.first_index())
491    }
492
493    /// Copy all files to the given path directory. directory should exist and be empty
494    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            // if file is locked by any Segment, call copy_to_path on it
524            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                // if file is not locked by any Segment, just copy it
532                fs::copy(entry.path(), &dst_path)?;
533            }
534        }
535        Ok(())
536    }
537
538    /// Set how many segments closed segments to retain on prefix truncation.
539    ///
540    /// Can't be less than 1. If 0 is provided, it will be set to 1.
541    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); // ignore non-files
589    }
590
591    let filename = entry.file_name().into_string().map_err(|_| error())?;
592    match filename.split_once('-') {
593        Some(("tmp", _)) => {
594            // remove temporary files.
595            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), // Ignore other files.
612    }
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    /// Windows has very slow IO
629    #[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        // Create empty wal with initial id 10.
650        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        // Create empty wal with initial id 10.
696        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        // read random entry back to make sure it is correct.
730        {
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    /// Check that entries appended to the write ahead log can be read back.
760    #[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    /// Check that the Wal will read previously written entries.
870    #[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                // Create fake temp file to simulate a crash.
895                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            // Check that all of the entries are present.
916            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 it ever reaches max, it should stay there
1014                if until < entry_count {
1015                    // If `until` is (much) lower we might retain more to satisfy prefix_truncate
1016                    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        // 2 entries should fit in each segment
1051        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; // Based on segment_capacity / entry_size
1096
1097        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            } // for num_entries=10 => 4 closed segments + 1 open segment (each having 2 entries)
1102
1103            let initial_closed_segments = wal.closed_segments.len();
1104
1105            // Do the prefix truncation
1106            wal.prefix_truncate(truncate_index).unwrap();
1107
1108            // Generalized logic
1109            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()); // truncate_index should be intact
1117
1118            for i in 0..expected_trimmed_until {
1119                // before should be trimmed
1120                assert!(wal.entry(i).is_none());
1121            }
1122
1123            for i in expected_trimmed_until..num_entries {
1124                // after should be intact
1125                assert!(wal.entry(i).is_some());
1126            }
1127
1128            assert_eq!(wal.closed_segments.len(), expected_closed_segments);
1129            wal.truncate(0).unwrap(); // Clean up for the next test case
1130        }
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        // 2 entries should fit in each segment
1145        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        // wal is empty
1157        assert!(wal.entry(0).is_none());
1158
1159        // add 10 entries
1160        for i in 0..10 {
1161            assert_eq!(i, wal.append(&&entry[..]).unwrap());
1162        }
1163
1164        // 4 closed segments
1165        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); // 4 x 2 entries
1169        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); // 1 x 2 entries
1174
1175        // first flush to set `flush_offset
1176        wal.flush_open_segment().unwrap();
1177
1178        // content unchanged after flushing
1179        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); // 4 x 2 entries
1183        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); // 1 x 2 entries
1188
1189        wal.truncate(9).unwrap();
1190
1191        assert_eq!(wal.open_segment.segment.len(), 1); // 1 x 2 entries
1192
1193        // truncate half of it
1194        wal.truncate(5).unwrap();
1195
1196        // assert truncation
1197        for i in 5..10 {
1198            assert!(wal.entry(i).is_none());
1199        }
1200
1201        // flush again with `flush_offset` > segment size
1202        wal.flush_open_segment().unwrap();
1203
1204        assert_eq!(wal.num_entries(), 5); // 5 entries removed
1205        assert_eq!(wal.first_index(), 0);
1206        assert_eq!(wal.last_index(), 4);
1207        assert_eq!(wal.closed_segments.len(), 3); // (0, 1) + (2, 3) + (4, empty slot)
1208        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); // empty open segment
1212
1213        // add 5 more entries
1214        for i in 0..5 {
1215            assert_eq!(i + 5, wal.append(&&entry[..]).unwrap());
1216        }
1217
1218        // 5 closed segments
1219        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); // 1,2
1224        assert_eq!(wal.closed_segments[1].segment.len(), 2); // 3
1225        assert_eq!(wal.closed_segments[2].segment.len(), 1); // 4 empty slot due to truncation
1226        assert_eq!(wal.closed_segments[3].segment.len(), 2); // 5, 6
1227        assert_eq!(wal.closed_segments[4].segment.len(), 2); // 7, 8
1228        assert_eq!(wal.open_segment.segment.len(), 1); // 9
1229
1230        eprintln!("wal: {wal:?}");
1231        eprintln!("wal open: {:?}", wal.open_segment);
1232        eprintln!("wal closed: {:?}", wal.closed_segments);
1233
1234        // test persistence
1235        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); // previously half truncated
1244        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    /// Tests that two Wal instances can not coexist for the same directory.
1250    #[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}