Skip to main content

qdrant_edge/wal/
mod.rs

1use std::cmp::Ordering;
2use std::io::{Error, ErrorKind, Result};
3use std::num::NonZeroUsize;
4use std::path::{Path, PathBuf};
5use std::str::FromStr;
6use std::{fmt, mem, ops, result, thread};
7
8use fs_err as fs;
9use fs_err::File;
10use log::{debug, info, trace};
11pub use segment::{Entry, Segment};
12
13use crate::wal::segment_creator::SegmentCreatorV2;
14
15mod mmap_view_sync;
16mod segment;
17mod segment_creator;
18pub mod test_utils;
19
20#[cfg(test)]
21mod test_segment_recovery;
22
23#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
24pub struct WalOptions {
25    /// The segment capacity. Defaults to 32MiB.
26    pub segment_capacity: usize,
27
28    /// The number of segments to create ahead of time, so that appends never
29    /// need to wait on creating a new segment.
30    pub segment_queue_len: usize,
31
32    /// The number of "closed-*" wal files to retain. Defaults to 1.
33    pub retain_closed: NonZeroUsize,
34}
35
36impl Default for WalOptions {
37    fn default() -> WalOptions {
38        WalOptions {
39            segment_capacity: 32 * 1024 * 1024,
40            segment_queue_len: 0,
41            retain_closed: NonZeroUsize::new(1).unwrap(),
42        }
43    }
44}
45
46/// An open segment and its ID.
47#[derive(Debug)]
48struct OpenSegment {
49    pub id: u64,
50    pub segment: Segment,
51}
52
53/// A closed segment, and the associated start and stop indices.
54#[derive(Debug)]
55struct ClosedSegment {
56    pub start_index: u64,
57    pub segment: Segment,
58}
59
60enum WalSegment {
61    Open(OpenSegment),
62    Closed(ClosedSegment),
63}
64
65/// A write ahead log.
66///
67/// ### Logging
68///
69/// Wal operations are logged. Metadata operations (open) are logged at `info`
70/// level. Segment operations (create, close, delete) are logged at `debug`
71/// level. Flush operations are logged at `debug` level. Entry operations
72/// (append, truncate) are logged at `trace` level. Long-running or multi-step
73/// operations will log a message at a lower level when beginning, and a final
74/// completion message.
75pub struct Wal {
76    /// The segment currently being appended to.
77    open_segment: OpenSegment,
78    closed_segments: Vec<ClosedSegment>,
79    creator: SegmentCreatorV2,
80
81    /// The number of closed segments to retain.
82    retain_closed: NonZeroUsize,
83
84    /// The directory which contains the write ahead log. Used to hold an open
85    /// file lock for the lifetime of the log.
86    #[expect(dead_code)]
87    dir: File,
88
89    /// The directory path.
90    path: PathBuf,
91
92    /// Tracks the flush status of recently closed segments between user calls
93    /// to `Wal::flush`.
94    flush: Option<thread::JoinHandle<Result<()>>>,
95}
96
97impl Wal {
98    pub fn open<P>(path: P) -> Result<Wal>
99    where
100        P: AsRef<Path>,
101    {
102        Wal::with_options(path, &WalOptions::default())
103    }
104
105    pub fn generate_empty_wal_starting_at_index(
106        path: impl Into<PathBuf>,
107        options: &WalOptions,
108        index: u64,
109    ) -> Result<()> {
110        let open_id = 0;
111        let mut path_buf = path.into();
112        path_buf.push(format!("open-{open_id}"));
113        let segment = OpenSegment {
114            id: index + 1,
115            segment: Segment::create(&path_buf, options.segment_capacity)?,
116        };
117
118        let mut close_segment = close_segment(segment, index + 1)?;
119
120        close_segment.segment.flush()
121    }
122
123    pub fn with_options<P>(path: P, options: &WalOptions) -> Result<Wal>
124    where
125        P: AsRef<Path>,
126    {
127        debug!("Wal {{ path: {:?} }}: opening", path.as_ref());
128
129        #[cfg(not(target_os = "windows"))]
130        let path = path.as_ref().to_path_buf();
131        #[cfg(not(target_os = "windows"))]
132        let dir = File::open(&path)?;
133
134        // Windows workaround. Directories cannot be exclusively held so we create a proxy file
135        // inside the tmp directory which is used for locking. This is done because:
136        // - A Windows directory is not a file unlike in Linux, so we cannot open it with
137        //   `File::open` nor lock it with `try_lock`
138        // - We want this to be auto-deleted together with the `TempDir`
139        #[cfg(target_os = "windows")]
140        let mut path = path.as_ref().to_path_buf();
141        #[cfg(target_os = "windows")]
142        let dir = {
143            path.push(".wal");
144            let dir = File::options()
145                .create(true)
146                .read(true)
147                .write(true)
148                .truncate(true)
149                .open(&path)?;
150            path.pop();
151            dir
152        };
153
154        // Use `fs4`'s `flock(2)`-based lock rather than `dir.try_lock()`.
155        //
156        // `dir.try_lock()` resolves to the inherent `fs_err`/`std` `File::try_lock`,
157        // which is gated to a fixed list of targets in stdlib and returns
158        // `ErrorKind::Unsupported` ("try_lock() not supported") on others — notably
159        // Android. `fs4::FileExt::try_lock` issues a direct `flock(LOCK_EX | LOCK_NB)`
160        // syscall, which Android supports. We call it via UFCS on the underlying
161        // `std::fs::File` because the trait method collides with the inherent one
162        // (which would otherwise win method resolution).
163        fs4::FileExt::try_lock(dir.file())?;
164
165        // Holds open segments in the directory.
166        let mut open_segments: Vec<OpenSegment> = Vec::new();
167        let mut closed_segments: Vec<ClosedSegment> = Vec::new();
168
169        for entry in fs::read_dir(&path)? {
170            match open_dir_entry(entry?)? {
171                Some(WalSegment::Open(open_segment)) => open_segments.push(open_segment),
172                Some(WalSegment::Closed(closed_segment)) => closed_segments.push(closed_segment),
173                None => {}
174            }
175        }
176
177        // Validate the closed segments. They must be non-overlapping, and contiguous.
178        closed_segments.sort_by_key(|s| s.start_index);
179        let mut next_start_index = closed_segments
180            .first()
181            .map_or(0, |segment| segment.start_index);
182        for &ClosedSegment {
183            start_index,
184            ref segment,
185        } in &closed_segments
186        {
187            match start_index.cmp(&next_start_index) {
188                Ordering::Less => {
189                    return Err(Error::new(
190                        ErrorKind::InvalidData,
191                        format!(
192                            "overlapping segments: segment at {start_index} overlaps with already-covered range up to {next_start_index}"
193                        ),
194                    ));
195                }
196                Ordering::Equal => {
197                    next_start_index = start_index + segment.len() as u64;
198                }
199                Ordering::Greater => {
200                    return Err(Error::new(
201                        ErrorKind::InvalidData,
202                        format!(
203                            "missing segment(s) containing wal entries {next_start_index} to {start_index}"
204                        ),
205                    ));
206                }
207            }
208        }
209
210        // Validate the open segments.
211        open_segments.sort_by_key(|s| s.id);
212
213        // The latest open segment, may already have segments.
214        let mut open_segment: Option<OpenSegment> = None;
215        // Unused open segments.
216        let mut unused_segments: Vec<OpenSegment> = Vec::new();
217
218        for segment in open_segments {
219            if !segment.segment.is_empty() {
220                // This segment has already been written to. If a previous open
221                // segment has also already been written to, we close it out and
222                // replace it with this new one. This may happen because when a
223                // segment is closed it is renamed, but the directory is not
224                // sync'd, so the operation is not guaranteed to be durable.
225                let stranded_segment = open_segment.take();
226                open_segment = Some(segment);
227                if let Some(segment) = stranded_segment {
228                    let closed_segment = close_segment(segment, next_start_index)?;
229                    next_start_index += closed_segment.segment.len() as u64;
230                    closed_segments.push(closed_segment);
231                }
232            } else if open_segment.is_none() {
233                open_segment = Some(segment);
234            } else {
235                unused_segments.push(segment);
236            }
237        }
238
239        let mut creator = SegmentCreatorV2::new(
240            &path,
241            open_segment.as_ref(),
242            unused_segments,
243            options.segment_capacity,
244            options.segment_queue_len,
245        );
246
247        let open_segment = match open_segment {
248            Some(segment) => segment,
249            None => creator.next()?,
250        };
251
252        let wal = Wal {
253            open_segment,
254            closed_segments,
255            retain_closed: options.retain_closed,
256            creator,
257            dir,
258            path,
259            flush: None,
260        };
261        info!("{wal:?}: opened");
262        Ok(wal)
263    }
264
265    fn retire_open_segment(&mut self) -> Result<()> {
266        trace!("{self:?}: retiring open segment");
267        let mut segment = self.creator.next()?;
268        mem::swap(&mut self.open_segment, &mut segment);
269
270        if let Some(flush) = self.flush.take() {
271            flush
272                .join()
273                .map_err(|err| Error::other(format!("wal flush thread panicked: {err:?}")))??;
274        };
275
276        self.flush = Some(segment.segment.flush_async());
277
278        let start_index = self.open_segment_start_index();
279
280        // If there is an empty closed segment, remove it before adding the new one.
281        if let Some(empty_segment) = self
282            .closed_segments
283            .pop_if(|last_closed| last_closed.segment.is_empty())
284        {
285            empty_segment.segment.delete()?;
286        }
287
288        self.closed_segments
289            .push(close_segment(segment, start_index)?);
290        debug!("{self:?}: open segment retired. start_index: {start_index}");
291        Ok(())
292    }
293
294    pub fn append<T>(&mut self, entry: &T) -> Result<u64>
295    where
296        T: ops::Deref<Target = [u8]>,
297    {
298        trace!("{:?}: appending entry of length {}", self, entry.len());
299        if !self.open_segment.segment.sufficient_capacity(entry.len()) {
300            if !self.open_segment.segment.is_empty() {
301                self.retire_open_segment()?;
302            }
303            self.open_segment.segment.ensure_capacity(entry.len())?;
304        }
305
306        Ok(self.open_segment_start_index()
307            + self.open_segment.segment.append(entry).unwrap() as u64)
308    }
309
310    pub fn flush_open_segment(&mut self) -> Result<()> {
311        trace!("{self:?}: flushing open segments");
312        self.open_segment.segment.flush()?;
313        Ok(())
314    }
315
316    pub fn flush_open_segment_async(&mut self) -> thread::JoinHandle<Result<()>> {
317        trace!("{self:?}: flushing open segments");
318        self.open_segment.segment.flush_async()
319    }
320
321    /// Retrieve the entry with the provided index from the log.
322    pub fn entry(&self, index: u64) -> Option<Entry> {
323        let open_start_index = self.open_segment_start_index();
324        if index >= open_start_index {
325            return self
326                .open_segment
327                .segment
328                .entry((index - open_start_index) as usize);
329        }
330
331        match self.find_closed_segment(index) {
332            Ok(segment_index) => {
333                let segment = &self.closed_segments[segment_index];
334                segment
335                    .segment
336                    .entry((index - segment.start_index) as usize)
337            }
338            Err(i) => {
339                // Sanity check that the missing index is less than the start of the log.
340                assert_eq!(0, i);
341                None
342            }
343        }
344    }
345
346    /// Truncates entries in the log beginning with `from`.
347    ///
348    /// Entries can be immediately appended to the log once this method returns,
349    /// but the truncated entries are not guaranteed to be removed until the
350    /// wal is flushed.
351    pub fn truncate(&mut self, from: u64) -> Result<()> {
352        trace!("{self:?}: truncate from entry {from}");
353        let open_start_index = self.open_segment_start_index();
354        if from >= open_start_index {
355            self.open_segment
356                .segment
357                .truncate((from - open_start_index) as usize);
358        } else {
359            // Truncate the open segment completely.
360            self.open_segment.segment.truncate(0);
361
362            match self.find_closed_segment(from) {
363                Ok(index) => {
364                    if from == self.closed_segments[index].start_index {
365                        for segment in self.closed_segments.drain(index..) {
366                            // TODO: this should be async
367                            segment.segment.delete()?;
368                        }
369                    } else {
370                        {
371                            let segment = &mut self.closed_segments[index];
372                            segment
373                                .segment
374                                .truncate((from - segment.start_index) as usize);
375                            // flushing closed segment after truncation
376                            segment.segment.flush()?;
377                        }
378                        if index + 1 < self.closed_segments.len() {
379                            for segment in self.closed_segments.drain(index + 1..) {
380                                // TODO: this should be async
381                                segment.segment.delete()?;
382                            }
383                        }
384                    }
385                }
386                Err(index) => {
387                    // The truncate index is before the first entry of the wal
388                    assert!(
389                        from <= self
390                            .closed_segments
391                            .get(index)
392                            .map_or(0, |segment| segment.start_index)
393                    );
394                    for segment in self.closed_segments.drain(..) {
395                        // TODO: this should be async
396                        segment.segment.delete()?;
397                    }
398                }
399            }
400        }
401        Ok(())
402    }
403
404    /// Possibly removes entries from the beginning of the log before the given index.
405    ///
406    /// After calling this method, the `first_index` will be between the current
407    /// `first_index` (inclusive), and `until` (exclusive).
408    ///
409    /// This always keeps at least one closed segment.
410    pub fn prefix_truncate(&mut self, until: u64) -> Result<()> {
411        trace!("{self:?}: prefix_truncate until entry {until}");
412
413        // Return early if everything up to `until` has already been truncated
414        if until
415            <= self
416                .closed_segments
417                .first()
418                .map_or(0, |segment| segment.start_index)
419        {
420            return Ok(());
421        }
422
423        // Retain closed segments starting from this index.
424        // If `retain_closed` is 1 (default), only the last closed segment will be preserved.
425        let retain_start_index = self
426            .closed_segments
427            .len()
428            .saturating_sub(self.retain_closed.get());
429
430        // If `until` goes into or above our open segment, delete till preserved closed segment index
431        if until >= self.open_segment_start_index() {
432            for segment in self.closed_segments.drain(..retain_start_index) {
433                segment.segment.delete()?
434            }
435            return Ok(());
436        }
437
438        // Delete all closed segments before the one `until` is in
439        let index = self.find_closed_segment(until).unwrap();
440        let truncate_until_index = index.min(retain_start_index);
441        trace!("{self:?}: prefix truncating until closed segment {truncate_until_index}");
442        for segment in self.closed_segments.drain(..truncate_until_index) {
443            segment.segment.delete()?
444        }
445        Ok(())
446    }
447
448    /// Returns the start index of the open segment.
449    fn open_segment_start_index(&self) -> u64 {
450        self.closed_segments
451            .last()
452            .map_or(0, |segment: &ClosedSegment| {
453                segment.start_index + segment.segment.len() as u64
454            })
455    }
456
457    fn find_closed_segment(&self, index: u64) -> result::Result<usize, usize> {
458        self.closed_segments.binary_search_by(|segment| {
459            if index < segment.start_index {
460                Ordering::Greater
461            } else if index >= segment.start_index + segment.segment.len() as u64 {
462                Ordering::Less
463            } else {
464                Ordering::Equal
465            }
466        })
467    }
468
469    pub fn path(&self) -> &Path {
470        &self.path
471    }
472
473    pub fn num_segments(&self) -> usize {
474        self.closed_segments.len() + 1
475    }
476
477    pub fn num_entries(&self) -> u64 {
478        self.open_segment_start_index()
479            - self
480                .closed_segments
481                .first()
482                .map_or(0, |segment| segment.start_index)
483            + self.open_segment.segment.len() as u64
484    }
485
486    /// The index of the first entry.
487    pub fn first_index(&self) -> u64 {
488        self.closed_segments
489            .first()
490            .map_or(0, |segment| segment.start_index)
491    }
492
493    /// The index of the last entry
494    pub fn last_index(&self) -> u64 {
495        let num_entries = self.num_entries();
496        self.first_index() + num_entries.saturating_sub(1)
497    }
498
499    /// Remove all entries
500    pub fn clear(&mut self) -> Result<()> {
501        self.truncate(self.first_index())
502    }
503
504    /// Set how many segments closed segments to retain on prefix truncation.
505    ///
506    /// Can't be less than 1. If 0 is provided, it will be set to 1.
507    pub fn set_retention(&mut self, retain_closed: usize) {
508        self.retain_closed = NonZeroUsize::new(retain_closed.max(1)).unwrap();
509    }
510}
511
512impl fmt::Debug for Wal {
513    fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
514        let Self {
515            open_segment,
516            closed_segments,
517            path,
518            creator: _,
519            retain_closed: _,
520            dir: _,
521            flush: _,
522        } = self;
523        let start_index = closed_segments
524            .first()
525            .map_or(0, |segment| segment.start_index);
526        let end_index = self.open_segment_start_index() + open_segment.segment.len() as u64;
527        f.debug_struct("Wal")
528            .field("path", path)
529            .field("segment-count", &(closed_segments.len() + 1))
530            .field("entries", &format_args!("[{start_index}, {end_index})"))
531            .finish_non_exhaustive()
532    }
533}
534
535fn close_segment(mut segment: OpenSegment, start_index: u64) -> Result<ClosedSegment> {
536    let new_path = segment
537        .segment
538        .path()
539        .with_file_name(format!("closed-{start_index}"));
540    segment.segment.rename(new_path)?;
541    segment.segment.close();
542    Ok(ClosedSegment {
543        start_index,
544        segment: segment.segment,
545    })
546}
547
548fn open_dir_entry(entry: fs::DirEntry) -> Result<Option<WalSegment>> {
549    let metadata = entry.metadata()?;
550
551    let error = || {
552        Error::new(
553            ErrorKind::InvalidData,
554            format!("unexpected entry in wal directory: {:?}", entry.path()),
555        )
556    };
557
558    if !metadata.is_file() {
559        return Ok(None); // ignore non-files
560    }
561
562    let filename = entry.file_name().into_string().map_err(|_| error())?;
563    match filename.split_once('-') {
564        Some(("tmp", _)) => {
565            // remove temporary files.
566            fs::remove_file(entry.path())?;
567            Ok(None)
568        }
569        Some(("open", id)) => {
570            let id = u64::from_str(id).map_err(|_| error())?;
571            let segment = Segment::open(entry.path())?;
572            Ok(Some(WalSegment::Open(OpenSegment { id, segment })))
573        }
574        Some(("closed", start)) => {
575            let start = u64::from_str(start).map_err(|_| error())?;
576            let segment = Segment::open(entry.path())?;
577            Ok(Some(WalSegment::Closed(ClosedSegment {
578                start_index: start,
579                segment,
580            })))
581        }
582        _ => Ok(None), // Ignore other files.
583    }
584}
585
586#[cfg(test)]
587mod test {
588    use std::io::{ErrorKind, Write};
589    use std::num::{NonZeroU8, NonZeroUsize};
590
591    use fs_err as fs;
592    use log::trace;
593    use quickcheck::{QuickCheck, TestResult};
594    use tempfile::Builder;
595
596    use super::{Wal, WalOptions};
597    use crate::wal::test_utils::EntryGenerator;
598
599    /// Windows has very slow IO
600    #[cfg(target_os = "windows")]
601    const QC_TESTS: u64 = 3;
602
603    #[cfg(not(target_os = "windows"))]
604    const QC_TESTS: u64 = 50;
605
606    fn init_logger() {
607        let _ = env_logger::builder().is_test(true).try_init();
608    }
609
610    #[test]
611    fn test_generate_empty_wal() {
612        init_logger();
613        let dir = Builder::new().prefix("wal").tempdir().unwrap();
614        let options = WalOptions {
615            segment_capacity: 80,
616            segment_queue_len: 3,
617            retain_closed: NonZeroUsize::new(1).unwrap(),
618        };
619
620        // Create empty wal with initial id 10.
621        let init_offset = 10;
622        Wal::generate_empty_wal_starting_at_index(dir.path(), &options, init_offset).unwrap();
623
624        let mut wal = Wal::with_options(dir.path(), &options).unwrap();
625
626        let first_index = wal.first_index();
627        let last_index = wal.last_index();
628        let num_entries = wal.num_entries();
629
630        assert!(first_index <= last_index);
631        assert_eq!(num_entries, 0);
632
633        let next_entry: Vec<u8> = vec![1, 2, 3];
634        let op = wal.append(&next_entry).unwrap();
635
636        assert!(op > init_offset);
637
638        let first_index = wal.first_index();
639        let last_index = wal.last_index();
640        let num_entries = wal.num_entries();
641
642        assert!(first_index <= last_index);
643        assert_eq!(num_entries, 1);
644
645        wal.append(&next_entry).unwrap();
646        wal.append(&next_entry).unwrap();
647
648        let first_index = wal.first_index();
649        let last_index = wal.last_index();
650        let num_entries = wal.num_entries();
651
652        assert!(first_index <= last_index);
653        assert_eq!(num_entries, 3);
654    }
655
656    #[test]
657    fn test_create_empty_wal_with_initial_id() {
658        init_logger();
659        let dir = Builder::new().prefix("wal").tempdir().unwrap();
660        let options = WalOptions {
661            segment_capacity: 80,
662            segment_queue_len: 3,
663            retain_closed: NonZeroUsize::new(1).unwrap(),
664        };
665
666        // Create empty wal with initial id 10.
667        let init_offset = 10;
668        Wal::generate_empty_wal_starting_at_index(dir.path(), &options, init_offset).unwrap();
669
670        let mut wal = Wal::with_options(dir.path(), &options).unwrap();
671
672        let last_index = wal.last_index();
673
674        assert!(last_index > init_offset);
675
676        assert_eq!(wal.num_entries(), 0);
677
678        let next_entry: Vec<u8> = vec![1, 2, 3];
679
680        wal.append(&next_entry).unwrap();
681
682        let last_index = wal.last_index();
683        assert_eq!(last_index, init_offset + 1);
684
685        assert_eq!(wal.num_entries(), 1);
686
687        let entry_count = 50;
688
689        let entries = EntryGenerator::new().take(entry_count).collect::<Vec<_>>();
690
691        for entry in &entries {
692            wal.append(entry).unwrap();
693        }
694
695        let last_index = wal.last_index();
696        assert_eq!(last_index, init_offset + 1 + entry_count as u64);
697
698        assert_eq!(wal.num_entries(), 1 + entry_count as u64);
699
700        // read random entry back to make sure it is correct.
701        {
702            let entry_index = init_offset + 1;
703            let entry = wal.entry(entry_index).unwrap();
704            assert_eq!(next_entry[..], entry[..]);
705
706            let entry_index = init_offset + 1 + entry_count as u64;
707            let entry = wal.entry(entry_index).unwrap();
708            assert_eq!(entries[entry_count - 1][..], entry[..]);
709
710            let entry_index = init_offset + 1 + 10;
711            let entry = wal.entry(entry_index).unwrap();
712            assert_eq!(entries[9][..], entry[..]);
713        }
714
715        wal.prefix_truncate(init_offset).unwrap();
716
717        assert_eq!(wal.num_entries(), entry_count as u64 + 1);
718
719        wal.prefix_truncate(init_offset + 20).unwrap();
720
721        assert!(wal.num_entries() < entry_count as u64 + 1);
722
723        let truncate_index = init_offset + 30;
724        wal.truncate(truncate_index).unwrap();
725
726        let last_index = wal.last_index();
727        assert_eq!(last_index, truncate_index - 1);
728    }
729
730    /// Check that entries appended to the write ahead log can be read back.
731    #[test]
732    fn check_wal() {
733        init_logger();
734        fn wal(entry_count: u8) -> TestResult {
735            let dir = Builder::new().prefix("wal").tempdir().unwrap();
736            let mut wal = Wal::with_options(
737                dir.path(),
738                &WalOptions {
739                    segment_capacity: 80,
740                    segment_queue_len: 3,
741                    retain_closed: NonZeroUsize::new(1).unwrap(),
742                },
743            )
744            .unwrap();
745            let entries = EntryGenerator::new()
746                .take(entry_count as usize)
747                .collect::<Vec<_>>();
748
749            for entry in &entries {
750                wal.append(entry).unwrap();
751            }
752
753            for (index, expected) in entries.iter().enumerate() {
754                match wal.entry(index as u64) {
755                    Some(ref entry) if entry[..] != expected[..] => return TestResult::failed(),
756                    None => return TestResult::failed(),
757                    _ => (),
758                }
759            }
760            TestResult::passed()
761        }
762
763        QuickCheck::new()
764            .tests(QC_TESTS)
765            .quickcheck(wal as fn(u8) -> TestResult);
766    }
767
768    #[test]
769    fn check_last_index() {
770        init_logger();
771        fn check(entry_count: u8) -> TestResult {
772            let dir = Builder::new().prefix("wal").tempdir().unwrap();
773            let mut wal = Wal::with_options(
774                dir.path(),
775                &WalOptions {
776                    segment_capacity: 80,
777                    segment_queue_len: 3,
778                    retain_closed: NonZeroUsize::new(1).unwrap(),
779                },
780            )
781            .unwrap();
782            let entries = EntryGenerator::new()
783                .take(entry_count as usize)
784                .collect::<Vec<_>>();
785
786            for entry in &entries {
787                wal.append(entry).unwrap();
788            }
789            if entries.is_empty() {
790                assert_eq!(wal.last_index(), 0);
791            } else {
792                assert_eq!(wal.last_index(), entries.len() as u64 - 1);
793            }
794
795            let last_index = wal.last_index();
796            if wal.entry(last_index).is_none() && wal.num_entries() != 0 {
797                return TestResult::failed();
798            }
799            if wal.entry(last_index + 1).is_some() {
800                return TestResult::failed();
801            }
802            TestResult::passed()
803        }
804
805        QuickCheck::new()
806            .tests(QC_TESTS)
807            .quickcheck(check as fn(u8) -> TestResult)
808    }
809
810    #[test]
811    fn check_clear() {
812        init_logger();
813        fn check(entry_count: u8) -> TestResult {
814            let dir = Builder::new().prefix("wal").tempdir().unwrap();
815            let mut wal = Wal::with_options(
816                dir.path(),
817                &WalOptions {
818                    segment_capacity: 80,
819                    segment_queue_len: 3,
820                    retain_closed: NonZeroUsize::new(1).unwrap(),
821                },
822            )
823            .unwrap();
824            let entries = EntryGenerator::new()
825                .take(entry_count as usize)
826                .collect::<Vec<_>>();
827
828            for entry in &entries {
829                wal.append(entry).unwrap();
830            }
831            wal.clear().unwrap();
832            TestResult::from_bool(wal.num_entries() == 0)
833        }
834
835        QuickCheck::new()
836            .tests(QC_TESTS)
837            .quickcheck(check as fn(u8) -> TestResult)
838    }
839
840    /// Check that the Wal will read previously written entries.
841    #[test]
842    fn check_reopen() {
843        init_logger();
844        fn wal(entry_count: u8) -> TestResult {
845            let entries = EntryGenerator::new()
846                .take(entry_count as usize)
847                .collect::<Vec<_>>();
848            let dir = Builder::new().prefix("wal").tempdir().unwrap();
849            {
850                let mut wal = Wal::with_options(
851                    dir.path(),
852                    &WalOptions {
853                        segment_capacity: 80,
854                        segment_queue_len: 3,
855                        retain_closed: NonZeroUsize::new(1).unwrap(),
856                    },
857                )
858                .unwrap();
859                for entry in &entries {
860                    let _ = wal.append(entry);
861                }
862            }
863
864            {
865                // Create fake temp file to simulate a crash.
866                let mut file = fs::OpenOptions::new()
867                    .read(true)
868                    .write(true)
869                    .create(true)
870                    .truncate(true)
871                    .open(dir.path().join("tmp-open-123"))
872                    .unwrap();
873
874                let _ = file.write(b"123").unwrap();
875            }
876
877            let wal = Wal::with_options(
878                dir.path(),
879                &WalOptions {
880                    segment_capacity: 80,
881                    segment_queue_len: 3,
882                    retain_closed: NonZeroUsize::new(1).unwrap(),
883                },
884            )
885            .unwrap();
886            // Check that all of the entries are present.
887            for (index, expected) in entries.iter().enumerate() {
888                match wal.entry(index as u64) {
889                    Some(ref entry) if entry[..] != expected[..] => return TestResult::failed(),
890                    None => return TestResult::failed(),
891                    _ => (),
892                }
893            }
894            TestResult::passed()
895        }
896
897        QuickCheck::new()
898            .tests(QC_TESTS)
899            .quickcheck(wal as fn(u8) -> TestResult);
900    }
901
902    #[test]
903    fn check_truncate() {
904        init_logger();
905        fn truncate(entry_count: u8, truncate: u8) -> TestResult {
906            if truncate > entry_count {
907                return TestResult::discard();
908            }
909            let dir = Builder::new().prefix("wal").tempdir().unwrap();
910            let mut wal = Wal::with_options(
911                dir.path(),
912                &WalOptions {
913                    segment_capacity: 80,
914                    segment_queue_len: 3,
915                    retain_closed: NonZeroUsize::new(1).unwrap(),
916                },
917            )
918            .unwrap();
919            let entries = EntryGenerator::new()
920                .take(entry_count as usize)
921                .collect::<Vec<_>>();
922
923            for entry in &entries {
924                if let Err(error) = wal.append(entry) {
925                    return TestResult::error(error.to_string());
926                }
927            }
928
929            wal.truncate(u64::from(truncate)).unwrap();
930
931            for (index, expected) in entries.iter().take(truncate as usize).enumerate() {
932                match wal.entry(index as u64) {
933                    Some(ref entry) if entry[..] != expected[..] => return TestResult::failed(),
934                    None => return TestResult::failed(),
935                    _ => (),
936                }
937            }
938
939            TestResult::from_bool(wal.entry(u64::from(truncate)).is_none())
940        }
941
942        QuickCheck::new()
943            .tests(QC_TESTS)
944            .quickcheck(truncate as fn(u8, u8) -> TestResult);
945    }
946
947    #[test]
948    fn check_prefix_truncate() {
949        init_logger();
950        fn prefix_truncate(entry_count: u8, until: u8, retain_closed: NonZeroU8) -> TestResult {
951            trace!(
952                "prefix truncate; entry_count: {entry_count}, until: {until}, retain_closed: {retain_closed}",
953            );
954            if until > entry_count {
955                return TestResult::discard();
956            }
957            let dir = Builder::new().prefix("wal").tempdir().unwrap();
958            let mut wal = Wal::with_options(
959                dir.path(),
960                &WalOptions {
961                    segment_capacity: 80,
962                    segment_queue_len: 3,
963                    retain_closed: NonZeroUsize::from(retain_closed),
964                },
965            )
966            .unwrap();
967            let entries = EntryGenerator::new()
968                .take(entry_count as usize)
969                .collect::<Vec<_>>();
970
971            let mut has_ever_reached_max = false;
972            let retain_closed = retain_closed.get() as usize;
973
974            for entry in &entries {
975                wal.append(entry).unwrap();
976                if wal.closed_segments.len() >= retain_closed {
977                    has_ever_reached_max = true;
978                }
979            }
980
981            wal.prefix_truncate(u64::from(until)).unwrap();
982
983            let retained = if has_ever_reached_max {
984                // If it ever reaches max, it should stay there
985                if until < entry_count {
986                    // If `until` is (much) lower we might retain more to satisfy prefix_truncate
987                    wal.closed_segments.len() >= retain_closed
988                } else {
989                    wal.closed_segments.len() == retain_closed
990                }
991            } else {
992                wal.closed_segments.len() < retain_closed
993            };
994
995            let num_entries = wal.num_entries() as u8;
996            TestResult::from_bool(
997                num_entries <= entry_count && num_entries >= entry_count - until && retained,
998            )
999        }
1000        QuickCheck::new()
1001            .tests(QC_TESTS)
1002            .quickcheck(prefix_truncate as fn(u8, u8, NonZeroU8) -> TestResult);
1003    }
1004
1005    #[test]
1006    fn test_append() {
1007        init_logger();
1008        let dir = Builder::new().prefix("wal").tempdir().unwrap();
1009        let mut wal = Wal::open(dir.path()).unwrap();
1010
1011        let entry: &[u8] = &[42u8; 4096];
1012        for _ in 1..10 {
1013            wal.append(&entry).unwrap();
1014        }
1015    }
1016
1017    #[test]
1018    fn test_truncate() {
1019        init_logger();
1020        let dir = Builder::new().prefix("wal").tempdir().unwrap();
1021        // 2 entries should fit in each segment
1022        let mut wal = Wal::with_options(
1023            dir.path(),
1024            &WalOptions {
1025                segment_capacity: 4096,
1026                segment_queue_len: 3,
1027                retain_closed: NonZeroUsize::new(1).unwrap(),
1028            },
1029        )
1030        .unwrap();
1031
1032        let entry: [u8; 2000] = [42u8; 2000];
1033
1034        for truncate_index in 0..10 {
1035            assert!(wal.entry(0).is_none());
1036            for i in 0..10 {
1037                assert_eq!(i, wal.append(&&entry[..]).unwrap());
1038            }
1039
1040            wal.truncate(truncate_index).unwrap();
1041
1042            assert!(wal.entry(truncate_index).is_none());
1043
1044            if truncate_index > 0 {
1045                assert!(wal.entry(truncate_index - 1).is_some());
1046            }
1047            wal.truncate(0).unwrap();
1048        }
1049    }
1050
1051    fn run_test_with_retain_closed(retain_closed: usize) {
1052        init_logger();
1053        let num_entries = 10;
1054        let dir = Builder::new().prefix("wal").tempdir().unwrap();
1055        let entry: [u8; 2000] = [42u8; 2000];
1056        let mut wal = Wal::with_options(
1057            dir.path(),
1058            &WalOptions {
1059                segment_capacity: 4096,
1060                segment_queue_len: 3,
1061                retain_closed: NonZeroUsize::new(retain_closed).unwrap(),
1062            },
1063        )
1064        .unwrap();
1065
1066        let entries_per_segment = 2; // Based on segment_capacity / entry_size
1067
1068        for truncate_index in 0..num_entries {
1069            assert!(wal.entry(0).is_none());
1070            for i in 0..num_entries {
1071                assert_eq!(i, wal.append(&&entry[..]).unwrap());
1072            } // for num_entries=10 => 4 closed segments + 1 open segment (each having 2 entries)
1073
1074            let initial_closed_segments = wal.closed_segments.len();
1075
1076            // Do the prefix truncation
1077            wal.prefix_truncate(truncate_index).unwrap();
1078
1079            // Generalized logic
1080            let segments_that_can_be_truncated = truncate_index / entries_per_segment;
1081            let expected_closed_segments = (initial_closed_segments
1082                - segments_that_can_be_truncated as usize)
1083                .max(retain_closed);
1084            let expected_trimmed_until =
1085                (initial_closed_segments - expected_closed_segments) as u64 * entries_per_segment;
1086
1087            assert!(wal.entry(truncate_index).is_some()); // truncate_index should be intact
1088
1089            for i in 0..expected_trimmed_until {
1090                // before should be trimmed
1091                assert!(wal.entry(i).is_none());
1092            }
1093
1094            for i in expected_trimmed_until..num_entries {
1095                // after should be intact
1096                assert!(wal.entry(i).is_some());
1097            }
1098
1099            assert_eq!(wal.closed_segments.len(), expected_closed_segments);
1100            wal.truncate(0).unwrap(); // Clean up for the next test case
1101        }
1102    }
1103
1104    #[test]
1105    fn test_prefix_truncate_parametric() {
1106        run_test_with_retain_closed(1);
1107        run_test_with_retain_closed(2);
1108        run_test_with_retain_closed(3);
1109    }
1110
1111    #[test]
1112    fn test_truncate_flush() {
1113        init_logger();
1114        let dir = Builder::new().prefix("wal").tempdir().unwrap();
1115        // 2 entries should fit in each segment
1116        let mut wal = Wal::with_options(
1117            dir.path(),
1118            &WalOptions {
1119                segment_capacity: 4096,
1120                segment_queue_len: 3,
1121                retain_closed: NonZeroUsize::new(1).unwrap(),
1122            },
1123        )
1124        .unwrap();
1125
1126        let entry: [u8; 2000] = [42u8; 2000];
1127        // wal is empty
1128        assert!(wal.entry(0).is_none());
1129
1130        // add 10 entries
1131        for i in 0..10 {
1132            assert_eq!(i, wal.append(&&entry[..]).unwrap());
1133        }
1134
1135        // 4 closed segments
1136        assert_eq!(wal.num_entries(), 10);
1137        assert_eq!(wal.first_index(), 0);
1138        assert_eq!(wal.last_index(), 9);
1139        assert_eq!(wal.closed_segments.len(), 4); // 4 x 2 entries
1140        assert_eq!(wal.closed_segments[0].segment.len(), 2);
1141        assert_eq!(wal.closed_segments[1].segment.len(), 2);
1142        assert_eq!(wal.closed_segments[2].segment.len(), 2);
1143        assert_eq!(wal.closed_segments[3].segment.len(), 2);
1144        assert_eq!(wal.open_segment.segment.len(), 2); // 1 x 2 entries
1145
1146        // first flush to set `flush_offset
1147        wal.flush_open_segment().unwrap();
1148
1149        // content unchanged after flushing
1150        assert_eq!(wal.num_entries(), 10);
1151        assert_eq!(wal.first_index(), 0);
1152        assert_eq!(wal.last_index(), 9);
1153        assert_eq!(wal.closed_segments.len(), 4); // 4 x 2 entries
1154        assert_eq!(wal.closed_segments[0].segment.len(), 2);
1155        assert_eq!(wal.closed_segments[1].segment.len(), 2);
1156        assert_eq!(wal.closed_segments[2].segment.len(), 2);
1157        assert_eq!(wal.closed_segments[3].segment.len(), 2);
1158        assert_eq!(wal.open_segment.segment.len(), 2); // 1 x 2 entries
1159
1160        wal.truncate(9).unwrap();
1161
1162        assert_eq!(wal.open_segment.segment.len(), 1); // 1 x 2 entries
1163
1164        // truncate half of it
1165        wal.truncate(5).unwrap();
1166
1167        // assert truncation
1168        for i in 5..10 {
1169            assert!(wal.entry(i).is_none());
1170        }
1171
1172        // flush again with `flush_offset` > segment size
1173        wal.flush_open_segment().unwrap();
1174
1175        assert_eq!(wal.num_entries(), 5); // 5 entries removed
1176        assert_eq!(wal.first_index(), 0);
1177        assert_eq!(wal.last_index(), 4);
1178        assert_eq!(wal.closed_segments.len(), 3); // (0, 1) + (2, 3) + (4, empty slot)
1179        assert_eq!(wal.closed_segments[0].segment.len(), 2);
1180        assert_eq!(wal.closed_segments[1].segment.len(), 2);
1181        assert_eq!(wal.closed_segments[2].segment.len(), 1);
1182        assert_eq!(wal.open_segment.segment.len(), 0); // empty open segment
1183
1184        // add 5 more entries
1185        for i in 0..5 {
1186            assert_eq!(i + 5, wal.append(&&entry[..]).unwrap());
1187        }
1188
1189        // 5 closed segments
1190        assert_eq!(wal.num_entries(), 10);
1191        assert_eq!(wal.first_index(), 0);
1192        assert_eq!(wal.last_index(), 9);
1193        assert_eq!(wal.closed_segments.len(), 5);
1194        assert_eq!(wal.closed_segments[0].segment.len(), 2); // 1,2
1195        assert_eq!(wal.closed_segments[1].segment.len(), 2); // 3
1196        assert_eq!(wal.closed_segments[2].segment.len(), 1); // 4 empty slot due to truncation
1197        assert_eq!(wal.closed_segments[3].segment.len(), 2); // 5, 6
1198        assert_eq!(wal.closed_segments[4].segment.len(), 2); // 7, 8
1199        assert_eq!(wal.open_segment.segment.len(), 1); // 9
1200
1201        eprintln!("wal: {wal:?}");
1202        eprintln!("wal open: {:?}", wal.open_segment);
1203        eprintln!("wal closed: {:?}", wal.closed_segments);
1204
1205        // test persistence
1206        drop(wal);
1207        let wal = Wal::open(dir.path()).unwrap();
1208        assert_eq!(wal.num_entries(), 10);
1209        assert_eq!(wal.first_index(), 0);
1210        assert_eq!(wal.last_index(), 9);
1211        assert_eq!(wal.closed_segments.len(), 5);
1212        assert_eq!(wal.closed_segments[0].segment.len(), 2);
1213        assert_eq!(wal.closed_segments[1].segment.len(), 2);
1214        assert_eq!(wal.closed_segments[2].segment.len(), 1); // previously half truncated
1215        assert_eq!(wal.closed_segments[3].segment.len(), 2);
1216        assert_eq!(wal.closed_segments[4].segment.len(), 2);
1217        assert_eq!(wal.open_segment.segment.len(), 1);
1218    }
1219
1220    /// Tests that two Wal instances can not coexist for the same directory.
1221    #[test]
1222    fn test_exclusive_lock() {
1223        init_logger();
1224        let dir = Builder::new().prefix("wal").tempdir().unwrap();
1225        let wal = Wal::open(dir.path()).unwrap();
1226        assert_eq!(
1227            ErrorKind::WouldBlock,
1228            Wal::open(dir.path()).unwrap_err().kind()
1229        );
1230        drop(wal);
1231        assert!(Wal::open(dir.path()).is_ok());
1232    }
1233
1234    #[test]
1235    fn test_record_id_preserving() {
1236        init_logger();
1237        let entry_count = 55;
1238        let dir = Builder::new().prefix("wal").tempdir().unwrap();
1239        let options = WalOptions {
1240            segment_capacity: 512,
1241            segment_queue_len: 3,
1242            retain_closed: NonZeroUsize::new(1).unwrap(),
1243        };
1244
1245        let mut wal = Wal::with_options(dir.path(), &options).unwrap();
1246        // Fixed-size entries make the segment boundaries deterministic.
1247        let entries = vec![vec![0; 32]; entry_count];
1248
1249        for entry in &entries {
1250            wal.append(entry).unwrap();
1251        }
1252        let closed_segments = wal.closed_segments.len();
1253        let start_index = wal.open_segment_start_index();
1254
1255        wal.prefix_truncate(25).unwrap();
1256        let half_trunk_closed_segments = wal.closed_segments.len();
1257        let half_trunk_start_index = wal.open_segment_start_index();
1258
1259        wal.prefix_truncate((entry_count - 2) as u64).unwrap();
1260        let full_trunk_closed_segments = wal.closed_segments.len();
1261        let full_trunk_start_index = wal.open_segment_start_index();
1262
1263        assert!(closed_segments > half_trunk_closed_segments);
1264        assert!(half_trunk_closed_segments > full_trunk_closed_segments);
1265
1266        assert_eq!(start_index, half_trunk_start_index);
1267        assert_eq!(start_index, full_trunk_start_index);
1268    }
1269
1270    #[test]
1271    fn test_offset_after_open() {
1272        init_logger();
1273        let entry_count = 55;
1274        let dir = Builder::new().prefix("wal").tempdir().unwrap();
1275        let options = WalOptions {
1276            segment_capacity: 512,
1277            segment_queue_len: 3,
1278            retain_closed: NonZeroUsize::new(1).unwrap(),
1279        };
1280        let start_index;
1281        {
1282            let mut wal = Wal::with_options(dir.path(), &options).unwrap();
1283            let entries = EntryGenerator::new().take(entry_count).collect::<Vec<_>>();
1284
1285            for entry in &entries {
1286                wal.append(entry).unwrap();
1287            }
1288            start_index = wal.open_segment_start_index();
1289            wal.prefix_truncate(25).unwrap();
1290            assert_eq!(start_index, wal.open_segment_start_index());
1291        }
1292        {
1293            let wal2 = Wal::with_options(dir.path(), &options).unwrap();
1294            assert_eq!(start_index, wal2.open_segment_start_index());
1295        }
1296    }
1297}