1use 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}