Skip to main content

reifydb_store_multi/tier/commit/memory/
storage.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use 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(&current, &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(&current, &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(&current, &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		// An out-of-order older commit must stay resolvable: a read takes the largest version <= the
972		// snapshot, so the v2 snapshot resolves to v1.
973		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		// The storage metric only ever decrements from what compact reports back, so a version removed
1273		// physically but omitted from the return value inflates historical_count forever.
1274		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		// Dropping the live version does not promote the newest survivor, so the removal is only visible
1306		// to the metric through the returned record; the survivors must stay in historical, unreported
1307		// and still readable, which is what makes skipping the promotion safe.
1308		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		// The metric routes a current removal to the current counters and a historical one to the
1342		// historical counters, so mislabelling the live entry moves rows between the two columns
1343		// instead of clearing them.
1344		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		// v2 is what a reader in [2, 3) resolves to, so it is the value that must be persisted.
1421		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		// Only the latest-<=cutoff value may be persisted: it is the single value a reader at the cutoff
1461		// snapshot resolves to, so persisting an older one corrupts that resolution.
1462		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		// A tombstone that is the latest-<=cutoff version must be carried to the persistent tier; dropping
1491		// it lets a later read resurrect the pre-delete value.
1492		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		// A key that is actively written while old snapshots age out: only the historical version may be
1512		// evicted, the current one is still hot and must not be persisted.
1513		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		// The cutoff applies per version, not per key: a key whose only version is above it must stay
1552		// fully resident even while a sibling key is evicted.
1553		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		// The budget bounds a flush slice so one transaction never persists the whole evictable set;
1577		// looping bounded calls must still drain exactly the below-cutoff set, no more, no less.
1578		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		// A key indexed anywhere but its smallest stored version can never be selected, so its old
1628		// versions leak in the buffer forever.
1629		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(&current, &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		// A stale index entry (key in neither map) would be re-selected on every sweep and never clear.
1642		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(&current, &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		// The eviction index is what keeps collect_evictable_below O(evictable) instead of O(table);
1660		// drift from the maps either strands a key forever or churns a ghost, so every maintenance path
1661		// is cross-checked against a full walk of both maps.
1662		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		// A late or replayed commit landing below the current version becomes the key's oldest; an index
1711		// that tracked only first-seen versions would strand it in the buffer forever.
1712		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		// The tally is incremental, so any drift from the true map contents misreports memory forever
1740		// after; every mutation shape is exercised then compared against an exhaustive walk.
1741		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		// Eviction drains the buffer continuously, so a leak in any release path accumulates into a
1798		// permanently inflated memory report; the live-version drop is included in the sequence.
1799		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		// None means "nothing is waiting to be flushed", which is what lets a retention floor sit at
1866		// the permitted watermark. Reporting a version here would peg the floor to a write that
1867		// does not exist.
1868		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		// The retention floor is global, so a single lagging keyspace has to hold it down. Taking a
1876		// per-kind minimum, or the newest version instead of the oldest, is what lets the tombstone
1877		// reaper delete rows an un-flushed write is about to rewrite.
1878		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		// One keyspace holds many keys, each indexed under its own oldest resident version. Reading
1913		// the newest bucket instead of the oldest reports a frontier above writes that are still
1914		// buffered - and the tombstone reaper deletes under exactly that frontier.
1915		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		// Overwriting a key moves the old version into the historical map; it is still resident and
1940		// still un-flushed. If the index followed the current version instead, the frontier would
1941		// jump past a version the flusher has not written yet.
1942		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		// A sweep that empties a keyspace has to retire its bucket from the index. A stale bucket
1970		// pins the retention floor at a version that is already durable, and reclamation stalls
1971		// forever with no failing symptom.
1972		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}