Skip to main content

graft/core/
commit.rs

1use std::{
2    ops::{Deref, DerefMut, Range, RangeInclusive},
3    time::SystemTime,
4};
5
6use bilrost::Message;
7use smallvec::SmallVec;
8use splinter_rs::Splinter;
9
10use crate::core::{
11    LogId, PageCount, PageIdx, SegmentId, commit_hash::CommitHash, logref::LogRef, lsn::LSN,
12    pageset::PageSet,
13};
14
15/// A Commit tracks which pages have changed in a volume at a particular point in time (LSN).
16/// A commit's `SegmentIdx` may be omitted if only the Volume's `PageCount` has changed.
17#[derive(Debug, Clone, Message, PartialEq, Eq, Default)]
18pub struct Commit {
19    /// The ID of the Log which this Commit is part of.
20    #[bilrost(1)]
21    pub log: LogId,
22
23    /// The LSN of the Commit.
24    #[bilrost(2)]
25    pub lsn: LSN,
26
27    /// The Volume's `PageCount` as of this Commit.
28    #[bilrost(3)]
29    pub page_count: PageCount,
30
31    /// An optional `CommitHash` for this Commit.
32    /// Always present on Commits uploaded to a remote Log.
33    /// May be omitted on commits to a local Log.
34    #[bilrost(4)]
35    pub commit_hash: Option<CommitHash>,
36
37    /// If this Commit contains any pages, `segment_idx` forms an index of those
38    /// pages within associated Segments.
39    #[bilrost(5)]
40    pub segment_idx: Option<SegmentIdx>,
41
42    /// If this commit is a checkpoint, this timestamp is set and records the time
43    /// the commit was made a checkpoint.
44    #[bilrost(6)]
45    pub checkpointed_at: Option<SystemTime>,
46}
47
48impl Commit {
49    /// Creates a new Commit for the given snapshot info
50    pub fn new(log: LogId, lsn: LSN, page_count: PageCount) -> Self {
51        Self {
52            log,
53            lsn,
54            page_count,
55            commit_hash: None,
56            segment_idx: None,
57            checkpointed_at: None,
58        }
59    }
60
61    pub fn with_log_id(self, log: LogId) -> Self {
62        Self { log, ..self }
63    }
64
65    pub fn with_lsn(self, lsn: LSN) -> Self {
66        Self { lsn, ..self }
67    }
68
69    pub fn with_commit_hash(self, commit_hash: Option<CommitHash>) -> Self {
70        Self { commit_hash, ..self }
71    }
72
73    /// Sets the segment index for this commit.
74    pub fn with_segment_idx(self, segment_idx: Option<SegmentIdx>) -> Self {
75        Self { segment_idx, ..self }
76    }
77
78    /// Sets the checkpointed timestamp for this commit.
79    pub fn with_checkpointed_at(self, checkpointed_at: Option<SystemTime>) -> Self {
80        Self { checkpointed_at, ..self }
81    }
82
83    pub fn log(&self) -> &LogId {
84        &self.log
85    }
86
87    pub fn lsn(&self) -> LSN {
88        self.lsn
89    }
90
91    pub fn logref(&self) -> LogRef {
92        LogRef::new(self.log.clone(), self.lsn)
93    }
94
95    pub fn page_count(&self) -> PageCount {
96        self.page_count
97    }
98
99    pub fn commit_hash(&self) -> Option<&CommitHash> {
100        self.commit_hash.as_ref()
101    }
102
103    pub fn segment_idx(&self) -> Option<&SegmentIdx> {
104        self.segment_idx.as_ref()
105    }
106
107    pub fn checkpointed_at(&self) -> Option<&SystemTime> {
108        self.checkpointed_at.as_ref()
109    }
110
111    pub fn is_checkpoint(&self) -> bool {
112        self.checkpointed_at.is_some()
113    }
114}
115
116#[derive(Debug, Clone, Message, PartialEq, Eq)]
117pub struct SegmentIdx {
118    /// The Segment ID
119    #[bilrost(1)]
120    pub sid: SegmentId,
121
122    /// The set of `PageIdxs` contained by this Segment.
123    #[bilrost(2)]
124    pub pageset: PageSet,
125
126    /// An index of `SegmentFrameIdxs` contained by this Segment.
127    /// Empty on local Segments which have not been encoded and uploaded to object storage.
128    #[bilrost(3)]
129    pub frames: SmallVec<[SegmentFrameIdx; 1]>,
130}
131
132impl SegmentIdx {
133    pub fn new(sid: SegmentId, pageset: PageSet) -> Self {
134        SegmentIdx { sid, pageset, frames: SmallVec::new() }
135    }
136
137    pub fn with_frames(self, frames: SmallVec<[SegmentFrameIdx; 1]>) -> Self {
138        Self { frames, ..self }
139    }
140
141    pub fn sid(&self) -> &SegmentId {
142        &self.sid
143    }
144
145    pub fn pageset(&self) -> &PageSet {
146        &self.pageset
147    }
148
149    pub fn iter_frames(
150        &self,
151        mut filter: impl FnMut(&RangeInclusive<PageIdx>) -> bool,
152    ) -> impl Iterator<Item = SegmentRangeRef> {
153        let first_page = self.pageset.iter().next().unwrap_or(PageIdx::FIRST);
154        self.frames
155            .iter()
156            .scan((0, first_page), |(bytes_acc, pages_acc), frame| {
157                let bytes = *bytes_acc..(*bytes_acc + frame.frame_size);
158                let pages = *pages_acc..=frame.last_pageidx;
159
160                *bytes_acc += frame.frame_size;
161                *pages_acc = frame.last_pageidx.saturating_next();
162
163                Some((bytes, pages))
164            })
165            .filter(move |(_, pages)| filter(pages))
166            .map(|(bytes, pages)| {
167                let pages = pages.start().to_u32()..=pages.end().to_u32();
168                let graft = (Splinter::from(pages) & self.pageset.splinter()).into();
169                SegmentRangeRef {
170                    sid: self.sid.clone(),
171                    bytes,
172                    pageset: graft,
173                }
174            })
175    }
176
177    pub fn frame_for_pageidx(&self, pageidx: PageIdx) -> Option<SegmentRangeRef> {
178        if !self.pageset.contains(pageidx) {
179            return None;
180        }
181        self.iter_frames(|pages| pages.contains(&pageidx)).next()
182    }
183}
184
185impl Deref for SegmentIdx {
186    type Target = PageSet;
187
188    fn deref(&self) -> &Self::Target {
189        &self.pageset
190    }
191}
192
193impl DerefMut for SegmentIdx {
194    fn deref_mut(&mut self) -> &mut Self::Target {
195        &mut self.pageset
196    }
197}
198
199#[derive(Debug, Clone, Message, PartialEq, Eq, Default)]
200pub struct SegmentFrameIdx {
201    /// The length of the compressed frame in bytes.
202    #[bilrost(1)]
203    frame_size: usize,
204
205    /// The last `PageIdx` contained by this `SegmentFrame`.
206    #[bilrost(2)]
207    last_pageidx: PageIdx,
208}
209
210impl SegmentFrameIdx {
211    pub fn new(frame_size: usize, last_pageidx: PageIdx) -> Self {
212        Self { frame_size, last_pageidx }
213    }
214
215    pub fn frame_size(&self) -> usize {
216        self.frame_size
217    }
218
219    pub fn last_pageidx(&self) -> PageIdx {
220        self.last_pageidx
221    }
222}
223
224/// A `SegmentRangeRef` contains the byte range and corresponding pages for a
225/// subset of a segment. The subset must correspond to one or more entire
226/// `SegmentFrames`.
227#[derive(Clone, PartialEq, Eq)]
228pub struct SegmentRangeRef {
229    pub sid: SegmentId,
230    pub bytes: Range<usize>,
231    pub pageset: PageSet,
232}
233
234impl std::fmt::Debug for SegmentRangeRef {
235    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
236        f.debug_struct("SegmentRangeRef")
237            .field("sid", &self.sid)
238            .field("bytes", &self.bytes)
239            .field("pages", &self.pageset.cardinality())
240            .finish()
241    }
242}
243
244impl SegmentRangeRef {
245    /// The size of the frame in bytes
246    pub fn size(&self) -> usize {
247        self.bytes.end - self.bytes.start
248    }
249
250    /// Attempt to coalesce two frame refs together.
251    /// Returns the two frame refs unmodified if coalescing is impossible.
252    #[allow(clippy::result_large_err)]
253    pub fn coalesce(self, other: Self) -> Result<Self, (Self, Self)> {
254        if self.sid != other.sid {
255            return Err((self, other));
256        }
257
258        let (left, right) = if self.bytes.end == other.bytes.start {
259            (self, other)
260        } else if other.bytes.end == self.bytes.start {
261            (other, self)
262        } else {
263            return Err((self, other));
264        };
265
266        let left_splinter: Splinter = left.pageset.into();
267        let right_splinter: Splinter = right.pageset.into();
268        Ok(Self {
269            sid: left.sid,
270            bytes: left.bytes.start..right.bytes.end,
271            pageset: (left_splinter | right_splinter).into(),
272        })
273    }
274}
275
276#[cfg(test)]
277mod tests {
278    use crate::pageidx;
279
280    use super::*;
281
282    #[test]
283    fn test_frame_for_pageidx() {
284        let pageset = PageSet::from_range(pageidx!(5)..=pageidx!(25));
285        let mut frames = SmallVec::new();
286        frames.push(SegmentFrameIdx {
287            frame_size: 100,
288            last_pageidx: pageidx!(10),
289        });
290        frames.push(SegmentFrameIdx {
291            frame_size: 200,
292            last_pageidx: pageidx!(20),
293        });
294        frames.push(SegmentFrameIdx {
295            frame_size: 150,
296            last_pageidx: pageidx!(25),
297        });
298
299        let sid = SegmentId::random();
300        let segment_idx = SegmentIdx { sid: sid.clone(), pageset, frames };
301
302        let tests = [
303            (pageidx!(4), None),
304            (
305                pageidx!(5),
306                Some(SegmentRangeRef {
307                    sid: sid.clone(),
308                    bytes: 0..100,
309                    pageset: PageSet::from_range(pageidx!(5)..=pageidx!(10)),
310                }),
311            ),
312            (
313                pageidx!(10),
314                Some(SegmentRangeRef {
315                    sid: sid.clone(),
316                    bytes: 0..100,
317                    pageset: PageSet::from_range(pageidx!(5)..=pageidx!(10)),
318                }),
319            ),
320            (
321                pageidx!(11),
322                Some(SegmentRangeRef {
323                    sid: sid.clone(),
324                    bytes: 100..300,
325                    pageset: PageSet::from_range(pageidx!(11)..=pageidx!(20)),
326                }),
327            ),
328            (
329                pageidx!(20),
330                Some(SegmentRangeRef {
331                    sid: sid.clone(),
332                    bytes: 100..300,
333                    pageset: PageSet::from_range(pageidx!(11)..=pageidx!(20)),
334                }),
335            ),
336            (
337                pageidx!(25),
338                Some(SegmentRangeRef {
339                    sid: sid.clone(),
340                    bytes: 300..450,
341                    pageset: PageSet::from_range(pageidx!(21)..=pageidx!(25)),
342                }),
343            ),
344            (pageidx!(26), None),
345        ];
346
347        for (pageidx, expected) in tests {
348            assert_eq!(
349                segment_idx.frame_for_pageidx(pageidx),
350                expected,
351                "wrong frame for pageidx {pageidx}"
352            );
353        }
354    }
355
356    #[test]
357    fn test_frame_for_pageidx_empty_frames() {
358        let segment_idx = SegmentIdx {
359            sid: SegmentId::random(),
360            pageset: PageSet::EMPTY,
361            frames: SmallVec::new(),
362        };
363
364        let result = segment_idx.frame_for_pageidx(pageidx!(1));
365        assert!(result.is_none());
366    }
367
368    #[test]
369    fn test_segment_range_ref_coalesce_adjacent() {
370        let sid = SegmentId::random();
371        // Test coalescing two adjacent ranges (first before second)
372        let frame1 = SegmentRangeRef {
373            sid: sid.clone(),
374            bytes: 0..100,
375            pageset: PageSet::from_range(pageidx!(5)..=pageidx!(10)),
376        };
377        let frame2 = SegmentRangeRef {
378            sid: sid.clone(),
379            bytes: 100..200,
380            pageset: PageSet::from_range(pageidx!(11)..=pageidx!(20)),
381        };
382
383        let result = frame1.clone().coalesce(frame2.clone()).unwrap();
384        assert_eq!(result.bytes, 0..200);
385        assert_eq!(
386            result.pageset,
387            PageSet::from_range(pageidx!(5)..=pageidx!(20))
388        );
389
390        // Test coalescing in reverse order (second before first)
391        let result = frame2.coalesce(frame1).unwrap();
392        assert_eq!(result.bytes, 0..200);
393        assert_eq!(
394            result.pageset,
395            PageSet::from_range(pageidx!(5)..=pageidx!(20))
396        );
397    }
398
399    #[test]
400    fn test_segment_range_ref_coalesce_non_adjacent() {
401        let sid = SegmentId::random();
402        // Test that non-adjacent ranges cannot be coalesced
403        let frame1 = SegmentRangeRef {
404            sid: sid.clone(),
405            bytes: 0..100,
406            pageset: PageSet::from_range(pageidx!(5)..=pageidx!(10)),
407        };
408        let frame2 = SegmentRangeRef {
409            sid: sid.clone(),
410            bytes: 150..250,
411            pageset: PageSet::from_range(pageidx!(20)..=pageidx!(30)),
412        };
413
414        let result = frame1.clone().coalesce(frame2.clone());
415        assert!(result.is_err());
416        let (f1, f2) = result.unwrap_err();
417        assert_eq!(f1, frame1);
418        assert_eq!(f2, frame2);
419    }
420
421    #[test]
422    fn test_segment_range_ref_coalesce_diff_segment() {
423        // Test that adjacent ranges from different segments don't combine
424        let frame1 = SegmentRangeRef {
425            sid: SegmentId::random(),
426            bytes: 0..100,
427            pageset: PageSet::from_range(pageidx!(5)..=pageidx!(10)),
428        };
429        let frame2 = SegmentRangeRef {
430            sid: SegmentId::random(),
431            bytes: 100..200,
432            pageset: PageSet::from_range(pageidx!(11)..=pageidx!(20)),
433        };
434
435        let result = frame1.clone().coalesce(frame2.clone());
436        assert!(result.is_err());
437        let (f1, f2) = result.unwrap_err();
438        assert_eq!(f1, frame1);
439        assert_eq!(f2, frame2);
440    }
441
442    #[test]
443    fn test_iter_frames_no_filter() {
444        let pageset = PageSet::from_range(pageidx!(5)..=pageidx!(25));
445        let mut frames = SmallVec::new();
446        frames.push(SegmentFrameIdx {
447            frame_size: 100,
448            last_pageidx: pageidx!(10),
449        });
450        frames.push(SegmentFrameIdx {
451            frame_size: 200,
452            last_pageidx: pageidx!(20),
453        });
454        frames.push(SegmentFrameIdx {
455            frame_size: 150,
456            last_pageidx: pageidx!(25),
457        });
458
459        let segment_idx = SegmentIdx {
460            sid: SegmentId::random(),
461            pageset,
462            frames,
463        };
464
465        // Collect all frames
466        let all_frames: Vec<_> = segment_idx.iter_frames(|_| true).collect();
467        assert_eq!(all_frames.len(), 3);
468
469        assert_eq!(all_frames[0].bytes, 0..100);
470        assert_eq!(
471            all_frames[0].pageset,
472            PageSet::from_range(pageidx!(5)..=pageidx!(10))
473        );
474
475        assert_eq!(all_frames[1].bytes, 100..300);
476        assert_eq!(
477            all_frames[1].pageset,
478            PageSet::from_range(pageidx!(11)..=pageidx!(20))
479        );
480
481        assert_eq!(all_frames[2].bytes, 300..450);
482        assert_eq!(
483            all_frames[2].pageset,
484            PageSet::from_range(pageidx!(21)..=pageidx!(25))
485        );
486    }
487
488    #[test]
489    fn test_iter_frames_with_filter() {
490        let pageset = PageSet::from_range(pageidx!(5)..=pageidx!(25));
491        let mut frames = SmallVec::new();
492        frames.push(SegmentFrameIdx {
493            frame_size: 100,
494            last_pageidx: pageidx!(10),
495        });
496        frames.push(SegmentFrameIdx {
497            frame_size: 200,
498            last_pageidx: pageidx!(20),
499        });
500        frames.push(SegmentFrameIdx {
501            frame_size: 150,
502            last_pageidx: pageidx!(25),
503        });
504
505        let segment_idx = SegmentIdx {
506            sid: SegmentId::random(),
507            pageset,
508            frames,
509        };
510
511        // Filter for frames containing page 15
512        let filtered_frames: Vec<_> = segment_idx
513            .iter_frames(|pages| pages.contains(&pageidx!(15)))
514            .collect();
515        assert_eq!(filtered_frames.len(), 1);
516        assert_eq!(filtered_frames[0].bytes, 100..300);
517        assert_eq!(
518            filtered_frames[0].pageset,
519            PageSet::from_range(pageidx!(11)..=pageidx!(20))
520        );
521    }
522
523    #[test]
524    fn test_iter_frames_empty() {
525        let segment_idx = SegmentIdx {
526            sid: SegmentId::random(),
527            pageset: PageSet::EMPTY,
528            frames: SmallVec::new(),
529        };
530
531        let frames: Vec<_> = segment_idx.iter_frames(|_| true).collect();
532        assert_eq!(frames.len(), 0);
533    }
534}