1use 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#[derive(Debug, Clone, PartialEq)]
23pub struct CommittedTx {
24 pub commit_lsn: Lsn,
27 pub events: Vec<MutationEvent>,
29}
30
31pub 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 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
137pub 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}