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, HashSet},
6	ops::{Bound, RangeBounds},
7};
8
9use reifydb_codec::{
10	encoded::row::EncodedRow,
11	key::encoded::{EncodedKey, EncodedKeyRange},
12};
13use reifydb_core::{
14	actors::drop::{DropMessage, DropRequest},
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		is_single_version_semantics_key,
22	},
23};
24use reifydb_store::row::page::PageId;
25use reifydb_value::{
26	reifydb_assertions,
27	util::{cowvec::CowVec, hex},
28};
29use tracing::{Span, field, instrument, warn};
30
31use super::StandardMultiStore;
32use crate::{
33	MultiVersionScope, Result,
34	tier::{
35		RangeBatch, RangeCursor, TierBatch, TierStorage, VersionedGetResult,
36		commit::buffer::MultiCommitBufferTier,
37		persistent::MultiPersistentTier,
38		read::{MultiReadBufferTier, ServedChunk},
39	},
40};
41
42const TIER_SCAN_CHUNK_SIZE: usize = 32;
43
44const OPERATOR_PAGE_WARM_CAP: usize = 131_072;
45
46pub(crate) const WARM_THRESHOLD: u64 = 4 * TIER_SCAN_CHUNK_SIZE as u64;
47
48impl MultiVersionGet for StandardMultiStore {
49	fn get(&self, key: &EncodedKey, version: CommitVersion) -> Result<Option<MultiVersionRow>> {
50		match classify_key(key) {
51			EntryKind::Operator(_) => self.get_operator(key, version),
52			EntryKind::OperatorInternal(_) => self.get_operator_internal(key, version),
53			EntryKind::Source(_) => self.get_source(key, version),
54			_ => self.get_multi(key, version),
55		}
56	}
57}
58
59impl StandardMultiStore {
60	#[instrument(name = "store::multi::get::operator", level = "trace", skip(self, key), fields(version = version.0))]
61	fn get_operator(&self, key: &EncodedKey, version: CommitVersion) -> Result<Option<MultiVersionRow>> {
62		self.get_impl(key, version)
63	}
64
65	#[instrument(name = "store::multi::get::operator_internal", level = "trace", skip(self, key), fields(version = version.0))]
66	fn get_operator_internal(&self, key: &EncodedKey, version: CommitVersion) -> Result<Option<MultiVersionRow>> {
67		self.get_impl(key, version)
68	}
69
70	#[instrument(name = "store::multi::get::source", level = "trace", skip(self, key), fields(version = version.0))]
71	fn get_source(&self, key: &EncodedKey, version: CommitVersion) -> Result<Option<MultiVersionRow>> {
72		self.get_impl(key, version)
73	}
74
75	#[instrument(name = "store::multi::get::multi", level = "trace", skip(self, key), fields(version = version.0))]
76	fn get_multi(&self, key: &EncodedKey, version: CommitVersion) -> Result<Option<MultiVersionRow>> {
77		self.get_impl(key, version)
78	}
79
80	#[inline]
81	fn get_impl(&self, key: &EncodedKey, version: CommitVersion) -> Result<Option<MultiVersionRow>> {
82		let table = classify_key(key);
83
84		if let Some(found) = self.get_probe_commit(table, key, version)? {
85			return Ok(found);
86		}
87		if let Some(found) = self.get_probe_read(key, version) {
88			return Ok(self.unmask_dropped(key, found, version));
89		}
90		if matches!(table, EntryKind::Operator(_) | EntryKind::OperatorInternal(_))
91			&& let Some(read) = &self.read
92			&& self.warm_operator_page(read.page_of_key(key))?
93			&& let Some(found) = self.get_probe_read(key, version)
94		{
95			return Ok(self.unmask_dropped(key, found, version));
96		}
97		if let Some(found) = self.get_probe_persistent(table, key, version)? {
98			return Ok(self.unmask_dropped(key, found, version));
99		}
100
101		Ok(None)
102	}
103
104	#[inline]
105	fn unmask_dropped(
106		&self,
107		key: &EncodedKey,
108		found: Option<MultiVersionRow>,
109		read: CommitVersion,
110	) -> Option<MultiVersionRow> {
111		match found {
112			Some(row) if self.pending_drops.masks(key, row.version, read) => None,
113			other => other,
114		}
115	}
116}
117
118impl StandardMultiStore {
119	#[inline]
120	fn get_probe_commit(
121		&self,
122		table: EntryKind,
123		key: &EncodedKey,
124		version: CommitVersion,
125	) -> Result<Option<Option<MultiVersionRow>>> {
126		let Some(commit) = &self.commit else {
127			return Ok(None);
128		};
129		Ok(match commit.get(table, key.as_ref(), version)? {
130			VersionedGetResult::Value {
131				value,
132				version: v,
133			} => Some(Some(MultiVersionRow {
134				key: key.clone(),
135				row: EncodedRow(value),
136				version: v,
137			})),
138			VersionedGetResult::Tombstone => Some(None),
139			VersionedGetResult::NotFound => None,
140		})
141	}
142
143	#[inline]
144	fn get_probe_read(&self, key: &EncodedKey, version: CommitVersion) -> Option<Option<MultiVersionRow>> {
145		let read = self.read.as_ref()?;
146		match read.get(key, version) {
147			VersionedGetResult::Value {
148				value,
149				version: v,
150			} => Some(Some(MultiVersionRow {
151				key: key.clone(),
152				row: EncodedRow(value),
153				version: v,
154			})),
155			VersionedGetResult::Tombstone => Some(None),
156			VersionedGetResult::NotFound => None,
157		}
158	}
159
160	#[inline]
161	fn get_probe_persistent(
162		&self,
163		table: EntryKind,
164		key: &EncodedKey,
165		version: CommitVersion,
166	) -> Result<Option<Option<MultiVersionRow>>> {
167		let Some(persistent) = &self.persistent else {
168			return Ok(None);
169		};
170		Ok(match persistent.get(table, key.as_ref(), version)? {
171			VersionedGetResult::Value {
172				value,
173				version: v,
174			} => {
175				if let Some(read) = &self.read {
176					read.insert(key.clone(), v, Some(value.clone()));
177				}
178				Some(Some(MultiVersionRow {
179					key: key.clone(),
180					row: EncodedRow(value),
181					version: v,
182				}))
183			}
184			VersionedGetResult::Tombstone => Some(None),
185			VersionedGetResult::NotFound => None,
186		})
187	}
188
189	#[instrument(name = "store::multi::warm_operator", level = "debug", skip(self), fields(node = ?page.kind, outcome = field::Empty, loaded = field::Empty))]
190	fn warm_operator_page(&self, page: PageId) -> Result<bool> {
191		let span = Span::current();
192		let (Some(read), Some(persistent)) = (&self.read, &self.persistent) else {
193			span.record("outcome", "no_tiers");
194			return Ok(false);
195		};
196		if !matches!(page.kind, EntryKind::Operator(_) | EntryKind::OperatorInternal(_)) {
197			span.record("outcome", "not_operator");
198			return Ok(false);
199		}
200		if read.page_is_complete(page) {
201			span.record("outcome", "already_complete");
202			return Ok(true);
203		}
204		if !read.page_is_warm_candidate(page) {
205			span.record("outcome", "blocked");
206			return Ok(false);
207		}
208		let Some(range) = read.page_key_range(page) else {
209			span.record("outcome", "no_range");
210			return Ok(false);
211		};
212		if !read.begin_warm(page) {
213			span.record("outcome", "busy");
214			return Ok(false);
215		}
216		let loaded = persistent.load_range_consistent(
217			page.kind,
218			bound_as_slice(&range.start),
219			bound_as_slice(&range.end),
220			CommitVersion(u64::MAX),
221			Some(OPERATOR_PAGE_WARM_CAP + 1),
222		);
223		let entries = match loaded {
224			Ok(entries) => entries,
225			Err(e) => {
226				read.abort_warm(page);
227				span.record("outcome", "load_error");
228				return Err(e);
229			}
230		};
231		span.record("loaded", entries.len());
232		if entries.len() > OPERATOR_PAGE_WARM_CAP {
233			read.abort_warm(page);
234			read.set_warm_blocked(page);
235			span.record("outcome", "over_cap");
236			return Ok(false);
237		}
238		if read.finish_warm(page, entries) {
239			span.record("outcome", "completed");
240			Ok(true)
241		} else {
242			span.record("outcome", "dirty_abort");
243			Ok(false)
244		}
245	}
246}
247
248#[inline]
249fn bound_as_slice(bound: &Bound<EncodedKey>) -> Bound<&[u8]> {
250	match bound {
251		Bound::Included(k) => Bound::Included(k.as_slice()),
252		Bound::Excluded(k) => Bound::Excluded(k.as_slice()),
253		Bound::Unbounded => Bound::Unbounded,
254	}
255}
256
257impl MultiVersionContains for StandardMultiStore {
258	#[instrument(name = "store::multi::contains", level = "trace", skip(self), fields(key_hex = %hex::display(key.as_ref()), version = version.0), ret)]
259	fn contains(&self, key: &EncodedKey, version: CommitVersion) -> Result<bool> {
260		Ok(MultiVersionGet::get(self, key, version)?.is_some())
261	}
262}
263
264impl MultiVersionCommit for StandardMultiStore {
265	#[instrument(name = "store::multi::commit", level = "debug", skip(self, deltas), fields(delta_count = deltas.len(), version = version.0, drop_count = field::Empty))]
266	fn commit(&self, deltas: CowVec<Delta>, version: CommitVersion) -> Result<()> {
267		let classified = classify_deltas(&deltas);
268
269		let (operator_drops, source_drops) = partition_drops(classified.explicit_drops);
270		Span::current().record("drop_count", operator_drops.len() + source_drops.len());
271		self.dispatch_drops(build_drop_batch(source_drops, &classified.pending_set_keys, version));
272
273		self.update_read_cache_on_commit(version, &classified.batches);
274
275		if !self.write_batches(version, classified.batches)? {
276			return Ok(());
277		}
278
279		self.evict_operator_state(&operator_drops, version)?;
280		self.emit_commit_metrics(classified.writes, classified.deletes, version);
281
282		Ok(())
283	}
284}
285
286type DropPartition = (Vec<(EntryKind, EncodedKey)>, Vec<(EntryKind, EncodedKey)>);
287
288#[inline]
289fn partition_drops(explicit_drops: Vec<(EntryKind, EncodedKey)>) -> DropPartition {
290	explicit_drops
291		.into_iter()
292		.partition(|(table, _)| matches!(table, EntryKind::Operator(_) | EntryKind::OperatorInternal(_)))
293}
294
295struct ClassifiedDeltas {
296	pending_set_keys: HashSet<EncodedKey>,
297	writes: Vec<MultiWrite>,
298	deletes: Vec<MultiDelete>,
299	batches: TierBatch,
300	explicit_drops: Vec<(EntryKind, EncodedKey)>,
301}
302
303#[inline]
304fn classify_deltas(deltas: &CowVec<Delta>) -> ClassifiedDeltas {
305	let mut pending_set_keys: HashSet<EncodedKey> = HashSet::new();
306	let mut writes: Vec<MultiWrite> = Vec::new();
307	let mut deletes: Vec<MultiDelete> = Vec::new();
308	let mut batches: TierBatch = HashMap::new();
309	let mut explicit_drops: Vec<(EntryKind, EncodedKey)> = Vec::new();
310
311	for delta in deltas.iter() {
312		let key = delta.key();
313		let table = classify_key(key);
314		let is_single_version = is_single_version_semantics_key(key);
315
316		match delta {
317			Delta::Set {
318				key,
319				row,
320			} => {
321				if is_single_version {
322					pending_set_keys.insert(key.clone());
323				}
324				writes.push(MultiWrite {
325					key: key.clone(),
326					value_bytes: row.len() as u64,
327				});
328				batches.entry(table).or_default().push((key.clone(), Some(row.0.clone())));
329			}
330			Delta::Unset {
331				key,
332				row,
333			} => {
334				deletes.push(MultiDelete {
335					key: key.clone(),
336					value_bytes: row.len() as u64,
337				});
338				batches.entry(table).or_default().push((key.clone(), None));
339			}
340			Delta::Remove {
341				key,
342			} => {
343				deletes.push(MultiDelete {
344					key: key.clone(),
345					value_bytes: 0,
346				});
347				batches.entry(table).or_default().push((key.clone(), None));
348			}
349			Delta::Drop {
350				key,
351			} => {
352				explicit_drops.push((table, key.clone()));
353			}
354		}
355	}
356
357	ClassifiedDeltas {
358		pending_set_keys,
359		writes,
360		deletes,
361		batches,
362		explicit_drops,
363	}
364}
365
366#[inline]
367fn build_drop_batch(
368	explicit_drops: Vec<(EntryKind, EncodedKey)>,
369	pending_set_keys: &HashSet<EncodedKey>,
370	version: CommitVersion,
371) -> Vec<DropRequest> {
372	let mut drop_batch = Vec::with_capacity(explicit_drops.len() + pending_set_keys.len());
373	for (table, key) in explicit_drops {
374		let pending_version = if pending_set_keys.contains(key.as_ref()) {
375			Some(version)
376		} else {
377			None
378		};
379		drop_batch.push(DropRequest {
380			table,
381			key,
382			commit_version: version,
383			pending_version,
384		});
385	}
386	for key in pending_set_keys.iter() {
387		let encoded = EncodedKey::new(key.to_vec());
388		let table = classify_key(&encoded);
389		drop_batch.push(DropRequest {
390			table,
391			key: encoded,
392			commit_version: version,
393			pending_version: Some(version),
394		});
395	}
396	drop_batch
397}
398
399impl StandardMultiStore {
400	pub fn get_many(
401		&self,
402		keys: &[EncodedKey],
403		version: CommitVersion,
404	) -> Result<HashMap<EncodedKey, MultiVersionRow>> {
405		let mut by_table: HashMap<EntryKind, Vec<&EncodedKey>> = HashMap::new();
406		for key in keys {
407			by_table.entry(classify_key(key)).or_default().push(key);
408		}
409
410		let mut out: HashMap<EncodedKey, MultiVersionRow> = HashMap::new();
411		for (table, table_keys) in by_table {
412			self.get_many_for_table(table, &table_keys, version, &mut out)?;
413		}
414
415		Ok(out)
416	}
417
418	#[inline]
419	fn get_many_for_table(
420		&self,
421		table: EntryKind,
422		table_keys: &[&EncodedKey],
423		version: CommitVersion,
424		out: &mut HashMap<EncodedKey, MultiVersionRow>,
425	) -> Result<()> {
426		let key_slices: Vec<&[u8]> = table_keys.iter().map(|k| k.as_ref()).collect();
427
428		let commit_results = self.probe_commit_batch(table, &key_slices, version)?;
429		let (read_aligned, persistent_aligned) = self.resolve_misses_through_read_and_persistent(
430			table,
431			table_keys,
432			&key_slices,
433			&commit_results,
434			version,
435		)?;
436
437		reifydb_assertions! {
438			let n = key_slices.len();
439			assert!(
440				commit_results.len() == n && read_aligned.len() == n && persistent_aligned.len() == n,
441				"per-tier result vectors must stay index-aligned with the table's keys, otherwise collect_resolved_rows \
442				 reads a tier result for the wrong key and returns mismatched rows (keys={n}, commit={}, read={}, persistent={})",
443				commit_results.len(),
444				read_aligned.len(),
445				persistent_aligned.len()
446			);
447		}
448
449		self.collect_resolved_rows(table_keys, &commit_results, &read_aligned, &persistent_aligned, out);
450		Ok(())
451	}
452
453	#[inline]
454	fn probe_commit_batch(
455		&self,
456		table: EntryKind,
457		key_slices: &[&[u8]],
458		version: CommitVersion,
459	) -> Result<Vec<VersionedGetResult>> {
460		match &self.commit {
461			Some(commit) => commit.get_many(table, key_slices, version),
462			None => Ok(vec![VersionedGetResult::NotFound; key_slices.len()]),
463		}
464	}
465
466	#[inline]
467	fn resolve_misses_through_read_and_persistent(
468		&self,
469		table: EntryKind,
470		table_keys: &[&EncodedKey],
471		key_slices: &[&[u8]],
472		commit_results: &[VersionedGetResult],
473		version: CommitVersion,
474	) -> Result<(Vec<VersionedGetResult>, Vec<VersionedGetResult>)> {
475		let mut read_aligned = vec![VersionedGetResult::NotFound; key_slices.len()];
476		let mut persistent_idx: Vec<usize> = Vec::new();
477		let mut persistent_slices: Vec<&[u8]> = Vec::new();
478		for (i, result) in commit_results.iter().enumerate() {
479			if !matches!(result, VersionedGetResult::NotFound) {
480				continue;
481			}
482			let read_hit = self
483				.read
484				.as_ref()
485				.map(|c| c.get(table_keys[i], version))
486				.unwrap_or(VersionedGetResult::NotFound);
487			match read_hit {
488				VersionedGetResult::Value {
489					value,
490					version: v,
491				} => {
492					read_aligned[i] = if self.pending_drops.masks(table_keys[i], v, version) {
493						VersionedGetResult::Tombstone
494					} else {
495						VersionedGetResult::Value {
496							value,
497							version: v,
498						}
499					};
500				}
501				VersionedGetResult::Tombstone => {
502					read_aligned[i] = VersionedGetResult::Tombstone;
503				}
504				VersionedGetResult::NotFound => {
505					persistent_idx.push(i);
506					persistent_slices.push(key_slices[i]);
507				}
508			}
509		}
510
511		if matches!(table, EntryKind::Operator(_) | EntryKind::OperatorInternal(_))
512			&& !persistent_idx.is_empty()
513			&& let Some(read) = &self.read
514		{
515			let mut pages: Vec<PageId> = Vec::new();
516			for &i in &persistent_idx {
517				let page = read.page_of_key(table_keys[i]);
518				if !pages.contains(&page) {
519					pages.push(page);
520				}
521			}
522			let mut warmed_any = false;
523			for page in pages {
524				warmed_any |= self.warm_operator_page(page)?;
525			}
526			if warmed_any {
527				let mut remaining_idx = Vec::new();
528				let mut remaining_slices = Vec::new();
529				for &i in &persistent_idx {
530					match read.get(table_keys[i], version) {
531						VersionedGetResult::Value {
532							value,
533							version: v,
534						} => {
535							read_aligned[i] = if self.pending_drops.masks(
536								table_keys[i],
537								v,
538								version,
539							) {
540								VersionedGetResult::Tombstone
541							} else {
542								VersionedGetResult::Value {
543									value,
544									version: v,
545								}
546							};
547						}
548						VersionedGetResult::Tombstone => {
549							read_aligned[i] = VersionedGetResult::Tombstone;
550						}
551						VersionedGetResult::NotFound => {
552							remaining_idx.push(i);
553							remaining_slices.push(key_slices[i]);
554						}
555					}
556				}
557				persistent_idx = remaining_idx;
558				persistent_slices = remaining_slices;
559			}
560		}
561
562		let mut persistent_aligned = vec![VersionedGetResult::NotFound; key_slices.len()];
563		if !persistent_slices.is_empty()
564			&& let Some(persistent) = &self.persistent
565		{
566			let persistent_results = persistent.get_many(table, &persistent_slices, version)?;
567			for (slot, result) in persistent_idx.into_iter().zip(persistent_results) {
568				if let VersionedGetResult::Value {
569					version: v,
570					..
571				} = &result && self.pending_drops.masks(table_keys[slot], *v, version)
572				{
573					persistent_aligned[slot] = VersionedGetResult::Tombstone;
574					continue;
575				}
576				if let (
577					Some(read),
578					VersionedGetResult::Value {
579						value,
580						version: v,
581					},
582				) = (&self.read, &result)
583				{
584					read.insert(table_keys[slot].clone(), *v, Some(value.clone()));
585				}
586				persistent_aligned[slot] = result;
587			}
588		}
589
590		Ok((read_aligned, persistent_aligned))
591	}
592
593	#[inline]
594	fn collect_resolved_rows(
595		&self,
596		table_keys: &[&EncodedKey],
597		commit_results: &[VersionedGetResult],
598		read_aligned: &[VersionedGetResult],
599		persistent_aligned: &[VersionedGetResult],
600		out: &mut HashMap<EncodedKey, MultiVersionRow>,
601	) {
602		for (i, key) in table_keys.iter().enumerate() {
603			let resolved = match &commit_results[i] {
604				VersionedGetResult::Value {
605					value,
606					version: v,
607				} => Some((value.clone(), *v)),
608				VersionedGetResult::Tombstone => None,
609				VersionedGetResult::NotFound => match &read_aligned[i] {
610					VersionedGetResult::Value {
611						value,
612						version: v,
613					} => Some((value.clone(), *v)),
614					VersionedGetResult::Tombstone => None,
615					VersionedGetResult::NotFound => match &persistent_aligned[i] {
616						VersionedGetResult::Value {
617							value,
618							version: v,
619						} => Some((value.clone(), *v)),
620						_ => None,
621					},
622				},
623			};
624
625			if let Some((value, v)) = resolved {
626				out.insert(
627					(*key).clone(),
628					MultiVersionRow {
629						key: (*key).clone(),
630						row: EncodedRow(value),
631						version: v,
632					},
633				);
634			}
635		}
636	}
637
638	#[inline]
639	fn dispatch_drops(&self, drop_batch: Vec<DropRequest>) {
640		if drop_batch.is_empty() {
641			return;
642		}
643		if let Some(actor) = &self.drop_actor
644			&& actor.send_blocking(DropMessage::Batch(drop_batch)).is_err()
645		{
646			warn!("Failed to send drop batch");
647		}
648	}
649
650	#[inline]
651	fn update_read_cache_on_commit(&self, version: CommitVersion, batches: &TierBatch) {
652		let Some(read) = &self.read else {
653			return;
654		};
655		for (table, entries) in batches {
656			match table {
657				EntryKind::Operator(_) | EntryKind::OperatorInternal(_) => {
658					for (key, value) in entries {
659						match value {
660							Some(value) => {
661								read.insert(key.clone(), version, Some(value.clone()))
662							}
663							None => read.insert(key.clone(), version, None),
664						}
665					}
666				}
667				_ => {
668					for (key, _) in entries {
669						read.invalidate(key);
670					}
671				}
672			}
673		}
674	}
675
676	#[inline]
677	fn write_batches(&self, version: CommitVersion, batches: TierBatch) -> Result<bool> {
678		if let Some(commit) = &self.commit {
679			commit.set(version, batches)?;
680		} else if let Some(persistent) = &self.persistent {
681			persistent.set(version, batches)?;
682		} else {
683			return Ok(false);
684		}
685		Ok(true)
686	}
687
688	#[instrument(name = "store::multi::evict_drops", level = "debug", skip_all, fields(drop_count = field::Empty))]
689	fn evict_operator_state(&self, drops: &[(EntryKind, EncodedKey)], version: CommitVersion) -> Result<()> {
690		if drops.is_empty() {
691			return Ok(());
692		}
693		Span::current().record("drop_count", drops.len());
694
695		self.record_pending_drops(drops, version);
696		self.evict_drops_from_commit(drops)?;
697		self.remove_drops_from_read(drops);
698		if !self.nudge_drop_purge() {
699			self.pending_drops.purge(self.persistent.as_ref(), self.read.as_ref());
700		}
701
702		Ok(())
703	}
704
705	#[inline]
706	fn nudge_drop_purge(&self) -> bool {
707		if self.persistent.is_none() {
708			return true;
709		}
710		let Some(actor) = &self.drop_actor else {
711			return false;
712		};
713		if actor.send_blocking(DropMessage::PurgePending).is_err() {
714			warn!("Failed to nudge drop purge, purging synchronously");
715			return false;
716		}
717		true
718	}
719
720	#[inline]
721	fn record_pending_drops(&self, drops: &[(EntryKind, EncodedKey)], version: CommitVersion) {
722		if self.persistent.is_none() {
723			return;
724		}
725		for (_, key) in drops {
726			self.pending_drops.record(key.clone(), version);
727		}
728	}
729
730	#[inline]
731	fn remove_drops_from_read(&self, drops: &[(EntryKind, EncodedKey)]) {
732		let Some(read) = &self.read else {
733			return;
734		};
735		for (_, key) in drops {
736			read.remove_dropped(key);
737		}
738	}
739
740	#[inline]
741	fn evict_drops_from_commit(&self, drops: &[(EntryKind, EncodedKey)]) -> Result<()> {
742		let Some(commit) = &self.commit else {
743			return Ok(());
744		};
745		let mut batches: HashMap<EntryKind, Vec<(EncodedKey, CommitVersion)>> = HashMap::new();
746		for (table, key) in drops {
747			for (entry_version, _) in commit.get_all_versions(*table, key.as_ref())? {
748				batches.entry(*table).or_default().push((key.clone(), entry_version));
749			}
750		}
751		if !batches.is_empty() {
752			commit.drop(batches)?;
753		}
754		Ok(())
755	}
756
757	#[inline]
758	fn emit_commit_metrics(&self, writes: Vec<MultiWrite>, deletes: Vec<MultiDelete>, version: CommitVersion) {
759		if writes.is_empty() && deletes.is_empty() {
760			return;
761		}
762		self.event_bus.emit(MultiCommittedEvent::new(writes, deletes, vec![], version));
763	}
764}
765
766#[derive(Debug, Clone, Default)]
767pub struct MultiVersionRangeCursor {
768	pub commit: RangeCursor,
769
770	pub persistent: RangeCursor,
771
772	pub exhausted: bool,
773
774	warm_bucket: Option<PageId>,
775
776	warm_consumed: u64,
777}
778
779impl MultiVersionRangeCursor {
780	pub fn new() -> Self {
781		Self::default()
782	}
783
784	pub fn is_exhausted(&self) -> bool {
785		self.exhausted
786	}
787}
788
789pub struct TierScanQuery<'a> {
790	pub table: EntryKind,
791	pub start: &'a [u8],
792	pub end: &'a [u8],
793	pub scope: MultiVersionScope,
794	pub range: &'a EncodedKeyRange,
795}
796
797pub fn scan_tier_chunk<S: TierStorage>(
798	storage: &S,
799	cursor: &mut RangeCursor,
800	scan: &TierScanQuery,
801	collected: &mut BTreeMap<Vec<u8>, (CommitVersion, Option<CowVec<u8>>)>,
802) -> Result<bool> {
803	let batch = storage.range_next(
804		scan.table,
805		cursor,
806		Bound::Included(scan.start),
807		Bound::Included(scan.end),
808		scan.scope,
809		TIER_SCAN_CHUNK_SIZE,
810	)?;
811	merge_tier_batch(batch, scan.range, collected)
812}
813
814pub fn scan_tier_chunk_rev<S: TierStorage>(
815	storage: &S,
816	cursor: &mut RangeCursor,
817	scan: &TierScanQuery,
818	collected: &mut BTreeMap<Vec<u8>, (CommitVersion, Option<CowVec<u8>>)>,
819) -> Result<bool> {
820	let batch = storage.range_rev_next(
821		scan.table,
822		cursor,
823		Bound::Included(scan.start),
824		Bound::Included(scan.end),
825		scan.scope,
826		TIER_SCAN_CHUNK_SIZE,
827	)?;
828	merge_tier_batch(batch, scan.range, collected)
829}
830
831#[inline]
832fn merge_tier_batch(
833	batch: RangeBatch,
834	range: &EncodedKeyRange,
835	collected: &mut BTreeMap<Vec<u8>, (CommitVersion, Option<CowVec<u8>>)>,
836) -> Result<bool> {
837	if batch.entries.is_empty() {
838		return Ok(false);
839	}
840
841	for entry in batch.entries {
842		let original_key = entry.key.as_slice().to_vec();
843		let entry_version = entry.version;
844
845		let original_key_encoded = EncodedKey::new(original_key.clone());
846		if !range.contains(&original_key_encoded) {
847			continue;
848		}
849
850		let should_update = match collected.get(&original_key) {
851			None => true,
852			Some((existing_version, _)) => entry_version > *existing_version,
853		};
854
855		if should_update {
856			collected.insert(original_key, (entry_version, entry.value));
857		}
858	}
859
860	Ok(true)
861}
862
863#[inline]
864pub fn collected_to_batch(
865	collected: BTreeMap<Vec<u8>, (CommitVersion, Option<CowVec<u8>>)>,
866	has_more: bool,
867) -> MultiVersionBatch {
868	let items: Vec<MultiVersionRow> = collected
869		.into_iter()
870		.filter_map(|(key_bytes, (v, value))| {
871			value.map(|val| MultiVersionRow {
872				key: EncodedKey::new(key_bytes),
873				row: EncodedRow(val),
874				version: v,
875			})
876		})
877		.collect();
878
879	MultiVersionBatch {
880		items,
881		has_more,
882	}
883}
884
885#[inline]
886fn step_all_tiers(
887	buffer: Option<&MultiCommitBufferTier>,
888	buffer_cursor: &mut RangeCursor,
889	persistent: Option<&MultiPersistentTier>,
890	persistent_cursor: &mut RangeCursor,
891	scan: &TierScanQuery,
892	collected: &mut BTreeMap<Vec<u8>, (CommitVersion, Option<CowVec<u8>>)>,
893) -> Result<bool> {
894	let mut any_progress = false;
895	if let Some(s) = buffer
896		&& !buffer_cursor.exhausted
897	{
898		any_progress |= scan_tier_chunk(s, buffer_cursor, scan, collected)?;
899	}
900	if let Some(s) = persistent
901		&& !persistent_cursor.exhausted
902	{
903		any_progress |= scan_tier_chunk(s, persistent_cursor, scan, collected)?;
904	}
905	Ok(any_progress)
906}
907
908pub fn scan_tiers_latest(
909	buffer: Option<&MultiCommitBufferTier>,
910	persistent: Option<&MultiPersistentTier>,
911	range: EncodedKeyRange,
912	scope: MultiVersionScope,
913	max_keys: usize,
914) -> Result<MultiVersionBatch> {
915	let table = classify_key_range(&range);
916	let (start, end) = make_range_bounds(&range);
917	let scan = TierScanQuery {
918		table,
919		start: &start,
920		end: &end,
921		scope,
922		range: &range,
923	};
924
925	let mut collected: BTreeMap<Vec<u8>, (CommitVersion, Option<CowVec<u8>>)> = BTreeMap::new();
926	let mut buffer_cursor = RangeCursor::default();
927	let mut persistent_cursor = RangeCursor::default();
928	let mut exhausted = false;
929
930	while collected.len() < max_keys {
931		let progress = step_all_tiers(
932			buffer,
933			&mut buffer_cursor,
934			persistent,
935			&mut persistent_cursor,
936			&scan,
937			&mut collected,
938		)?;
939		if !progress {
940			exhausted = true;
941			break;
942		}
943	}
944
945	Ok(collected_to_batch(collected, !exhausted))
946}
947
948impl StandardMultiStore {
949	pub fn range_next(
950		&self,
951		cursor: &mut MultiVersionRangeCursor,
952		range: EncodedKeyRange,
953		scope: MultiVersionScope,
954		batch_size: u64,
955	) -> Result<MultiVersionBatch> {
956		if cursor.exhausted {
957			return Ok(MultiVersionBatch {
958				items: Vec::new(),
959				has_more: false,
960			});
961		}
962
963		mark_unconfigured_exhausted(self, cursor);
964
965		let table = classify_key_range(&range);
966		let (start, end) = make_range_bounds(&range);
967		let batch_size = batch_size as usize;
968		let scan = TierScanQuery {
969			table,
970			start: &start,
971			end: &end,
972			scope,
973			range: &range,
974		};
975
976		let mut collected: BTreeMap<Vec<u8>, (CommitVersion, Option<CowVec<u8>>)> = BTreeMap::new();
977
978		while collected.len() < batch_size {
979			let mut any_progress = false;
980
981			if let Some(commit) = &self.commit
982				&& !cursor.commit.exhausted
983			{
984				any_progress |= scan_tier_chunk(commit, &mut cursor.commit, &scan, &mut collected)?;
985			}
986
987			if self.persistent.is_some() && !cursor.persistent.exhausted {
988				any_progress |= self.step_persistent_cached(&scan, cursor, &mut collected, false)?;
989			}
990
991			if !any_progress {
992				cursor.exhausted = true;
993				break;
994			}
995		}
996
997		apply_forward_horizon(cursor, &mut collected);
998
999		let items: Vec<MultiVersionRow> = collected
1000			.into_iter()
1001			.filter_map(|(key_bytes, (v, value))| {
1002				value.map(|val| MultiVersionRow {
1003					key: EncodedKey::new(key_bytes),
1004					row: EncodedRow(val),
1005					version: v,
1006				})
1007			})
1008			.collect();
1009
1010		let has_more = !cursor.exhausted;
1011
1012		Ok(MultiVersionBatch {
1013			items,
1014			has_more,
1015		})
1016	}
1017
1018	pub fn range(
1019		&self,
1020		range: EncodedKeyRange,
1021		scope: MultiVersionScope,
1022		batch_size: usize,
1023	) -> MultiVersionRangeIter {
1024		MultiVersionRangeIter {
1025			store: self.clone(),
1026			cursor: MultiVersionRangeCursor::new(),
1027			range,
1028			scope,
1029			batch_size,
1030			current_batch: Vec::new(),
1031			current_index: 0,
1032		}
1033	}
1034
1035	pub fn range_rev(
1036		&self,
1037		range: EncodedKeyRange,
1038		scope: MultiVersionScope,
1039		batch_size: usize,
1040	) -> MultiVersionRangeRevIter {
1041		MultiVersionRangeRevIter {
1042			store: self.clone(),
1043			cursor: MultiVersionRangeCursor::new(),
1044			range,
1045			scope,
1046			batch_size,
1047			current_batch: Vec::new(),
1048			current_index: 0,
1049		}
1050	}
1051
1052	fn range_rev_next(
1053		&self,
1054		cursor: &mut MultiVersionRangeCursor,
1055		range: EncodedKeyRange,
1056		scope: MultiVersionScope,
1057		batch_size: u64,
1058	) -> Result<MultiVersionBatch> {
1059		if cursor.exhausted {
1060			return Ok(MultiVersionBatch {
1061				items: Vec::new(),
1062				has_more: false,
1063			});
1064		}
1065
1066		mark_unconfigured_exhausted(self, cursor);
1067
1068		let table = classify_key_range(&range);
1069		let (start, end) = make_range_bounds(&range);
1070		let batch_size = batch_size as usize;
1071		let scan = TierScanQuery {
1072			table,
1073			start: &start,
1074			end: &end,
1075			scope,
1076			range: &range,
1077		};
1078
1079		let mut collected: BTreeMap<Vec<u8>, (CommitVersion, Option<CowVec<u8>>)> = BTreeMap::new();
1080
1081		while collected.len() < batch_size {
1082			let mut any_progress = false;
1083
1084			if let Some(commit) = &self.commit
1085				&& !cursor.commit.exhausted
1086			{
1087				any_progress |= scan_tier_chunk_rev(commit, &mut cursor.commit, &scan, &mut collected)?;
1088			}
1089
1090			if self.persistent.is_some() && !cursor.persistent.exhausted {
1091				any_progress |= self.step_persistent_cached(&scan, cursor, &mut collected, true)?;
1092			}
1093
1094			if !any_progress {
1095				cursor.exhausted = true;
1096				break;
1097			}
1098		}
1099
1100		apply_reverse_horizon(cursor, &mut collected);
1101
1102		let items: Vec<MultiVersionRow> = collected
1103			.into_iter()
1104			.rev()
1105			.filter_map(|(key_bytes, (v, value))| {
1106				value.map(|val| MultiVersionRow {
1107					key: EncodedKey::new(key_bytes),
1108					row: EncodedRow(val),
1109					version: v,
1110				})
1111			})
1112			.collect();
1113
1114		let has_more = !cursor.exhausted;
1115
1116		Ok(MultiVersionBatch {
1117			items,
1118			has_more,
1119		})
1120	}
1121
1122	fn step_persistent_cached(
1123		&self,
1124		scan: &TierScanQuery,
1125		cursor: &mut MultiVersionRangeCursor,
1126		collected: &mut BTreeMap<Vec<u8>, (CommitVersion, Option<CowVec<u8>>)>,
1127		descending: bool,
1128	) -> Result<bool> {
1129		let Some(persistent) = &self.persistent else {
1130			return Ok(false);
1131		};
1132
1133		if let Some(served) = self.serve_from_read_cache(scan, cursor, collected, descending) {
1134			return served;
1135		}
1136
1137		if matches!(scan.table, EntryKind::Operator(_) | EntryKind::OperatorInternal(_))
1138			&& let Some(read) = &self.read
1139			&& self.warm_operator_page(read.page_of_key(&EncodedKey::new(scan.start.to_vec())))?
1140			&& let Some(served) = self.serve_from_read_cache(scan, cursor, collected, descending)
1141		{
1142			return served;
1143		}
1144
1145		let (consumed, progressed) =
1146			self.scan_persistent_chunk(persistent, scan, cursor, collected, descending)?;
1147		self.warm_read_bucket_after_scan(persistent, scan, cursor, consumed)?;
1148
1149		Ok(progressed)
1150	}
1151
1152	#[inline]
1153	fn serve_from_read_cache(
1154		&self,
1155		scan: &TierScanQuery,
1156		cursor: &mut MultiVersionRangeCursor,
1157		collected: &mut BTreeMap<Vec<u8>, (CommitVersion, Option<CowVec<u8>>)>,
1158		descending: bool,
1159	) -> Option<Result<bool>> {
1160		let (Some(read), EntryKind::Source(_) | EntryKind::Operator(_) | EntryKind::OperatorInternal(_)) =
1161			(&self.read, scan.table)
1162		else {
1163			return None;
1164		};
1165		match read.serve_persistent_chunk(
1166			scan.table,
1167			&mut cursor.persistent,
1168			scan.start,
1169			scan.end,
1170			scan.scope,
1171			TIER_SCAN_CHUNK_SIZE,
1172			descending,
1173		) {
1174			ServedChunk::Served(batch) => {
1175				let batch = self.mask_dropped_persistent_rows(scan, batch);
1176				Some(merge_tier_batch(batch, scan.range, collected))
1177			}
1178			ServedChunk::Gap => None,
1179		}
1180	}
1181
1182	#[inline]
1183	fn scan_persistent_chunk(
1184		&self,
1185		persistent: &MultiPersistentTier,
1186		scan: &TierScanQuery,
1187		cursor: &mut MultiVersionRangeCursor,
1188		collected: &mut BTreeMap<Vec<u8>, (CommitVersion, Option<CowVec<u8>>)>,
1189		descending: bool,
1190	) -> Result<(usize, bool)> {
1191		let batch = if descending {
1192			persistent.range_rev_next(
1193				scan.table,
1194				&mut cursor.persistent,
1195				Bound::Included(scan.start),
1196				Bound::Included(scan.end),
1197				scan.scope,
1198				TIER_SCAN_CHUNK_SIZE,
1199			)?
1200		} else {
1201			persistent.range_next(
1202				scan.table,
1203				&mut cursor.persistent,
1204				Bound::Included(scan.start),
1205				Bound::Included(scan.end),
1206				scan.scope,
1207				TIER_SCAN_CHUNK_SIZE,
1208			)?
1209		};
1210		let consumed = batch.entries.len();
1211		let batch = self.mask_dropped_persistent_rows(scan, batch);
1212		let progressed = merge_tier_batch(batch, scan.range, collected)?;
1213		Ok((consumed, progressed))
1214	}
1215
1216	#[inline]
1217	fn mask_dropped_persistent_rows(&self, scan: &TierScanQuery, mut batch: RangeBatch) -> RangeBatch {
1218		if !matches!(scan.table, EntryKind::Operator(_) | EntryKind::OperatorInternal(_))
1219			|| self.pending_drops.is_empty()
1220		{
1221			return batch;
1222		}
1223		for entry in batch.entries.iter_mut() {
1224			if entry.value.is_some()
1225				&& self.pending_drops.masks(&entry.key, entry.version, scan.scope.read())
1226			{
1227				entry.value = None;
1228			}
1229		}
1230		batch
1231	}
1232
1233	#[inline]
1234	fn warm_read_bucket_after_scan(
1235		&self,
1236		persistent: &MultiPersistentTier,
1237		scan: &TierScanQuery,
1238		cursor: &mut MultiVersionRangeCursor,
1239		consumed: usize,
1240	) -> Result<()> {
1241		if let (Some(read), EntryKind::Source(_)) = (&self.read, scan.table) {
1242			maybe_warm_bucket(read, persistent, cursor, scan.table, consumed)?;
1243		}
1244		Ok(())
1245	}
1246}
1247
1248fn maybe_warm_bucket(
1249	read: &MultiReadBufferTier,
1250	persistent: &MultiPersistentTier,
1251	cursor: &mut MultiVersionRangeCursor,
1252	table: EntryKind,
1253	consumed: usize,
1254) -> Result<()> {
1255	let page = {
1256		let Some(last) = cursor.persistent.last_key.as_ref() else {
1257			return Ok(());
1258		};
1259		read.page_of_key(last)
1260	};
1261	if !matches!(page.kind, EntryKind::Source(_)) {
1262		return Ok(());
1263	}
1264
1265	if cursor.warm_bucket == Some(page) {
1266		cursor.warm_consumed = cursor.warm_consumed.saturating_add(consumed as u64);
1267	} else {
1268		cursor.warm_bucket = Some(page);
1269		cursor.warm_consumed = consumed as u64;
1270	}
1271
1272	if cursor.warm_consumed <= WARM_THRESHOLD {
1273		return Ok(());
1274	}
1275
1276	let Some(range) = read.page_key_range(page) else {
1277		return Ok(());
1278	};
1279	let (Bound::Included(lo), Bound::Included(hi)) = (range.start, range.end) else {
1280		return Ok(());
1281	};
1282	let entries = persistent.load_range_consistent(
1283		table,
1284		Bound::Included(lo.as_slice()),
1285		Bound::Included(hi.as_slice()),
1286		CommitVersion(u64::MAX),
1287		None,
1288	)?;
1289	read.populate_page(page, entries, true);
1290	cursor.warm_bucket = None;
1291	cursor.warm_consumed = 0;
1292	Ok(())
1293}
1294
1295fn mark_unconfigured_exhausted(store: &StandardMultiStore, cursor: &mut MultiVersionRangeCursor) {
1296	if store.commit.is_none() {
1297		cursor.commit.exhausted = true;
1298	}
1299	if store.persistent.is_none() {
1300		cursor.persistent.exhausted = true;
1301	}
1302}
1303
1304fn apply_forward_horizon(
1305	cursor: &mut MultiVersionRangeCursor,
1306	collected: &mut BTreeMap<Vec<u8>, (CommitVersion, Option<CowVec<u8>>)>,
1307) {
1308	let horizon = forward_horizon(cursor);
1309	if let Some(h) = horizon {
1310		collected.retain(|k, _| k.as_slice() <= h.as_slice());
1311		rewind_over_advanced_forward(cursor, &h);
1312	}
1313}
1314
1315fn apply_reverse_horizon(
1316	cursor: &mut MultiVersionRangeCursor,
1317	collected: &mut BTreeMap<Vec<u8>, (CommitVersion, Option<CowVec<u8>>)>,
1318) {
1319	let horizon = reverse_horizon(cursor);
1320	if let Some(h) = horizon {
1321		collected.retain(|k, _| k.as_slice() >= h.as_slice());
1322		rewind_over_advanced_reverse(cursor, &h);
1323	}
1324}
1325
1326fn forward_horizon(cursor: &MultiVersionRangeCursor) -> Option<EncodedKey> {
1327	let mut horizon: Option<EncodedKey> = None;
1328	for tier in [&cursor.commit, &cursor.persistent] {
1329		if tier.exhausted {
1330			continue;
1331		}
1332		let last = match &tier.last_key {
1333			Some(k) => k.clone(),
1334
1335			None => return None,
1336		};
1337		horizon = Some(match horizon {
1338			None => last,
1339			Some(prev) => {
1340				if last.as_slice() < prev.as_slice() {
1341					last
1342				} else {
1343					prev
1344				}
1345			}
1346		});
1347	}
1348	horizon
1349}
1350
1351fn reverse_horizon(cursor: &MultiVersionRangeCursor) -> Option<EncodedKey> {
1352	let mut horizon: Option<EncodedKey> = None;
1353	for tier in [&cursor.commit, &cursor.persistent] {
1354		if tier.exhausted {
1355			continue;
1356		}
1357		let last = match &tier.last_key {
1358			Some(k) => k.clone(),
1359			None => return None,
1360		};
1361		horizon = Some(match horizon {
1362			None => last,
1363			Some(prev) => {
1364				if last.as_slice() > prev.as_slice() {
1365					last
1366				} else {
1367					prev
1368				}
1369			}
1370		});
1371	}
1372	horizon
1373}
1374
1375fn rewind_over_advanced_forward(cursor: &mut MultiVersionRangeCursor, horizon: &EncodedKey) {
1376	for tier in [&mut cursor.commit, &mut cursor.persistent] {
1377		if let Some(last) = &tier.last_key
1378			&& last.as_slice() > horizon.as_slice()
1379		{
1380			tier.last_key = Some(horizon.clone());
1381			tier.exhausted = false;
1382		}
1383	}
1384}
1385
1386fn rewind_over_advanced_reverse(cursor: &mut MultiVersionRangeCursor, horizon: &EncodedKey) {
1387	for tier in [&mut cursor.commit, &mut cursor.persistent] {
1388		if let Some(last) = &tier.last_key
1389			&& last.as_slice() < horizon.as_slice()
1390		{
1391			tier.last_key = Some(horizon.clone());
1392			tier.exhausted = false;
1393		}
1394	}
1395}
1396
1397impl MultiVersionGetPrevious for StandardMultiStore {
1398	fn get_previous_version(
1399		&self,
1400		key: &EncodedKey,
1401		before_version: CommitVersion,
1402	) -> Result<Option<MultiVersionRow>> {
1403		if before_version.0 == 0 {
1404			return Ok(None);
1405		}
1406
1407		let table = classify_key(key);
1408		reifydb_assertions! {
1409			assert!(
1410				before_version.0 >= 1,
1411				"the before_version==0 guard must precede this subtraction, otherwise before_version.0 - 1 \
1412				 wraps to u64::MAX and the probe reads the latest version instead of the previous one \
1413				 (before_version={})",
1414				before_version.0
1415			);
1416		}
1417		let prev_version = CommitVersion(before_version.0 - 1);
1418
1419		if let Some(found) = self.previous_probe_commit(table, key, prev_version)? {
1420			return Ok(found);
1421		}
1422		if let Some(found) = self.previous_probe_read(key, prev_version) {
1423			return Ok(found);
1424		}
1425		if let Some(found) = self.previous_probe_persistent(table, key, prev_version)? {
1426			return Ok(found);
1427		}
1428
1429		Ok(None)
1430	}
1431}
1432
1433impl StandardMultiStore {
1434	#[inline]
1435	fn previous_probe_commit(
1436		&self,
1437		table: EntryKind,
1438		key: &EncodedKey,
1439		prev_version: CommitVersion,
1440	) -> Result<Option<Option<MultiVersionRow>>> {
1441		let Some(commit) = &self.commit else {
1442			return Ok(None);
1443		};
1444		Ok(match commit.get(table, key.as_ref(), prev_version)? {
1445			VersionedGetResult::Value {
1446				value,
1447				version,
1448			} => Some(Some(MultiVersionRow {
1449				key: key.clone(),
1450				row: EncodedRow(CowVec::new(value.to_vec())),
1451				version,
1452			})),
1453			VersionedGetResult::Tombstone => Some(None),
1454			VersionedGetResult::NotFound => None,
1455		})
1456	}
1457
1458	#[inline]
1459	fn previous_probe_read(
1460		&self,
1461		key: &EncodedKey,
1462		prev_version: CommitVersion,
1463	) -> Option<Option<MultiVersionRow>> {
1464		let read = self.read.as_ref()?;
1465		match read.get(key, prev_version) {
1466			VersionedGetResult::Value {
1467				value,
1468				version,
1469			} => Some(Some(MultiVersionRow {
1470				key: key.clone(),
1471				row: EncodedRow(CowVec::new(value.to_vec())),
1472				version,
1473			})),
1474			VersionedGetResult::Tombstone => Some(None),
1475			VersionedGetResult::NotFound => None,
1476		}
1477	}
1478
1479	#[inline]
1480	fn previous_probe_persistent(
1481		&self,
1482		table: EntryKind,
1483		key: &EncodedKey,
1484		prev_version: CommitVersion,
1485	) -> Result<Option<Option<MultiVersionRow>>> {
1486		let Some(persistent) = &self.persistent else {
1487			return Ok(None);
1488		};
1489		Ok(match persistent.get(table, key.as_ref(), prev_version)? {
1490			VersionedGetResult::Value {
1491				value,
1492				version,
1493			} => {
1494				if let Some(read) = &self.read {
1495					read.insert(key.clone(), version, Some(value.clone()));
1496				}
1497				Some(Some(MultiVersionRow {
1498					key: key.clone(),
1499					row: EncodedRow(CowVec::new(value.to_vec())),
1500					version,
1501				}))
1502			}
1503			VersionedGetResult::Tombstone => Some(None),
1504			VersionedGetResult::NotFound => None,
1505		})
1506	}
1507}
1508
1509impl MultiVersionStore for StandardMultiStore {}
1510
1511pub struct MultiVersionRangeIter {
1512	store: StandardMultiStore,
1513	cursor: MultiVersionRangeCursor,
1514	range: EncodedKeyRange,
1515	scope: MultiVersionScope,
1516	batch_size: usize,
1517	current_batch: Vec<MultiVersionRow>,
1518	current_index: usize,
1519}
1520
1521impl Iterator for MultiVersionRangeIter {
1522	type Item = Result<MultiVersionRow>;
1523
1524	fn next(&mut self) -> Option<Self::Item> {
1525		if self.current_index < self.current_batch.len() {
1526			let item = self.current_batch[self.current_index].clone();
1527			self.current_index += 1;
1528			return Some(Ok(item));
1529		}
1530
1531		if self.cursor.exhausted {
1532			return None;
1533		}
1534
1535		match self.store.range_next(&mut self.cursor, self.range.clone(), self.scope, self.batch_size as u64) {
1536			Ok(batch) => {
1537				if batch.items.is_empty() {
1538					if self.cursor.exhausted {
1539						return None;
1540					}
1541					return self.next();
1542				}
1543				self.current_batch = batch.items;
1544				self.current_index = 0;
1545				self.next()
1546			}
1547			Err(e) => Some(Err(e)),
1548		}
1549	}
1550}
1551
1552pub struct MultiVersionRangeRevIter {
1553	store: StandardMultiStore,
1554	cursor: MultiVersionRangeCursor,
1555	range: EncodedKeyRange,
1556	scope: MultiVersionScope,
1557	batch_size: usize,
1558	current_batch: Vec<MultiVersionRow>,
1559	current_index: usize,
1560}
1561
1562impl Iterator for MultiVersionRangeRevIter {
1563	type Item = Result<MultiVersionRow>;
1564
1565	fn next(&mut self) -> Option<Self::Item> {
1566		if self.current_index < self.current_batch.len() {
1567			let item = self.current_batch[self.current_index].clone();
1568			self.current_index += 1;
1569			return Some(Ok(item));
1570		}
1571
1572		if self.cursor.exhausted {
1573			return None;
1574		}
1575
1576		match self.store.range_rev_next(
1577			&mut self.cursor,
1578			self.range.clone(),
1579			self.scope,
1580			self.batch_size as u64,
1581		) {
1582			Ok(batch) => {
1583				if batch.items.is_empty() {
1584					if self.cursor.exhausted {
1585						return None;
1586					}
1587					return self.next();
1588				}
1589				self.current_batch = batch.items;
1590				self.current_index = 0;
1591				self.next()
1592			}
1593			Err(e) => Some(Err(e)),
1594		}
1595	}
1596}
1597
1598fn classify_key_range(range: &EncodedKeyRange) -> EntryKind {
1599	classify_range(range).unwrap_or(EntryKind::Multi)
1600}
1601
1602fn make_range_bounds(range: &EncodedKeyRange) -> (Vec<u8>, Vec<u8>) {
1603	let start = match &range.start {
1604		Bound::Included(key) => key.as_ref().to_vec(),
1605		Bound::Excluded(key) => key.as_ref().to_vec(),
1606		Bound::Unbounded => vec![],
1607	};
1608
1609	let end = match &range.end {
1610		Bound::Included(key) => key.as_ref().to_vec(),
1611		Bound::Excluded(key) => key.as_ref().to_vec(),
1612		Bound::Unbounded => vec![0xFFu8; 256],
1613	};
1614
1615	(start, end)
1616}
1617
1618#[cfg(all(test, feature = "sqlite", not(target_arch = "wasm32")))]
1619mod cache_tests {
1620	use std::collections::HashMap;
1621
1622	use reifydb_codec::{encoded::row::EncodedRow, key::encoded::EncodedKey};
1623	use reifydb_core::{
1624		common::CommitVersion,
1625		delta::Delta,
1626		interface::{
1627			catalog::{flow::FlowNodeId, id::TableId, shape::ShapeId},
1628			store::{EntryKind, MultiVersionCommit},
1629		},
1630		key::{
1631			EncodableKey, flow_node_internal_state::FlowNodeInternalStateKey,
1632			flow_node_state::FlowNodeStateKey, row::RowKey,
1633		},
1634	};
1635	use reifydb_value::{cow_vec, util::cowvec::CowVec};
1636
1637	use crate::{
1638		MultiVersionScope,
1639		store::{StandardMultiStore, multi::WARM_THRESHOLD},
1640		tier::{RawEntry, TierStorage, VersionedGetResult, commit::buffer::MultiCommitBufferTier},
1641	};
1642
1643	const SHAPE: ShapeId = ShapeId::Table(TableId(1));
1644
1645	fn commit_row(store: &StandardMultiStore, n: u64, version: u64) {
1646		MultiVersionCommit::commit(
1647			store,
1648			cow_vec![Delta::Set {
1649				key: RowKey::encoded(SHAPE, n),
1650				row: EncodedRow(CowVec::new(format!("v{n}").into_bytes())),
1651			}],
1652			CommitVersion(version),
1653		)
1654		.unwrap();
1655	}
1656
1657	fn flush(store: &StandardMultiStore, cutoff: CommitVersion) {
1658		let commit = store.commit().expect("commit tier");
1659		for kind in commit.list_all_entry_kinds().unwrap() {
1660			let (to_persist, to_drop) = match commit {
1661				MultiCommitBufferTier::Memory(s) => s.collect_evictable_below(kind, cutoff),
1662			};
1663			if to_drop.is_empty() {
1664				continue;
1665			}
1666			if !to_persist.is_empty() {
1667				let persistent = store.persistent().expect("persistent tier");
1668				let mut by_version: HashMap<
1669					CommitVersion,
1670					HashMap<EntryKind, Vec<(EncodedKey, Option<CowVec<u8>>)>>,
1671				> = HashMap::new();
1672				for (key, version, value) in to_persist {
1673					by_version
1674						.entry(version)
1675						.or_default()
1676						.entry(kind)
1677						.or_default()
1678						.push((key, value));
1679				}
1680				for (version, batch) in by_version {
1681					persistent.set(version, batch).unwrap();
1682				}
1683			}
1684			for (key, _) in &to_drop {
1685				store.invalidate_read_key(key);
1686			}
1687			commit.drop(HashMap::from([(kind, to_drop)])).unwrap();
1688		}
1689	}
1690
1691	#[test]
1692	fn operator_drop_fully_removes_state_leaving_no_tombstone() {
1693		let (store, _g) = StandardMultiStore::testing_memory_with_persistent_sqlite();
1694		let node = FlowNodeId(7);
1695		let table = EntryKind::Operator(node);
1696		let internal_table = EntryKind::OperatorInternal(node);
1697		let data_key = FlowNodeStateKey::encoded(node, vec![1u8]);
1698		let internal_key = FlowNodeInternalStateKey::encoded(node, vec![2u8]);
1699
1700		for v in [1u64, 2] {
1701			MultiVersionCommit::commit(
1702				&store,
1703				cow_vec![Delta::Set {
1704					key: data_key.clone(),
1705					row: EncodedRow(CowVec::new(vec![v as u8])),
1706				}],
1707				CommitVersion(v),
1708			)
1709			.unwrap();
1710		}
1711		for v in [3u64, 4] {
1712			MultiVersionCommit::commit(
1713				&store,
1714				cow_vec![Delta::Set {
1715					key: internal_key.clone(),
1716					row: EncodedRow(CowVec::new(vec![v as u8])),
1717				}],
1718				CommitVersion(v),
1719			)
1720			.unwrap();
1721		}
1722
1723		let commit = store.commit().expect("commit tier");
1724		assert!(!commit.get_all_versions(table, data_key.as_ref()).unwrap().is_empty());
1725		assert!(!commit.get_all_versions(internal_table, internal_key.as_ref()).unwrap().is_empty());
1726
1727		MultiVersionCommit::commit(
1728			&store,
1729			cow_vec![
1730				Delta::Drop {
1731					key: data_key.clone(),
1732				},
1733				Delta::Drop {
1734					key: internal_key.clone(),
1735				}
1736			],
1737			CommitVersion(5),
1738		)
1739		.unwrap();
1740
1741		assert!(
1742			commit.get_all_versions(table, data_key.as_ref()).unwrap().is_empty(),
1743			"operator data-state Drop must remove every version, not leave a tombstone"
1744		);
1745		assert!(
1746			commit.get_all_versions(internal_table, internal_key.as_ref()).unwrap().is_empty(),
1747			"operator internal-state Drop must remove every version, not leave a tombstone"
1748		);
1749	}
1750
1751	#[test]
1752	fn operator_remove_leaves_a_tombstone_in_commit_tier() {
1753		let (store, _g) = StandardMultiStore::testing_memory_with_persistent_sqlite();
1754		let node = FlowNodeId(8);
1755		let table = EntryKind::OperatorInternal(node);
1756		let key = FlowNodeInternalStateKey::encoded(node, vec![9u8]);
1757
1758		MultiVersionCommit::commit(
1759			&store,
1760			cow_vec![Delta::Set {
1761				key: key.clone(),
1762				row: EncodedRow(CowVec::new(vec![1u8])),
1763			}],
1764			CommitVersion(1),
1765		)
1766		.unwrap();
1767		MultiVersionCommit::commit(
1768			&store,
1769			cow_vec![Delta::Remove {
1770				key: key.clone(),
1771			}],
1772			CommitVersion(2),
1773		)
1774		.unwrap();
1775
1776		let commit = store.commit().expect("commit tier");
1777		let versions = commit.get_all_versions(table, key.as_ref()).unwrap();
1778		assert!(
1779			versions.iter().any(|(_, value)| value.is_none()),
1780			"Remove leaves a tombstone in the commit tier (the path Drop must avoid); versions={versions:?}"
1781		);
1782	}
1783
1784	#[test]
1785	fn operator_state_drop_keeps_keyspace_bounded_under_churn() {
1786		const ROUNDS: u64 = 200;
1787
1788		fn current_count(store: &StandardMultiStore, table: EntryKind) -> u64 {
1789			match store.commit().expect("commit tier") {
1790				MultiCommitBufferTier::Memory(s) => s.count_current(table).unwrap(),
1791			}
1792		}
1793
1794		fn churn(evict_with_drop: bool) -> u64 {
1795			let (store, _g) = StandardMultiStore::testing_memory_with_persistent_sqlite();
1796			let node = FlowNodeId(21);
1797			let table = EntryKind::OperatorInternal(node);
1798			let key_at = |round: u64| FlowNodeInternalStateKey::encoded(node, round.to_be_bytes().to_vec());
1799
1800			let mut version = 0u64;
1801			for round in 0..ROUNDS {
1802				version += 1;
1803				MultiVersionCommit::commit(
1804					&store,
1805					cow_vec![Delta::Set {
1806						key: key_at(round),
1807						row: EncodedRow(CowVec::new(vec![1u8])),
1808					}],
1809					CommitVersion(version),
1810				)
1811				.unwrap();
1812
1813				if round > 0 {
1814					version += 1;
1815					let prev = key_at(round - 1);
1816					let delta = if evict_with_drop {
1817						Delta::Drop {
1818							key: prev,
1819						}
1820					} else {
1821						Delta::Remove {
1822							key: prev,
1823						}
1824					};
1825					MultiVersionCommit::commit(&store, cow_vec![delta], CommitVersion(version))
1826						.unwrap();
1827				}
1828			}
1829			current_count(&store, table)
1830		}
1831
1832		let drop_live = churn(true);
1833		let remove_live = churn(false);
1834
1835		assert!(
1836			drop_live <= 2,
1837			"Drop must keep the operator keyspace bounded to the live set; got {drop_live}"
1838		);
1839		assert!(
1840			remove_live >= ROUNDS - 1,
1841			"Remove leaves a tombstone per round (the path Drop avoids); got {remove_live} after {ROUNDS} rounds"
1842		);
1843	}
1844
1845	#[test]
1846	fn warm_threshold_warms_only_buckets_above_threshold() {
1847		const HEAVY: u64 = WARM_THRESHOLD + 64;
1848		const LIGHT: u64 = 20;
1849		let (store, _g) = StandardMultiStore::testing_memory_with_persistent_sqlite();
1850
1851		for n in 1..=HEAVY {
1852			commit_row(&store, n, 1);
1853		}
1854		for n in 0..LIGHT {
1855			commit_row(&store, (1u64 << 16) + n, 1);
1856		}
1857		flush(&store, CommitVersion(1));
1858
1859		let read = store.read.clone().expect("read tier configured");
1860		let heavy_bucket = read.page_of_key(&RowKey::encoded(SHAPE, 1));
1861		let light_bucket = read.page_of_key(&RowKey::encoded(SHAPE, 1u64 << 16));
1862		assert_ne!(heavy_bucket, light_bucket, "the two row groups must land in different buckets");
1863		assert!(!read.page_is_complete(heavy_bucket), "nothing is warm before the scan");
1864
1865		let scanned = store
1866			.range(
1867				RowKey::full_scan(SHAPE),
1868				MultiVersionScope::AsOf {
1869					read: CommitVersion(10),
1870				},
1871				32,
1872			)
1873			.collect::<Result<Vec<_>, _>>()
1874			.unwrap();
1875		assert_eq!(scanned.len() as u64, HEAVY + LIGHT, "the scan returns every row regardless of warming");
1876
1877		assert!(read.page_is_complete(heavy_bucket), "a bucket scanned past the threshold must be warmed");
1878		assert!(
1879			!read.page_is_complete(light_bucket),
1880			"a bucket scanned below the threshold must not be warmed"
1881		);
1882	}
1883
1884	#[test]
1885	fn operator_state_write_through_keeps_read_cache_warm() {
1886		let (store, _g) = StandardMultiStore::testing_memory_with_persistent_sqlite();
1887		let read = store.read.clone().expect("read tier configured");
1888
1889		let opkey = FlowNodeStateKey::new(FlowNodeId(7), vec![1, 2, 3]).encode();
1890		MultiVersionCommit::commit(
1891			&store,
1892			cow_vec![Delta::Set {
1893				key: opkey.clone(),
1894				row: EncodedRow(CowVec::new(b"state-v10".to_vec())),
1895			}],
1896			CommitVersion(10),
1897		)
1898		.unwrap();
1899
1900		match read.get(&opkey, CommitVersion(10)) {
1901			VersionedGetResult::Value {
1902				value,
1903				version,
1904			} => {
1905				assert_eq!(
1906					value.as_ref(),
1907					b"state-v10",
1908					"the cached operator state must be the committed value"
1909				);
1910				assert_eq!(
1911					version,
1912					CommitVersion(10),
1913					"the cached entry must carry the commit version"
1914				);
1915			}
1916			other => {
1917				panic!("operator state must be served from the read cache after commit, got {other:?}")
1918			}
1919		}
1920
1921		assert!(
1922			matches!(read.get(&opkey, CommitVersion(9)), VersionedGetResult::NotFound),
1923			"a pre-write snapshot read must miss the write-through entry, not see the newer value"
1924		);
1925	}
1926
1927	#[test]
1928	fn source_row_write_clears_range_complete_on_its_page() {
1929		let (store, _g) = StandardMultiStore::testing_memory_with_persistent_sqlite();
1930		let read = store.read.clone().expect("read tier configured");
1931
1932		let neighbor = RowKey::encoded(SHAPE, 1);
1933		let page = read.page_of_key(&neighbor);
1934		assert_eq!(
1935			read.page_of_key(&RowKey::encoded(SHAPE, 2)),
1936			page,
1937			"both source rows must share a page for this test to exercise flag-clearing"
1938		);
1939		read.populate_page(
1940			page,
1941			vec![RawEntry {
1942				key: neighbor,
1943				version: CommitVersion(1),
1944				value: Some(CowVec::new(b"neighbor".to_vec())),
1945			}],
1946			true,
1947		);
1948		assert!(read.page_is_complete(page), "the page must start range-complete");
1949
1950		commit_row(&store, 2, 5);
1951
1952		assert!(
1953			!read.page_is_complete(page),
1954			"writing a source row into a range-complete page must clear the flag so the range cache re-warms"
1955		);
1956	}
1957}