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