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#[derive(Debug, Clone, Message, PartialEq, Eq, Default)]
18pub struct Commit {
19 #[bilrost(1)]
21 pub log: LogId,
22
23 #[bilrost(2)]
25 pub lsn: LSN,
26
27 #[bilrost(3)]
29 pub page_count: PageCount,
30
31 #[bilrost(4)]
35 pub commit_hash: Option<CommitHash>,
36
37 #[bilrost(5)]
40 pub segment_idx: Option<SegmentIdx>,
41
42 #[bilrost(6)]
45 pub checkpointed_at: Option<SystemTime>,
46}
47
48impl Commit {
49 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 pub fn with_segment_idx(self, segment_idx: Option<SegmentIdx>) -> Self {
75 Self { segment_idx, ..self }
76 }
77
78 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 #[bilrost(1)]
120 pub sid: SegmentId,
121
122 #[bilrost(2)]
124 pub pageset: PageSet,
125
126 #[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 #[bilrost(1)]
203 frame_size: usize,
204
205 #[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#[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 pub fn size(&self) -> usize {
247 self.bytes.end - self.bytes.start
248 }
249
250 #[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 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 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 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 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 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 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}