1use crate::epoch::{Epoch, Snapshot};
21use crate::memtable::Row;
22use crate::pma::Pma;
23use crate::rowid::RowId;
24use std::cmp::Ordering;
25use std::collections::{BTreeMap, BinaryHeap};
26use std::sync::Arc;
27
28type VersionKey = (RowId, Epoch);
31
32#[derive(Clone)]
35struct MutableRunSegment {
36 pma: Pma<VersionKey, Row>,
37 byte_size: u64,
38}
39
40#[derive(Clone)]
41pub struct MutableRun {
42 frozen: Arc<Vec<Arc<MutableRunSegment>>>,
43 active: MutableRunSegment,
44 byte_size: u64,
45}
46
47impl Default for MutableRun {
48 fn default() -> Self {
49 Self::new()
50 }
51}
52
53impl MutableRun {
54 pub fn new() -> Self {
55 Self {
56 frozen: Arc::new(Vec::new()),
57 active: MutableRunSegment {
58 pma: Pma::new(),
59 byte_size: 0,
60 },
61 byte_size: 0,
62 }
63 }
64
65 pub fn insert_many(&mut self, rows: Vec<Row>) {
69 let batch: Vec<(VersionKey, Row)> = rows
70 .into_iter()
71 .map(|r| {
72 let bytes = r.estimated_bytes();
73 self.byte_size = self.byte_size.saturating_add(bytes);
74 self.active.byte_size = self.active.byte_size.saturating_add(bytes);
75 ((r.row_id, r.committed_epoch), r)
76 })
77 .collect();
78 self.active.pma.extend_sorted(batch);
79 }
80
81 pub fn len(&self) -> usize {
83 self.active.pma.len()
84 + self
85 .frozen
86 .iter()
87 .map(|segment| segment.pma.len())
88 .sum::<usize>()
89 }
90
91 pub fn is_empty(&self) -> bool {
92 self.active.pma.is_empty() && self.frozen.is_empty()
93 }
94
95 pub fn approx_bytes(&self) -> u64 {
97 self.byte_size
98 }
99
100 pub fn get_version(&self, row_id: RowId, snapshot_epoch: Epoch) -> Option<(Epoch, Row)> {
104 self.get_version_at(row_id, Snapshot::at(snapshot_epoch))
105 }
106
107 pub fn get_version_at(&self, row_id: RowId, snapshot: Snapshot) -> Option<(Epoch, Row)> {
111 let mut best: Option<Row> = None;
112 for pma in self
113 .frozen
114 .iter()
115 .map(|segment| &segment.pma)
116 .chain(std::iter::once(&self.active.pma))
117 {
118 let end_epoch = if snapshot.uses_hlc_authority() {
121 Epoch(u64::MAX)
122 } else {
123 snapshot.epoch
124 };
125 for ((rid, _epoch), row) in pma.iter_from(&(row_id, Epoch::ZERO)) {
126 if *rid != row_id {
127 break;
128 }
129 if !snapshot.uses_hlc_authority() && row.committed_epoch > end_epoch {
131 break;
132 }
133 if !snapshot.observes_row(row.committed_epoch, row.commit_ts) {
134 continue;
135 }
136 if best.as_ref().is_none_or(|current| {
137 crate::epoch::version_supersedes(
138 row.committed_epoch,
139 row.commit_ts,
140 current.committed_epoch,
141 current.commit_ts,
142 )
143 }) {
144 best = Some(row.clone());
145 }
146 }
147 }
148 best.map(|row| (row.committed_epoch, row))
149 }
150
151 pub fn visible_versions(&self, snapshot_epoch: Epoch) -> Vec<Row> {
155 self.visible_versions_at(Snapshot::at(snapshot_epoch))
156 }
157
158 pub fn visible_versions_at(&self, snapshot: Snapshot) -> Vec<Row> {
161 self.newest_visible_map(snapshot).into_values().collect()
162 }
163
164 pub(crate) fn newest_visible_map(&self, snapshot: Snapshot) -> BTreeMap<RowId, Row> {
168 let mut by_row: BTreeMap<RowId, Row> = BTreeMap::new();
169 for pma in self
170 .frozen
171 .iter()
172 .map(|segment| &segment.pma)
173 .chain(std::iter::once(&self.active.pma))
174 {
175 for ((_rid, _epoch), row) in pma.iter() {
176 if !snapshot.observes_row(row.committed_epoch, row.commit_ts) {
177 continue;
178 }
179 by_row
180 .entry(row.row_id)
181 .and_modify(|existing| {
182 if crate::epoch::version_supersedes(
183 row.committed_epoch,
184 row.commit_ts,
185 existing.committed_epoch,
186 existing.commit_ts,
187 ) {
188 *existing = row.clone();
189 }
190 })
191 .or_insert_with(|| row.clone());
192 }
193 }
194 by_row
195 }
196
197 pub fn newest_visible_iter<'a>(
205 &'a self,
206 snapshot: &Snapshot,
207 ) -> MutableRunVisibleVersionCursor<'a> {
208 let snap = *snapshot;
209 let mut sources: Vec<Box<dyn Iterator<Item = &'a (VersionKey, Row)> + 'a>> = self
210 .frozen
211 .iter()
212 .map(|segment| Box::new(segment.pma.iter()) as Box<_>)
213 .chain(std::iter::once(Box::new(self.active.pma.iter()) as Box<_>))
214 .collect();
215 let mut heap = BinaryHeap::new();
216 for (index, source) in sources.iter_mut().enumerate() {
217 if let Some(((rid, _epoch), row)) = source.next() {
218 heap.push(PmaHead {
219 rid: *rid,
220 index,
221 row,
222 });
223 }
224 }
225 MutableRunVisibleVersionCursor {
226 sources,
227 heap,
228 snapshot: snap,
229 current: None,
230 last_examined: 0,
231 peak_examined: 0,
232 }
233 }
234
235 pub(crate) fn seal(&mut self) {
236 if self.active.pma.is_empty() {
237 return;
238 }
239 let active = std::mem::replace(
240 &mut self.active,
241 MutableRunSegment {
242 pma: Pma::new(),
243 byte_size: 0,
244 },
245 );
246 Arc::make_mut(&mut self.frozen).push(Arc::new(active));
247 if self.frozen.len() >= crate::MAX_READ_GENERATION_LAYERS {
248 self.consolidate();
249 }
250 }
251
252 fn consolidate(&mut self) {
253 let mut rows = self
254 .frozen
255 .iter()
256 .flat_map(|segment| segment.pma.iter().map(|(_, row)| row.clone()))
257 .collect::<Vec<_>>();
258 rows.sort_by_key(|row| (row.row_id, row.committed_epoch));
259 let mut pma = Pma::new();
260 pma.extend_sorted(
261 rows.into_iter()
262 .map(|row| ((row.row_id, row.committed_epoch), row))
263 .collect(),
264 );
265 self.frozen = Arc::new(vec![Arc::new(MutableRunSegment {
266 pma,
267 byte_size: self.byte_size,
268 })]);
269 }
270
271 #[cfg(test)]
272 pub(crate) fn frozen_layer_count(&self) -> usize {
273 self.frozen.len()
274 }
275
276 pub fn drain_sorted(&mut self) -> Vec<Row> {
279 let mut out = self
280 .frozen
281 .iter()
282 .flat_map(|segment| segment.pma.iter().map(|(_, row)| row.clone()))
283 .chain(self.active.pma.iter().map(|(_, row)| row.clone()))
284 .collect::<Vec<_>>();
285 out.sort_by_key(|row| (row.row_id, row.committed_epoch));
286 self.frozen = Arc::new(Vec::new());
287 self.active = MutableRunSegment {
288 pma: Pma::new(),
289 byte_size: 0,
290 };
291 self.byte_size = 0;
292 out
293 }
294}
295
296#[derive(Clone, Copy)]
297struct PmaHead<'a> {
298 rid: RowId,
299 index: usize,
300 row: &'a Row,
301}
302
303impl PartialEq for PmaHead<'_> {
304 fn eq(&self, other: &Self) -> bool {
305 (self.rid, self.index) == (other.rid, other.index)
306 }
307}
308impl Eq for PmaHead<'_> {}
309impl PartialOrd for PmaHead<'_> {
310 fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
311 Some(self.cmp(other))
312 }
313}
314impl Ord for PmaHead<'_> {
315 fn cmp(&self, other: &Self) -> Ordering {
316 (other.rid, other.index).cmp(&(self.rid, self.index))
317 }
318}
319
320pub struct MutableRunVisibleVersionCursor<'a> {
326 sources: Vec<Box<dyn Iterator<Item = &'a (VersionKey, Row)> + 'a>>,
327 heap: BinaryHeap<PmaHead<'a>>,
328 snapshot: Snapshot,
329 current: Option<(RowId, Epoch, &'a Row)>,
330 pub(crate) last_examined: usize,
332 pub(crate) peak_examined: usize,
334}
335
336impl<'a> Iterator for MutableRunVisibleVersionCursor<'a> {
337 type Item = (RowId, Epoch, &'a Row);
338
339 fn next(&mut self) -> Option<Self::Item> {
340 if let Some(current) = self.current.take() {
341 return Some(current);
342 }
343 self.last_examined = 0;
344 let PmaHead { rid, index, row } = self.heap.pop()?;
345 self.last_examined += 1;
346 if let Some(((next_rid, _), next_row)) = self.sources[index].next() {
347 self.heap.push(PmaHead {
348 rid: *next_rid,
349 index,
350 row: next_row,
351 });
352 }
353 let mut best = if self
354 .snapshot
355 .observes_row(row.committed_epoch, row.commit_ts)
356 {
357 Some(row)
358 } else {
359 None
360 };
361 while self.heap.peek().is_some_and(|h| h.rid == rid) {
362 let head = self.heap.pop().unwrap();
363 self.last_examined += 1;
364 if let Some(((next_rid, _), next_row)) = self.sources[head.index].next() {
365 self.heap.push(PmaHead {
366 rid: *next_rid,
367 index: head.index,
368 row: next_row,
369 });
370 }
371 if self
372 .snapshot
373 .observes_row(head.row.committed_epoch, head.row.commit_ts)
374 && best.is_none_or(|current| {
375 crate::epoch::version_supersedes(
376 head.row.committed_epoch,
377 head.row.commit_ts,
378 current.committed_epoch,
379 current.commit_ts,
380 )
381 })
382 {
383 best = Some(head.row);
384 }
385 }
386 self.peak_examined = self.peak_examined.max(self.last_examined);
387 best.map(|row| (row.row_id, row.committed_epoch, row))
388 }
389}
390
391#[cfg(test)]
392mod tests {
393 use super::*;
394 use crate::memtable::Value;
395
396 fn row(id: u64, epoch: u64, v: i64) -> Row {
397 Row::new(RowId(id), Epoch(epoch)).with_column(1, Value::Int64(v))
398 }
399
400 fn tomb(id: u64, epoch: u64) -> Row {
401 Row {
402 row_id: RowId(id),
403 committed_epoch: Epoch(epoch),
404 columns: std::collections::HashMap::new(),
405 deleted: true,
406 commit_ts: None,
407 }
408 }
409
410 fn int_of(r: &Row) -> i64 {
411 match r.columns.get(&1) {
412 Some(Value::Int64(x)) => *x,
413 _ => panic!("expected Int64 column"),
414 }
415 }
416
417 #[test]
418 fn get_version_returns_newest_visible() {
419 let mut mr = MutableRun::new();
420 mr.insert_many(vec![row(1, 1, 10), row(1, 3, 30), row(1, 9, 90)]);
421 assert_eq!(int_of(&mr.get_version(RowId(1), Epoch(5)).unwrap().1), 30);
423 assert_eq!(int_of(&mr.get_version(RowId(1), Epoch(9)).unwrap().1), 90);
425 assert!(mr.get_version(RowId(1), Epoch(0)).is_none());
427 assert!(mr.get_version(RowId(2), Epoch(100)).is_none());
429 }
430
431 #[test]
432 fn tombstone_is_returned_as_a_version() {
433 let mut mr = MutableRun::new();
434 mr.insert_many(vec![row(1, 1, 10), tomb(1, 2)]);
435 let v = mr.get_version(RowId(1), Epoch(5)).unwrap().1;
436 assert!(v.deleted);
437 let v0 = mr.get_version(RowId(1), Epoch(1)).unwrap().1;
439 assert!(!v0.deleted);
440 }
441
442 #[test]
443 fn sealed_generations_share_rows_and_consolidate() {
444 let mut writer = MutableRun::new();
445 for id in 0..crate::MAX_READ_GENERATION_LAYERS as u64 + 2 {
446 writer.insert_many(vec![row(id, id + 1, id as i64)]);
447 writer.seal();
448 }
449 assert!(writer.frozen_layer_count() < crate::MAX_READ_GENERATION_LAYERS);
450 let generation = writer.clone();
451 writer.insert_many(vec![row(99, 99, 99)]);
452 assert!(generation.get_version(RowId(99), Epoch(99)).is_none());
453 assert!(writer.get_version(RowId(99), Epoch(99)).is_some());
454 }
455
456 #[test]
457 fn visible_versions_dedups_to_newest_ascending() {
458 let mut mr = MutableRun::new();
459 mr.insert_many(vec![
460 row(3, 1, 30),
461 row(1, 1, 10),
462 row(2, 9, 20), row(1, 3, 11), row(3, 2, 31),
465 ]);
466 let out = mr.visible_versions(Epoch(5));
467 let got: Vec<(u64, i64)> = out.iter().map(|r| (r.row_id.0, int_of(r))).collect();
468 assert_eq!(got, vec![(1, 11), (3, 31)], "row 2 hidden, newest wins");
469 }
470
471 #[test]
472 fn drain_sorted_is_ascending_version_order_and_empties() {
473 let mut mr = MutableRun::new();
474 mr.insert_many(vec![row(3, 1, 0), row(1, 2, 0), row(1, 1, 0), row(2, 1, 0)]);
475 let out = mr.drain_sorted();
476 let keys: Vec<(u64, u64)> = out
477 .iter()
478 .map(|r| (r.row_id.0, r.committed_epoch.0))
479 .collect();
480 assert_eq!(keys, vec![(1, 1), (1, 2), (2, 1), (3, 1)]);
481 assert!(mr.is_empty());
482 assert_eq!(mr.approx_bytes(), 0);
483 }
484
485 #[test]
486 fn many_inserts_stay_queryable() {
487 let mut mr = MutableRun::new();
488 let mut rows = Vec::new();
489 for i in 0..500u64 {
490 rows.push(row(i, 1, i as i64));
491 }
492 mr.insert_many(rows);
493 assert_eq!(mr.len(), 500);
494 for i in 0..500u64 {
495 assert_eq!(
496 int_of(&mr.get_version(RowId(i), Epoch(1)).unwrap().1),
497 i as i64
498 );
499 }
500 assert!(mr.approx_bytes() > 0);
501 }
502
503 fn hlc(physical_micros: u64) -> mongreldb_types::hlc::HlcTimestamp {
504 mongreldb_types::hlc::HlcTimestamp {
505 physical_micros,
506 logical: 0,
507 node_tiebreaker: 1,
508 }
509 }
510
511 fn hlc_row(id: u64, epoch: u64, ts: mongreldb_types::hlc::HlcTimestamp, v: i64) -> Row {
512 Row::new_with_hlc(RowId(id), Epoch(epoch), ts).with_column(1, Value::Int64(v))
513 }
514
515 #[test]
516 fn hlc_visibility_is_authoritative_when_stamped() {
517 let mut mr = MutableRun::new();
518 let early = hlc(100);
519 let late = hlc(200);
520 mr.insert_many(vec![hlc_row(1, 1, early, 1), hlc_row(1, 2, late, 2)]);
521 let snap = Snapshot::at_hlc(Epoch(99), early);
522 let versions = mr.visible_versions_at(snap);
523 assert_eq!(versions.len(), 1);
524 assert_eq!(int_of(&versions[0]), 1);
525 assert_eq!(int_of(&mr.get_version_at(RowId(1), snap).unwrap().1), 1);
526 let snap2 = Snapshot::at_hlc(Epoch(1), late);
527 assert_eq!(int_of(&mr.visible_versions_at(snap2)[0]), 2);
528 assert_eq!(int_of(&mr.get_version_at(RowId(1), snap2).unwrap().1), 2);
529 }
530
531 #[test]
532 fn snapshot_hlc_hides_later_commit_ts_even_if_epoch_higher() {
533 let mut mr = MutableRun::new();
534 let early = hlc(100);
535 let late = hlc(200);
536 mr.insert_many(vec![hlc_row(1, 1, late, 99), hlc_row(1, 50, early, 1)]);
539 let snap = Snapshot::at_hlc(Epoch(99), early);
540 let versions = mr.visible_versions_at(snap);
541 assert_eq!(versions.len(), 1);
542 assert_eq!(int_of(&versions[0]), 1);
543 assert_eq!(versions[0].commit_ts, Some(early));
544 assert_eq!(int_of(&mr.get_version_at(RowId(1), snap).unwrap().1), 1);
545 }
546
547 #[test]
548 fn epoch_only_snapshot_sees_hlc_stamped_rows_by_epoch() {
549 let mut mr = MutableRun::new();
550 mr.insert_many(vec![hlc_row(1, 1, hlc(50), 1), row(2, 1, 2)]);
551 let legacy = Snapshot::at(Epoch(99));
552 let versions = mr.visible_versions_at(legacy);
553 assert_eq!(
554 versions.len(),
555 2,
556 "dual-model: epoch pin sees HLC rows by epoch"
557 );
558 assert!(mr.get_version_at(RowId(1), legacy).is_some());
559 assert!(mr.get_version_at(RowId(2), legacy).is_some());
560 }
561
562 fn hlc_from_raw(raw: u64) -> mongreldb_types::hlc::HlcTimestamp {
565 mongreldb_types::hlc::HlcTimestamp {
566 physical_micros: raw,
567 logical: 0,
568 node_tiebreaker: 1,
569 }
570 }
571
572 #[test]
573 fn newest_visible_iter_empty_yields_nothing() {
574 let mr = MutableRun::new();
575 let snap = Snapshot::at(Epoch(99));
576 let got: Vec<_> = mr.newest_visible_iter(&snap).collect();
577 assert!(got.is_empty());
578 }
579
580 #[test]
581 fn newest_visible_iter_single_insert_yields_one() {
582 let mut mr = MutableRun::new();
583 mr.insert_many(vec![row(7, 1, 70)]);
584 let snap = Snapshot::at(Epoch(1));
585 let got: Vec<_> = mr.newest_visible_iter(&snap).collect();
586 assert_eq!(got.len(), 1);
587 assert_eq!(got[0].0, RowId(7));
588 assert_eq!(got[0].1, Epoch(1));
589 assert_eq!(int_of(got[0].2), 70);
590 }
591
592 #[test]
593 fn newest_visible_iter_newer_epoch_wins_for_same_rowid() {
594 let mut mr = MutableRun::new();
595 mr.insert_many(vec![row(1, 1, 10), row(1, 3, 30), row(1, 9, 90)]);
596 let snap = Snapshot::at(Epoch(99));
597 let got: Vec<(u64, u64, i64)> = mr
598 .newest_visible_iter(&snap)
599 .map(|(rid, epoch, row)| (rid.0, epoch.0, int_of(row)))
600 .collect();
601 assert_eq!(got, vec![(1, 9, 90)]);
602 }
603
604 #[test]
605 fn newest_visible_iter_tombstone_suppresses_older_live_version() {
606 let mut mr = MutableRun::new();
607 mr.insert_many(vec![row(1, 1, 10), tomb(1, 2)]);
608 let snap = Snapshot::at(Epoch(5));
609 let got: Vec<(u64, u64, bool)> = mr
610 .newest_visible_iter(&snap)
611 .map(|(rid, epoch, row)| (rid.0, epoch.0, row.deleted))
612 .collect();
613 assert_eq!(got, vec![(1, 2, true)], "tombstone is the newest");
614 let snap_early = Snapshot::at(Epoch(1));
616 let got_early: Vec<(u64, u64, bool)> = mr
617 .newest_visible_iter(&snap_early)
618 .map(|(rid, epoch, row)| (rid.0, epoch.0, row.deleted))
619 .collect();
620 assert_eq!(got_early, vec![(1, 1, false)]);
621 }
622
623 #[test]
624 fn newest_visible_iter_hlc_vs_epoch_newness() {
625 let mut mr = MutableRun::new();
626 mr.insert_many(vec![
629 hlc_row(1, 50, hlc_from_raw(100), 7),
630 hlc_row(1, 1, hlc_from_raw(200), 42),
631 ]);
632 let snap = Snapshot::at_hlc(Epoch(99), hlc_from_raw(250));
633 let got: Vec<(u64, u64, i64, Option<mongreldb_types::hlc::HlcTimestamp>)> = mr
634 .newest_visible_iter(&snap)
635 .map(|(rid, epoch, row)| (rid.0, epoch.0, int_of(row), row.commit_ts))
636 .collect();
637 assert_eq!(
638 got,
639 vec![(1, 1, 42, Some(hlc_from_raw(200)))],
640 "HLC newer wins over epoch-newer when both stamped"
641 );
642 }
643}