1use std::{
5 cmp::{Ordering, Reverse},
6 collections::{HashMap, HashSet},
7 ops::Bound,
8 sync::Arc,
9};
10
11use reifydb_codec::key::encoded::EncodedKey;
12use reifydb_core::{common::CommitVersion, interface::store::EntryKind};
13use reifydb_value::{Result, byte_size::ByteSize, reifydb_assertions, util::cowvec::CowVec};
14use tracing::{Span, field, instrument};
15
16use crate::{
17 MultiVersionScope,
18 tier::{
19 DisplacedValues, HistoricalCursor, RangeBatch, RangeCursor, RawEntry, TierBackend, TierBatch,
20 TierStorage, VersionedGetResult,
21 commit::memory::entry::{
22 CurrentMap, Entries, Entry, HistoricalMap, OldestIndex, entry_bytes, entry_bytes_with,
23 oldest_version, reconcile_oldest,
24 },
25 },
26};
27
28type EvictablePersist = Vec<(EncodedKey, CommitVersion, Option<CowVec<u8>>)>;
29type EvictableDrop = Vec<EvictedVersion>;
30
31#[derive(Clone, Debug)]
32pub struct EvictedVersion {
33 pub key: EncodedKey,
34 pub version: CommitVersion,
35 pub value_bytes: ByteSize,
36 pub current: bool,
37}
38
39fn value_bytes_of(value: &Option<CowVec<u8>>) -> ByteSize {
40 ByteSize::from_bytes(value.as_ref().map(|v| v.len() as u64).unwrap_or(0))
41}
42
43#[derive(Clone)]
44pub struct MemoryRowStorage {
45 inner: Arc<MemoryRowStorageInner>,
46}
47
48struct MemoryRowStorageInner {
49 entries: Entries,
50}
51
52impl Default for MemoryRowStorage {
53 fn default() -> Self {
54 Self::new()
55 }
56}
57
58impl MemoryRowStorage {
59 #[instrument(name = "store::multi::memory::new", level = "debug")]
60 pub fn new() -> Self {
61 Self {
62 inner: Arc::new(MemoryRowStorageInner {
63 entries: Entries::default(),
64 }),
65 }
66 }
67
68 pub fn count_current(&self, table: EntryKind) -> Result<u64> {
69 Ok(self.inner.entries.data.get(&table).map(|e| e.current.read().len() as u64).unwrap_or(0))
70 }
71
72 pub fn list_all_entry_kinds(&self) -> Result<Vec<EntryKind>> {
73 Ok(self.inner.entries.data.keys())
74 }
75
76 fn collect_oldest_pending(&self) -> Vec<(EntryKind, CommitVersion)> {
77 self.inner
78 .entries
79 .data
80 .keys()
81 .into_iter()
82 .filter_map(|kind| {
83 let entry = self.inner.entries.data.get(&kind)?;
84 let oldest = entry.oldest.read();
85 let version = *oldest.keys().next()?;
86 Some((kind, version))
87 })
88 .collect()
89 }
90
91 pub fn list_entry_kinds_by_oldest_pending(&self) -> Result<Vec<EntryKind>> {
92 let mut pending = self.collect_oldest_pending();
93 pending.sort_by_key(|(_, version)| *version);
94 Ok(pending.into_iter().map(|(kind, _)| kind).collect())
95 }
96
97 pub fn oldest_pending_for(&self, kind: EntryKind) -> Option<CommitVersion> {
98 let entry = self.inner.entries.data.get(&kind)?;
99 let oldest = entry.oldest.read();
100 oldest.keys().next().copied()
101 }
102
103 pub fn count_historical(&self, table: EntryKind) -> Result<u64> {
104 Ok(self.inner
105 .entries
106 .data
107 .get(&table)
108 .map(|e| {
109 let hist = e.historical.read();
110 hist.values().map(|m| m.len() as u64).sum()
111 })
112 .unwrap_or(0))
113 }
114
115 pub fn current_resident_bytes(&self) -> ByteSize {
116 let total = self
117 .inner
118 .entries
119 .data
120 .keys()
121 .into_iter()
122 .filter_map(|kind| self.inner.entries.data.get(&kind))
123 .map(|entry| entry.bytes.current())
124 .sum();
125 ByteSize::from_bytes(total)
126 }
127
128 pub fn historical_resident_bytes(&self) -> ByteSize {
129 let total = self
130 .inner
131 .entries
132 .data
133 .keys()
134 .into_iter()
135 .filter_map(|kind| self.inner.entries.data.get(&kind))
136 .map(|entry| entry.bytes.historical())
137 .sum();
138 ByteSize::from_bytes(total)
139 }
140
141 #[inline]
142 #[instrument(name = "store::multi::memory::get_or_create_table", level = "trace", skip(self), fields(table = ?table))]
143 fn get_or_create_table(&self, table: EntryKind) -> Entry {
144 self.inner.entries.data.get_or_insert_with(table, Entry::new)
145 }
146
147 #[inline]
148 #[instrument(name = "store::multi::memory::set::table", level = "trace", skip(self, entries), fields(
149 table = ?table,
150 entry_count = entries.len(),
151 ))]
152 fn process_table(
153 &self,
154 table: EntryKind,
155 version: CommitVersion,
156 entries: Vec<(EncodedKey, Option<CowVec<u8>>)>,
157 displaced: &mut DisplacedValues,
158 ) {
159 let table_entry = self.get_or_create_table(table);
160 let (mut current, mut historical) = table_entry.write_pair();
161 let mut oldest = table_entry.oldest.write();
162
163 for (key, value) in entries {
164 if let Some((pre_version, pre_value)) = current.get(&key) {
165 if *pre_version < version {
166 let pre_version = *pre_version;
167 let pre_value = pre_value.clone();
168 reifydb_assertions! {
169 assert!(
170 version.0 > pre_version.0,
171 "promoting current entry to historical requires the incoming version to exceed it, otherwise the same version appears in both tiers and point-reads return the wrong entry (version={} pre_version={})",
172 version.0,
173 pre_version.0
174 );
175 }
176 displaced.push((key.clone(), pre_value.as_ref().map_or(0, |v| v.len() as u64)));
177 let pre_bytes = entry_bytes(&key, &pre_value);
178 let new_bytes = entry_bytes(&key, &value);
179 let replaced = historical
180 .entry(key.clone())
181 .or_default()
182 .insert(Reverse(pre_version), pre_value);
183 table_entry.bytes.add_historical(pre_bytes);
184 if let Some(replaced) = replaced {
185 table_entry.bytes.sub_historical(entry_bytes(&key, &replaced));
186 }
187
188 current.insert(key, (version, value));
189 table_entry.bytes.add_current(new_bytes);
190 table_entry.bytes.sub_current(pre_bytes);
191 } else {
192 let key_heap = key.heap_bytes();
193 let new_bytes = entry_bytes_with(key_heap, &value);
194 let old_oldest = oldest_version(¤t, &historical, &key);
195 let new_oldest = Some(old_oldest.map_or(version, |o| o.min(version)));
196 let index_key = key.clone();
197 let replaced =
198 historical.entry(key).or_default().insert(Reverse(version), value);
199 table_entry.bytes.add_historical(new_bytes);
200 if let Some(replaced) = replaced {
201 table_entry.bytes.sub_historical(entry_bytes_with(key_heap, &replaced));
202 }
203 reconcile_oldest(&mut oldest, &index_key, old_oldest, new_oldest);
204 }
205 } else {
206 let new_bytes = entry_bytes(&key, &value);
207 oldest.entry(version).or_default().insert(key.clone());
208 current.insert(key, (version, value));
209 table_entry.bytes.add_current(new_bytes);
210 }
211 }
212 }
213
214 pub fn oldest_pending_version(&self) -> Option<CommitVersion> {
215 self.collect_oldest_pending().into_iter().map(|(_, version)| version).min()
216 }
217
218 pub fn collect_evictable_below(
219 &self,
220 table: EntryKind,
221 cutoff: CommitVersion,
222 budget: usize,
223 ) -> (EvictablePersist, EvictableDrop, bool) {
224 let entry = match self.inner.entries.data.get(&table) {
225 Some(e) => e,
226 None => return (Vec::new(), Vec::new(), false),
227 };
228 let current = entry.current.read();
229 let historical = entry.historical.read();
230 let oldest = entry.oldest.read();
231
232 let mut selected: Vec<EncodedKey> = Vec::with_capacity(budget.min(current.len()));
233 let mut more = false;
234 'select: for (_bucket, keys) in oldest.range(..=cutoff) {
235 for key in keys {
236 if selected.len() >= budget {
237 more = true;
238 break 'select;
239 }
240 selected.push(key.clone());
241 }
242 }
243
244 let mut latest: HashMap<EncodedKey, (CommitVersion, Option<CowVec<u8>>)> =
245 HashMap::with_capacity(selected.len());
246 let mut to_drop: EvictableDrop = Vec::new();
247 for key in &selected {
248 if let Some((v, val)) = current.get(key)
249 && *v <= cutoff
250 {
251 to_drop.push(EvictedVersion {
252 key: key.clone(),
253 version: *v,
254 value_bytes: value_bytes_of(val),
255 current: true,
256 });
257 latest.insert(key.clone(), (*v, val.clone()));
258 }
259 if let Some(versions) = historical.get(key) {
260 for (Reverse(v), val) in versions.iter() {
261 if *v <= cutoff {
262 to_drop.push(EvictedVersion {
263 key: key.clone(),
264 version: *v,
265 value_bytes: value_bytes_of(val),
266 current: false,
267 });
268 match latest.get(key) {
269 Some((best, _)) if *best >= *v => {}
270 _ => {
271 latest.insert(key.clone(), (*v, val.clone()));
272 }
273 }
274 }
275 }
276 }
277 }
278
279 let to_persist = latest.into_iter().map(|(key, (v, val))| (key, v, val)).collect();
280 (to_persist, to_drop, more)
281 }
282}
283
284impl TierStorage for MemoryRowStorage {
285 #[instrument(name = "store::multi::memory::get", level = "trace", skip(self, key), fields(table = ?table, key_len = key.len(), version = version.0))]
286 fn get(&self, table: EntryKind, key: &[u8], version: CommitVersion) -> Result<VersionedGetResult> {
287 let entry = match self.inner.entries.data.get(&table) {
288 Some(e) => e,
289 None => return Ok(VersionedGetResult::NotFound),
290 };
291
292 let current = entry.current.read();
293 if let Some((cur_version, value)) = current.get(key)
294 && *cur_version <= version
295 {
296 return Ok(match value {
297 Some(v) => VersionedGetResult::Value {
298 value: v.clone(),
299 version: *cur_version,
300 },
301 None => VersionedGetResult::Tombstone,
302 });
303 }
304 drop(current);
305
306 let historical = entry.historical.read();
307 if let Some(versions) = historical.get(key) {
308 for (Reverse(v), value) in versions.range(Reverse(version)..) {
309 if *v <= version {
310 return Ok(match value {
311 Some(val) => VersionedGetResult::Value {
312 value: val.clone(),
313 version: *v,
314 },
315 None => VersionedGetResult::Tombstone,
316 });
317 }
318 }
319 }
320
321 Ok(VersionedGetResult::NotFound)
322 }
323
324 #[instrument(name = "store::multi::memory::contains", level = "trace", skip(self, key), fields(table = ?table, key_len = key.len(), version = version.0), ret)]
325 fn contains(&self, table: EntryKind, key: &[u8], version: CommitVersion) -> Result<bool> {
326 let entry = match self.inner.entries.data.get(&table) {
327 Some(e) => e,
328 None => return Ok(false),
329 };
330
331 let current = entry.current.read();
332 if let Some((cur_version, value)) = current.get(key)
333 && *cur_version <= version
334 {
335 return Ok(value.is_some());
336 }
337 drop(current);
338
339 let historical = entry.historical.read();
340 if let Some(versions) = historical.get(key) {
341 for (Reverse(v), value) in versions.range(Reverse(version)..) {
342 if *v <= version {
343 return Ok(value.is_some());
344 }
345 }
346 }
347
348 Ok(false)
349 }
350
351 #[instrument(name = "store::multi::memory::set", level = "trace", skip(self, batches), fields(
352 table_count = batches.len(),
353 total_entry_count = field::Empty,
354 version = version.0
355 ))]
356 fn set(&self, version: CommitVersion, batches: TierBatch) -> Result<DisplacedValues> {
357 let total_entries: usize = batches.values().map(|v| v.len()).sum();
358
359 let mut displaced = DisplacedValues::with_capacity(total_entries);
360 batches.into_iter().for_each(|(table, entries)| {
361 self.process_table(table, version, entries, &mut displaced);
362 });
363
364 Span::current().record("total_entry_count", total_entries);
365 Ok(displaced)
366 }
367
368 #[instrument(name = "store::multi::memory::range_next", level = "trace", skip(self, cursor, start, end), fields(table = ?table, batch_size = batch_size, scope = ?scope))]
369 fn range_next(
370 &self,
371 table: EntryKind,
372 cursor: &mut RangeCursor,
373 start: Bound<&[u8]>,
374 end: Bound<&[u8]>,
375 scope: MultiVersionScope,
376 batch_size: usize,
377 ) -> Result<RangeBatch> {
378 if cursor.exhausted {
379 return Ok(RangeBatch::empty());
380 }
381
382 let entry = match self.inner.entries.data.get(&table) {
383 Some(e) => e,
384 None => {
385 cursor.exhausted = true;
386 return Ok(RangeBatch::empty());
387 }
388 };
389
390 let cursor_key = cursor.last_key.clone();
391
392 let current = entry.current.read();
393 let historical = entry.historical.read();
394
395 let mut entries: Vec<RawEntry> = Vec::with_capacity(batch_size + 1);
396
397 let iter_start: Bound<&[u8]> = match &cursor_key {
398 Some(last) => Bound::Excluded(last.as_slice()),
399 None => start,
400 };
401
402 let iter_end: Bound<&[u8]> = end;
403
404 let mut cur_iter = current.range::<[u8], _>((iter_start, iter_end)).peekable();
405 let mut hist_iter = historical.range::<[u8], _>((iter_start, iter_end)).peekable();
406
407 while entries.len() <= batch_size {
408 let (take_cur, take_hist) = match (cur_iter.peek(), hist_iter.peek()) {
409 (None, None) => break,
410 (Some(_), None) => (true, false),
411 (None, Some(_)) => (false, true),
412 (Some((kc, _)), Some((kh, _))) => match kc.cmp(kh) {
413 Ordering::Less => (true, false),
414 Ordering::Greater => (false, true),
415 Ordering::Equal => (true, true),
416 },
417 };
418
419 if take_cur && take_hist {
420 let (key, (cur_version, cur_value)) = cur_iter.next().unwrap();
421 let (_, versions) = hist_iter.next().unwrap();
422 if scope.contains(*cur_version) {
423 entries.push(RawEntry {
424 key: key.clone(),
425 version: *cur_version,
426 value: cur_value.clone(),
427 });
428 } else if *cur_version > scope.read() {
429 for (Reverse(v), value) in versions.range(Reverse(scope.read())..) {
430 if scope.contains(*v) {
431 entries.push(RawEntry {
432 key: key.clone(),
433 version: *v,
434 value: value.clone(),
435 });
436 break;
437 }
438 if let MultiVersionScope::Between {
439 after,
440 ..
441 } = scope && *v <= after
442 {
443 break;
444 }
445 }
446 }
447 } else if take_cur {
448 let (key, (cur_version, cur_value)) = cur_iter.next().unwrap();
449 if scope.contains(*cur_version) {
450 entries.push(RawEntry {
451 key: key.clone(),
452 version: *cur_version,
453 value: cur_value.clone(),
454 });
455 }
456 } else {
457 let (key, versions) = hist_iter.next().unwrap();
458 for (Reverse(v), value) in versions.range(Reverse(scope.read())..) {
459 if scope.contains(*v) {
460 entries.push(RawEntry {
461 key: key.clone(),
462 version: *v,
463 value: value.clone(),
464 });
465 break;
466 }
467 if let MultiVersionScope::Between {
468 after,
469 ..
470 } = scope && *v <= after
471 {
472 break;
473 }
474 }
475 }
476 }
477
478 let has_more = entries.len() > batch_size;
479 if has_more {
480 entries.truncate(batch_size);
481 }
482
483 if let Some(last_entry) = entries.last() {
484 cursor.last_key = Some(last_entry.key.clone());
485 }
486 if !has_more {
487 cursor.exhausted = true;
488 }
489
490 Ok(RangeBatch {
491 entries,
492 has_more,
493 })
494 }
495
496 #[instrument(name = "store::multi::memory::range_rev_next", level = "trace", skip(self, cursor, start, end), fields(table = ?table, batch_size = batch_size, scope = ?scope))]
497 fn range_rev_next(
498 &self,
499 table: EntryKind,
500 cursor: &mut RangeCursor,
501 start: Bound<&[u8]>,
502 end: Bound<&[u8]>,
503 scope: MultiVersionScope,
504 batch_size: usize,
505 ) -> Result<RangeBatch> {
506 if cursor.exhausted {
507 return Ok(RangeBatch::empty());
508 }
509
510 let entry = match self.inner.entries.data.get(&table) {
511 Some(e) => e,
512 None => {
513 cursor.exhausted = true;
514 return Ok(RangeBatch::empty());
515 }
516 };
517
518 let cursor_key = cursor.last_key.clone();
519
520 let current = entry.current.read();
521 let historical = entry.historical.read();
522
523 let mut entries: Vec<RawEntry> = Vec::with_capacity(batch_size + 1);
524
525 let iter_start: Bound<&[u8]> = start;
526
527 let iter_end: Bound<&[u8]> = match &cursor_key {
528 Some(last) => Bound::Excluded(last.as_slice()),
529 None => end,
530 };
531
532 let mut cur_iter = current.range::<[u8], _>((iter_start, iter_end)).rev().peekable();
533 let mut hist_iter = historical.range::<[u8], _>((iter_start, iter_end)).rev().peekable();
534
535 while entries.len() <= batch_size {
536 let (take_cur, take_hist) = match (cur_iter.peek(), hist_iter.peek()) {
537 (None, None) => break,
538 (Some(_), None) => (true, false),
539 (None, Some(_)) => (false, true),
540 (Some((kc, _)), Some((kh, _))) => match kc.cmp(kh) {
541 Ordering::Greater => (true, false),
542 Ordering::Less => (false, true),
543 Ordering::Equal => (true, true),
544 },
545 };
546
547 if take_cur && take_hist {
548 let (key, (cur_version, cur_value)) = cur_iter.next().unwrap();
549 let (_, versions) = hist_iter.next().unwrap();
550 if scope.contains(*cur_version) {
551 entries.push(RawEntry {
552 key: key.clone(),
553 version: *cur_version,
554 value: cur_value.clone(),
555 });
556 } else if *cur_version > scope.read() {
557 for (Reverse(v), value) in versions.range(Reverse(scope.read())..) {
558 if scope.contains(*v) {
559 entries.push(RawEntry {
560 key: key.clone(),
561 version: *v,
562 value: value.clone(),
563 });
564 break;
565 }
566 if let MultiVersionScope::Between {
567 after,
568 ..
569 } = scope && *v <= after
570 {
571 break;
572 }
573 }
574 }
575 } else if take_cur {
576 let (key, (cur_version, cur_value)) = cur_iter.next().unwrap();
577 if scope.contains(*cur_version) {
578 entries.push(RawEntry {
579 key: key.clone(),
580 version: *cur_version,
581 value: cur_value.clone(),
582 });
583 }
584 } else {
585 let (key, versions) = hist_iter.next().unwrap();
586 for (Reverse(v), value) in versions.range(Reverse(scope.read())..) {
587 if scope.contains(*v) {
588 entries.push(RawEntry {
589 key: key.clone(),
590 version: *v,
591 value: value.clone(),
592 });
593 break;
594 }
595 if let MultiVersionScope::Between {
596 after,
597 ..
598 } = scope && *v <= after
599 {
600 break;
601 }
602 }
603 }
604 }
605
606 let has_more = entries.len() > batch_size;
607 if has_more {
608 entries.truncate(batch_size);
609 }
610
611 if let Some(last_entry) = entries.last() {
612 cursor.last_key = Some(last_entry.key.clone());
613 }
614 if !has_more {
615 cursor.exhausted = true;
616 }
617
618 Ok(RangeBatch {
619 entries,
620 has_more,
621 })
622 }
623
624 #[instrument(name = "store::multi::memory::ensure_table", level = "trace", skip(self), fields(table = ?table))]
625 fn ensure_table(&self, table: EntryKind) -> Result<()> {
626 let _ = self.get_or_create_table(table);
627 Ok(())
628 }
629
630 #[instrument(name = "store::multi::memory::clear_table", level = "debug", skip(self), fields(table = ?table))]
631 fn clear_table(&self, table: EntryKind) -> Result<()> {
632 if let Some(entry) = self.inner.entries.data.get(&table) {
633 *entry.current.write() = CurrentMap::new();
634 *entry.historical.write() = HistoricalMap::new();
635 *entry.oldest.write() = OldestIndex::new();
636 entry.bytes.reset();
637 }
638 Ok(())
639 }
640}
641
642impl MemoryRowStorage {
643 #[instrument(name = "store::multi::memory::drop", level = "debug", skip(self, batches), fields(
644 table_count = batches.len(),
645 total_entry_count = field::Empty
646 ))]
647 pub fn compact(
648 &self,
649 batches: HashMap<EntryKind, Vec<(EncodedKey, CommitVersion)>>,
650 ) -> Result<Vec<EvictedVersion>> {
651 let total_entries: usize = batches.values().map(|v| v.len()).sum();
652 let mut removed: Vec<EvictedVersion> = Vec::with_capacity(total_entries);
653
654 for (table, entries) in batches {
655 let table_entry = self.get_or_create_table(table);
656 let (mut current, mut historical) = table_entry.write_pair();
657 let mut oldest = table_entry.oldest.write();
658
659 let mut by_key: HashMap<EncodedKey, Vec<CommitVersion>> = HashMap::new();
660 for (key, version) in entries {
661 by_key.entry(key).or_default().push(version);
662 }
663
664 for (key, dropped_versions) in by_key {
665 let old_oldest = oldest_version(¤t, &historical, &key);
666 let dropped_set: HashSet<CommitVersion> = dropped_versions.iter().copied().collect();
667
668 let cur_version = current.get(&key).map(|(v, _)| *v);
669 let stored_hist_covered = historical
670 .get(&key)
671 .map(|m| m.keys().all(|Reverse(v)| dropped_set.contains(v)))
672 .unwrap_or(true);
673 let stored_cur_covered = cur_version.is_none_or(|v| dropped_set.contains(&v));
674
675 if stored_cur_covered && stored_hist_covered {
676 if let Some((version, value)) = current.remove(&key) {
677 table_entry.bytes.sub_current(entry_bytes(&key, &value));
678 removed.push(EvictedVersion {
679 key: key.clone(),
680 version,
681 value_bytes: value_bytes_of(&value),
682 current: true,
683 });
684 }
685 if let Some(versions) = historical.remove(&key) {
686 for (Reverse(version), value) in versions.iter() {
687 table_entry.bytes.sub_historical(entry_bytes(&key, value));
688 removed.push(EvictedVersion {
689 key: key.clone(),
690 version: *version,
691 value_bytes: value_bytes_of(value),
692 current: false,
693 });
694 }
695 }
696 reconcile_oldest(&mut oldest, &key, old_oldest, None);
697 continue;
698 }
699
700 for version in dropped_versions {
701 let cur_matches = current.get(&key).map(|(v, _)| *v) == Some(version);
702 if cur_matches {
703 if let Some((version, value)) = current.remove(&key) {
704 table_entry.bytes.sub_current(entry_bytes(&key, &value));
705 removed.push(EvictedVersion {
706 key: key.clone(),
707 version,
708 value_bytes: value_bytes_of(&value),
709 current: true,
710 });
711 }
712 } else {
713 let now_empty = if let Some(versions) = historical.get_mut(&key) {
714 if let Some(value) = versions.remove(&Reverse(version)) {
715 table_entry
716 .bytes
717 .sub_historical(entry_bytes(&key, &value));
718 removed.push(EvictedVersion {
719 key: key.clone(),
720 version,
721 value_bytes: value_bytes_of(&value),
722 current: false,
723 });
724 }
725 versions.is_empty()
726 } else {
727 false
728 };
729 if now_empty {
730 historical.remove(&key);
731 }
732 }
733 }
734
735 let new_oldest = oldest_version(¤t, &historical, &key);
736 reconcile_oldest(&mut oldest, &key, old_oldest, new_oldest);
737 }
738 }
739
740 Span::current().record("total_entry_count", total_entries);
741 Ok(removed)
742 }
743
744 #[instrument(name = "store::multi::memory::get_all_versions", level = "trace", skip(self, key), fields(table = ?table, key_len = key.len()))]
745 pub fn get_all_versions(
746 &self,
747 table: EntryKind,
748 key: &[u8],
749 ) -> Result<Vec<(CommitVersion, Option<CowVec<u8>>)>> {
750 let entry = match self.inner.entries.data.get(&table) {
751 Some(e) => e,
752 None => return Ok(Vec::new()),
753 };
754
755 let current = entry.current.read();
756 let current_hit = current.get(key).map(|(cur_version, value)| (*cur_version, value.clone()));
757 drop(current);
758
759 let historical = entry.historical.read();
760 let hist_versions = historical.get(key);
761
762 let mut versions: Vec<(CommitVersion, Option<CowVec<u8>>)> =
763 Vec::with_capacity(current_hit.is_some() as usize + hist_versions.map_or(0, |v| v.len()));
764 if let Some(hit) = current_hit {
765 versions.push(hit);
766 }
767 if let Some(hist_versions) = hist_versions {
768 for (Reverse(v), value) in hist_versions.iter() {
769 versions.push((*v, value.clone()));
770 }
771 }
772
773 versions.sort_by(|a, b| b.0.cmp(&a.0));
774
775 Ok(versions)
776 }
777
778 #[instrument(name = "store::multi::memory::scan_historical_below", level = "trace", skip(self, cursor), fields(table = ?table, cutoff = cutoff.0, batch_size = batch_size))]
779 pub fn scan_historical_below(
780 &self,
781 table: EntryKind,
782 cutoff: CommitVersion,
783 cursor: &mut HistoricalCursor,
784 batch_size: usize,
785 ) -> Result<Vec<(EncodedKey, CommitVersion)>> {
786 if cursor.exhausted || batch_size == 0 {
787 return Ok(Vec::new());
788 }
789
790 let entry = match self.inner.entries.data.get(&table) {
791 Some(e) => e,
792 None => {
793 cursor.exhausted = true;
794 return Ok(Vec::new());
795 }
796 };
797
798 let historical = entry.historical.read();
799
800 let mut collected: Vec<(EncodedKey, CommitVersion)> = Vec::new();
801 let mut over_limit = false;
802
803 for (key, versions) in historical.iter() {
804 match (cursor.last_key.as_ref(), cursor.last_version) {
805 (Some(lk), _) if key < lk => continue,
806 (Some(lk), Some(lv)) if key == lk => {
807 for (Reverse(v), _value) in versions.iter().rev() {
808 if *v <= lv {
809 continue;
810 }
811 if *v >= cutoff {
812 continue;
813 }
814 collected.push((key.clone(), *v));
815 if collected.len() > batch_size {
816 over_limit = true;
817 break;
818 }
819 }
820 }
821 _ => {
822 for (Reverse(v), _value) in versions.iter().rev() {
823 if *v >= cutoff {
824 continue;
825 }
826 collected.push((key.clone(), *v));
827 if collected.len() > batch_size {
828 over_limit = true;
829 break;
830 }
831 }
832 }
833 }
834
835 if over_limit {
836 break;
837 }
838 }
839
840 collected.sort_by(|a, b| a.0.as_slice().cmp(b.0.as_slice()).then(a.1.0.cmp(&b.1.0)));
841
842 let has_more = collected.len() > batch_size;
843 if has_more {
844 collected.truncate(batch_size);
845 }
846
847 if let Some(last) = collected.last() {
848 cursor.last_key = Some(last.0.clone());
849 cursor.last_version = Some(last.1);
850 }
851 if !has_more {
852 cursor.exhausted = true;
853 }
854
855 Ok(collected)
856 }
857}
858
859impl TierBackend for MemoryRowStorage {}
860
861#[cfg(test)]
862pub mod tests {
863 use reifydb_core::interface::catalog::{id::TableId, storage::StorageId};
864
865 use super::*;
866
867 #[test]
868 fn test_basic_operations() {
869 let storage = MemoryRowStorage::new();
870
871 let key = EncodedKey::new(b"key1");
872 let version = CommitVersion(1);
873
874 storage.set(
875 version,
876 HashMap::from([(EntryKind::Multi, vec![(key.clone(), Some(CowVec::new(b"value1".to_vec())))])]),
877 )
878 .unwrap();
879
880 let value = storage.get(EntryKind::Multi, &key, version).unwrap().value();
881 assert_eq!(value.as_deref(), Some(b"value1".as_slice()));
882
883 assert!(storage.contains(EntryKind::Multi, &key, version).unwrap());
884
885 assert!(!storage.contains(EntryKind::Multi, b"nonexistent", version).unwrap());
886
887 let version2 = CommitVersion(2);
888 storage.set(version2, HashMap::from([(EntryKind::Multi, vec![(key.clone(), None)])])).unwrap();
889 assert!(!storage.contains(EntryKind::Multi, &key, version2).unwrap());
890 }
891
892 #[test]
893 fn test_source_tables() {
894 let storage = MemoryRowStorage::new();
895
896 let source1 = StorageId::Table(TableId(1));
897 let source2 = StorageId::Table(TableId(2));
898
899 let key = EncodedKey::new(b"key");
900 let version = CommitVersion(1);
901
902 storage.set(
903 version,
904 HashMap::from([(
905 EntryKind::Source(source1),
906 vec![(key.clone(), Some(CowVec::new(b"table1".to_vec())))],
907 )]),
908 )
909 .unwrap();
910 storage.set(
911 version,
912 HashMap::from([(
913 EntryKind::Source(source2),
914 vec![(key.clone(), Some(CowVec::new(b"table2".to_vec())))],
915 )]),
916 )
917 .unwrap();
918
919 assert_eq!(
920 storage.get(EntryKind::Source(source1), &key, version).unwrap().value().as_deref(),
921 Some(b"table1".as_slice())
922 );
923 assert_eq!(
924 storage.get(EntryKind::Source(source2), &key, version).unwrap().value().as_deref(),
925 Some(b"table2".as_slice())
926 );
927 }
928
929 #[test]
930 fn test_version_promotion_to_historical() {
931 let storage = MemoryRowStorage::new();
932
933 let key = EncodedKey::new(b"key1");
934
935 storage.set(
936 CommitVersion(1),
937 HashMap::from([(EntryKind::Multi, vec![(key.clone(), Some(CowVec::new(b"v1".to_vec())))])]),
938 )
939 .unwrap();
940
941 storage.set(
942 CommitVersion(2),
943 HashMap::from([(EntryKind::Multi, vec![(key.clone(), Some(CowVec::new(b"v2".to_vec())))])]),
944 )
945 .unwrap();
946
947 storage.set(
948 CommitVersion(3),
949 HashMap::from([(EntryKind::Multi, vec![(key.clone(), Some(CowVec::new(b"v3".to_vec())))])]),
950 )
951 .unwrap();
952
953 assert_eq!(
954 storage.get(EntryKind::Multi, &key, CommitVersion(3)).unwrap().value().as_deref(),
955 Some(b"v3".as_slice())
956 );
957
958 assert_eq!(
959 storage.get(EntryKind::Multi, &key, CommitVersion(2)).unwrap().value().as_deref(),
960 Some(b"v2".as_slice())
961 );
962
963 assert_eq!(
964 storage.get(EntryKind::Multi, &key, CommitVersion(1)).unwrap().value().as_deref(),
965 Some(b"v1".as_slice())
966 );
967 }
968
969 #[test]
970 fn test_insert_older_version() {
971 let storage = MemoryRowStorage::new();
974
975 let key = EncodedKey::new(b"key1");
976
977 storage.set(
978 CommitVersion(3),
979 HashMap::from([(EntryKind::Multi, vec![(key.clone(), Some(CowVec::new(b"v3".to_vec())))])]),
980 )
981 .unwrap();
982
983 storage.set(
984 CommitVersion(1),
985 HashMap::from([(EntryKind::Multi, vec![(key.clone(), Some(CowVec::new(b"v1".to_vec())))])]),
986 )
987 .unwrap();
988
989 assert_eq!(
990 storage.get(EntryKind::Multi, &key, CommitVersion(3)).unwrap().value().as_deref(),
991 Some(b"v3".as_slice())
992 );
993
994 assert_eq!(
995 storage.get(EntryKind::Multi, &key, CommitVersion(1)).unwrap().value().as_deref(),
996 Some(b"v1".as_slice())
997 );
998
999 assert_eq!(
1000 storage.get(EntryKind::Multi, &key, CommitVersion(2)).unwrap().value().as_deref(),
1001 Some(b"v1".as_slice())
1002 );
1003 }
1004
1005 #[test]
1006 fn test_range_next() {
1007 let storage = MemoryRowStorage::new();
1008
1009 let version = CommitVersion(1);
1010 storage.set(
1011 version,
1012 HashMap::from([(
1013 EntryKind::Multi,
1014 vec![
1015 (EncodedKey::new(b"a"), Some(CowVec::new(b"1".to_vec()))),
1016 (EncodedKey::new(b"b"), Some(CowVec::new(b"2".to_vec()))),
1017 (EncodedKey::new(b"c"), Some(CowVec::new(b"3".to_vec()))),
1018 ],
1019 )]),
1020 )
1021 .unwrap();
1022
1023 let mut cursor = RangeCursor::new();
1024 let batch = storage
1025 .range_next(
1026 EntryKind::Multi,
1027 &mut cursor,
1028 Bound::Unbounded,
1029 Bound::Unbounded,
1030 MultiVersionScope::AsOf {
1031 read: version,
1032 },
1033 100,
1034 )
1035 .unwrap();
1036
1037 assert_eq!(batch.entries.len(), 3);
1038 assert!(!batch.has_more);
1039 assert!(cursor.exhausted);
1040
1041 assert_eq!(&*batch.entries[0].key, b"a");
1042 assert_eq!(&*batch.entries[1].key, b"b");
1043 assert_eq!(&*batch.entries[2].key, b"c");
1044 }
1045
1046 #[test]
1047 fn test_range_rev_next() {
1048 let storage = MemoryRowStorage::new();
1049
1050 let version = CommitVersion(1);
1051 storage.set(
1052 version,
1053 HashMap::from([(
1054 EntryKind::Multi,
1055 vec![
1056 (EncodedKey::new(b"a"), Some(CowVec::new(b"1".to_vec()))),
1057 (EncodedKey::new(b"b"), Some(CowVec::new(b"2".to_vec()))),
1058 (EncodedKey::new(b"c"), Some(CowVec::new(b"3".to_vec()))),
1059 ],
1060 )]),
1061 )
1062 .unwrap();
1063
1064 let mut cursor = RangeCursor::new();
1065 let batch = storage
1066 .range_rev_next(
1067 EntryKind::Multi,
1068 &mut cursor,
1069 Bound::Unbounded,
1070 Bound::Unbounded,
1071 MultiVersionScope::AsOf {
1072 read: version,
1073 },
1074 100,
1075 )
1076 .unwrap();
1077
1078 assert_eq!(batch.entries.len(), 3);
1079 assert!(!batch.has_more);
1080 assert!(cursor.exhausted);
1081
1082 assert_eq!(&*batch.entries[0].key, b"c");
1083 assert_eq!(&*batch.entries[1].key, b"b");
1084 assert_eq!(&*batch.entries[2].key, b"a");
1085 }
1086
1087 #[test]
1088 fn test_range_streaming_pagination() {
1089 let storage = MemoryRowStorage::new();
1090
1091 let version = CommitVersion(1);
1092
1093 let entries: Vec<_> =
1094 (0..10u8).map(|i| (EncodedKey::new(vec![i]), Some(CowVec::new(vec![i * 10])))).collect();
1095 storage.set(version, HashMap::from([(EntryKind::Multi, entries)])).unwrap();
1096
1097 let mut cursor = RangeCursor::new();
1098
1099 let batch1 = storage
1100 .range_next(
1101 EntryKind::Multi,
1102 &mut cursor,
1103 Bound::Unbounded,
1104 Bound::Unbounded,
1105 MultiVersionScope::AsOf {
1106 read: version,
1107 },
1108 3,
1109 )
1110 .unwrap();
1111 assert_eq!(batch1.entries.len(), 3);
1112 assert!(batch1.has_more);
1113 assert!(!cursor.exhausted);
1114
1115 assert_eq!(&*batch1.entries[0].key, &[0]);
1116 assert_eq!(&*batch1.entries[2].key, &[2]);
1117
1118 let batch2 = storage
1119 .range_next(
1120 EntryKind::Multi,
1121 &mut cursor,
1122 Bound::Unbounded,
1123 Bound::Unbounded,
1124 MultiVersionScope::AsOf {
1125 read: version,
1126 },
1127 3,
1128 )
1129 .unwrap();
1130 assert_eq!(batch2.entries.len(), 3);
1131 assert!(batch2.has_more);
1132 assert!(!cursor.exhausted);
1133
1134 assert_eq!(&*batch2.entries[0].key, &[3]);
1135 assert_eq!(&*batch2.entries[2].key, &[5]);
1136
1137 let batch3 = storage
1138 .range_next(
1139 EntryKind::Multi,
1140 &mut cursor,
1141 Bound::Unbounded,
1142 Bound::Unbounded,
1143 MultiVersionScope::AsOf {
1144 read: version,
1145 },
1146 3,
1147 )
1148 .unwrap();
1149 assert_eq!(batch3.entries.len(), 3);
1150 assert!(batch3.has_more);
1151 assert!(!cursor.exhausted);
1152
1153 assert_eq!(&*batch3.entries[0].key, &[6]);
1154 assert_eq!(&*batch3.entries[2].key, &[8]);
1155
1156 let batch4 = storage
1157 .range_next(
1158 EntryKind::Multi,
1159 &mut cursor,
1160 Bound::Unbounded,
1161 Bound::Unbounded,
1162 MultiVersionScope::AsOf {
1163 read: version,
1164 },
1165 3,
1166 )
1167 .unwrap();
1168 assert_eq!(batch4.entries.len(), 1);
1169 assert!(!batch4.has_more);
1170 assert!(cursor.exhausted);
1171
1172 assert_eq!(&*batch4.entries[0].key, &[9]);
1173
1174 let batch5 = storage
1175 .range_next(
1176 EntryKind::Multi,
1177 &mut cursor,
1178 Bound::Unbounded,
1179 Bound::Unbounded,
1180 MultiVersionScope::AsOf {
1181 read: version,
1182 },
1183 3,
1184 )
1185 .unwrap();
1186 assert!(batch5.entries.is_empty());
1187 }
1188
1189 #[test]
1190 fn test_range_reving_pagination() {
1191 let storage = MemoryRowStorage::new();
1192
1193 let version = CommitVersion(1);
1194
1195 let entries: Vec<_> =
1196 (0..10u8).map(|i| (EncodedKey::new(vec![i]), Some(CowVec::new(vec![i * 10])))).collect();
1197 storage.set(version, HashMap::from([(EntryKind::Multi, entries)])).unwrap();
1198
1199 let mut cursor = RangeCursor::new();
1200
1201 let batch1 = storage
1202 .range_rev_next(
1203 EntryKind::Multi,
1204 &mut cursor,
1205 Bound::Unbounded,
1206 Bound::Unbounded,
1207 MultiVersionScope::AsOf {
1208 read: version,
1209 },
1210 3,
1211 )
1212 .unwrap();
1213 assert_eq!(batch1.entries.len(), 3);
1214 assert!(batch1.has_more);
1215 assert!(!cursor.exhausted);
1216
1217 assert_eq!(&*batch1.entries[0].key, &[9]);
1218 assert_eq!(&*batch1.entries[2].key, &[7]);
1219
1220 let batch2 = storage
1221 .range_rev_next(
1222 EntryKind::Multi,
1223 &mut cursor,
1224 Bound::Unbounded,
1225 Bound::Unbounded,
1226 MultiVersionScope::AsOf {
1227 read: version,
1228 },
1229 3,
1230 )
1231 .unwrap();
1232 assert_eq!(batch2.entries.len(), 3);
1233 assert!(batch2.has_more);
1234 assert!(!cursor.exhausted);
1235
1236 assert_eq!(&*batch2.entries[0].key, &[6]);
1237 assert_eq!(&*batch2.entries[2].key, &[4]);
1238 }
1239
1240 #[test]
1241 fn test_drop_from_historical() {
1242 let storage = MemoryRowStorage::new();
1243
1244 let key = EncodedKey::new(b"key1");
1245
1246 for v in 1..=3u64 {
1247 storage.set(
1248 CommitVersion(v),
1249 HashMap::from([(
1250 EntryKind::Multi,
1251 vec![(key.clone(), Some(CowVec::new(format!("v{}", v).into_bytes())))],
1252 )]),
1253 )
1254 .unwrap();
1255 }
1256
1257 storage.compact(HashMap::from([(EntryKind::Multi, vec![(key.clone(), CommitVersion(1))])])).unwrap();
1258
1259 assert!(storage.get(EntryKind::Multi, &key, CommitVersion(1)).unwrap().value().is_none());
1260
1261 assert_eq!(
1262 storage.get(EntryKind::Multi, &key, CommitVersion(2)).unwrap().value().as_deref(),
1263 Some(b"v2".as_slice())
1264 );
1265 assert_eq!(
1266 storage.get(EntryKind::Multi, &key, CommitVersion(3)).unwrap().value().as_deref(),
1267 Some(b"v3".as_slice())
1268 );
1269 }
1270 #[test]
1271 fn compact_returns_each_removed_historical_version_flagged_as_not_current() {
1272 let storage = MemoryRowStorage::new();
1275
1276 let key = EncodedKey::new(b"key1");
1277
1278 for v in 1..=3u64 {
1279 storage.set(
1280 CommitVersion(v),
1281 HashMap::from([(
1282 EntryKind::Multi,
1283 vec![(key.clone(), Some(CowVec::new(format!("v{}", v).into_bytes())))],
1284 )]),
1285 )
1286 .unwrap();
1287 }
1288
1289 let removed = storage
1290 .compact(HashMap::from([(
1291 EntryKind::Multi,
1292 vec![(key.clone(), CommitVersion(1)), (key.clone(), CommitVersion(2))],
1293 )]))
1294 .unwrap();
1295
1296 let mut versions: Vec<u64> = removed.iter().map(|entry| entry.version.0).collect();
1297 versions.sort_unstable();
1298 assert_eq!(versions, vec![1, 2]);
1299 assert!(removed.iter().all(|entry| !entry.current));
1300 assert!(removed.iter().all(|entry| entry.value_bytes == ByteSize::from_bytes(2)));
1301 }
1302
1303 #[test]
1304 fn compact_reports_the_live_version_and_leaves_surviving_history_in_place() {
1305 let storage = MemoryRowStorage::new();
1309
1310 let key = EncodedKey::new(b"key1");
1311
1312 for v in 1..=3u64 {
1313 storage.set(
1314 CommitVersion(v),
1315 HashMap::from([(
1316 EntryKind::Multi,
1317 vec![(key.clone(), Some(CowVec::new(format!("v{}", v).into_bytes())))],
1318 )]),
1319 )
1320 .unwrap();
1321 }
1322
1323 let removed = storage
1324 .compact(HashMap::from([(EntryKind::Multi, vec![(key.clone(), CommitVersion(3))])]))
1325 .unwrap();
1326
1327 assert_eq!(removed.len(), 1, "only the live version was dropped");
1328 assert_eq!(removed[0].version, CommitVersion(3));
1329 assert!(removed[0].current, "the dropped version was the live one");
1330 assert_eq!(removed[0].value_bytes, ByteSize::from_bytes(2));
1331
1332 assert_eq!(
1333 storage.get(EntryKind::Multi, &key, CommitVersion(3)).unwrap().value().as_deref(),
1334 Some(b"v2".as_slice()),
1335 "the newest survivor is still readable from historical without being promoted"
1336 );
1337 }
1338
1339 #[test]
1340 fn compact_reports_the_live_entry_when_every_version_of_a_key_is_removed() {
1341 let storage = MemoryRowStorage::new();
1345
1346 let key = EncodedKey::new(b"key1");
1347
1348 for v in 1..=3u64 {
1349 storage.set(
1350 CommitVersion(v),
1351 HashMap::from([(
1352 EntryKind::Multi,
1353 vec![(key.clone(), Some(CowVec::new(format!("v{}", v).into_bytes())))],
1354 )]),
1355 )
1356 .unwrap();
1357 }
1358
1359 let removed = storage
1360 .compact(HashMap::from([(
1361 EntryKind::Multi,
1362 vec![
1363 (key.clone(), CommitVersion(1)),
1364 (key.clone(), CommitVersion(2)),
1365 (key.clone(), CommitVersion(3)),
1366 ],
1367 )]))
1368 .unwrap();
1369
1370 assert_eq!(removed.len(), 3);
1371
1372 let live: Vec<u64> =
1373 removed.iter().filter(|entry| entry.current).map(|entry| entry.version.0).collect();
1374 assert_eq!(live, vec![3]);
1375
1376 let mut historical: Vec<u64> =
1377 removed.iter().filter(|entry| !entry.current).map(|entry| entry.version.0).collect();
1378 historical.sort_unstable();
1379 assert_eq!(historical, vec![1, 2]);
1380 }
1381
1382 #[test]
1383 fn test_tombstones() {
1384 let storage = MemoryRowStorage::new();
1385
1386 let key = EncodedKey::new(b"key1");
1387
1388 storage.set(
1389 CommitVersion(1),
1390 HashMap::from([(EntryKind::Multi, vec![(key.clone(), Some(CowVec::new(b"value".to_vec())))])]),
1391 )
1392 .unwrap();
1393
1394 storage.set(CommitVersion(2), HashMap::from([(EntryKind::Multi, vec![(key.clone(), None)])])).unwrap();
1395
1396 assert!(storage.get(EntryKind::Multi, &key, CommitVersion(2)).unwrap().value().is_none());
1397 assert!(!storage.contains(EntryKind::Multi, &key, CommitVersion(2)).unwrap());
1398
1399 assert_eq!(
1400 storage.get(EntryKind::Multi, &key, CommitVersion(1)).unwrap().value().as_deref(),
1401 Some(b"value".as_slice())
1402 );
1403 }
1404
1405 #[test]
1406 fn test_collect_evictable_below_keeps_versions_above_cutoff() {
1407 let storage = MemoryRowStorage::new();
1408 let key = EncodedKey::new(b"k");
1409 for v in 1..=3u64 {
1410 storage.set(
1411 CommitVersion(v),
1412 HashMap::from([(
1413 EntryKind::Multi,
1414 vec![(key.clone(), Some(CowVec::new(format!("v{v}").into_bytes())))],
1415 )]),
1416 )
1417 .unwrap();
1418 }
1419
1420 let (to_persist, to_drop, _more) =
1422 storage.collect_evictable_below(EntryKind::Multi, CommitVersion(2), usize::MAX);
1423 assert_eq!(to_persist.len(), 1);
1424 assert_eq!(to_persist[0].0, key);
1425 assert_eq!(to_persist[0].1, CommitVersion(2));
1426 assert_eq!(to_persist[0].2.as_deref(), Some(b"v2".as_slice()));
1427 let dropped: HashSet<CommitVersion> = to_drop.iter().map(|e| e.version).collect();
1428 assert_eq!(dropped, HashSet::from([CommitVersion(1), CommitVersion(2)]));
1429
1430 storage.compact(HashMap::from([(
1431 EntryKind::Multi,
1432 to_drop.into_iter().map(|e| (e.key, e.version)).collect(),
1433 )]))
1434 .unwrap();
1435 assert_eq!(
1436 storage.get(EntryKind::Multi, &key, CommitVersion(3)).unwrap().value().as_deref(),
1437 Some(b"v3".as_slice())
1438 );
1439 assert!(storage.get(EntryKind::Multi, &key, CommitVersion(2)).unwrap().value().is_none());
1440 assert!(storage.get(EntryKind::Multi, &key, CommitVersion(1)).unwrap().value().is_none());
1441 }
1442
1443 #[test]
1444 fn test_collect_evictable_below_empty_when_all_above_cutoff() {
1445 let storage = MemoryRowStorage::new();
1446 let key = EncodedKey::new(b"k");
1447 storage.set(
1448 CommitVersion(5),
1449 HashMap::from([(EntryKind::Multi, vec![(key.clone(), Some(CowVec::new(b"v".to_vec())))])]),
1450 )
1451 .unwrap();
1452 let (to_persist, to_drop, _more) =
1453 storage.collect_evictable_below(EntryKind::Multi, CommitVersion(3), usize::MAX);
1454 assert!(to_persist.is_empty());
1455 assert!(to_drop.is_empty());
1456 }
1457
1458 #[test]
1459 fn test_collect_evictable_below_persists_exactly_one_value_per_key() {
1460 let storage = MemoryRowStorage::new();
1463 let key = EncodedKey::new(b"k");
1464 for v in 1..=5u64 {
1465 storage.set(
1466 CommitVersion(v),
1467 HashMap::from([(
1468 EntryKind::Multi,
1469 vec![(key.clone(), Some(CowVec::new(format!("v{v}").into_bytes())))],
1470 )]),
1471 )
1472 .unwrap();
1473 }
1474
1475 let (to_persist, to_drop, _more) =
1476 storage.collect_evictable_below(EntryKind::Multi, CommitVersion(4), usize::MAX);
1477 assert_eq!(to_persist.len(), 1, "exactly one value persisted per key");
1478 assert_eq!(to_persist[0].1, CommitVersion(4), "the latest version <= cutoff");
1479 assert_eq!(to_persist[0].2.as_deref(), Some(b"v4".as_slice()));
1480
1481 let dropped: HashSet<CommitVersion> = to_drop.iter().map(|e| e.version).collect();
1482 assert_eq!(
1483 dropped,
1484 HashSet::from([CommitVersion(1), CommitVersion(2), CommitVersion(3), CommitVersion(4)])
1485 );
1486 }
1487
1488 #[test]
1489 fn test_collect_evictable_below_persists_tombstone_when_it_is_the_latest() {
1490 let storage = MemoryRowStorage::new();
1493 let key = EncodedKey::new(b"k");
1494 storage.set(
1495 CommitVersion(1),
1496 HashMap::from([(EntryKind::Multi, vec![(key.clone(), Some(CowVec::new(b"v1".to_vec())))])]),
1497 )
1498 .unwrap();
1499 storage.set(CommitVersion(2), HashMap::from([(EntryKind::Multi, vec![(key.clone(), None)])])).unwrap();
1500
1501 let (to_persist, to_drop, _more) =
1502 storage.collect_evictable_below(EntryKind::Multi, CommitVersion(2), usize::MAX);
1503 assert_eq!(to_persist.len(), 1);
1504 assert_eq!(to_persist[0].1, CommitVersion(2), "the tombstone is the latest version");
1505 assert!(to_persist[0].2.is_none(), "the persisted latest value must be the tombstone, not v1");
1506 assert_eq!(to_drop.len(), 2, "both v1 and the tombstone are dropped from the buffer");
1507 }
1508
1509 #[test]
1510 fn test_collect_evictable_below_only_drops_historical_when_current_is_above_cutoff() {
1511 let storage = MemoryRowStorage::new();
1514 let key = EncodedKey::new(b"k");
1515 storage.set(
1516 CommitVersion(2),
1517 HashMap::from([(EntryKind::Multi, vec![(key.clone(), Some(CowVec::new(b"v2".to_vec())))])]),
1518 )
1519 .unwrap();
1520 storage.set(
1521 CommitVersion(5),
1522 HashMap::from([(EntryKind::Multi, vec![(key.clone(), Some(CowVec::new(b"v5".to_vec())))])]),
1523 )
1524 .unwrap();
1525
1526 let (to_persist, to_drop, _more) =
1527 storage.collect_evictable_below(EntryKind::Multi, CommitVersion(3), usize::MAX);
1528 assert_eq!(to_persist.len(), 1);
1529 assert_eq!(to_persist[0].1, CommitVersion(2), "only the aged-out historical version is persisted");
1530 assert_eq!(to_persist[0].2.as_deref(), Some(b"v2".as_slice()));
1531 let dropped: HashSet<CommitVersion> = to_drop.iter().map(|e| e.version).collect();
1532 assert_eq!(dropped, HashSet::from([CommitVersion(2)]), "v5 (current, > cutoff) is never dropped");
1533
1534 storage.compact(HashMap::from([(
1535 EntryKind::Multi,
1536 to_drop.into_iter().map(|e| (e.key, e.version)).collect(),
1537 )]))
1538 .unwrap();
1539 assert_eq!(
1540 storage.get(EntryKind::Multi, &key, CommitVersion(5)).unwrap().value().as_deref(),
1541 Some(b"v5".as_slice())
1542 );
1543 assert!(
1544 storage.get(EntryKind::Multi, &key, CommitVersion(3)).unwrap().value().is_none(),
1545 "the v2 a reader at snapshot 3 used to see is gone from the buffer after eviction"
1546 );
1547 }
1548
1549 #[test]
1550 fn test_collect_evictable_below_handles_multiple_keys_independently() {
1551 let storage = MemoryRowStorage::new();
1554 let cold = EncodedKey::new(b"cold");
1555 let hot = EncodedKey::new(b"hot");
1556 storage.set(
1557 CommitVersion(1),
1558 HashMap::from([(EntryKind::Multi, vec![(cold.clone(), Some(CowVec::new(b"cold1".to_vec())))])]),
1559 )
1560 .unwrap();
1561 storage.set(
1562 CommitVersion(9),
1563 HashMap::from([(EntryKind::Multi, vec![(hot.clone(), Some(CowVec::new(b"hot9".to_vec())))])]),
1564 )
1565 .unwrap();
1566
1567 let (to_persist, to_drop, _more) =
1568 storage.collect_evictable_below(EntryKind::Multi, CommitVersion(5), usize::MAX);
1569 assert_eq!(to_persist.len(), 1, "only the cold key is evictable below the cutoff");
1570 assert_eq!(to_persist[0].0, cold);
1571 assert!(to_drop.iter().all(|e| e.key == cold), "the hot key must not be scheduled for drop");
1572 }
1573
1574 #[test]
1575 fn test_collect_evictable_below_bounds_to_budget_and_drains_across_calls() {
1576 let storage = MemoryRowStorage::new();
1579 for i in 0..5u64 {
1580 let key = EncodedKey::new(format!("k{i}").into_bytes());
1581 storage.set(
1582 CommitVersion(1),
1583 HashMap::from([(EntryKind::Multi, vec![(key, Some(CowVec::new(vec![i as u8])))])]),
1584 )
1585 .unwrap();
1586 }
1587
1588 let (to_persist, to_drop, more) =
1589 storage.collect_evictable_below(EntryKind::Multi, CommitVersion(1), 2);
1590 assert_eq!(to_persist.len(), 2, "budget caps the collected key count");
1591 assert_eq!(to_drop.len(), 2);
1592 assert!(more, "three keys remain below the cutoff");
1593
1594 let mut drained = to_persist.len();
1595 let mut compaction_batch: HashMap<EntryKind, Vec<(EncodedKey, CommitVersion)>> = HashMap::new();
1596 compaction_batch.insert(EntryKind::Multi, to_drop.into_iter().map(|e| (e.key, e.version)).collect());
1597 storage.compact(compaction_batch).unwrap();
1598 loop {
1599 let (p, d, more) = storage.collect_evictable_below(EntryKind::Multi, CommitVersion(1), 2);
1600 if p.is_empty() {
1601 assert!(!more, "an empty collect must not claim more remains");
1602 break;
1603 }
1604 drained += p.len();
1605 let mut batch: HashMap<EntryKind, Vec<(EncodedKey, CommitVersion)>> = HashMap::new();
1606 batch.insert(EntryKind::Multi, d.into_iter().map(|e| (e.key, e.version)).collect());
1607 storage.compact(batch).unwrap();
1608 if !more {
1609 break;
1610 }
1611 }
1612 assert_eq!(drained, 5, "every below-cutoff key is drained exactly once");
1613 }
1614
1615 fn indexed_oldest(storage: &MemoryRowStorage, table: EntryKind, key: &EncodedKey) -> Option<CommitVersion> {
1616 let entry = storage.inner.entries.data.get(&table)?;
1617 let oldest = entry.oldest.read();
1618 oldest.iter().find(|(_, keys)| keys.contains(key)).map(|(v, _)| *v)
1619 }
1620
1621 fn assert_index_consistent(storage: &MemoryRowStorage, table: EntryKind) {
1622 let entry = storage.inner.entries.data.get(&table).expect("table exists");
1623 let current = entry.current.read();
1624 let historical = entry.historical.read();
1625 let oldest = entry.oldest.read();
1626
1627 let mut resident: HashSet<EncodedKey> = HashSet::new();
1630 resident.extend(current.keys().cloned());
1631 resident.extend(historical.keys().cloned());
1632 for key in &resident {
1633 let expected = oldest_version(¤t, &historical, key);
1634 let indexed = oldest.iter().find(|(_, keys)| keys.contains(key)).map(|(v, _)| *v);
1635 assert_eq!(
1636 indexed, expected,
1637 "a resident key must sit in the index bucket of its smallest version"
1638 );
1639 }
1640
1641 for (bucket, keys) in oldest.iter() {
1643 for key in keys {
1644 assert!(
1645 resident.contains(key),
1646 "index holds a key that is in neither map (stale entry)"
1647 );
1648 assert_eq!(
1649 oldest_version(¤t, &historical, key),
1650 Some(*bucket),
1651 "index bucket must equal the key's smallest stored version"
1652 );
1653 }
1654 }
1655 }
1656
1657 #[test]
1658 fn index_stays_consistent_across_new_monotonic_out_of_order_and_drops() {
1659 let storage = MemoryRowStorage::new();
1663 let kind = EntryKind::Multi;
1664 let a = EncodedKey::new(b"a");
1665 let b = EncodedKey::new(b"b");
1666 let c = EncodedKey::new(b"c");
1667
1668 let set = |v: u64, key: &EncodedKey, val: &str| {
1669 storage.set(
1670 CommitVersion(v),
1671 HashMap::from([(
1672 kind,
1673 vec![(key.clone(), Some(CowVec::new(val.as_bytes().to_vec())))],
1674 )]),
1675 )
1676 .unwrap();
1677 };
1678
1679 set(10, &a, "a10");
1680 set(20, &a, "a20");
1681 set(3, &a, "a3");
1682 set(5, &b, "b5");
1683 set(7, &c, "c7");
1684 assert_eq!(
1685 indexed_oldest(&storage, kind, &a),
1686 Some(CommitVersion(3)),
1687 "an out-of-order write below the current version must lower a's bucket to 3"
1688 );
1689 assert_index_consistent(&storage, kind);
1690
1691 storage.compact(HashMap::from([(kind, vec![(a.clone(), CommitVersion(3))])])).unwrap();
1692 assert_eq!(
1693 indexed_oldest(&storage, kind, &a),
1694 Some(CommitVersion(10)),
1695 "dropping the oldest version must raise the bucket to the next-smallest stored version"
1696 );
1697 assert_index_consistent(&storage, kind);
1698
1699 storage.compact(HashMap::from([(kind, vec![(b.clone(), CommitVersion(5))])])).unwrap();
1700 assert_eq!(
1701 indexed_oldest(&storage, kind, &b),
1702 None,
1703 "a fully dropped key must leave the index entirely"
1704 );
1705 assert_index_consistent(&storage, kind);
1706 }
1707
1708 #[test]
1709 fn out_of_order_landing_is_selected_for_eviction() {
1710 let storage = MemoryRowStorage::new();
1713 let key = EncodedKey::new(b"k");
1714 storage.set(
1715 CommitVersion(20),
1716 HashMap::from([(EntryKind::Multi, vec![(key.clone(), Some(CowVec::new(b"v20".to_vec())))])]),
1717 )
1718 .unwrap();
1719 storage.set(
1720 CommitVersion(3),
1721 HashMap::from([(EntryKind::Multi, vec![(key.clone(), Some(CowVec::new(b"v3".to_vec())))])]),
1722 )
1723 .unwrap();
1724
1725 let (to_persist, to_drop, _more) =
1726 storage.collect_evictable_below(EntryKind::Multi, CommitVersion(5), usize::MAX);
1727 let dropped: HashSet<CommitVersion> = to_drop.iter().map(|e| e.version).collect();
1728 assert_eq!(
1729 dropped,
1730 HashSet::from([CommitVersion(3)]),
1731 "the out-of-order v3 must be selected; v20 stays resident"
1732 );
1733 assert_eq!(to_persist.len(), 1);
1734 assert_eq!(to_persist[0].1, CommitVersion(3), "the aged-out v3 is the value persisted");
1735 }
1736
1737 #[test]
1738 fn byte_tally_matches_a_full_walk_across_mixed_mutations() {
1739 let storage = MemoryRowStorage::new();
1742 let k1 = EncodedKey::new(b"key-one");
1743 let k2 = EncodedKey::new(b"key-two");
1744
1745 for v in 1..=3u64 {
1746 storage.set(
1747 CommitVersion(v),
1748 HashMap::from([(
1749 EntryKind::Multi,
1750 vec![(k1.clone(), Some(CowVec::new(format!("value-{v}").into_bytes())))],
1751 )]),
1752 )
1753 .unwrap();
1754 }
1755 storage.set(
1756 CommitVersion(5),
1757 HashMap::from([(EntryKind::Multi, vec![(k2.clone(), Some(CowVec::new(b"x".to_vec())))])]),
1758 )
1759 .unwrap();
1760 storage.set(
1761 CommitVersion(2),
1762 HashMap::from([(EntryKind::Multi, vec![(k2.clone(), Some(CowVec::new(b"older".to_vec())))])]),
1763 )
1764 .unwrap();
1765 storage.set(CommitVersion(6), HashMap::from([(EntryKind::Multi, vec![(k2.clone(), None)])])).unwrap();
1766
1767 storage.compact(HashMap::from([(EntryKind::Multi, vec![(k1.clone(), CommitVersion(3))])])).unwrap();
1768 storage.compact(HashMap::from([(EntryKind::Multi, vec![(k2.clone(), CommitVersion(2))])])).unwrap();
1769
1770 let entry = storage.inner.entries.data.get(&EntryKind::Multi).unwrap();
1771 let current = entry.current.read();
1772 let walked_current: u64 = current.iter().map(|(k, (_, v))| entry_bytes(k, v)).sum();
1773 drop(current);
1774 let historical = entry.historical.read();
1775 let walked_historical: u64 = historical
1776 .iter()
1777 .map(|(k, versions)| versions.values().map(|v| entry_bytes(k, v)).sum::<u64>())
1778 .sum();
1779 drop(historical);
1780
1781 assert!(walked_current > 0, "precondition: the scenario must leave current entries behind");
1782 assert!(walked_historical > 0, "precondition: the scenario must leave historical entries behind");
1783 assert_eq!(
1784 storage.current_resident_bytes().as_bytes(),
1785 walked_current,
1786 "the incremental current tally must equal an exhaustive walk of the current map"
1787 );
1788 assert_eq!(
1789 storage.historical_resident_bytes().as_bytes(),
1790 walked_historical,
1791 "the incremental historical tally must equal an exhaustive walk of the historical map"
1792 );
1793 }
1794
1795 #[test]
1796 fn byte_tally_nets_to_zero_when_the_buffer_is_fully_drained() {
1797 let storage = MemoryRowStorage::new();
1800 let key = EncodedKey::new(b"k");
1801 for v in 1..=4u64 {
1802 storage.set(
1803 CommitVersion(v),
1804 HashMap::from([(
1805 EntryKind::Multi,
1806 vec![(key.clone(), Some(CowVec::new(format!("v{v}").into_bytes())))],
1807 )]),
1808 )
1809 .unwrap();
1810 }
1811 assert!(storage.current_resident_bytes().as_bytes() > 0);
1812 assert!(storage.historical_resident_bytes().as_bytes() > 0);
1813
1814 storage.compact(HashMap::from([(EntryKind::Multi, vec![(key.clone(), CommitVersion(4))])])).unwrap();
1815 storage.compact(HashMap::from([(
1816 EntryKind::Multi,
1817 vec![
1818 (key.clone(), CommitVersion(1)),
1819 (key.clone(), CommitVersion(2)),
1820 (key.clone(), CommitVersion(3)),
1821 ],
1822 )]))
1823 .unwrap();
1824
1825 assert_eq!(
1826 storage.current_resident_bytes(),
1827 ByteSize::ZERO,
1828 "draining every entry must return the current tally to zero"
1829 );
1830 assert_eq!(
1831 storage.historical_resident_bytes(),
1832 ByteSize::ZERO,
1833 "draining every entry must return the historical tally to zero"
1834 );
1835 }
1836
1837 #[test]
1838 fn clear_table_resets_the_byte_tally() {
1839 let storage = MemoryRowStorage::new();
1840 let key = EncodedKey::new(b"k");
1841 storage.set(
1842 CommitVersion(1),
1843 HashMap::from([(EntryKind::Multi, vec![(key.clone(), Some(CowVec::new(b"v".to_vec())))])]),
1844 )
1845 .unwrap();
1846 storage.set(
1847 CommitVersion(2),
1848 HashMap::from([(EntryKind::Multi, vec![(key.clone(), Some(CowVec::new(b"w".to_vec())))])]),
1849 )
1850 .unwrap();
1851 assert!(storage.current_resident_bytes().as_bytes() > 0);
1852 assert!(storage.historical_resident_bytes().as_bytes() > 0);
1853
1854 storage.clear_table(EntryKind::Multi).unwrap();
1855 assert_eq!(
1856 storage.current_resident_bytes(),
1857 ByteSize::ZERO,
1858 "clearing a table must zero its byte tally, not leak it"
1859 );
1860 assert_eq!(storage.historical_resident_bytes(), ByteSize::ZERO);
1861 }
1862
1863 #[test]
1864 fn an_empty_buffer_has_no_oldest_pending_version() {
1865 let storage = MemoryRowStorage::new();
1869
1870 assert_eq!(storage.oldest_pending_version(), None);
1871 }
1872
1873 #[test]
1874 fn oldest_pending_version_is_the_minimum_across_every_entry_kind() {
1875 let storage = MemoryRowStorage::new();
1879 let key = EncodedKey::new(b"k");
1880
1881 storage.set(
1882 CommitVersion(40),
1883 HashMap::from([(EntryKind::Multi, vec![(key.clone(), Some(CowVec::new(b"late".to_vec())))])]),
1884 )
1885 .unwrap();
1886 storage.set(
1887 CommitVersion(7),
1888 HashMap::from([(
1889 EntryKind::Source(StorageId::Table(TableId(1))),
1890 vec![(key.clone(), Some(CowVec::new(b"early".to_vec())))],
1891 )]),
1892 )
1893 .unwrap();
1894 storage.set(
1895 CommitVersion(19),
1896 HashMap::from([(
1897 EntryKind::Source(StorageId::Table(TableId(2))),
1898 vec![(key.clone(), Some(CowVec::new(b"mid".to_vec())))],
1899 )]),
1900 )
1901 .unwrap();
1902
1903 assert_eq!(
1904 storage.oldest_pending_version(),
1905 Some(CommitVersion(7)),
1906 "the oldest un-flushed write in any keyspace is what bounds the durable frontier"
1907 );
1908 }
1909
1910 #[test]
1911 fn the_oldest_bucket_in_a_keyspace_wins_over_its_newer_ones() {
1912 let storage = MemoryRowStorage::new();
1916
1917 storage.set(
1918 CommitVersion(11),
1919 HashMap::from([(
1920 EntryKind::Multi,
1921 vec![(EncodedKey::new(b"late"), Some(CowVec::new(b"v".to_vec())))],
1922 )]),
1923 )
1924 .unwrap();
1925 storage.set(
1926 CommitVersion(4),
1927 HashMap::from([(
1928 EntryKind::Multi,
1929 vec![(EncodedKey::new(b"early"), Some(CowVec::new(b"v".to_vec())))],
1930 )]),
1931 )
1932 .unwrap();
1933
1934 assert_eq!(storage.oldest_pending_version(), Some(CommitVersion(4)));
1935 }
1936
1937 #[test]
1938 fn a_superseded_version_still_counts_as_pending_until_it_is_compacted_away() {
1939 let storage = MemoryRowStorage::new();
1943 let key = EncodedKey::new(b"k");
1944
1945 storage.set(
1946 CommitVersion(3),
1947 HashMap::from([(EntryKind::Multi, vec![(key.clone(), Some(CowVec::new(b"v3".to_vec())))])]),
1948 )
1949 .unwrap();
1950 storage.set(
1951 CommitVersion(9),
1952 HashMap::from([(EntryKind::Multi, vec![(key.clone(), Some(CowVec::new(b"v9".to_vec())))])]),
1953 )
1954 .unwrap();
1955
1956 assert_eq!(storage.oldest_pending_version(), Some(CommitVersion(3)));
1957
1958 storage.compact(HashMap::from([(EntryKind::Multi, vec![(key.clone(), CommitVersion(3))])])).unwrap();
1959
1960 assert_eq!(
1961 storage.oldest_pending_version(),
1962 Some(CommitVersion(9)),
1963 "once the sweep drains v3 the frontier may advance to the next un-flushed write"
1964 );
1965 }
1966
1967 #[test]
1968 fn draining_the_buffer_clears_the_oldest_pending_version() {
1969 let storage = MemoryRowStorage::new();
1973 let key = EncodedKey::new(b"k");
1974
1975 storage.set(
1976 CommitVersion(5),
1977 HashMap::from([(EntryKind::Multi, vec![(key.clone(), Some(CowVec::new(b"v".to_vec())))])]),
1978 )
1979 .unwrap();
1980 storage.compact(HashMap::from([(EntryKind::Multi, vec![(key.clone(), CommitVersion(5))])])).unwrap();
1981
1982 assert_eq!(storage.oldest_pending_version(), None);
1983 }
1984}