Skip to main content

reifydb_store_multi/store/
multi.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use std::{
5	collections::{BTreeMap, HashMap, btree_map::Entry},
6	ops::{Bound, RangeBounds},
7	vec,
8};
9
10use reifydb_codec::{
11	key::encoded::{EncodedKey, EncodedKeyRange},
12	row::bytes::EncodedBytes,
13};
14use reifydb_core::{
15	common::CommitVersion,
16	delta::Delta,
17	event::metric::{MultiCommittedEvent, MultiDelete, MultiWrite},
18	interface::store::{
19		EntryKind, MultiVersionBatch, MultiVersionCommit, MultiVersionContains, MultiVersionGet,
20		MultiVersionGetPrevious, MultiVersionRow, MultiVersionStore, classify_key, classify_range,
21	},
22};
23use reifydb_store::row::page::PageId;
24use reifydb_value::{
25	reifydb_assertions,
26	util::{cowvec::CowVec, hex},
27};
28use tracing::instrument;
29
30use super::StandardMultiStore;
31use crate::{
32	MultiVersionScope, Result,
33	tier::{
34		DisplacedValues, RangeBatch, RangeCursor, TierBatch, TierStorage, VersionedGetResult,
35		commit::buffer::MultiCommitBufferTier,
36		persistent::MultiPersistentTier,
37		read::{MultiReadBufferTier, ServedChunk},
38	},
39};
40
41const TIER_SCAN_CHUNK_SIZE: usize = 32;
42
43pub(crate) const WARM_THRESHOLD: u64 = 4 * TIER_SCAN_CHUNK_SIZE as u64;
44
45impl MultiVersionGet for StandardMultiStore {
46	fn get(&self, key: &EncodedKey, version: CommitVersion) -> Result<Option<MultiVersionRow>> {
47		match classify_key(key) {
48			EntryKind::Source(_) => self.get_source(key, version),
49			_ => self.get_multi(key, version),
50		}
51	}
52}
53
54impl StandardMultiStore {
55	#[instrument(name = "store::multi::get::source", level = "trace", skip(self, key), fields(version = version.0))]
56	fn get_source(&self, key: &EncodedKey, version: CommitVersion) -> Result<Option<MultiVersionRow>> {
57		self.get_impl(key, version)
58	}
59
60	#[instrument(name = "store::multi::get::multi", level = "trace", skip(self, key), fields(version = version.0))]
61	fn get_multi(&self, key: &EncodedKey, version: CommitVersion) -> Result<Option<MultiVersionRow>> {
62		self.get_impl(key, version)
63	}
64
65	#[inline]
66	fn get_impl(&self, key: &EncodedKey, version: CommitVersion) -> Result<Option<MultiVersionRow>> {
67		let table = classify_key(key);
68
69		if let Some(found) = self.get_probe_commit(table, key, version)? {
70			return Ok(found);
71		}
72		if let Some(found) = self.get_probe_read(key, version) {
73			return Ok(found);
74		}
75		if let Some(found) = self.get_probe_persistent(table, key, version)? {
76			return Ok(found);
77		}
78
79		Ok(None)
80	}
81}
82
83impl StandardMultiStore {
84	#[inline]
85	fn get_probe_commit(
86		&self,
87		table: EntryKind,
88		key: &EncodedKey,
89		version: CommitVersion,
90	) -> Result<Option<Option<MultiVersionRow>>> {
91		Ok(match self.commit.get(table, key.as_ref(), version)? {
92			VersionedGetResult::Value {
93				value,
94				version: v,
95			} => Some(Some(MultiVersionRow {
96				key: key.clone(),
97				bytes: EncodedBytes(value),
98				version: v,
99			})),
100			VersionedGetResult::Tombstone => Some(None),
101			VersionedGetResult::NotFound => None,
102		})
103	}
104
105	#[inline]
106	fn get_probe_read(&self, key: &EncodedKey, version: CommitVersion) -> Option<Option<MultiVersionRow>> {
107		let read = self.read.as_ref()?;
108		match read.get(key, version) {
109			VersionedGetResult::Value {
110				value,
111				version: v,
112			} => Some(Some(MultiVersionRow {
113				key: key.clone(),
114				bytes: EncodedBytes(value),
115				version: v,
116			})),
117			VersionedGetResult::Tombstone => Some(None),
118			VersionedGetResult::NotFound => None,
119		}
120	}
121
122	#[inline]
123	fn get_probe_persistent(
124		&self,
125		table: EntryKind,
126		key: &EncodedKey,
127		version: CommitVersion,
128	) -> Result<Option<Option<MultiVersionRow>>> {
129		let Some(persistent) = &self.persistent else {
130			return Ok(None);
131		};
132		Ok(match persistent.get(table, key.as_ref(), version)? {
133			VersionedGetResult::Value {
134				value,
135				version: v,
136			} => {
137				if let Some(read) = &self.read {
138					read.insert(key.clone(), v, Some(value.clone()));
139				}
140				Some(Some(MultiVersionRow {
141					key: key.clone(),
142					bytes: EncodedBytes(value),
143					version: v,
144				}))
145			}
146			VersionedGetResult::Tombstone => Some(None),
147			VersionedGetResult::NotFound => None,
148		})
149	}
150}
151
152impl MultiVersionContains for StandardMultiStore {
153	#[instrument(name = "store::multi::contains", level = "trace", skip(self), fields(key_hex = %hex::display(key.as_ref()), version = version.0), ret)]
154	fn contains(&self, key: &EncodedKey, version: CommitVersion) -> Result<bool> {
155		Ok(MultiVersionGet::get(self, key, version)?.is_some())
156	}
157}
158
159impl MultiVersionCommit for StandardMultiStore {
160	#[instrument(name = "store::multi::commit", level = "debug", skip(self, deltas), fields(delta_count = deltas.len(), version = version.0))]
161	fn commit(&self, deltas: CowVec<Delta>, version: CommitVersion) -> Result<()> {
162		let classified = classify_deltas(&deltas);
163
164		self.update_read_cache_on_commit(&classified.batches);
165
166		let displaced = self.write_batches(version, classified.batches)?;
167
168		self.emit_commit_metrics(classified.writes, classified.deletes, displaced, version);
169
170		Ok(())
171	}
172}
173
174struct ClassifiedDeltas {
175	writes: Vec<MultiWrite>,
176	deletes: Vec<MultiDelete>,
177	batches: TierBatch,
178}
179
180#[inline]
181fn classify_deltas(deltas: &CowVec<Delta>) -> ClassifiedDeltas {
182	let mut writes: Vec<MultiWrite> = Vec::new();
183	let mut deletes: Vec<MultiDelete> = Vec::new();
184	let mut batches: TierBatch = HashMap::new();
185
186	for delta in deltas.iter() {
187		let key = delta.key();
188		let table = classify_key(key);
189
190		match delta {
191			Delta::Set {
192				key,
193				bytes,
194			} => {
195				writes.push(MultiWrite {
196					key: key.clone(),
197					value_bytes: bytes.len() as u64,
198				});
199				batches.entry(table).or_default().push((key.clone(), Some(bytes.0.clone())));
200			}
201			Delta::Remove {
202				key,
203				..
204			} => {
205				deletes.push(MultiDelete {
206					key: key.clone(),
207					value_bytes: 0,
208				});
209				batches.entry(table).or_default().push((key.clone(), None));
210			}
211		}
212	}
213
214	ClassifiedDeltas {
215		writes,
216		deletes,
217		batches,
218	}
219}
220
221impl StandardMultiStore {
222	pub fn get_many(
223		&self,
224		keys: &[EncodedKey],
225		version: CommitVersion,
226	) -> Result<HashMap<EncodedKey, MultiVersionRow>> {
227		let mut by_table: HashMap<EntryKind, Vec<&EncodedKey>> = HashMap::new();
228		for key in keys {
229			by_table.entry(classify_key(key)).or_default().push(key);
230		}
231
232		let mut out: HashMap<EncodedKey, MultiVersionRow> = HashMap::new();
233		for (table, table_keys) in by_table {
234			self.get_many_for_table(table, &table_keys, version, &mut out)?;
235		}
236
237		Ok(out)
238	}
239
240	#[inline]
241	fn get_many_for_table(
242		&self,
243		table: EntryKind,
244		table_keys: &[&EncodedKey],
245		version: CommitVersion,
246		out: &mut HashMap<EncodedKey, MultiVersionRow>,
247	) -> Result<()> {
248		let key_slices: Vec<&[u8]> = table_keys.iter().map(|k| k.as_ref()).collect();
249
250		let commit_results = self.probe_commit_batch(table, &key_slices, version)?;
251		let (read_aligned, persistent_aligned) = self.resolve_misses_through_read_and_persistent(
252			table,
253			table_keys,
254			&key_slices,
255			&commit_results,
256			version,
257		)?;
258
259		reifydb_assertions! {
260			let n = key_slices.len();
261			assert!(
262				commit_results.len() == n && read_aligned.len() == n && persistent_aligned.len() == n,
263				"per-tier result vectors must stay index-aligned with the table's keys, otherwise collect_resolved_rows \
264				 reads a tier result for the wrong key and returns mismatched rows (keys={n}, commit={}, read={}, persistent={})",
265				commit_results.len(),
266				read_aligned.len(),
267				persistent_aligned.len()
268			);
269		}
270
271		self.collect_resolved_rows(table_keys, &commit_results, &read_aligned, &persistent_aligned, out);
272		Ok(())
273	}
274
275	#[inline]
276	fn probe_commit_batch(
277		&self,
278		table: EntryKind,
279		key_slices: &[&[u8]],
280		version: CommitVersion,
281	) -> Result<Vec<VersionedGetResult>> {
282		self.commit.get_many(table, key_slices, version)
283	}
284
285	#[inline]
286	fn resolve_misses_through_read_and_persistent(
287		&self,
288		table: EntryKind,
289		table_keys: &[&EncodedKey],
290		key_slices: &[&[u8]],
291		commit_results: &[VersionedGetResult],
292		version: CommitVersion,
293	) -> Result<(Vec<VersionedGetResult>, Vec<VersionedGetResult>)> {
294		let mut read_aligned = vec![VersionedGetResult::NotFound; key_slices.len()];
295		let mut persistent_idx: Vec<usize> = Vec::new();
296		let mut persistent_slices: Vec<&[u8]> = Vec::new();
297		for (i, result) in commit_results.iter().enumerate() {
298			if !matches!(result, VersionedGetResult::NotFound) {
299				continue;
300			}
301			let read_hit = self
302				.read
303				.as_ref()
304				.map(|c| c.get(table_keys[i], version))
305				.unwrap_or(VersionedGetResult::NotFound);
306			match read_hit {
307				VersionedGetResult::Value {
308					value,
309					version: v,
310				} => {
311					read_aligned[i] = VersionedGetResult::Value {
312						value,
313						version: v,
314					};
315				}
316				VersionedGetResult::Tombstone => {
317					read_aligned[i] = VersionedGetResult::Tombstone;
318				}
319				VersionedGetResult::NotFound => {
320					persistent_idx.push(i);
321					persistent_slices.push(key_slices[i]);
322				}
323			}
324		}
325
326		let mut persistent_aligned = vec![VersionedGetResult::NotFound; key_slices.len()];
327		if !persistent_slices.is_empty()
328			&& let Some(persistent) = &self.persistent
329		{
330			let persistent_results = persistent.get_many(table, &persistent_slices, version)?;
331			for (slot, result) in persistent_idx.into_iter().zip(persistent_results) {
332				if let (
333					Some(read),
334					VersionedGetResult::Value {
335						value,
336						version: v,
337					},
338				) = (&self.read, &result)
339				{
340					read.insert(table_keys[slot].clone(), *v, Some(value.clone()));
341				}
342				persistent_aligned[slot] = result;
343			}
344		}
345
346		Ok((read_aligned, persistent_aligned))
347	}
348
349	#[inline]
350	fn collect_resolved_rows(
351		&self,
352		table_keys: &[&EncodedKey],
353		commit_results: &[VersionedGetResult],
354		read_aligned: &[VersionedGetResult],
355		persistent_aligned: &[VersionedGetResult],
356		out: &mut HashMap<EncodedKey, MultiVersionRow>,
357	) {
358		for (i, key) in table_keys.iter().enumerate() {
359			let resolved = match &commit_results[i] {
360				VersionedGetResult::Value {
361					value,
362					version: v,
363				} => Some((value.clone(), *v)),
364				VersionedGetResult::Tombstone => None,
365				VersionedGetResult::NotFound => match &read_aligned[i] {
366					VersionedGetResult::Value {
367						value,
368						version: v,
369					} => Some((value.clone(), *v)),
370					VersionedGetResult::Tombstone => None,
371					VersionedGetResult::NotFound => match &persistent_aligned[i] {
372						VersionedGetResult::Value {
373							value,
374							version: v,
375						} => Some((value.clone(), *v)),
376						_ => None,
377					},
378				},
379			};
380
381			if let Some((value, v)) = resolved {
382				out.insert(
383					(*key).clone(),
384					MultiVersionRow {
385						key: (*key).clone(),
386						bytes: EncodedBytes(value),
387						version: v,
388					},
389				);
390			}
391		}
392	}
393
394	#[inline]
395	fn update_read_cache_on_commit(&self, batches: &TierBatch) {
396		let Some(read) = &self.read else {
397			return;
398		};
399		for entries in batches.values() {
400			for (key, _) in entries {
401				read.invalidate(key);
402			}
403		}
404	}
405
406	#[inline]
407	fn write_batches(&self, version: CommitVersion, batches: TierBatch) -> Result<DisplacedValues> {
408		self.commit.set(version, batches)
409	}
410
411	#[inline]
412	fn emit_commit_metrics(
413		&self,
414		writes: Vec<MultiWrite>,
415		mut deletes: Vec<MultiDelete>,
416		displaced: DisplacedValues,
417		version: CommitVersion,
418	) {
419		if writes.is_empty() && deletes.is_empty() {
420			return;
421		}
422		if !deletes.is_empty() {
423			let displaced: HashMap<&EncodedKey, u64> = displaced.iter().map(|(k, b)| (k, *b)).collect();
424			for delete in deletes.iter_mut() {
425				delete.value_bytes = displaced.get(&delete.key).copied().unwrap_or(0);
426			}
427		}
428		self.event_bus.emit(MultiCommittedEvent::new(writes, deletes, version));
429	}
430}
431
432#[derive(Debug, Clone, Default)]
433pub struct MultiVersionRangeCursor {
434	pub commit: RangeCursor,
435
436	pub persistent: RangeCursor,
437
438	pub exhausted: bool,
439
440	warm: bool,
441
442	warm_bucket: Option<PageId>,
443
444	warm_consumed: u64,
445}
446
447impl MultiVersionRangeCursor {
448	pub fn new() -> Self {
449		Self {
450			warm: true,
451			..Default::default()
452		}
453	}
454
455	pub fn cold() -> Self {
456		Self {
457			warm: false,
458			..Default::default()
459		}
460	}
461
462	pub fn is_exhausted(&self) -> bool {
463		self.exhausted
464	}
465}
466
467pub struct TierScanQuery<'a> {
468	pub table: EntryKind,
469	pub start: &'a [u8],
470	pub end: &'a [u8],
471	pub scope: MultiVersionScope,
472	pub range: &'a EncodedKeyRange,
473}
474
475pub fn scan_tier_chunk<S: TierStorage>(
476	storage: &S,
477	cursor: &mut RangeCursor,
478	scan: &TierScanQuery,
479	collected: &mut BTreeMap<EncodedKey, (CommitVersion, Option<CowVec<u8>>)>,
480) -> Result<bool> {
481	let batch = storage.range_next(
482		scan.table,
483		cursor,
484		Bound::Included(scan.start),
485		Bound::Included(scan.end),
486		scan.scope,
487		TIER_SCAN_CHUNK_SIZE,
488	)?;
489	merge_tier_batch(batch, scan.range, collected)
490}
491
492pub fn scan_tier_chunk_rev<S: TierStorage>(
493	storage: &S,
494	cursor: &mut RangeCursor,
495	scan: &TierScanQuery,
496	collected: &mut BTreeMap<EncodedKey, (CommitVersion, Option<CowVec<u8>>)>,
497) -> Result<bool> {
498	let batch = storage.range_rev_next(
499		scan.table,
500		cursor,
501		Bound::Included(scan.start),
502		Bound::Included(scan.end),
503		scan.scope,
504		TIER_SCAN_CHUNK_SIZE,
505	)?;
506	merge_tier_batch(batch, scan.range, collected)
507}
508
509#[inline]
510fn merge_tier_batch(
511	batch: RangeBatch,
512	range: &EncodedKeyRange,
513	collected: &mut BTreeMap<EncodedKey, (CommitVersion, Option<CowVec<u8>>)>,
514) -> Result<bool> {
515	if batch.entries.is_empty() {
516		return Ok(false);
517	}
518
519	for entry in batch.entries {
520		if !range.contains(&entry.key) {
521			continue;
522		}
523
524		match collected.entry(entry.key) {
525			Entry::Vacant(slot) => {
526				slot.insert((entry.version, entry.value));
527			}
528			Entry::Occupied(mut slot) => {
529				if entry.version > slot.get().0 {
530					slot.insert((entry.version, entry.value));
531				}
532			}
533		}
534	}
535
536	Ok(true)
537}
538
539#[inline]
540pub fn collected_to_batch(
541	collected: BTreeMap<EncodedKey, (CommitVersion, Option<CowVec<u8>>)>,
542	has_more: bool,
543) -> MultiVersionBatch {
544	let items: Vec<MultiVersionRow> = collected
545		.into_iter()
546		.filter_map(|(key, (v, value))| {
547			value.map(|val| MultiVersionRow {
548				key,
549				bytes: EncodedBytes(val),
550				version: v,
551			})
552		})
553		.collect();
554
555	MultiVersionBatch {
556		items,
557		has_more,
558	}
559}
560
561#[inline]
562fn step_all_tiers(
563	buffer: Option<&MultiCommitBufferTier>,
564	buffer_cursor: &mut RangeCursor,
565	persistent: Option<&MultiPersistentTier>,
566	persistent_cursor: &mut RangeCursor,
567	scan: &TierScanQuery,
568	collected: &mut BTreeMap<EncodedKey, (CommitVersion, Option<CowVec<u8>>)>,
569) -> Result<bool> {
570	let mut any_progress = false;
571	if let Some(s) = buffer
572		&& !buffer_cursor.exhausted
573	{
574		any_progress |= scan_tier_chunk(s, buffer_cursor, scan, collected)?;
575	}
576	if let Some(s) = persistent
577		&& !persistent_cursor.exhausted
578	{
579		any_progress |= scan_tier_chunk(s, persistent_cursor, scan, collected)?;
580	}
581	Ok(any_progress)
582}
583
584pub fn scan_tiers_latest(
585	buffer: Option<&MultiCommitBufferTier>,
586	persistent: Option<&MultiPersistentTier>,
587	range: EncodedKeyRange,
588	scope: MultiVersionScope,
589	max_keys: usize,
590) -> Result<MultiVersionBatch> {
591	let table = classify_key_range(&range);
592	let (start, end) = make_range_bounds(&range);
593	let scan = TierScanQuery {
594		table,
595		start: &start,
596		end: &end,
597		scope,
598		range: &range,
599	};
600
601	let mut collected: BTreeMap<EncodedKey, (CommitVersion, Option<CowVec<u8>>)> = BTreeMap::new();
602	let mut buffer_cursor = RangeCursor::default();
603	let mut persistent_cursor = RangeCursor::default();
604	let mut exhausted = false;
605
606	while collected.len() < max_keys {
607		let progress = step_all_tiers(
608			buffer,
609			&mut buffer_cursor,
610			persistent,
611			&mut persistent_cursor,
612			&scan,
613			&mut collected,
614		)?;
615		if !progress {
616			exhausted = true;
617			break;
618		}
619	}
620
621	Ok(collected_to_batch(collected, !exhausted))
622}
623
624impl StandardMultiStore {
625	pub fn range_next(
626		&self,
627		cursor: &mut MultiVersionRangeCursor,
628		range: EncodedKeyRange,
629		scope: MultiVersionScope,
630		batch_size: u64,
631	) -> Result<MultiVersionBatch> {
632		if cursor.exhausted {
633			return Ok(MultiVersionBatch {
634				items: Vec::new(),
635				has_more: false,
636			});
637		}
638
639		mark_unconfigured_exhausted(self, cursor);
640
641		let table = classify_key_range(&range);
642		let (start, end) = make_range_bounds(&range);
643		let batch_size = batch_size as usize;
644		let scan = TierScanQuery {
645			table,
646			start: &start,
647			end: &end,
648			scope,
649			range: &range,
650		};
651
652		let mut collected: BTreeMap<EncodedKey, (CommitVersion, Option<CowVec<u8>>)> = BTreeMap::new();
653
654		while collected.len() < batch_size {
655			let mut any_progress = false;
656
657			if !cursor.commit.exhausted {
658				any_progress |=
659					scan_tier_chunk(&self.commit, &mut cursor.commit, &scan, &mut collected)?;
660			}
661
662			if self.persistent.is_some() && !cursor.persistent.exhausted {
663				any_progress |= self.step_persistent_cached(&scan, cursor, &mut collected, false)?;
664			}
665
666			if !any_progress {
667				cursor.exhausted = true;
668				break;
669			}
670		}
671
672		apply_forward_horizon(cursor, &mut collected);
673
674		let items: Vec<MultiVersionRow> = collected
675			.into_iter()
676			.filter_map(|(key_bytes, (v, value))| {
677				value.map(|val| MultiVersionRow {
678					key: EncodedKey::new(key_bytes),
679					bytes: EncodedBytes(val),
680					version: v,
681				})
682			})
683			.collect();
684
685		let has_more = !cursor.exhausted;
686
687		Ok(MultiVersionBatch {
688			items,
689			has_more,
690		})
691	}
692
693	pub fn range(
694		&self,
695		range: EncodedKeyRange,
696		scope: MultiVersionScope,
697		batch_size: usize,
698	) -> MultiVersionRangeIter {
699		MultiVersionRangeIter {
700			store: self.clone(),
701			cursor: MultiVersionRangeCursor::new(),
702			range,
703			scope,
704			batch_size,
705			current_batch: Vec::new().into_iter(),
706		}
707	}
708
709	pub fn range_persistence(
710		&self,
711		range: EncodedKeyRange,
712		scope: MultiVersionScope,
713		batch_size: usize,
714	) -> MultiVersionRangeIter {
715		MultiVersionRangeIter {
716			store: self.clone(),
717			cursor: MultiVersionRangeCursor::cold(),
718			range,
719			scope,
720			batch_size,
721			current_batch: Vec::new().into_iter(),
722		}
723	}
724
725	pub fn range_rev(
726		&self,
727		range: EncodedKeyRange,
728		scope: MultiVersionScope,
729		batch_size: usize,
730	) -> MultiVersionRangeRevIter {
731		MultiVersionRangeRevIter {
732			store: self.clone(),
733			cursor: MultiVersionRangeCursor::new(),
734			range,
735			scope,
736			batch_size,
737			current_batch: Vec::new().into_iter(),
738		}
739	}
740
741	pub fn range_rev_persistence(
742		&self,
743		range: EncodedKeyRange,
744		scope: MultiVersionScope,
745		batch_size: usize,
746	) -> MultiVersionRangeRevIter {
747		MultiVersionRangeRevIter {
748			store: self.clone(),
749			cursor: MultiVersionRangeCursor::cold(),
750			range,
751			scope,
752			batch_size,
753			current_batch: Vec::new().into_iter(),
754		}
755	}
756
757	fn range_rev_next(
758		&self,
759		cursor: &mut MultiVersionRangeCursor,
760		range: EncodedKeyRange,
761		scope: MultiVersionScope,
762		batch_size: u64,
763	) -> Result<MultiVersionBatch> {
764		if cursor.exhausted {
765			return Ok(MultiVersionBatch {
766				items: Vec::new(),
767				has_more: false,
768			});
769		}
770
771		mark_unconfigured_exhausted(self, cursor);
772
773		let table = classify_key_range(&range);
774		let (start, end) = make_range_bounds(&range);
775		let batch_size = batch_size as usize;
776		let scan = TierScanQuery {
777			table,
778			start: &start,
779			end: &end,
780			scope,
781			range: &range,
782		};
783
784		let mut collected: BTreeMap<EncodedKey, (CommitVersion, Option<CowVec<u8>>)> = BTreeMap::new();
785
786		while collected.len() < batch_size {
787			let mut any_progress = false;
788
789			if !cursor.commit.exhausted {
790				any_progress |=
791					scan_tier_chunk_rev(&self.commit, &mut cursor.commit, &scan, &mut collected)?;
792			}
793
794			if self.persistent.is_some() && !cursor.persistent.exhausted {
795				any_progress |= self.step_persistent_cached(&scan, cursor, &mut collected, true)?;
796			}
797
798			if !any_progress {
799				cursor.exhausted = true;
800				break;
801			}
802		}
803
804		apply_reverse_horizon(cursor, &mut collected);
805
806		let items: Vec<MultiVersionRow> = collected
807			.into_iter()
808			.rev()
809			.filter_map(|(key_bytes, (v, value))| {
810				value.map(|val| MultiVersionRow {
811					key: EncodedKey::new(key_bytes),
812					bytes: EncodedBytes(val),
813					version: v,
814				})
815			})
816			.collect();
817
818		let has_more = !cursor.exhausted;
819
820		Ok(MultiVersionBatch {
821			items,
822			has_more,
823		})
824	}
825
826	fn step_persistent_cached(
827		&self,
828		scan: &TierScanQuery,
829		cursor: &mut MultiVersionRangeCursor,
830		collected: &mut BTreeMap<EncodedKey, (CommitVersion, Option<CowVec<u8>>)>,
831		descending: bool,
832	) -> Result<bool> {
833		let Some(persistent) = &self.persistent else {
834			return Ok(false);
835		};
836
837		if let Some(served) = self.serve_from_read_cache(scan, cursor, collected, descending) {
838			return served;
839		}
840
841		let (consumed, progressed) =
842			self.scan_persistent_chunk(persistent, scan, cursor, collected, descending)?;
843		self.warm_read_bucket_after_scan(persistent, scan, cursor, consumed)?;
844
845		Ok(progressed)
846	}
847
848	#[inline]
849	fn serve_from_read_cache(
850		&self,
851		scan: &TierScanQuery,
852		cursor: &mut MultiVersionRangeCursor,
853		collected: &mut BTreeMap<EncodedKey, (CommitVersion, Option<CowVec<u8>>)>,
854		descending: bool,
855	) -> Option<Result<bool>> {
856		let (Some(read), EntryKind::Source(_)) = (&self.read, scan.table) else {
857			return None;
858		};
859		match read.serve_persistent_chunk(
860			scan.table,
861			&mut cursor.persistent,
862			scan.start,
863			scan.end,
864			scan.scope,
865			TIER_SCAN_CHUNK_SIZE,
866			descending,
867		) {
868			ServedChunk::Served(batch) => Some(merge_tier_batch(batch, scan.range, collected)),
869			ServedChunk::Gap => None,
870		}
871	}
872
873	#[inline]
874	fn scan_persistent_chunk(
875		&self,
876		persistent: &MultiPersistentTier,
877		scan: &TierScanQuery,
878		cursor: &mut MultiVersionRangeCursor,
879		collected: &mut BTreeMap<EncodedKey, (CommitVersion, Option<CowVec<u8>>)>,
880		descending: bool,
881	) -> Result<(usize, bool)> {
882		let batch = if descending {
883			persistent.range_rev_next(
884				scan.table,
885				&mut cursor.persistent,
886				Bound::Included(scan.start),
887				Bound::Included(scan.end),
888				scan.scope,
889				TIER_SCAN_CHUNK_SIZE,
890			)?
891		} else {
892			persistent.range_next(
893				scan.table,
894				&mut cursor.persistent,
895				Bound::Included(scan.start),
896				Bound::Included(scan.end),
897				scan.scope,
898				TIER_SCAN_CHUNK_SIZE,
899			)?
900		};
901		let consumed = batch.entries.len();
902		let progressed = merge_tier_batch(batch, scan.range, collected)?;
903		Ok((consumed, progressed))
904	}
905
906	#[inline]
907	fn warm_read_bucket_after_scan(
908		&self,
909		persistent: &MultiPersistentTier,
910		scan: &TierScanQuery,
911		cursor: &mut MultiVersionRangeCursor,
912		consumed: usize,
913	) -> Result<()> {
914		if !cursor.warm {
915			return Ok(());
916		}
917		if let (Some(read), EntryKind::Source(_)) = (&self.read, scan.table) {
918			maybe_warm_bucket(read, persistent, cursor, scan.table, consumed)?;
919		}
920		Ok(())
921	}
922}
923
924fn maybe_warm_bucket(
925	read: &MultiReadBufferTier,
926	persistent: &MultiPersistentTier,
927	cursor: &mut MultiVersionRangeCursor,
928	table: EntryKind,
929	consumed: usize,
930) -> Result<()> {
931	let page = {
932		let Some(last) = cursor.persistent.last_key.as_ref() else {
933			return Ok(());
934		};
935		read.page_of_key(last)
936	};
937	if !matches!(page.kind, EntryKind::Source(_)) {
938		return Ok(());
939	}
940
941	if cursor.warm_bucket == Some(page) {
942		cursor.warm_consumed = cursor.warm_consumed.saturating_add(consumed as u64);
943	} else {
944		cursor.warm_bucket = Some(page);
945		cursor.warm_consumed = consumed as u64;
946	}
947
948	if cursor.warm_consumed <= WARM_THRESHOLD {
949		return Ok(());
950	}
951
952	let settle = |cursor: &mut MultiVersionRangeCursor| {
953		cursor.warm_bucket = None;
954		cursor.warm_consumed = 0;
955	};
956
957	if read.page_is_complete(page) {
958		settle(cursor);
959		return Ok(());
960	}
961
962	let Some(range) = read.page_key_range(page) else {
963		return Ok(());
964	};
965	let (Bound::Included(lo), Bound::Included(hi)) = (range.start, range.end) else {
966		return Ok(());
967	};
968
969	if !read.begin_warm(page) {
970		settle(cursor);
971		return Ok(());
972	}
973
974	let loaded = persistent.load_range_consistent(
975		table,
976		Bound::Included(lo.as_slice()),
977		Bound::Included(hi.as_slice()),
978		CommitVersion(u64::MAX),
979		None,
980	);
981	let entries = match loaded {
982		Ok(entries) => entries,
983		Err(e) => {
984			read.abort_warm(page);
985			settle(cursor);
986			return Err(e);
987		}
988	};
989
990	read.finish_warm(page, entries);
991	settle(cursor);
992	Ok(())
993}
994
995fn mark_unconfigured_exhausted(store: &StandardMultiStore, cursor: &mut MultiVersionRangeCursor) {
996	if store.persistent.is_none() {
997		cursor.persistent.exhausted = true;
998	}
999}
1000
1001fn apply_forward_horizon(
1002	cursor: &mut MultiVersionRangeCursor,
1003	collected: &mut BTreeMap<EncodedKey, (CommitVersion, Option<CowVec<u8>>)>,
1004) {
1005	let horizon = forward_horizon(cursor);
1006	if let Some(h) = horizon {
1007		collected.retain(|k, _| k.as_slice() <= h.as_slice());
1008		rewind_over_advanced_forward(cursor, &h);
1009	}
1010}
1011
1012fn apply_reverse_horizon(
1013	cursor: &mut MultiVersionRangeCursor,
1014	collected: &mut BTreeMap<EncodedKey, (CommitVersion, Option<CowVec<u8>>)>,
1015) {
1016	let horizon = reverse_horizon(cursor);
1017	if let Some(h) = horizon {
1018		collected.retain(|k, _| k.as_slice() >= h.as_slice());
1019		rewind_over_advanced_reverse(cursor, &h);
1020	}
1021}
1022
1023fn forward_horizon(cursor: &MultiVersionRangeCursor) -> Option<EncodedKey> {
1024	let mut horizon: Option<EncodedKey> = None;
1025	for tier in [&cursor.commit, &cursor.persistent] {
1026		if tier.exhausted {
1027			continue;
1028		}
1029		let last = match &tier.last_key {
1030			Some(k) => k.clone(),
1031
1032			None => return None,
1033		};
1034		horizon = Some(match horizon {
1035			None => last,
1036			Some(prev) => {
1037				if last.as_slice() < prev.as_slice() {
1038					last
1039				} else {
1040					prev
1041				}
1042			}
1043		});
1044	}
1045	horizon
1046}
1047
1048fn reverse_horizon(cursor: &MultiVersionRangeCursor) -> Option<EncodedKey> {
1049	let mut horizon: Option<EncodedKey> = None;
1050	for tier in [&cursor.commit, &cursor.persistent] {
1051		if tier.exhausted {
1052			continue;
1053		}
1054		let last = match &tier.last_key {
1055			Some(k) => k.clone(),
1056			None => return None,
1057		};
1058		horizon = Some(match horizon {
1059			None => last,
1060			Some(prev) => {
1061				if last.as_slice() > prev.as_slice() {
1062					last
1063				} else {
1064					prev
1065				}
1066			}
1067		});
1068	}
1069	horizon
1070}
1071
1072fn rewind_over_advanced_forward(cursor: &mut MultiVersionRangeCursor, horizon: &EncodedKey) {
1073	for tier in [&mut cursor.commit, &mut cursor.persistent] {
1074		if let Some(last) = &tier.last_key
1075			&& last.as_slice() > horizon.as_slice()
1076		{
1077			tier.last_key = Some(horizon.clone());
1078			tier.exhausted = false;
1079		}
1080	}
1081}
1082
1083fn rewind_over_advanced_reverse(cursor: &mut MultiVersionRangeCursor, horizon: &EncodedKey) {
1084	for tier in [&mut cursor.commit, &mut cursor.persistent] {
1085		if let Some(last) = &tier.last_key
1086			&& last.as_slice() < horizon.as_slice()
1087		{
1088			tier.last_key = Some(horizon.clone());
1089			tier.exhausted = false;
1090		}
1091	}
1092}
1093
1094impl MultiVersionGetPrevious for StandardMultiStore {
1095	fn get_previous_version(
1096		&self,
1097		key: &EncodedKey,
1098		before_version: CommitVersion,
1099	) -> Result<Option<MultiVersionRow>> {
1100		if before_version.0 == 0 {
1101			return Ok(None);
1102		}
1103
1104		let table = classify_key(key);
1105		reifydb_assertions! {
1106			assert!(
1107				before_version.0 >= 1,
1108				"the before_version==0 guard must precede this subtraction, otherwise before_version.0 - 1 \
1109				 wraps to u64::MAX and the probe reads the latest version instead of the previous one \
1110				 (before_version={})",
1111				before_version.0
1112			);
1113		}
1114		let prev_version = CommitVersion(before_version.0 - 1);
1115
1116		if let Some(found) = self.previous_probe_commit(table, key, prev_version)? {
1117			return Ok(found);
1118		}
1119		if let Some(found) = self.previous_probe_read(key, prev_version) {
1120			return Ok(found);
1121		}
1122		if let Some(found) = self.previous_probe_persistent(table, key, prev_version)? {
1123			return Ok(found);
1124		}
1125
1126		Ok(None)
1127	}
1128}
1129
1130impl StandardMultiStore {
1131	#[inline]
1132	fn previous_probe_commit(
1133		&self,
1134		table: EntryKind,
1135		key: &EncodedKey,
1136		prev_version: CommitVersion,
1137	) -> Result<Option<Option<MultiVersionRow>>> {
1138		Ok(match self.commit.get(table, key.as_ref(), prev_version)? {
1139			VersionedGetResult::Value {
1140				value,
1141				version,
1142			} => Some(Some(MultiVersionRow {
1143				key: key.clone(),
1144				bytes: EncodedBytes(CowVec::new(value.to_vec())),
1145				version,
1146			})),
1147			VersionedGetResult::Tombstone => Some(None),
1148			VersionedGetResult::NotFound => None,
1149		})
1150	}
1151
1152	#[inline]
1153	fn previous_probe_read(
1154		&self,
1155		key: &EncodedKey,
1156		prev_version: CommitVersion,
1157	) -> Option<Option<MultiVersionRow>> {
1158		let read = self.read.as_ref()?;
1159		match read.get(key, prev_version) {
1160			VersionedGetResult::Value {
1161				value,
1162				version,
1163			} => Some(Some(MultiVersionRow {
1164				key: key.clone(),
1165				bytes: EncodedBytes(CowVec::new(value.to_vec())),
1166				version,
1167			})),
1168			VersionedGetResult::Tombstone => Some(None),
1169			VersionedGetResult::NotFound => None,
1170		}
1171	}
1172
1173	#[inline]
1174	fn previous_probe_persistent(
1175		&self,
1176		table: EntryKind,
1177		key: &EncodedKey,
1178		prev_version: CommitVersion,
1179	) -> Result<Option<Option<MultiVersionRow>>> {
1180		let Some(persistent) = &self.persistent else {
1181			return Ok(None);
1182		};
1183		Ok(match persistent.get(table, key.as_ref(), prev_version)? {
1184			VersionedGetResult::Value {
1185				value,
1186				version,
1187			} => {
1188				if let Some(read) = &self.read {
1189					read.insert(key.clone(), version, Some(value.clone()));
1190				}
1191				Some(Some(MultiVersionRow {
1192					key: key.clone(),
1193					bytes: EncodedBytes(CowVec::new(value.to_vec())),
1194					version,
1195				}))
1196			}
1197			VersionedGetResult::Tombstone => Some(None),
1198			VersionedGetResult::NotFound => None,
1199		})
1200	}
1201}
1202
1203impl MultiVersionStore for StandardMultiStore {}
1204
1205pub struct MultiVersionRangeIter {
1206	store: StandardMultiStore,
1207	cursor: MultiVersionRangeCursor,
1208	range: EncodedKeyRange,
1209	scope: MultiVersionScope,
1210	batch_size: usize,
1211	current_batch: vec::IntoIter<MultiVersionRow>,
1212}
1213
1214impl Iterator for MultiVersionRangeIter {
1215	type Item = Result<MultiVersionRow>;
1216
1217	fn next(&mut self) -> Option<Self::Item> {
1218		if let Some(item) = self.current_batch.next() {
1219			return Some(Ok(item));
1220		}
1221
1222		if self.cursor.exhausted {
1223			return None;
1224		}
1225
1226		match self.store.range_next(&mut self.cursor, self.range.clone(), self.scope, self.batch_size as u64) {
1227			Ok(batch) => {
1228				if batch.items.is_empty() {
1229					if self.cursor.exhausted {
1230						return None;
1231					}
1232					return self.next();
1233				}
1234				self.current_batch = batch.items.into_iter();
1235				self.next()
1236			}
1237			Err(e) => Some(Err(e)),
1238		}
1239	}
1240}
1241
1242pub struct MultiVersionRangeRevIter {
1243	store: StandardMultiStore,
1244	cursor: MultiVersionRangeCursor,
1245	range: EncodedKeyRange,
1246	scope: MultiVersionScope,
1247	batch_size: usize,
1248	current_batch: vec::IntoIter<MultiVersionRow>,
1249}
1250
1251impl Iterator for MultiVersionRangeRevIter {
1252	type Item = Result<MultiVersionRow>;
1253
1254	fn next(&mut self) -> Option<Self::Item> {
1255		if let Some(item) = self.current_batch.next() {
1256			return Some(Ok(item));
1257		}
1258
1259		if self.cursor.exhausted {
1260			return None;
1261		}
1262
1263		match self.store.range_rev_next(
1264			&mut self.cursor,
1265			self.range.clone(),
1266			self.scope,
1267			self.batch_size as u64,
1268		) {
1269			Ok(batch) => {
1270				if batch.items.is_empty() {
1271					if self.cursor.exhausted {
1272						return None;
1273					}
1274					return self.next();
1275				}
1276				self.current_batch = batch.items.into_iter();
1277				self.next()
1278			}
1279			Err(e) => Some(Err(e)),
1280		}
1281	}
1282}
1283
1284fn classify_key_range(range: &EncodedKeyRange) -> EntryKind {
1285	classify_range(range).unwrap_or(EntryKind::Multi)
1286}
1287
1288fn make_range_bounds(range: &EncodedKeyRange) -> (Vec<u8>, Vec<u8>) {
1289	let start = match &range.start {
1290		Bound::Included(key) => key.as_ref().to_vec(),
1291		Bound::Excluded(key) => key.as_ref().to_vec(),
1292		Bound::Unbounded => vec![],
1293	};
1294
1295	let end = match &range.end {
1296		Bound::Included(key) => key.as_ref().to_vec(),
1297		Bound::Excluded(key) => key.as_ref().to_vec(),
1298		Bound::Unbounded => vec![0xFFu8; 256],
1299	};
1300
1301	(start, end)
1302}
1303
1304#[cfg(all(test, feature = "sqlite", not(target_arch = "wasm32")))]
1305mod cache_tests {
1306	use std::collections::HashMap;
1307
1308	use reifydb_codec::{key::encoded::EncodedKey, row::bytes::EncodedBytes};
1309	use reifydb_core::{
1310		common::CommitVersion,
1311		delta::Delta,
1312		interface::{
1313			catalog::{flow::OperatorId, id::TableId, storage::StorageId},
1314			store::{EntryKind, MultiVersionCommit, MultiVersionGet},
1315		},
1316		key::{
1317			EncodableKey,
1318			operator_state::{GroupId, Keyspace, OperatorStateKey},
1319			row::RowKey,
1320		},
1321	};
1322	use reifydb_value::{cow_vec, util::cowvec::CowVec};
1323
1324	use crate::{
1325		MultiVersionScope,
1326		store::{StandardMultiStore, multi::WARM_THRESHOLD},
1327		tier::{RawEntry, TierStorage, VersionedGetResult, commit::buffer::MultiCommitBufferTier},
1328	};
1329
1330	const STORAGE: StorageId = StorageId::Table(TableId(1));
1331
1332	fn commit_row(store: &StandardMultiStore, n: u64, version: u64) {
1333		MultiVersionCommit::commit(
1334			store,
1335			cow_vec![Delta::Set {
1336				key: RowKey::encoded(STORAGE, n),
1337				bytes: EncodedBytes(CowVec::new(format!("v{n}").into_bytes())),
1338			}],
1339			CommitVersion(version),
1340		)
1341		.unwrap();
1342	}
1343
1344	fn flush(store: &StandardMultiStore, cutoff: CommitVersion) {
1345		let commit = store.commit();
1346		for kind in commit.list_all_entry_kinds().unwrap() {
1347			let (to_persist, to_compact, _more) = match commit {
1348				MultiCommitBufferTier::Memory(s) => s.collect_evictable_below(kind, cutoff, usize::MAX),
1349			};
1350			if to_compact.is_empty() {
1351				continue;
1352			}
1353			if !to_persist.is_empty() {
1354				let persistent = store.persistent().expect("persistent tier");
1355				let mut by_version: HashMap<
1356					CommitVersion,
1357					HashMap<EntryKind, Vec<(EncodedKey, Option<CowVec<u8>>)>>,
1358				> = HashMap::new();
1359				for (key, version, value) in to_persist {
1360					by_version
1361						.entry(version)
1362						.or_default()
1363						.entry(kind)
1364						.or_default()
1365						.push((key, value));
1366				}
1367				for (version, batch) in by_version {
1368					persistent.set(version, batch).unwrap();
1369				}
1370			}
1371			for evicted in &to_compact {
1372				store.invalidate_read_key(&evicted.key);
1373			}
1374			commit.compact(HashMap::from([(
1375				kind,
1376				to_compact.into_iter().map(|e| (e.key, e.version)).collect(),
1377			)]))
1378			.unwrap();
1379		}
1380	}
1381
1382	#[test]
1383	fn warm_threshold_warms_only_buckets_above_threshold() {
1384		const HEAVY: u64 = WARM_THRESHOLD + 64;
1385		const LIGHT: u64 = 20;
1386		let (store, _g) = StandardMultiStore::testing_memory_with_persistent_sqlite();
1387
1388		for n in 1..=HEAVY {
1389			commit_row(&store, n, 1);
1390		}
1391		for n in 0..LIGHT {
1392			commit_row(&store, (1u64 << 16) + n, 1);
1393		}
1394		flush(&store, CommitVersion(1));
1395
1396		let read = store.read.clone().expect("read tier configured");
1397		let heavy_bucket = read.page_of_key(&RowKey::encoded(STORAGE, 1));
1398		let light_bucket = read.page_of_key(&RowKey::encoded(STORAGE, 1u64 << 16));
1399		assert_ne!(heavy_bucket, light_bucket, "the two row groups must land in different buckets");
1400		assert!(!read.page_is_complete(heavy_bucket), "nothing is warm before the scan");
1401
1402		let scanned = store
1403			.range(
1404				RowKey::full_scan(STORAGE),
1405				MultiVersionScope::AsOf {
1406					read: CommitVersion(10),
1407				},
1408				32,
1409			)
1410			.collect::<Result<Vec<_>, _>>()
1411			.unwrap();
1412		assert_eq!(scanned.len() as u64, HEAVY + LIGHT, "the scan returns every row regardless of warming");
1413
1414		assert!(read.page_is_complete(heavy_bucket), "a bucket scanned past the threshold must be warmed");
1415		assert!(
1416			!read.page_is_complete(light_bucket),
1417			"a bucket scanned below the threshold must not be warmed"
1418		);
1419	}
1420
1421	#[test]
1422	fn operator_state_commit_does_not_populate_the_read_tier() {
1423		let (store, _g) = StandardMultiStore::testing_memory_with_persistent_sqlite();
1424		let read = store.read.clone().expect("read tier configured");
1425
1426		let opkey =
1427			OperatorStateKey::new(OperatorId(7), GroupId::ROOT, Keyspace::CUSTOM, vec![1, 2, 3]).encode();
1428		MultiVersionCommit::commit(
1429			&store,
1430			cow_vec![Delta::Set {
1431				key: opkey.clone(),
1432				bytes: EncodedBytes(CowVec::new(b"state-v10".to_vec())),
1433			}],
1434			CommitVersion(10),
1435		)
1436		.unwrap();
1437
1438		assert!(
1439			matches!(read.get(&opkey, CommitVersion(10)), VersionedGetResult::NotFound),
1440			"an operator commit must not write through into the read tier"
1441		);
1442		assert_eq!(read.resident_pages(), 0, "no operator page may become resident on commit");
1443
1444		let row = MultiVersionGet::get(&store, &opkey, CommitVersion(10))
1445			.unwrap()
1446			.expect("the committed operator state must still be readable through the store");
1447		assert_eq!(row.bytes.as_slice(), b"state-v10");
1448		assert_eq!(row.version, CommitVersion(10));
1449
1450		assert!(
1451			matches!(read.get(&opkey, CommitVersion(10)), VersionedGetResult::NotFound),
1452			"a store-level operator read must not back-populate the read tier"
1453		);
1454	}
1455
1456	#[test]
1457	fn source_row_write_clears_range_complete_on_its_page() {
1458		let (store, _g) = StandardMultiStore::testing_memory_with_persistent_sqlite();
1459		let read = store.read.clone().expect("read tier configured");
1460
1461		let neighbor = RowKey::encoded(STORAGE, 1);
1462		let page = read.page_of_key(&neighbor);
1463		assert_eq!(
1464			read.page_of_key(&RowKey::encoded(STORAGE, 2)),
1465			page,
1466			"both source rows must share a page for this test to exercise flag-clearing"
1467		);
1468		read.populate_page(
1469			page,
1470			vec![RawEntry {
1471				key: neighbor,
1472				version: CommitVersion(1),
1473				value: Some(CowVec::new(b"neighbor".to_vec())),
1474			}],
1475			true,
1476		);
1477		assert!(read.page_is_complete(page), "the page must start range-complete");
1478
1479		commit_row(&store, 2, 5);
1480
1481		assert!(
1482			!read.page_is_complete(page),
1483			"writing a source row into a range-complete page must clear the flag so the range cache re-warms"
1484		);
1485	}
1486
1487	#[test]
1488	fn source_warm_does_not_publish_a_page_another_warm_has_claimed() {
1489		const HEAVY: u64 = WARM_THRESHOLD + 64;
1490		let (store, _g) = StandardMultiStore::testing_memory_with_persistent_sqlite();
1491
1492		for n in 1..=HEAVY {
1493			commit_row(&store, n, 1);
1494		}
1495		flush(&store, CommitVersion(1));
1496
1497		let read = store.read.clone().expect("read tier configured");
1498		let page = read.page_of_key(&RowKey::encoded(STORAGE, 1));
1499		assert!(!read.page_is_complete(page), "nothing is warm before the scan");
1500
1501		assert!(read.begin_warm(page), "the page is unclaimed, so this claim must succeed");
1502
1503		let scanned = store
1504			.range(
1505				RowKey::full_scan(STORAGE),
1506				MultiVersionScope::AsOf {
1507					read: CommitVersion(10),
1508				},
1509				32,
1510			)
1511			.collect::<Result<Vec<_>, _>>()
1512			.unwrap();
1513		assert_eq!(scanned.len() as u64, HEAVY, "the scan still returns every row");
1514
1515		assert!(
1516			!read.page_is_complete(page),
1517			"a source range scan published a page that another warm had claimed. The operator warm path \
1518			 claims with begin_warm and publishes with finish_warm, which refuses a claim that a \
1519			 concurrent drop has dirtied; the source path claims nothing and publishes with \
1520			 populate_page, which sets range_complete unconditionally. So a drop landing during a source \
1521			 warm cannot invalidate it, and the stale pre-drop snapshot is republished as authoritative - \
1522			 resurrecting the dropped row in both point reads and range scans, permanently, because the \
1523			 persistent tier no longer holds anything to contradict the cache"
1524		);
1525	}
1526
1527	#[test]
1528	fn source_warm_releases_its_claim_when_it_publishes() {
1529		const HEAVY: u64 = WARM_THRESHOLD + 64;
1530		let (store, _g) = StandardMultiStore::testing_memory_with_persistent_sqlite();
1531
1532		for n in 1..=HEAVY {
1533			commit_row(&store, n, 1);
1534		}
1535		flush(&store, CommitVersion(1));
1536
1537		let read = store.read.clone().expect("read tier configured");
1538		let page = read.page_of_key(&RowKey::encoded(STORAGE, 1));
1539
1540		let _ = store
1541			.range(
1542				RowKey::full_scan(STORAGE),
1543				MultiVersionScope::AsOf {
1544					read: CommitVersion(10),
1545				},
1546				32,
1547			)
1548			.collect::<Result<Vec<_>, _>>()
1549			.unwrap();
1550		assert!(read.page_is_complete(page), "a bucket scanned past the threshold must be warmed");
1551
1552		assert!(
1553			read.begin_warm(page),
1554			"the source warm did not hand its claim back. Publishing through finish_warm consumes the \
1555			 claim; publishing through populate_page leaves it stranded in shard.warming, and then every \
1556			 later begin_warm on this page is refused - so once the page is invalidated it can never warm \
1557			 again for the life of the process"
1558		);
1559	}
1560}