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