Skip to main content

lora_wal/
history.rs

1//! Forward reader over committed WAL transactions.
2//!
3//! Recovery ([`crate::replay`]) flattens every committed transaction into
4//! one event stream. Change-feed consumers need the transaction boundaries
5//! and the commit LSN of each transaction instead, so they can hand out a
6//! resume token per commit. [`CommittedTxReader`] walks the same segments
7//! lazily, one committed transaction at a time, and stops at an upper LSN
8//! bound so it can run against a WAL that is still being appended to.
9
10use std::collections::BTreeMap;
11use std::path::{Path, PathBuf};
12
13use lora_store::MutationEvent;
14
15use crate::dir::SegmentDir;
16use crate::errors::WalError;
17use crate::lsn::Lsn;
18use crate::record::WalRecord;
19use crate::segment::SegmentReader;
20
21/// One committed transaction as read back from the log.
22#[derive(Debug, Clone, PartialEq)]
23pub struct CommittedTx {
24    /// LSN of the `TxCommit` record. Strictly increasing across
25    /// transactions.
26    pub commit_lsn: Lsn,
27    /// Every mutation the transaction committed, in append order.
28    pub events: Vec<MutationEvent>,
29}
30
31/// Lazily yields committed transactions whose `TxBegin` LSN is above
32/// `after` and whose commit LSN is at or below `upto`.
33///
34/// Records past `upto` are never decoded, so the reader is safe to use
35/// while a writer keeps appending to the active segment: everything at or
36/// below `upto` was written to the OS before `upto` was observed.
37pub struct CommittedTxReader {
38    paths: Vec<PathBuf>,
39    next_path: usize,
40    reader: Option<SegmentReader>,
41    after: Lsn,
42    upto: Lsn,
43    pending: BTreeMap<Lsn, Vec<MutationEvent>>,
44    done: bool,
45}
46
47impl CommittedTxReader {
48    pub fn open(dir: &Path, after: Lsn, upto: Lsn) -> Result<Self, WalError> {
49        let paths = SegmentDir::new(dir)
50            .list()?
51            .into_iter()
52            .map(|entry| entry.path)
53            .collect();
54        Ok(Self {
55            paths,
56            next_path: 0,
57            reader: None,
58            after,
59            upto,
60            pending: BTreeMap::new(),
61            done: after >= upto,
62        })
63    }
64
65    /// Next committed transaction, or `None` once `upto` is reached or the
66    /// log ends.
67    pub fn next_tx(&mut self) -> Result<Option<CommittedTx>, WalError> {
68        while !self.done {
69            let Some(record) = self.next_record()? else {
70                self.done = true;
71                break;
72            };
73            let lsn = record.lsn();
74            if lsn >= self.upto {
75                self.done = true;
76            }
77            if lsn > self.upto {
78                break;
79            }
80            match record {
81                WalRecord::TxBegin { lsn } if lsn > self.after => {
82                    self.pending.insert(lsn, Vec::new());
83                }
84                WalRecord::Mutation {
85                    tx_begin_lsn,
86                    event,
87                    ..
88                } => {
89                    if let Some(events) = self.pending.get_mut(&tx_begin_lsn) {
90                        events.push(event);
91                    }
92                }
93                WalRecord::MutationBatch {
94                    tx_begin_lsn,
95                    events,
96                    ..
97                } => {
98                    if let Some(pending) = self.pending.get_mut(&tx_begin_lsn) {
99                        pending.extend(events);
100                    }
101                }
102                WalRecord::TxCommit { lsn, tx_begin_lsn } => {
103                    if let Some(events) = self.pending.remove(&tx_begin_lsn) {
104                        return Ok(Some(CommittedTx {
105                            commit_lsn: lsn,
106                            events,
107                        }));
108                    }
109                }
110                WalRecord::TxAbort { tx_begin_lsn, .. } => {
111                    self.pending.remove(&tx_begin_lsn);
112                }
113                WalRecord::TxBegin { .. } | WalRecord::Checkpoint { .. } => {}
114            }
115        }
116        Ok(None)
117    }
118
119    fn next_record(&mut self) -> Result<Option<WalRecord>, WalError> {
120        loop {
121            if self.reader.is_none() {
122                let Some(path) = self.paths.get(self.next_path) else {
123                    return Ok(None);
124                };
125                self.next_path += 1;
126                self.reader = Some(SegmentReader::open(path)?);
127            }
128            let reader = self.reader.as_mut().expect("reader was just opened");
129            match reader.read_record()? {
130                Some(record) => return Ok(Some(record)),
131                None => self.reader = None,
132            }
133        }
134    }
135}
136
137/// Lowest LSN the WAL in `dir` still holds a record for: the `base_lsn` of
138/// its oldest segment. `None` when the directory has no segments.
139///
140/// Every record at or above this LSN is retained; everything below it was
141/// truncated after a checkpoint.
142pub fn oldest_retained_lsn(dir: &Path) -> Result<Option<Lsn>, WalError> {
143    let entries = SegmentDir::new(dir).list()?;
144    match entries.first() {
145        Some(first) => Ok(Some(SegmentDir::base_lsn(&first.path)?)),
146        None => Ok(None),
147    }
148}