1use std::collections::hash_map::HashMap;
16use std::collections::{HashSet, VecDeque};
17use std::future::poll_fn;
18use std::hash::Hash;
19use std::ops::{Bound, Range};
20use std::sync::{Arc, LazyLock};
21use std::task::Poll;
22use std::time::{Duration, Instant};
23
24use foyer::RangeBoundsExt;
25use futures::future::{join_all, try_join_all};
26use futures::{Future, FutureExt};
27use itertools::Itertools;
28use prometheus::core::{AtomicU64, GenericCounter, GenericCounterVec};
29use prometheus::{
30 Histogram, HistogramVec, IntGauge, Registry, register_histogram_vec_with_registry,
31 register_int_counter_vec_with_registry, register_int_gauge_with_registry,
32};
33use risingwave_common::bitmap::Bitmap;
34use risingwave_common::config::Role;
35use risingwave_common::config::streaming::CacheRefillPolicy;
36use risingwave_common::hash::VirtualNode;
37use risingwave_common::license::Feature;
38use risingwave_common::monitor::GLOBAL_METRICS_REGISTRY;
39use risingwave_common::util::iter_util::ZipEqFast;
40use risingwave_hummock_sdk::compaction_group::hummock_version_ext::SstDeltaInfo;
41use risingwave_hummock_sdk::key::{FullKey, vnode_range};
42use risingwave_hummock_sdk::{HummockSstableObjectId, KeyComparator};
43use risingwave_pb::id::TableId;
44use thiserror_ext::AsReport;
45use tokio::sync::Semaphore;
46use tokio::task::JoinHandle;
47
48use crate::hummock::local_version::pinned_version::PinnedVersion;
49use crate::hummock::{
50 Block, HummockError, HummockResult, RecentFilterTrait, Sstable, SstableBlockIndex,
51 SstableStoreRef, TableHolder,
52};
53use crate::monitor::StoreLocalStatistic;
54use crate::opts::StorageOpts;
55
56pub static GLOBAL_CACHE_REFILL_METRICS: LazyLock<CacheRefillMetrics> =
57 LazyLock::new(|| CacheRefillMetrics::new(&GLOBAL_METRICS_REGISTRY));
58
59pub struct CacheRefillMetrics {
60 pub refill_duration: HistogramVec,
61 pub refill_total: GenericCounterVec<AtomicU64>,
62 pub refill_bytes: GenericCounterVec<AtomicU64>,
63
64 pub data_refill_success_duration: Histogram,
65 pub meta_refill_success_duration: Histogram,
66
67 pub data_refill_filtered_total: GenericCounter<AtomicU64>,
68 pub data_refill_attempts_total: GenericCounter<AtomicU64>,
69 pub data_refill_started_total: GenericCounter<AtomicU64>,
70 pub meta_refill_attempts_total: GenericCounter<AtomicU64>,
71
72 pub data_refill_parent_meta_lookup_hit_total: GenericCounter<AtomicU64>,
73 pub data_refill_parent_meta_lookup_miss_total: GenericCounter<AtomicU64>,
74 pub data_refill_unit_inheritance_hit_total: GenericCounter<AtomicU64>,
75 pub data_refill_unit_inheritance_miss_total: GenericCounter<AtomicU64>,
76
77 pub data_refill_block_unfiltered_total: GenericCounter<AtomicU64>,
78 pub data_refill_block_success_total: GenericCounter<AtomicU64>,
79
80 pub data_refill_ideal_bytes: GenericCounter<AtomicU64>,
81 pub data_refill_success_bytes: GenericCounter<AtomicU64>,
82
83 pub refill_queue_total: IntGauge,
84}
85
86impl CacheRefillMetrics {
87 pub fn new(registry: &Registry) -> Self {
88 let refill_duration = register_histogram_vec_with_registry!(
89 "refill_duration",
90 "refill duration",
91 &["type", "op"],
92 registry,
93 )
94 .unwrap();
95 let refill_total = register_int_counter_vec_with_registry!(
96 "refill_total",
97 "refill total",
98 &["type", "op"],
99 registry,
100 )
101 .unwrap();
102 let refill_bytes = register_int_counter_vec_with_registry!(
103 "refill_bytes",
104 "refill bytes",
105 &["type", "op"],
106 registry,
107 )
108 .unwrap();
109
110 let data_refill_success_duration = refill_duration
111 .get_metric_with_label_values(&["data", "success"])
112 .unwrap();
113 let meta_refill_success_duration = refill_duration
114 .get_metric_with_label_values(&["meta", "success"])
115 .unwrap();
116
117 let data_refill_filtered_total = refill_total
118 .get_metric_with_label_values(&["data", "filtered"])
119 .unwrap();
120 let data_refill_attempts_total = refill_total
121 .get_metric_with_label_values(&["data", "attempts"])
122 .unwrap();
123 let data_refill_started_total = refill_total
124 .get_metric_with_label_values(&["data", "started"])
125 .unwrap();
126 let meta_refill_attempts_total = refill_total
127 .get_metric_with_label_values(&["meta", "attempts"])
128 .unwrap();
129
130 let data_refill_parent_meta_lookup_hit_total = refill_total
131 .get_metric_with_label_values(&["parent_meta", "hit"])
132 .unwrap();
133 let data_refill_parent_meta_lookup_miss_total = refill_total
134 .get_metric_with_label_values(&["parent_meta", "miss"])
135 .unwrap();
136 let data_refill_unit_inheritance_hit_total = refill_total
137 .get_metric_with_label_values(&["unit_inheritance", "hit"])
138 .unwrap();
139 let data_refill_unit_inheritance_miss_total = refill_total
140 .get_metric_with_label_values(&["unit_inheritance", "miss"])
141 .unwrap();
142
143 let data_refill_block_unfiltered_total = refill_total
144 .get_metric_with_label_values(&["block", "unfiltered"])
145 .unwrap();
146 let data_refill_block_success_total = refill_total
147 .get_metric_with_label_values(&["block", "success"])
148 .unwrap();
149
150 let data_refill_ideal_bytes = refill_bytes
151 .get_metric_with_label_values(&["data", "ideal"])
152 .unwrap();
153 let data_refill_success_bytes = refill_bytes
154 .get_metric_with_label_values(&["data", "success"])
155 .unwrap();
156
157 let refill_queue_total = register_int_gauge_with_registry!(
158 "refill_queue_total",
159 "refill queue total",
160 registry,
161 )
162 .unwrap();
163
164 Self {
165 refill_duration,
166 refill_total,
167 refill_bytes,
168
169 data_refill_success_duration,
170 meta_refill_success_duration,
171 data_refill_filtered_total,
172 data_refill_attempts_total,
173 data_refill_started_total,
174 meta_refill_attempts_total,
175
176 data_refill_parent_meta_lookup_hit_total,
177 data_refill_parent_meta_lookup_miss_total,
178 data_refill_unit_inheritance_hit_total,
179 data_refill_unit_inheritance_miss_total,
180
181 data_refill_block_unfiltered_total,
182 data_refill_block_success_total,
183
184 data_refill_ideal_bytes,
185 data_refill_success_bytes,
186
187 refill_queue_total,
188 }
189 }
190}
191
192#[derive(Debug)]
193pub struct CacheRefillConfig {
194 pub timeout: Duration,
196
197 pub data_refill_levels: HashSet<u32>,
199
200 pub meta_refill_concurrency: usize,
202
203 pub concurrency: usize,
205
206 pub unit: usize,
208
209 pub threshold: f64,
213
214 pub skip_recent_filter: bool,
216
217 pub skip_inheritance_filter: bool,
219
220 pub table_cache_refill_default_policy: CacheRefillPolicy,
222}
223
224impl CacheRefillConfig {
225 pub fn from_storage_opts(options: &StorageOpts) -> Self {
226 let data_refill_levels = match Feature::ElasticDiskCache.check_available() {
227 Ok(_) => options
228 .cache_refill_data_refill_levels
229 .iter()
230 .copied()
231 .collect(),
232 Err(e) => {
233 tracing::warn!(error = %e.as_report(), "ElasticDiskCache is not available.");
234 HashSet::new()
235 }
236 };
237
238 Self {
239 timeout: Duration::from_millis(options.cache_refill_timeout_ms),
240 data_refill_levels,
241 concurrency: options.cache_refill_concurrency,
242 meta_refill_concurrency: options.cache_refill_meta_refill_concurrency,
243 unit: options.cache_refill_unit,
244 threshold: options.cache_refill_threshold,
245 skip_recent_filter: options.cache_refill_skip_recent_filter,
246 skip_inheritance_filter: options.cache_refill_skip_inheritance_filter,
247 table_cache_refill_default_policy: options
248 .cache_refill_table_cache_refill_default_policy,
249 }
250 }
251}
252
253struct Item {
254 handle: JoinHandle<()>,
255 event: CacheRefillerEvent,
256}
257
258pub(crate) type SpawnRefillTask = Arc<
259 dyn Fn(Vec<SstDeltaInfo>, CacheRefillContext, PinnedVersion, PinnedVersion) -> JoinHandle<()>
261 + Send
262 + Sync
263 + 'static,
264>;
265
266pub type TableCacheRefillContextMap = HashMap<TableId, TableCacheRefillContext>;
267
268#[derive(Clone)]
271pub struct TableCacheRefillContext {
272 pub streaming_vnode_bitmap: Option<Bitmap>,
274 pub serving_vnode_bitmap: Option<Bitmap>,
276 pub policy: CacheRefillPolicy,
278}
279
280#[derive(Clone)]
284pub struct TableCacheRefillMonitorSnapshot {
285 pub contexts: TableCacheRefillContextMap,
286 pub policies: HashMap<TableId, CacheRefillPolicy>,
287 pub default_policy: CacheRefillPolicy,
288 pub streaming_table_vnode_mapping: HashMap<TableId, Bitmap>,
289 pub serving_table_vnode_mapping: HashMap<TableId, Bitmap>,
290}
291
292fn vnode_range_overlaps_bitmap(vnode_range: (usize, usize), bitmap: &Bitmap) -> bool {
293 assert!(vnode_range.0 <= vnode_range.1);
294 let start = vnode_range.0.min(bitmap.len());
295 let end = vnode_range.1.min(bitmap.len());
296 if start == end || !bitmap.any() {
297 return false;
298 }
299 if bitmap.all() {
300 return true;
301 }
302 (start..end).any(|vnode| bitmap.is_set(vnode))
303}
304
305impl TableCacheRefillContext {
306 fn allows_normal_data_refill_block(&self, sstable: &Sstable, block_index: usize) -> bool {
307 if self.policy.is_unscoped_enabled() {
308 return true;
309 }
310
311 (self.policy.is_streaming_scoped()
312 && self.check_table_refill_streaming_vnodes(sstable, block_index))
313 || (self.policy.is_serving_scoped()
314 && self.check_table_refill_serving_vnodes(sstable, block_index))
315 }
316
317 fn allows_insert_only_data_refill_block(&self, sstable: &Sstable, block_index: usize) -> bool {
318 (self.policy.is_unscoped_enabled() || self.policy.is_serving_scoped())
322 && self.check_table_refill_serving_vnodes(sstable, block_index)
323 }
324
325 fn check_table_refill_streaming_vnodes(&self, sstable: &Sstable, block_index: usize) -> bool {
326 self.streaming_vnode_bitmap.as_ref().is_some_and(|bitmap| {
327 let vnode_range = block_vnode_range(sstable, block_index);
328 vnode_range_overlaps_bitmap(vnode_range, bitmap)
329 })
330 }
331
332 fn check_table_refill_serving_vnodes(&self, sstable: &Sstable, block_index: usize) -> bool {
333 self.serving_vnode_bitmap.as_ref().is_some_and(|bitmap| {
334 let vnode_range = block_vnode_range(sstable, block_index);
335 vnode_range_overlaps_bitmap(vnode_range, bitmap)
336 })
337 }
338}
339
340fn block_vnode_range(sstable: &Sstable, block_index: usize) -> (usize, usize) {
341 let block_meta = &sstable.meta.block_metas[block_index];
342 let block_smallest_key = FullKey::decode(&block_meta.smallest_key);
343 let table_key_end = match sstable.meta.block_metas.get(block_index + 1) {
344 Some(next_block_meta) if next_block_meta.table_id() != block_meta.table_id() => {
347 Bound::Unbounded
348 }
349 Some(next_block_meta) => Bound::Included(
352 FullKey::decode(&next_block_meta.smallest_key)
353 .user_key
354 .table_key,
355 ),
356 None => Bound::Included(
360 FullKey::decode(&sstable.meta.largest_key)
361 .user_key
362 .table_key,
363 ),
364 };
365
366 let table_key_range = (
367 Bound::Included(block_smallest_key.user_key.table_key),
368 table_key_end,
369 );
370 if match &table_key_range.0 {
374 Bound::Included(key) | Bound::Excluded(key) => key.as_ref().len() < VirtualNode::SIZE,
375 Bound::Unbounded => false,
376 } || match &table_key_range.1 {
377 Bound::Included(key) | Bound::Excluded(key) => key.as_ref().len() < VirtualNode::SIZE,
378 Bound::Unbounded => false,
379 } {
380 return (0, VirtualNode::MAX_REPRESENTABLE.to_index() + 1);
381 }
382 vnode_range(&table_key_range)
383}
384
385pub(crate) struct CacheRefiller {
387 queue: VecDeque<Item>,
389
390 spawn_refill_task: SpawnRefillTask,
391
392 config: Arc<CacheRefillConfig>,
393 meta_refill_concurrency: Option<Arc<Semaphore>>,
394 concurrency: Arc<Semaphore>,
395 sstable_store: SstableStoreRef,
396
397 role: Role,
398 default_policy: CacheRefillPolicy,
399 table_cache_refill_policies: HashMap<TableId, CacheRefillPolicy>,
400 streaming_table_vnode_mapping: HashMap<TableId, Bitmap>,
401 serving_table_vnode_mapping: HashMap<TableId, Bitmap>,
402}
403
404impl CacheRefiller {
405 pub(crate) fn new(
406 role: Role,
407 config: CacheRefillConfig,
408 sstable_store: SstableStoreRef,
409 spawn_refill_task: SpawnRefillTask,
410 ) -> Self {
411 let config = Arc::new(config);
412 let concurrency = Arc::new(Semaphore::new(config.concurrency));
413 let default_policy = config.table_cache_refill_default_policy;
414 let meta_refill_concurrency = if config.meta_refill_concurrency == 0 {
415 None
416 } else {
417 Some(Arc::new(Semaphore::new(config.meta_refill_concurrency)))
418 };
419 Self {
420 queue: VecDeque::new(),
421 spawn_refill_task,
422 config,
423 meta_refill_concurrency,
424 concurrency,
425 sstable_store,
426 role,
427 default_policy,
428 table_cache_refill_policies: HashMap::new(),
429 streaming_table_vnode_mapping: HashMap::new(),
430 serving_table_vnode_mapping: HashMap::new(),
431 }
432 }
433
434 pub(crate) fn default_spawn_refill_task() -> SpawnRefillTask {
435 Arc::new(|deltas, context, _, _| {
436 let task = CacheRefillTask { deltas, context };
437 tokio::spawn(task.run())
438 })
439 }
440
441 pub(crate) fn start_cache_refill(
442 &mut self,
443 mut deltas: Vec<SstDeltaInfo>,
444 pinned_version: PinnedVersion,
445 new_pinned_version: PinnedVersion,
446 ) {
447 for delta in &mut deltas {
448 let for_serving = self.role.for_serving();
449 let for_streaming =
452 self.role.for_streaming() && !delta.delete_sst_object_ids.is_empty();
453
454 if !for_serving && !for_streaming {
455 delta.insert_sst_infos.clear();
456 continue;
457 }
458
459 delta.insert_sst_infos.retain(|sst| {
464 sst.table_ids.iter().any(|table_id| {
465 let policy = self
468 .table_cache_refill_policies
469 .get(table_id)
470 .copied()
471 .unwrap_or(self.default_policy);
472
473 match policy {
476 CacheRefillPolicy::Enabled => for_streaming || for_serving,
477 CacheRefillPolicy::Disabled => false,
478 CacheRefillPolicy::Streaming => {
479 for_streaming
480 && self.streaming_table_vnode_mapping.contains_key(table_id)
481 }
482 CacheRefillPolicy::Serving => {
483 for_serving && self.serving_table_vnode_mapping.contains_key(table_id)
484 }
485 CacheRefillPolicy::Both => {
486 (for_streaming
487 && self.streaming_table_vnode_mapping.contains_key(table_id))
488 || (for_serving
489 && self.serving_table_vnode_mapping.contains_key(table_id))
490 }
491 }
492 })
493 });
494 }
495 let context = self.new_cache_refill_context(&deltas);
496 let handle = (self.spawn_refill_task)(
497 deltas,
498 context,
499 pinned_version.clone(),
500 new_pinned_version.clone(),
501 );
502 let event = CacheRefillerEvent {
503 pinned_version,
504 new_pinned_version,
505 };
506 let item = Item { handle, event };
507 self.queue.push_back(item);
508 GLOBAL_CACHE_REFILL_METRICS.refill_queue_total.add(1);
509 }
510
511 fn new_cache_refill_context(&self, deltas: &[SstDeltaInfo]) -> CacheRefillContext {
512 let table_ids = deltas.iter().flat_map(|delta| {
513 delta
514 .insert_sst_infos
515 .iter()
516 .flat_map(|sst| sst.table_ids.iter().copied())
517 });
518 CacheRefillContext {
519 config: self.config.clone(),
520 meta_refill_concurrency: self.meta_refill_concurrency.clone(),
521 concurrency: self.concurrency.clone(),
522 sstable_store: self.sstable_store.clone(),
523 table_cache_refill_context_map: Arc::new(self.table_cache_refill_contexts(table_ids)),
524 }
525 }
526
527 pub(crate) fn last_new_pinned_version(&self) -> Option<&PinnedVersion> {
528 self.queue.back().map(|item| &item.event.new_pinned_version)
529 }
530
531 pub(crate) fn replace_table_cache_refill_policies(
533 &mut self,
534 policies: HashMap<TableId, CacheRefillPolicy>,
535 ) {
536 self.table_cache_refill_policies = policies;
537 }
538
539 pub(crate) fn replace_serving_table_vnode_mapping(
541 &mut self,
542 mapping: HashMap<TableId, Bitmap>,
543 ) {
544 self.serving_table_vnode_mapping = mapping;
545 }
546
547 pub(crate) fn update_streaming_table_vnodes(
548 &mut self,
549 table_id: TableId,
550 streaming_vnodes: Option<Bitmap>,
551 ) {
552 if let Some(streaming_vnodes) = streaming_vnodes {
553 self.streaming_table_vnode_mapping
554 .insert(table_id, streaming_vnodes);
555 } else {
556 self.streaming_table_vnode_mapping.remove(&table_id);
557 }
558 }
559
560 fn table_cache_refill_contexts(
561 &self,
562 table_ids: impl IntoIterator<Item = TableId>,
563 ) -> TableCacheRefillContextMap {
564 let for_streaming = self.role.for_streaming();
565 let for_serving = self.role.for_serving();
566 table_ids
567 .into_iter()
568 .filter_map(|table_id| {
569 if for_serving
570 && !for_streaming
571 && !self.serving_table_vnode_mapping.contains_key(&table_id)
572 {
573 return None;
574 }
575 let policy = self
576 .table_cache_refill_policies
577 .get(&table_id)
578 .copied()
579 .unwrap_or(self.default_policy);
580 let streaming_vnode_bitmap = (for_streaming && policy.is_streaming_scoped())
581 .then(|| self.streaming_table_vnode_mapping.get(&table_id).cloned())
582 .flatten();
583 let serving_vnode_bitmap = (for_serving
586 && (policy.is_serving_scoped() || policy.is_unscoped_enabled()))
587 .then(|| self.serving_table_vnode_mapping.get(&table_id).cloned())
588 .flatten();
589 Some((
590 table_id,
591 TableCacheRefillContext {
592 streaming_vnode_bitmap,
593 serving_vnode_bitmap,
594 policy,
595 },
596 ))
597 })
598 .collect()
599 }
600
601 pub(crate) fn table_cache_refill_monitor_snapshot(&self) -> TableCacheRefillMonitorSnapshot {
602 let table_ids = self
603 .table_cache_refill_policies
604 .keys()
605 .chain(self.streaming_table_vnode_mapping.keys())
606 .chain(self.serving_table_vnode_mapping.keys())
607 .copied();
608 TableCacheRefillMonitorSnapshot {
609 contexts: self.table_cache_refill_contexts(table_ids),
610 policies: self.table_cache_refill_policies.clone(),
611 default_policy: self.default_policy,
612 streaming_table_vnode_mapping: self.streaming_table_vnode_mapping.clone(),
613 serving_table_vnode_mapping: self.serving_table_vnode_mapping.clone(),
614 }
615 }
616}
617
618impl CacheRefiller {
619 pub(crate) fn next_events(&mut self) -> impl Future<Output = Vec<CacheRefillerEvent>> + '_ {
620 poll_fn(|cx| {
621 const MAX_BATCH_SIZE: usize = 16;
622 let mut events = None;
623 while let Some(item) = self.queue.front_mut()
624 && let Poll::Ready(result) = item.handle.poll_unpin(cx)
625 {
626 result.unwrap();
627 let item = self.queue.pop_front().unwrap();
628 GLOBAL_CACHE_REFILL_METRICS.refill_queue_total.sub(1);
629 let events = events.get_or_insert_with(|| Vec::with_capacity(MAX_BATCH_SIZE));
630 events.push(item.event);
631 if events.len() >= MAX_BATCH_SIZE {
632 break;
633 }
634 }
635 if let Some(events) = events {
636 Poll::Ready(events)
637 } else {
638 Poll::Pending
639 }
640 })
641 }
642}
643
644pub struct CacheRefillerEvent {
645 pub pinned_version: PinnedVersion,
646 pub new_pinned_version: PinnedVersion,
647}
648
649#[derive(Clone)]
650pub(crate) struct CacheRefillContext {
651 config: Arc<CacheRefillConfig>,
652 meta_refill_concurrency: Option<Arc<Semaphore>>,
653 concurrency: Arc<Semaphore>,
654 sstable_store: SstableStoreRef,
655 table_cache_refill_context_map: Arc<TableCacheRefillContextMap>,
656}
657
658struct DataCacheRefillTaskGenerator<'a> {
659 context: &'a CacheRefillContext,
660 delta: &'a SstDeltaInfo,
661 ssts: &'a [TableHolder],
662}
663
664impl DataCacheRefillTaskGenerator<'_> {
665 fn generate_unfiltered_tasks(&self) -> Vec<DataCacheRefillTask> {
666 let mut tasks = Vec::new();
667
668 if !self.context.sstable_store.block_cache().is_hybrid() {
670 return tasks;
671 }
672
673 if self.delta.insert_sst_infos.is_empty() {
674 return tasks;
675 }
676
677 let has_parent_ssts = !self.delta.delete_sst_object_ids.is_empty();
678 debug_assert!(has_parent_ssts || self.delta.insert_sst_level == 0);
681
682 if !self
684 .context
685 .config
686 .data_refill_levels
687 .contains(&self.delta.insert_sst_level)
688 {
689 return tasks;
690 }
691
692 let unit = self.context.config.unit;
695 assert!(unit > 0, "cache refill unit must be positive");
696 let table_cache_refill_context_map = &self.context.table_cache_refill_context_map;
697 for (sst_info, sst) in self.delta.insert_sst_infos.iter().zip_eq_fast(self.ssts) {
698 debug_assert_eq!(sst_info.object_id, sst.id);
699 debug_assert!(sst_info.table_ids.is_sorted());
700 let mut blk_start = 0;
701 while blk_start < sst.block_count() {
702 let table_id = sst.meta.block_metas[blk_start].table_id();
706 let mut blk_end = std::cmp::min(sst.block_count(), blk_start + unit);
707 if let Some(table_boundary) = (blk_start + 1..blk_end)
708 .find(|&block_index| sst.meta.block_metas[block_index].table_id() != table_id)
709 {
710 blk_end = table_boundary;
711 }
712
713 let should_refill = sst_info.table_ids.binary_search(&table_id).is_ok()
714 && (blk_start..blk_end).any(|block_index| {
715 table_cache_refill_context_map
716 .get(&table_id)
717 .is_some_and(|context| {
718 if has_parent_ssts {
719 context.allows_normal_data_refill_block(sst, block_index)
720 } else {
721 context.allows_insert_only_data_refill_block(sst, block_index)
722 }
723 })
724 });
725 if should_refill {
726 tasks.push(DataCacheRefillTask {
727 sst: sst.clone(),
728 blks: blk_start..blk_end,
729 });
730 }
731 blk_start = blk_end;
732 }
733 }
734
735 if tasks.is_empty() {
736 return tasks;
737 }
738
739 if has_parent_ssts
742 && !self.context.config.skip_recent_filter
743 && !self.filter_by_recent_filter()
744 {
745 GLOBAL_CACHE_REFILL_METRICS
746 .data_refill_filtered_total
747 .inc_by(self.delta.delete_sst_object_ids.len() as u64);
748 return vec![];
749 }
750
751 tasks
752 }
753
754 async fn filter_by_inheritance_if_needed(
755 &self,
756 tasks: Vec<DataCacheRefillTask>,
757 ) -> Vec<DataCacheRefillTask> {
758 let should_filter_by_inheritance = !tasks.is_empty()
761 && !self.delta.delete_sst_object_ids.is_empty()
762 && self.delta.insert_sst_level != 0
763 && !self.context.config.skip_recent_filter
764 && !self.context.config.skip_inheritance_filter;
765 if should_filter_by_inheritance {
766 self.filter_by_inheritance_filter(tasks).await
767 } else {
768 tasks
769 }
770 }
771
772 fn filter_by_recent_filter(&self) -> bool {
774 let recent_filter = self.context.sstable_store.recent_filter();
775 let targets = self
776 .delta
777 .delete_sst_object_ids
778 .iter()
779 .map(|id| (*id, usize::MAX))
780 .collect_vec();
781 recent_filter.contains_any(targets.iter())
782 }
783
784 async fn filter_by_inheritance_filter(
785 &self,
786 originals: Vec<DataCacheRefillTask>,
787 ) -> Vec<DataCacheRefillTask> {
788 let sstable_store = self.context.sstable_store.clone();
790 let futures = self.delta.delete_sst_object_ids.iter().map(|sst_obj_id| {
791 let store = &sstable_store;
792 async move {
793 let res = store.sstable_cached(*sst_obj_id).await;
794 match res {
795 Ok(Some(_)) => GLOBAL_CACHE_REFILL_METRICS
796 .data_refill_parent_meta_lookup_hit_total
797 .inc(),
798 Ok(None) => GLOBAL_CACHE_REFILL_METRICS
799 .data_refill_parent_meta_lookup_miss_total
800 .inc(),
801 _ => {}
802 }
803 res
804 }
805 });
806 let parent_ssts = match try_join_all(futures).await {
807 Ok(parent_ssts) => parent_ssts.into_iter().flatten(),
808 Err(e) => {
809 tracing::error!(error = %e.as_report(), "get old meta from cache error");
810 return vec![];
811 }
812 };
813
814 if cfg!(debug_assertions) {
816 originals.iter().tuple_windows().for_each(|(a, b)| {
817 debug_assert_ne!(
818 KeyComparator::compare_encoded_full_key(a.largest_key(), b.smallest_key()),
819 std::cmp::Ordering::Greater
820 )
821 });
822 }
823
824 let mut filtered: HashSet<SstableUnit> = HashSet::default();
825 let recent_filter = self.context.sstable_store.recent_filter();
826 for psst in parent_ssts {
827 for pblk in 0..psst.block_count() {
828 let pleft = &psst.meta.block_metas[pblk].smallest_key;
829 let pright = if pblk + 1 == psst.block_count() {
830 &psst.meta.largest_key
832 } else {
833 &psst.meta.block_metas[pblk + 1].smallest_key
834 };
835
836 let uleft = originals.partition_point(|task| {
838 KeyComparator::compare_encoded_full_key(task.largest_key(), pleft)
839 == std::cmp::Ordering::Less
840 });
841 let uright = originals.partition_point(|task| {
843 KeyComparator::compare_encoded_full_key(task.smallest_key(), pright)
844 != std::cmp::Ordering::Greater
845 });
846
847 for task in originals.iter().take(uright).skip(uleft) {
849 let unit = task.unit();
850 if filtered.contains(&unit) {
851 continue;
852 }
853 if recent_filter.contains(&(psst.id, pblk)) {
854 filtered.insert(unit);
855 }
856 }
857 }
858 }
859
860 let hit = filtered.len();
861 let miss = originals.len() - hit;
862 GLOBAL_CACHE_REFILL_METRICS
863 .data_refill_unit_inheritance_hit_total
864 .inc_by(hit as u64);
865 GLOBAL_CACHE_REFILL_METRICS
866 .data_refill_unit_inheritance_miss_total
867 .inc_by(miss as u64);
868
869 originals
870 .into_iter()
871 .filter(|task| filtered.contains(&task.unit()))
872 .collect()
873 }
874}
875
876#[derive(Debug)]
877struct DataCacheRefillTask {
878 sst: TableHolder,
879 blks: Range<usize>,
880}
881
882impl DataCacheRefillTask {
883 fn unit(&self) -> SstableUnit {
884 SstableUnit {
885 sst_obj_id: self.sst.id,
886 blks: self.blks.clone(),
887 }
888 }
889
890 fn smallest_key(&self) -> &[u8] {
891 &self.sst.meta.block_metas[self.blks.start].smallest_key
892 }
893
894 fn largest_key(&self) -> &[u8] {
895 if self.blks.end == self.sst.block_count() {
896 &self.sst.meta.largest_key
897 } else {
898 &self.sst.meta.block_metas[self.blks.end].smallest_key
899 }
900 }
901}
902
903struct CacheRefillTask {
904 deltas: Vec<SstDeltaInfo>,
905 context: CacheRefillContext,
906}
907
908impl CacheRefillTask {
909 async fn run(self) {
910 let tasks = self
911 .deltas
912 .iter()
913 .map(|delta| {
914 let context = self.context.clone();
915 async move {
916 let holders = match Self::meta_cache_refill(&context, delta).await {
917 Ok(holders) => holders,
918 Err(e) => {
919 tracing::warn!(error = %e.as_report(), "meta cache refill error");
920 return;
921 }
922 };
923 let generator = DataCacheRefillTaskGenerator {
924 context: &context,
925 delta,
926 ssts: &holders,
927 };
928 let tasks = generator.generate_unfiltered_tasks();
929
930 let unfiltered_block_count =
932 tasks.iter().map(|task| task.blks.len() as u64).sum();
933 GLOBAL_CACHE_REFILL_METRICS
934 .data_refill_block_unfiltered_total
935 .inc_by(unfiltered_block_count);
936
937 let tasks = generator.filter_by_inheritance_if_needed(tasks).await;
938 Self::data_cache_refill(&context, tasks).await;
939 }
940 })
941 .collect_vec();
942 let future = join_all(tasks);
943
944 let _ = tokio::time::timeout(self.context.config.timeout, future).await;
945 }
946
947 async fn meta_cache_refill(
948 context: &CacheRefillContext,
949 delta: &SstDeltaInfo,
950 ) -> HummockResult<Vec<TableHolder>> {
951 let tasks = delta
952 .insert_sst_infos
953 .iter()
954 .map(|info| async {
955 let mut stats = StoreLocalStatistic::default();
956 GLOBAL_CACHE_REFILL_METRICS.meta_refill_attempts_total.inc();
957
958 let permit = if let Some(c) = &context.meta_refill_concurrency {
959 Some(c.acquire().await.unwrap())
960 } else {
961 None
962 };
963
964 let now = Instant::now();
965 let res = context.sstable_store.sstable(info, &mut stats).await;
966 stats.discard();
967 if res.is_ok() {
968 GLOBAL_CACHE_REFILL_METRICS
969 .meta_refill_success_duration
970 .observe(now.elapsed().as_secs_f64());
971 }
972 drop(permit);
973
974 res
975 })
976 .collect_vec();
977 let holders = try_join_all(tasks).await?;
978 Ok(holders)
979 }
980
981 async fn data_cache_refill(context: &CacheRefillContext, tasks: Vec<DataCacheRefillTask>) {
982 let mut futures = Vec::with_capacity(tasks.len());
983 for task in tasks {
984 context
986 .sstable_store
987 .recent_filter()
988 .insert((task.sst.id, usize::MAX));
989
990 let blocks = task.blks.len();
991 let mut contexts = Vec::with_capacity(blocks);
992 let mut admits = 0;
993
994 let (range_first, _) = task.sst.calculate_block_info(task.blks.start);
995 let (range_last, _) = task.sst.calculate_block_info(task.blks.end - 1);
996 let range = range_first.start..range_last.end;
997
998 let size = range.size().unwrap();
999
1000 GLOBAL_CACHE_REFILL_METRICS
1001 .data_refill_ideal_bytes
1002 .inc_by(size as _);
1003
1004 for blk in task.blks {
1005 let (range, uncompressed_capacity) = task.sst.calculate_block_info(blk);
1006 let key = SstableBlockIndex {
1007 sst_id: task.sst.id,
1008 block_idx: blk as u64,
1009 };
1010
1011 let mut writer = context.sstable_store.block_cache().storage_writer(key);
1012
1013 if writer.filter(size).is_admitted() {
1014 admits += 1;
1015 }
1016
1017 contexts.push((writer, range, uncompressed_capacity))
1018 }
1019
1020 if admits as f64 / contexts.len() as f64 >= context.config.threshold {
1021 let sstable_store = context.sstable_store.clone();
1022 let context = context.clone();
1023 let future = async move {
1024 GLOBAL_CACHE_REFILL_METRICS.data_refill_attempts_total.inc();
1025
1026 let permit = context.concurrency.acquire().await.unwrap();
1027
1028 GLOBAL_CACHE_REFILL_METRICS.data_refill_started_total.inc();
1029
1030 let now = Instant::now();
1031
1032 let data = sstable_store
1033 .store()
1034 .read(&sstable_store.get_sst_data_path(task.sst.id), range.clone())
1035 .await?;
1036 let mut apply_disk_cache_futures = vec![];
1037 for (w, r, uc) in contexts {
1038 let offset = r.start - range.start;
1039 let len = r.end - r.start;
1040 let bytes = data.slice(offset..offset + len);
1041 let future = async move {
1042 let value = Box::new(Block::decode(bytes, uc)?);
1043 if let Some(_entry) = w.force().insert(value) {
1045 GLOBAL_CACHE_REFILL_METRICS
1046 .data_refill_success_bytes
1047 .inc_by(len as u64);
1048 GLOBAL_CACHE_REFILL_METRICS
1049 .data_refill_block_success_total
1050 .inc();
1051 }
1052 Ok::<_, HummockError>(())
1053 };
1054 apply_disk_cache_futures.push(future);
1055 }
1056 try_join_all(apply_disk_cache_futures)
1057 .await
1058 .map_err(HummockError::file_cache)?;
1059
1060 GLOBAL_CACHE_REFILL_METRICS
1061 .data_refill_success_duration
1062 .observe(now.elapsed().as_secs_f64());
1063 drop(permit);
1064
1065 Ok::<_, HummockError>(())
1066 };
1067 futures.push(future);
1068 }
1069 }
1070
1071 let futures = futures.into_iter().map(|future| async move {
1072 if let Err(e) = future.await {
1073 tracing::error!(error = %e.as_report(), "data cache refill task error");
1074 }
1075 });
1076
1077 join_all(futures).await;
1078 }
1079}
1080
1081#[derive(Debug)]
1082pub struct SstableBlock {
1083 pub sst_obj_id: HummockSstableObjectId,
1084 pub blk_idx: usize,
1085}
1086
1087#[derive(Debug, Hash, PartialEq, Eq)]
1088pub struct SstableUnit {
1089 pub sst_obj_id: HummockSstableObjectId,
1090 pub blks: Range<usize>,
1091}
1092
1093impl Ord for SstableUnit {
1094 fn cmp(&self, other: &Self) -> std::cmp::Ordering {
1095 match self.sst_obj_id.cmp(&other.sst_obj_id) {
1096 std::cmp::Ordering::Equal => {}
1097 ord => return ord,
1098 }
1099 match self.blks.start.cmp(&other.blks.start) {
1100 std::cmp::Ordering::Equal => {}
1101 ord => return ord,
1102 }
1103 self.blks.end.cmp(&other.blks.end)
1104 }
1105}
1106
1107impl PartialOrd for SstableUnit {
1108 fn partial_cmp(&self, other: &Self) -> Option<std::cmp::Ordering> {
1109 Some(self.cmp(other))
1110 }
1111}
1112
1113#[cfg(test)]
1114mod tests {
1115 use std::collections::{HashMap, HashSet};
1116 use std::sync::Arc;
1117 use std::time::Duration;
1118
1119 use bytes::Bytes;
1120 use foyer::{
1121 BlockEngineConfig, CacheBuilder, DeviceBuilder, FsDeviceBuilder, HybridCacheBuilder,
1122 PsyncIoEngineConfig,
1123 };
1124 use parking_lot::Mutex;
1125 use risingwave_common::bitmap::Bitmap;
1126 use risingwave_common::config::streaming::CacheRefillPolicy;
1127 use risingwave_common::config::{MetricLevel, Role};
1128 use risingwave_common::hash::VirtualNode;
1129 use risingwave_common::util::epoch::test_epoch;
1130 use risingwave_hummock_sdk::compaction_group::group_split::split_sst_with_table_ids;
1131 use risingwave_hummock_sdk::key::{FullKey, UserKey, prefix_slice_with_vnode};
1132 use risingwave_hummock_sdk::sstable_info::{SstableInfo, SstableInfoInner};
1133 use risingwave_hummock_sdk::version::HummockVersion;
1134 use risingwave_hummock_sdk::{EpochWithGap, HummockSstableObjectId};
1135 use risingwave_pb::hummock::PbHummockVersion;
1136 use risingwave_pb::id::TableId;
1137 use tokio::sync::mpsc::unbounded_channel;
1138
1139 use super::{
1140 CacheRefillConfig, CacheRefillContext, CacheRefiller, DataCacheRefillTaskGenerator,
1141 SpawnRefillTask, SstDeltaInfo, block_vnode_range, vnode_range_overlaps_bitmap,
1142 };
1143 use crate::hummock::iterator::test_utils::{iterator_test_table_key_of, mock_sstable_store};
1144 use crate::hummock::local_version::pinned_version::PinnedVersion;
1145 use crate::hummock::recent_filter::simple::SimpleRecentFilter;
1146 use crate::hummock::sstable_store::{SstableStore, SstableStoreConfig};
1147 use crate::hummock::test_utils::{
1148 default_builder_opt_for_test, gen_test_sstable_with_table_ids,
1149 };
1150 use crate::hummock::value::HummockValue;
1151 use crate::hummock::{RecentFilter, RecentFilterTrait, SstableStoreRef, TableHolder};
1152 use crate::monitor::global_hummock_state_store_metrics;
1153
1154 async fn mock_sstable_store_with_disk_cache(
1155 recent_filter: Option<Arc<RecentFilter<(HummockSstableObjectId, usize)>>>,
1156 ) -> (SstableStoreRef, tempfile::TempDir) {
1157 let cache_dir = tempfile::tempdir().unwrap();
1158 let device = FsDeviceBuilder::new(cache_dir.path())
1159 .with_capacity(16 << 20)
1160 .build()
1161 .unwrap();
1162 let block_cache = HybridCacheBuilder::new()
1163 .memory(64 << 20)
1164 .with_shards(2)
1165 .storage()
1166 .with_io_engine_config(PsyncIoEngineConfig::new())
1167 .with_engine_config(
1168 BlockEngineConfig::new(device)
1169 .with_block_size(1 << 20)
1170 .with_buffer_pool_size(4 << 20),
1171 )
1172 .build()
1173 .await
1174 .unwrap();
1175 assert!(block_cache.is_hybrid());
1176 let memory_store = mock_sstable_store().await;
1177 let sstable_store = Arc::new(SstableStore::new(SstableStoreConfig {
1178 store: memory_store.store(),
1179 path: "test".to_owned(),
1180 prefetch_buffer_capacity: 64 << 20,
1181 max_prefetch_block_number: 16,
1182 recent_filter: recent_filter.unwrap_or_else(|| memory_store.recent_filter().clone()),
1183 state_store_metrics: Arc::new(global_hummock_state_store_metrics(
1184 MetricLevel::Disabled,
1185 )),
1186 use_new_object_prefix_strategy: true,
1187 skip_bloom_filter_in_serde: false,
1188 meta_cache: memory_store.meta_cache().clone(),
1189 block_cache,
1190 vector_meta_cache: CacheBuilder::new(64 << 20).build(),
1191 vector_block_cache: CacheBuilder::new(64 << 20).build(),
1192 }));
1193 (sstable_store, cache_dir)
1194 }
1195
1196 fn test_refill_config(default_policy: CacheRefillPolicy) -> CacheRefillConfig {
1197 CacheRefillConfig {
1198 timeout: Duration::from_secs(1),
1199 data_refill_levels: HashSet::new(),
1200 meta_refill_concurrency: 1,
1201 concurrency: 1,
1202 unit: 1,
1203 threshold: 0.0,
1204 skip_recent_filter: true,
1205 skip_inheritance_filter: true,
1206 table_cache_refill_default_policy: default_policy,
1207 }
1208 }
1209
1210 fn pinned_version_for_test() -> PinnedVersion {
1211 PinnedVersion::new(
1212 HummockVersion::from(PbHummockVersion::default()),
1213 unbounded_channel().0,
1214 )
1215 }
1216
1217 async fn gen_test_sst_with_object_id(
1218 table_id: TableId,
1219 sstable_store: SstableStoreRef,
1220 object_id: u64,
1221 ) -> (TableHolder, SstableInfo) {
1222 gen_test_sstable_with_table_ids(
1223 default_builder_opt_for_test(),
1224 object_id,
1225 (0..2).map(|idx| {
1226 (
1227 FullKey {
1228 user_key: risingwave_hummock_sdk::key::UserKey::for_test(
1229 table_id,
1230 iterator_test_table_key_of(idx),
1231 ),
1232 epoch_with_gap: EpochWithGap::new_from_epoch(test_epoch(233)),
1233 },
1234 HummockValue::put(vec![idx as u8]),
1235 )
1236 }),
1237 sstable_store,
1238 vec![table_id.as_raw_id()],
1239 )
1240 .await
1241 }
1242
1243 struct DataRefillGeneratorTestFixture {
1244 table_id: TableId,
1245 sstable_store: SstableStoreRef,
1246 sst: TableHolder,
1247 sst_info: SstableInfo,
1248 deleted_sst_object_id: HummockSstableObjectId,
1249 _cache_dir: tempfile::TempDir,
1250 }
1251
1252 impl DataRefillGeneratorTestFixture {
1253 async fn new(
1254 recent_filter: Option<Arc<RecentFilter<(HummockSstableObjectId, usize)>>>,
1255 ) -> Self {
1256 let table_id = TableId::from(233);
1257 let (sstable_store, cache_dir) =
1258 mock_sstable_store_with_disk_cache(recent_filter).await;
1259 let (sst, sst_info) =
1260 gen_test_sst_with_object_id(table_id, sstable_store.clone(), 1).await;
1261 Self {
1262 table_id,
1263 sstable_store,
1264 sst,
1265 sst_info,
1266 deleted_sst_object_id: 2330.into(),
1267 _cache_dir: cache_dir,
1268 }
1269 }
1270
1271 fn context(
1272 &self,
1273 policy: CacheRefillPolicy,
1274 streaming_vnode_bitmap: Option<Bitmap>,
1275 serving_vnode_bitmap: Option<Bitmap>,
1276 configure: impl FnOnce(&mut CacheRefillConfig),
1277 ) -> CacheRefillContext {
1278 let mut config = test_refill_config(CacheRefillPolicy::Enabled);
1279 config.data_refill_levels.insert(0);
1280 configure(&mut config);
1281 CacheRefillContext {
1282 config: Arc::new(config),
1283 meta_refill_concurrency: None,
1284 concurrency: Arc::new(tokio::sync::Semaphore::new(1)),
1285 sstable_store: self.sstable_store.clone(),
1286 table_cache_refill_context_map: Arc::new(HashMap::from([(
1287 self.table_id,
1288 super::TableCacheRefillContext {
1289 streaming_vnode_bitmap,
1290 serving_vnode_bitmap,
1291 policy,
1292 },
1293 )])),
1294 }
1295 }
1296
1297 fn normal_delta(
1298 &self,
1299 insert_sst_level: u32,
1300 deleted_sst_object_id: HummockSstableObjectId,
1301 ) -> SstDeltaInfo {
1302 SstDeltaInfo {
1303 insert_sst_infos: vec![self.sst_info.clone()],
1304 delete_sst_object_ids: vec![deleted_sst_object_id],
1305 insert_sst_level,
1306 }
1307 }
1308
1309 fn normal_l0_delta(&self) -> SstDeltaInfo {
1310 self.normal_delta(0, self.deleted_sst_object_id)
1311 }
1312
1313 fn l0_insert_only_delta(&self) -> SstDeltaInfo {
1314 SstDeltaInfo {
1315 insert_sst_infos: vec![self.sst_info.clone()],
1316 delete_sst_object_ids: vec![],
1317 insert_sst_level: 0,
1318 }
1319 }
1320
1321 async fn generate(
1322 &self,
1323 context: &CacheRefillContext,
1324 delta: &SstDeltaInfo,
1325 ) -> Vec<super::DataCacheRefillTask> {
1326 let generator = DataCacheRefillTaskGenerator {
1327 context,
1328 delta,
1329 ssts: std::slice::from_ref(&self.sst),
1330 };
1331 let tasks = generator.generate_unfiltered_tasks();
1332 generator.filter_by_inheritance_if_needed(tasks).await
1333 }
1334 }
1335
1336 #[tokio::test]
1337 async fn test_table_cache_refill_contexts_by_role_and_policy() {
1338 struct Case {
1339 name: &'static str,
1340 role: Role,
1341 default_policy: CacheRefillPolicy,
1342 policy: Option<CacheRefillPolicy>,
1343 has_streaming_vnodes: bool,
1344 has_serving_vnodes: bool,
1345 expected: Option<(CacheRefillPolicy, bool, bool)>,
1346 }
1347
1348 let cases = [
1349 Case {
1350 name: "streaming role uses streaming side of Both",
1351 role: Role::Streaming,
1352 default_policy: CacheRefillPolicy::Disabled,
1353 policy: Some(CacheRefillPolicy::Both),
1354 has_streaming_vnodes: true,
1355 has_serving_vnodes: true,
1356 expected: Some((CacheRefillPolicy::Both, true, false)),
1357 },
1358 Case {
1359 name: "serving role uses serving side of Both",
1360 role: Role::Serving,
1361 default_policy: CacheRefillPolicy::Disabled,
1362 policy: Some(CacheRefillPolicy::Both),
1363 has_streaming_vnodes: true,
1364 has_serving_vnodes: true,
1365 expected: Some((CacheRefillPolicy::Both, false, true)),
1366 },
1367 Case {
1368 name: "both role keeps both sides",
1369 role: Role::Both,
1370 default_policy: CacheRefillPolicy::Disabled,
1371 policy: Some(CacheRefillPolicy::Both),
1372 has_streaming_vnodes: true,
1373 has_serving_vnodes: true,
1374 expected: Some((CacheRefillPolicy::Both, true, true)),
1375 },
1376 Case {
1377 name: "both role keeps streaming-only ownership",
1378 role: Role::Both,
1379 default_policy: CacheRefillPolicy::Disabled,
1380 policy: Some(CacheRefillPolicy::Both),
1381 has_streaming_vnodes: true,
1382 has_serving_vnodes: false,
1383 expected: Some((CacheRefillPolicy::Both, true, false)),
1384 },
1385 Case {
1386 name: "streaming scope without ownership has no usable bitmap",
1387 role: Role::Streaming,
1388 default_policy: CacheRefillPolicy::Disabled,
1389 policy: Some(CacheRefillPolicy::Streaming),
1390 has_streaming_vnodes: false,
1391 has_serving_vnodes: false,
1392 expected: Some((CacheRefillPolicy::Streaming, false, false)),
1393 },
1394 Case {
1395 name: "pure serving worker excludes unmapped table",
1396 role: Role::Serving,
1397 default_policy: CacheRefillPolicy::Disabled,
1398 policy: Some(CacheRefillPolicy::Serving),
1399 has_streaming_vnodes: true,
1400 has_serving_vnodes: false,
1401 expected: None,
1402 },
1403 Case {
1404 name: "default Enabled retains serving ownership",
1405 role: Role::Serving,
1406 default_policy: CacheRefillPolicy::Enabled,
1407 policy: None,
1408 has_streaming_vnodes: false,
1409 has_serving_vnodes: true,
1410 expected: Some((CacheRefillPolicy::Enabled, false, true)),
1411 },
1412 Case {
1413 name: "explicit policy overrides default",
1414 role: Role::Serving,
1415 default_policy: CacheRefillPolicy::Enabled,
1416 policy: Some(CacheRefillPolicy::Disabled),
1417 has_streaming_vnodes: false,
1418 has_serving_vnodes: true,
1419 expected: Some((CacheRefillPolicy::Disabled, false, false)),
1420 },
1421 ];
1422
1423 let table_id = TableId::from(233);
1424 let streaming_vnodes = Bitmap::from_indices(VirtualNode::COUNT_FOR_TEST, [1, 3]);
1425 let serving_vnodes = Bitmap::from_indices(VirtualNode::COUNT_FOR_TEST, [2, 4]);
1426 let sstable_store = mock_sstable_store().await;
1427 for case in cases {
1428 let mut refiller = CacheRefiller::new(
1429 case.role,
1430 test_refill_config(case.default_policy),
1431 sstable_store.clone(),
1432 CacheRefiller::default_spawn_refill_task(),
1433 );
1434 if let Some(policy) = case.policy {
1435 refiller.replace_table_cache_refill_policies(HashMap::from([(table_id, policy)]));
1436 }
1437 if case.has_streaming_vnodes {
1438 refiller.update_streaming_table_vnodes(table_id, Some(streaming_vnodes.clone()));
1439 }
1440 if case.has_serving_vnodes {
1441 refiller.replace_serving_table_vnode_mapping(HashMap::from([(
1442 table_id,
1443 serving_vnodes.clone(),
1444 )]));
1445 }
1446
1447 let contexts = refiller.table_cache_refill_contexts([table_id]);
1448 let actual = contexts.get(&table_id).map(|context| {
1449 (
1450 context.policy,
1451 context.streaming_vnode_bitmap.as_ref(),
1452 context.serving_vnode_bitmap.as_ref(),
1453 )
1454 });
1455 let expected = case.expected.map(|(policy, streaming, serving)| {
1456 (
1457 policy,
1458 streaming.then_some(&streaming_vnodes),
1459 serving.then_some(&serving_vnodes),
1460 )
1461 });
1462 assert_eq!(actual, expected, "{}", case.name);
1463 }
1464 }
1465
1466 #[tokio::test]
1467 async fn test_refill_task_captures_runtime_context_snapshot() {
1468 let table_id = TableId::from(233);
1469 let old_vnodes = Bitmap::ones(VirtualNode::COUNT_FOR_TEST);
1470 let new_vnodes = Bitmap::from_range(VirtualNode::COUNT_FOR_TEST, 0..8);
1471 let captured_context = Arc::new(Mutex::new(None::<CacheRefillContext>));
1472 let captured_context_clone = captured_context.clone();
1473 let spawn_refill_task: SpawnRefillTask = Arc::new(move |_, context, _, _| {
1474 *captured_context_clone.lock() = Some(context);
1475 tokio::spawn(async {})
1476 });
1477 let mut refiller = CacheRefiller::new(
1478 Role::Serving,
1479 test_refill_config(CacheRefillPolicy::Enabled),
1480 mock_sstable_store().await,
1481 spawn_refill_task,
1482 );
1483
1484 refiller.replace_table_cache_refill_policies(HashMap::from([(
1485 table_id,
1486 CacheRefillPolicy::Serving,
1487 )]));
1488 refiller
1489 .replace_serving_table_vnode_mapping(HashMap::from([(table_id, old_vnodes.clone())]));
1490
1491 refiller.start_cache_refill(
1492 vec![SstDeltaInfo {
1493 insert_sst_infos: vec![SstableInfo::from(SstableInfoInner {
1494 table_ids: vec![table_id],
1495 ..Default::default()
1496 })],
1497 ..Default::default()
1498 }],
1499 pinned_version_for_test(),
1500 pinned_version_for_test(),
1501 );
1502 refiller.replace_table_cache_refill_policies(HashMap::from([(
1503 table_id,
1504 CacheRefillPolicy::Disabled,
1505 )]));
1506 refiller.replace_serving_table_vnode_mapping(HashMap::from([(table_id, new_vnodes)]));
1507
1508 let captured_context = captured_context.lock();
1509 let context = captured_context
1510 .as_ref()
1511 .unwrap()
1512 .table_cache_refill_context_map
1513 .get(&table_id)
1514 .unwrap();
1515 assert_eq!(context.policy, CacheRefillPolicy::Serving);
1516 assert_eq!(context.serving_vnode_bitmap.as_ref(), Some(&old_vnodes));
1517 }
1518
1519 #[tokio::test]
1520 async fn test_cache_refill_prunes_whole_ssts_before_meta_load() {
1521 let streaming_table = TableId::from(1);
1522 let serving_table = TableId::from(2);
1523 let disabled_table = TableId::from(3);
1524 let internal_table = TableId::from(4);
1525 let fallback_table = TableId::from(5);
1526 let serving_vnodes = Bitmap::ones(VirtualNode::COUNT_FOR_TEST);
1527 let sstable_store = mock_sstable_store().await;
1528 let sst_info = |table_ids: Vec<TableId>| {
1529 SstableInfo::from(SstableInfoInner {
1530 table_ids,
1531 ..Default::default()
1532 })
1533 };
1534 let capture_pruned_table_ids = |role,
1535 default_policy,
1536 policies: HashMap<TableId, CacheRefillPolicy>,
1537 streaming_table_vnodes: HashMap<TableId, Bitmap>,
1538 serving_table_vnodes: HashMap<TableId, Bitmap>,
1539 delta| {
1540 let captured_deltas = Arc::new(Mutex::new(None::<Vec<SstDeltaInfo>>));
1541 let captured_deltas_clone = captured_deltas.clone();
1542 let spawn_refill_task: SpawnRefillTask = Arc::new(move |deltas, _, _, _| {
1543 *captured_deltas_clone.lock() = Some(deltas);
1544 tokio::spawn(async {})
1545 });
1546 let mut refiller = CacheRefiller::new(
1547 role,
1548 test_refill_config(default_policy),
1549 sstable_store.clone(),
1550 spawn_refill_task,
1551 );
1552 refiller.replace_table_cache_refill_policies(policies);
1553 for (table_id, vnodes) in streaming_table_vnodes {
1554 refiller.update_streaming_table_vnodes(table_id, Some(vnodes));
1555 }
1556 refiller.replace_serving_table_vnode_mapping(serving_table_vnodes);
1557 refiller.start_cache_refill(
1558 vec![delta],
1559 pinned_version_for_test(),
1560 pinned_version_for_test(),
1561 );
1562 captured_deltas
1563 .lock()
1564 .take()
1565 .unwrap()
1566 .pop()
1567 .unwrap()
1568 .insert_sst_infos
1569 .into_iter()
1570 .map(|sst| sst.table_ids.clone())
1571 .collect::<Vec<_>>()
1572 };
1573
1574 let normal_delta = |insert_sst_infos| SstDeltaInfo {
1575 insert_sst_infos,
1576 delete_sst_object_ids: vec![1.into()],
1577 insert_sst_level: 1,
1578 };
1579 let insert_only_delta = |insert_sst_infos| SstDeltaInfo {
1580 insert_sst_infos,
1581 delete_sst_object_ids: vec![],
1582 insert_sst_level: 0,
1583 };
1584
1585 assert_eq!(
1586 capture_pruned_table_ids(
1587 Role::Both,
1588 CacheRefillPolicy::Disabled,
1589 HashMap::from([
1590 (disabled_table, CacheRefillPolicy::Disabled),
1591 (streaming_table, CacheRefillPolicy::Enabled),
1592 ]),
1593 HashMap::new(),
1594 HashMap::new(),
1595 normal_delta(vec![
1596 sst_info(vec![disabled_table]),
1597 sst_info(vec![streaming_table]),
1598 ]),
1599 ),
1600 vec![vec![streaming_table]],
1601 );
1602
1603 assert_eq!(
1604 capture_pruned_table_ids(
1605 Role::Both,
1606 CacheRefillPolicy::Disabled,
1607 HashMap::from([
1608 (streaming_table, CacheRefillPolicy::Streaming),
1609 (internal_table, CacheRefillPolicy::Serving),
1610 ]),
1611 HashMap::from([(streaming_table, serving_vnodes.clone())]),
1612 HashMap::new(),
1613 normal_delta(vec![
1614 sst_info(vec![streaming_table]),
1615 sst_info(vec![internal_table]),
1616 ]),
1617 ),
1618 vec![vec![streaming_table]],
1619 );
1620
1621 assert_eq!(
1622 capture_pruned_table_ids(
1623 Role::Both,
1624 CacheRefillPolicy::Disabled,
1625 HashMap::from([
1626 (streaming_table, CacheRefillPolicy::Streaming),
1627 (serving_table, CacheRefillPolicy::Serving),
1628 (internal_table, CacheRefillPolicy::Serving),
1629 ]),
1630 HashMap::from([(streaming_table, serving_vnodes.clone())]),
1631 HashMap::from([(serving_table, serving_vnodes.clone())]),
1632 insert_only_delta(vec![
1633 sst_info(vec![streaming_table]),
1634 sst_info(vec![serving_table]),
1635 sst_info(vec![internal_table]),
1636 ]),
1637 ),
1638 vec![vec![serving_table]],
1639 );
1640
1641 assert_eq!(
1642 capture_pruned_table_ids(
1643 Role::Serving,
1644 CacheRefillPolicy::Disabled,
1645 HashMap::from([
1646 (streaming_table, CacheRefillPolicy::Streaming),
1647 (serving_table, CacheRefillPolicy::Serving),
1648 ]),
1649 HashMap::new(),
1650 HashMap::from([(serving_table, serving_vnodes.clone())]),
1651 normal_delta(vec![
1652 sst_info(vec![streaming_table]),
1653 sst_info(vec![serving_table]),
1654 ]),
1655 ),
1656 vec![vec![serving_table]],
1657 );
1658
1659 assert_eq!(
1660 capture_pruned_table_ids(
1661 Role::Serving,
1662 CacheRefillPolicy::Disabled,
1663 HashMap::from([
1664 (streaming_table, CacheRefillPolicy::Enabled),
1665 (serving_table, CacheRefillPolicy::Enabled),
1666 ]),
1667 HashMap::new(),
1668 HashMap::from([(serving_table, serving_vnodes.clone())]),
1669 normal_delta(vec![
1670 sst_info(vec![streaming_table]),
1671 sst_info(vec![serving_table]),
1672 sst_info(vec![streaming_table, serving_table]),
1673 ]),
1674 ),
1675 vec![
1676 vec![streaming_table],
1677 vec![serving_table],
1678 vec![streaming_table, serving_table],
1679 ],
1680 );
1681
1682 assert_eq!(
1683 capture_pruned_table_ids(
1684 Role::Both,
1685 CacheRefillPolicy::Disabled,
1686 HashMap::from([
1687 (streaming_table, CacheRefillPolicy::Both),
1688 (serving_table, CacheRefillPolicy::Both),
1689 ]),
1690 HashMap::new(),
1691 HashMap::new(),
1692 normal_delta(vec![
1693 sst_info(vec![streaming_table]),
1694 sst_info(vec![serving_table]),
1695 ]),
1696 ),
1697 Vec::<Vec<TableId>>::new(),
1698 );
1699
1700 assert_eq!(
1701 capture_pruned_table_ids(
1702 Role::Both,
1703 CacheRefillPolicy::Disabled,
1704 HashMap::from([
1705 (streaming_table, CacheRefillPolicy::Both),
1706 (serving_table, CacheRefillPolicy::Both),
1707 ]),
1708 HashMap::from([(streaming_table, serving_vnodes.clone())]),
1709 HashMap::new(),
1710 normal_delta(vec![
1711 sst_info(vec![streaming_table]),
1712 sst_info(vec![serving_table]),
1713 ]),
1714 ),
1715 vec![vec![streaming_table]],
1716 );
1717
1718 assert_eq!(
1719 capture_pruned_table_ids(
1720 Role::Both,
1721 CacheRefillPolicy::Disabled,
1722 HashMap::from([
1723 (streaming_table, CacheRefillPolicy::Both),
1724 (serving_table, CacheRefillPolicy::Both),
1725 ]),
1726 HashMap::new(),
1727 HashMap::from([(serving_table, serving_vnodes.clone())]),
1728 normal_delta(vec![
1729 sst_info(vec![streaming_table]),
1730 sst_info(vec![serving_table]),
1731 ]),
1732 ),
1733 vec![vec![serving_table]],
1734 );
1735
1736 assert_eq!(
1737 capture_pruned_table_ids(
1738 Role::Streaming,
1739 CacheRefillPolicy::Disabled,
1740 HashMap::from([
1741 (streaming_table, CacheRefillPolicy::Streaming),
1742 (serving_table, CacheRefillPolicy::Serving),
1743 ]),
1744 HashMap::new(),
1745 HashMap::new(),
1746 normal_delta(vec![
1747 sst_info(vec![streaming_table]),
1748 sst_info(vec![serving_table]),
1749 ]),
1750 ),
1751 Vec::<Vec<TableId>>::new(),
1752 );
1753
1754 assert_eq!(
1755 capture_pruned_table_ids(
1756 Role::Streaming,
1757 CacheRefillPolicy::Disabled,
1758 HashMap::from([
1759 (streaming_table, CacheRefillPolicy::Streaming),
1760 (serving_table, CacheRefillPolicy::Serving),
1761 ]),
1762 HashMap::new(),
1763 HashMap::new(),
1764 insert_only_delta(vec![
1765 sst_info(vec![streaming_table]),
1766 sst_info(vec![serving_table]),
1767 ]),
1768 ),
1769 Vec::<Vec<TableId>>::new(),
1770 );
1771
1772 assert_eq!(
1773 capture_pruned_table_ids(
1774 Role::Both,
1775 CacheRefillPolicy::Enabled,
1776 HashMap::from([(disabled_table, CacheRefillPolicy::Disabled)]),
1777 HashMap::new(),
1778 HashMap::new(),
1779 normal_delta(vec![
1780 sst_info(vec![disabled_table]),
1781 sst_info(vec![fallback_table]),
1782 sst_info(vec![disabled_table, fallback_table]),
1783 ]),
1784 ),
1785 vec![vec![fallback_table], vec![disabled_table, fallback_table]],
1786 );
1787 }
1788
1789 #[tokio::test]
1790 async fn test_data_refill_requires_disk_cache() {
1791 let fixture = DataRefillGeneratorTestFixture::new(None).await;
1792 let delta = fixture.normal_l0_delta();
1793 let mut context = fixture.context(CacheRefillPolicy::Enabled, None, None, |_| {});
1794 assert!(!fixture.generate(&context, &delta).await.is_empty());
1795
1796 context.sstable_store = mock_sstable_store().await;
1797 assert!(!context.sstable_store.block_cache().is_hybrid());
1798 assert!(fixture.generate(&context, &delta).await.is_empty());
1799 }
1800
1801 #[tokio::test]
1802 async fn test_normal_refill_applies_policy_and_vnode_ownership() {
1803 let fixture = DataRefillGeneratorTestFixture::new(None).await;
1804 let delta = fixture.normal_l0_delta();
1805 let owned = Bitmap::from_indices(VirtualNode::COUNT_FOR_TEST, [0]);
1806 let unowned = Bitmap::from_indices(VirtualNode::COUNT_FOR_TEST, [1]);
1807 let cases = vec![
1808 ("Enabled", CacheRefillPolicy::Enabled, None, None, true),
1809 (
1810 "Disabled",
1811 CacheRefillPolicy::Disabled,
1812 Some(owned.clone()),
1813 Some(owned.clone()),
1814 false,
1815 ),
1816 (
1817 "Streaming match",
1818 CacheRefillPolicy::Streaming,
1819 Some(owned.clone()),
1820 None,
1821 true,
1822 ),
1823 (
1824 "Streaming miss",
1825 CacheRefillPolicy::Streaming,
1826 Some(unowned.clone()),
1827 None,
1828 false,
1829 ),
1830 (
1831 "Streaming ownership missing",
1832 CacheRefillPolicy::Streaming,
1833 None,
1834 None,
1835 false,
1836 ),
1837 (
1838 "Serving match",
1839 CacheRefillPolicy::Serving,
1840 None,
1841 Some(owned.clone()),
1842 true,
1843 ),
1844 (
1845 "Serving miss",
1846 CacheRefillPolicy::Serving,
1847 None,
1848 Some(unowned.clone()),
1849 false,
1850 ),
1851 (
1852 "Serving ownership missing",
1853 CacheRefillPolicy::Serving,
1854 None,
1855 None,
1856 false,
1857 ),
1858 (
1859 "Both streaming match",
1860 CacheRefillPolicy::Both,
1861 Some(owned.clone()),
1862 Some(unowned.clone()),
1863 true,
1864 ),
1865 (
1866 "Both serving match",
1867 CacheRefillPolicy::Both,
1868 Some(unowned.clone()),
1869 Some(owned),
1870 true,
1871 ),
1872 (
1873 "Both misses",
1874 CacheRefillPolicy::Both,
1875 Some(unowned.clone()),
1876 Some(unowned),
1877 false,
1878 ),
1879 ];
1880
1881 for (name, policy, streaming_vnodes, serving_vnodes, should_refill) in cases {
1882 let context = fixture.context(policy, streaming_vnodes, serving_vnodes, |_| {});
1883 assert_eq!(
1884 !fixture.generate(&context, &delta).await.is_empty(),
1885 should_refill,
1886 "{name}"
1887 );
1888 }
1889 }
1890
1891 #[tokio::test]
1892 async fn test_normal_refill_applies_recent_and_inheritance_filters() {
1893 let recent_filter = SimpleRecentFilter::new(3, Duration::from_secs(60));
1894 let fixture =
1895 DataRefillGeneratorTestFixture::new(Some(Arc::new(recent_filter.clone().into()))).await;
1896 let delta = fixture.normal_l0_delta();
1897
1898 let serving_context = fixture.context(
1899 CacheRefillPolicy::Serving,
1900 None,
1901 Some(Bitmap::ones(VirtualNode::COUNT_FOR_TEST)),
1902 |config| {
1903 config.skip_recent_filter = false;
1904 },
1905 );
1906 assert!(
1907 fixture.generate(&serving_context, &delta).await.is_empty(),
1908 "explicit Serving policy is not an implicit skip_recent_filter"
1909 );
1910
1911 recent_filter.insert((fixture.deleted_sst_object_id, usize::MAX));
1912 assert!(
1913 !fixture.generate(&serving_context, &delta).await.is_empty(),
1914 "explicit Serving policy should produce tasks after recent admission hits"
1915 );
1916
1917 let (_, parent_sst_info) =
1918 gen_test_sst_with_object_id(fixture.table_id, fixture.sstable_store.clone(), 2).await;
1919 let non_l0_delta = fixture.normal_delta(1, parent_sst_info.object_id);
1920 let non_l0_context = fixture.context(CacheRefillPolicy::Enabled, None, None, |config| {
1921 config.data_refill_levels.insert(1);
1922 config.skip_recent_filter = false;
1923 config.skip_inheritance_filter = false;
1924 });
1925
1926 recent_filter.insert((parent_sst_info.object_id, usize::MAX));
1927 let generator = DataCacheRefillTaskGenerator {
1928 context: &non_l0_context,
1929 delta: &non_l0_delta,
1930 ssts: std::slice::from_ref(&fixture.sst),
1931 };
1932 let unfiltered_tasks = generator.generate_unfiltered_tasks();
1933 assert_eq!(
1934 unfiltered_tasks
1935 .iter()
1936 .map(|task| task.blks.len())
1937 .sum::<usize>(),
1938 fixture.sst.block_count(),
1939 "recent-admitted blocks should reach the inheritance stage"
1940 );
1941 assert!(
1942 generator
1943 .filter_by_inheritance_if_needed(unfiltered_tasks)
1944 .await
1945 .is_empty(),
1946 "after recent admission, parent block recent miss should filter non-L0 normal refill"
1947 );
1948
1949 recent_filter.insert((parent_sst_info.object_id, 0));
1950 let tasks = fixture.generate(&non_l0_context, &non_l0_delta).await;
1951 assert_eq!(tasks.len(), 1);
1952 assert_eq!(tasks[0].sst.id, fixture.sst.id);
1953 assert_eq!(tasks[0].blks, 0..1);
1954 }
1955
1956 #[tokio::test]
1957 async fn test_l0_insert_only_refill_policy_uses_serving_ownership() {
1958 let fixture = DataRefillGeneratorTestFixture::new(None).await;
1959 let delta = fixture.l0_insert_only_delta();
1960 let cases = [
1961 (
1962 "Enabled + serving overlap",
1963 CacheRefillPolicy::Enabled,
1964 None,
1965 Some(Bitmap::ones(VirtualNode::COUNT_FOR_TEST)),
1966 true,
1967 ),
1968 (
1969 "Enabled without serving ownership",
1970 CacheRefillPolicy::Enabled,
1971 None,
1972 None,
1973 false,
1974 ),
1975 (
1976 "Serving + serving overlap",
1977 CacheRefillPolicy::Serving,
1978 None,
1979 Some(Bitmap::ones(VirtualNode::COUNT_FOR_TEST)),
1980 true,
1981 ),
1982 (
1983 "Streaming + streaming overlap",
1984 CacheRefillPolicy::Streaming,
1985 Some(Bitmap::ones(VirtualNode::COUNT_FOR_TEST)),
1986 None,
1987 false,
1988 ),
1989 (
1990 "Both + streaming overlap + serving non-overlap",
1991 CacheRefillPolicy::Both,
1992 Some(Bitmap::ones(VirtualNode::COUNT_FOR_TEST)),
1993 Some(Bitmap::from_indices(VirtualNode::COUNT_FOR_TEST, [1])),
1994 false,
1995 ),
1996 (
1997 "Both + serving overlap",
1998 CacheRefillPolicy::Both,
1999 Some(Bitmap::ones(VirtualNode::COUNT_FOR_TEST)),
2000 Some(Bitmap::ones(VirtualNode::COUNT_FOR_TEST)),
2001 true,
2002 ),
2003 (
2004 "Disabled + serving overlap",
2005 CacheRefillPolicy::Disabled,
2006 None,
2007 Some(Bitmap::ones(VirtualNode::COUNT_FOR_TEST)),
2008 false,
2009 ),
2010 ];
2011
2012 for (name, policy, streaming_vnodes, serving_vnodes, should_refill) in cases {
2013 let context = fixture.context(policy, streaming_vnodes, serving_vnodes, |config| {
2014 config.skip_recent_filter = false;
2015 config.skip_inheritance_filter = false;
2016 });
2017 assert_eq!(
2018 !fixture.generate(&context, &delta).await.is_empty(),
2019 should_refill,
2020 "{name}"
2021 );
2022 }
2023 }
2024
2025 #[tokio::test]
2026 async fn test_refill_units_do_not_cross_table_projection_boundaries() {
2027 let table_a = TableId::from(233);
2028 let table_b = TableId::from(234);
2029 let (sstable_store, _cache_dir) = mock_sstable_store_with_disk_cache(None).await;
2030 let (sst, sst_info) = gen_test_sstable_with_table_ids(
2031 default_builder_opt_for_test(),
2032 1,
2033 [table_a, table_b].into_iter().map(|table_id| {
2034 (
2035 FullKey {
2036 user_key: UserKey::for_test(table_id, iterator_test_table_key_of(0)),
2037 epoch_with_gap: EpochWithGap::new_from_epoch(test_epoch(233)),
2038 },
2039 HummockValue::put(b"value".to_vec()),
2040 )
2041 }),
2042 sstable_store.clone(),
2043 vec![table_a.as_raw_id(), table_b.as_raw_id()],
2044 )
2045 .await;
2046 assert_eq!(sst.block_count(), 2, "table switch must form a new block");
2047
2048 let mut next_sst_id = 100.into();
2049 let (table_a_projection, table_b_projection) =
2050 split_sst_with_table_ids(&sst_info, &mut next_sst_id, 1, 1, vec![table_b]);
2051 assert_eq!(table_a_projection.object_id, sst_info.object_id);
2052 assert_eq!(table_b_projection.object_id, sst_info.object_id);
2053 assert_ne!(table_a_projection.sst_id, table_b_projection.sst_id);
2054 assert_eq!(table_a_projection.table_ids, vec![table_a]);
2055 assert_eq!(table_b_projection.table_ids, vec![table_b]);
2056
2057 let deltas = [table_a_projection, table_b_projection].map(|projection| SstDeltaInfo {
2058 insert_sst_infos: vec![projection],
2059 delete_sst_object_ids: vec![],
2060 insert_sst_level: 0,
2061 });
2062 let normal_deltas = deltas.clone().map(|mut delta| {
2063 delta.delete_sst_object_ids = vec![999.into()];
2066 delta
2067 });
2068 let serving_vnodes = Bitmap::ones(VirtualNode::COUNT_FOR_TEST);
2069 let table_cache_refill_context_map = Arc::new(
2070 [table_a, table_b]
2071 .into_iter()
2072 .map(|table_id| {
2073 (
2074 table_id,
2075 super::TableCacheRefillContext {
2076 streaming_vnode_bitmap: None,
2077 serving_vnode_bitmap: Some(serving_vnodes.clone()),
2078 policy: CacheRefillPolicy::Serving,
2079 },
2080 )
2081 })
2082 .collect::<super::TableCacheRefillContextMap>(),
2083 );
2084 let make_context = |unit| {
2085 let mut config = test_refill_config(CacheRefillPolicy::Disabled);
2086 config.data_refill_levels.insert(0);
2087 config.unit = unit;
2088 CacheRefillContext {
2089 config: Arc::new(config),
2090 meta_refill_concurrency: None,
2091 concurrency: Arc::new(tokio::sync::Semaphore::new(1)),
2092 sstable_store: sstable_store.clone(),
2093 table_cache_refill_context_map: table_cache_refill_context_map.clone(),
2094 }
2095 };
2096 let generated_tasks = |context: &CacheRefillContext| {
2097 deltas
2098 .iter()
2099 .map(|delta| {
2100 DataCacheRefillTaskGenerator {
2101 context,
2102 delta,
2103 ssts: std::slice::from_ref(&sst),
2104 }
2105 .generate_unfiltered_tasks()
2106 })
2107 .collect::<Vec<_>>()
2108 };
2109 let generated_ranges = |context: &CacheRefillContext| {
2110 generated_tasks(context)
2111 .into_iter()
2112 .map(|tasks| tasks.into_iter().map(|task| task.blks).collect::<Vec<_>>())
2113 .collect::<Vec<_>>()
2114 };
2115
2116 assert_eq!(
2117 generated_ranges(&make_context(1)),
2118 vec![vec![0..1], vec![1..2]],
2119 "each logical projection must select only its own block"
2120 );
2121
2122 let wide_unit_context = make_context(2);
2123 assert_eq!(
2124 generated_ranges(&wide_unit_context),
2125 vec![vec![0..1], vec![1..2]],
2126 "units are clipped at table boundaries even when unit is larger than a table run"
2127 );
2128
2129 let normal_ranges = normal_deltas
2130 .iter()
2131 .map(|delta| {
2132 DataCacheRefillTaskGenerator {
2133 context: &wide_unit_context,
2134 delta,
2135 ssts: std::slice::from_ref(&sst),
2136 }
2137 .generate_unfiltered_tasks()
2138 .into_iter()
2139 .map(|task| task.blks)
2140 .collect::<Vec<_>>()
2141 })
2142 .collect::<Vec<_>>();
2143 assert_eq!(
2144 normal_ranges,
2145 vec![vec![0..1], vec![1..2]],
2146 "normal refill uses the same table-boundary geometry"
2147 );
2148
2149 for task in generated_tasks(&wide_unit_context).into_iter().flatten() {
2150 assert!(task.blks.len() <= wide_unit_context.config.unit);
2151 assert_eq!(
2152 task.sst.meta.block_metas[task.blks.start].table_id(),
2153 task.sst.meta.block_metas[task.blks.end - 1].table_id(),
2154 "a refill unit must not cross a table boundary"
2155 );
2156 }
2157 }
2158
2159 #[tokio::test]
2160 async fn test_scoped_refill_handles_multi_table_vnode_boundary() {
2161 let table_a = TableId::from(233);
2162 let table_b = TableId::from(234);
2163 let vnode_a = VirtualNode::COUNT_FOR_TEST - 1;
2164 let (sstable_store, _cache_dir) = mock_sstable_store_with_disk_cache(None).await;
2165 let (sst, sst_info) = gen_test_sstable_with_table_ids(
2166 default_builder_opt_for_test(),
2167 1,
2168 [
2169 (
2170 FullKey {
2171 user_key: UserKey::for_test(
2172 table_a,
2173 prefix_slice_with_vnode(VirtualNode::from_index(vnode_a), b"table_a"),
2174 ),
2175 epoch_with_gap: EpochWithGap::new_from_epoch(test_epoch(233)),
2176 },
2177 HummockValue::put(Bytes::from_static(b"a")),
2178 ),
2179 (
2180 FullKey {
2181 user_key: UserKey::for_test(
2182 table_b,
2183 prefix_slice_with_vnode(VirtualNode::ZERO, b"table_b"),
2184 ),
2185 epoch_with_gap: EpochWithGap::new_from_epoch(test_epoch(233)),
2186 },
2187 HummockValue::put(Bytes::from_static(b"b")),
2188 ),
2189 ]
2190 .into_iter(),
2191 sstable_store.clone(),
2192 vec![table_a.as_raw_id(), table_b.as_raw_id()],
2193 )
2194 .await;
2195 assert_eq!(sst.block_count(), 2, "table switch must form a new block");
2196
2197 let generate = |streaming_vnodes| {
2198 let sstable_store = sstable_store.clone();
2199 let sst = sst.clone();
2200 let sst_info = sst_info.clone();
2201 let mut config = test_refill_config(CacheRefillPolicy::Streaming);
2202 config.data_refill_levels.insert(0);
2203 let context = CacheRefillContext {
2204 config: Arc::new(config),
2205 meta_refill_concurrency: None,
2206 concurrency: Arc::new(tokio::sync::Semaphore::new(1)),
2207 sstable_store,
2208 table_cache_refill_context_map: Arc::new(HashMap::from([(
2209 table_a,
2210 super::TableCacheRefillContext {
2211 streaming_vnode_bitmap: Some(streaming_vnodes),
2212 serving_vnode_bitmap: None,
2213 policy: CacheRefillPolicy::Streaming,
2214 },
2215 )])),
2216 };
2217 async move {
2218 let generator = DataCacheRefillTaskGenerator {
2219 context: &context,
2220 delta: &SstDeltaInfo {
2221 insert_sst_infos: vec![sst_info.clone()],
2222 delete_sst_object_ids: vec![2330.into()],
2223 insert_sst_level: 0,
2224 },
2225 ssts: std::slice::from_ref(&sst),
2226 };
2227 let tasks = generator.generate_unfiltered_tasks();
2228 generator.filter_by_inheritance_if_needed(tasks).await
2229 }
2230 };
2231
2232 let matching_tasks =
2233 generate(Bitmap::from_indices(VirtualNode::COUNT_FOR_TEST, [vnode_a])).await;
2234 assert_eq!(matching_tasks.len(), 1);
2235 assert_eq!(matching_tasks[0].blks, 0..1);
2236
2237 let non_matching_tasks = generate(Bitmap::from_indices(
2238 VirtualNode::COUNT_FOR_TEST,
2239 [VirtualNode::ZERO.to_index()],
2240 ))
2241 .await;
2242 assert!(non_matching_tasks.is_empty());
2243 }
2244
2245 #[tokio::test]
2246 async fn test_block_vnode_range_handles_vnode_only_block_boundaries() {
2247 let table_id = TableId::from(233);
2248 let vnode = VirtualNode::ZERO;
2249 let sstable_store = mock_sstable_store().await;
2250 let mut builder_options = default_builder_opt_for_test();
2251 builder_options.block_capacity = 1;
2252 let (sst, _) = gen_test_sstable_with_table_ids(
2253 builder_options,
2254 1,
2255 [234, 233].into_iter().map(|epoch| {
2256 (
2257 FullKey {
2258 user_key: UserKey::for_test(table_id, prefix_slice_with_vnode(vnode, b"")),
2259 epoch_with_gap: EpochWithGap::new_from_epoch(test_epoch(epoch)),
2260 },
2261 HummockValue::put(Bytes::from_static(b"value")),
2262 )
2263 }),
2264 sstable_store.clone(),
2265 vec![table_id.as_raw_id()],
2266 )
2267 .await;
2268 assert_eq!(sst.block_count(), 2);
2269 let expected = (vnode.to_index(), vnode.to_index() + 1);
2270 assert_eq!(block_vnode_range(&sst, 0), expected);
2271 assert_eq!(block_vnode_range(&sst, 1), expected);
2272 }
2273
2274 #[tokio::test]
2275 async fn test_block_vnode_range_fails_open_for_shortened_meta_keys() {
2276 let table_id = TableId::from(233);
2277 let sstable_store = mock_sstable_store().await;
2278 let mut builder_options = default_builder_opt_for_test();
2279 builder_options.block_capacity = 1;
2280 builder_options.shorten_block_meta_key_threshold = Some(0);
2281 let (sst, _) = gen_test_sstable_with_table_ids(
2282 builder_options,
2283 1,
2284 [255, 256].into_iter().map(|vnode| {
2285 (
2286 FullKey {
2287 user_key: UserKey::for_test(
2288 table_id,
2289 prefix_slice_with_vnode(VirtualNode::from_index(vnode), b"long-key"),
2290 ),
2291 epoch_with_gap: EpochWithGap::new_from_epoch(test_epoch(233)),
2292 },
2293 HummockValue::put(Bytes::from_static(b"value")),
2294 )
2295 }),
2296 sstable_store,
2297 vec![table_id.as_raw_id()],
2298 )
2299 .await;
2300 assert_eq!(sst.block_count(), 2);
2301 assert!(
2302 FullKey::decode(&sst.meta.block_metas[1].smallest_key)
2303 .user_key
2304 .table_key
2305 .as_ref()
2306 .len()
2307 < VirtualNode::SIZE
2308 );
2309 let full_range = (0, VirtualNode::MAX_REPRESENTABLE.to_index() + 1);
2310 assert_eq!(block_vnode_range(&sst, 0), full_range);
2311 assert_eq!(block_vnode_range(&sst, 1), full_range);
2312 }
2313
2314 #[test]
2315 fn test_vnode_range_overlaps_bitmap_uses_right_exclusive_end() {
2316 let right_exclusive = Bitmap::from_indices(VirtualNode::COUNT_FOR_TEST, [12]);
2317 assert!(!vnode_range_overlaps_bitmap((10, 12), &right_exclusive));
2318
2319 let inside_range = Bitmap::from_indices(VirtualNode::COUNT_FOR_TEST, [11]);
2320 assert!(vnode_range_overlaps_bitmap((10, 12), &inside_range));
2321
2322 let last_vnode = Bitmap::from_indices(
2323 VirtualNode::COUNT_FOR_TEST,
2324 [VirtualNode::COUNT_FOR_TEST - 1],
2325 );
2326 assert!(vnode_range_overlaps_bitmap(
2327 (VirtualNode::COUNT_FOR_TEST - 1, VirtualNode::COUNT_FOR_TEST),
2328 &last_vnode
2329 ));
2330 assert!(!vnode_range_overlaps_bitmap(
2331 (VirtualNode::COUNT_FOR_TEST, VirtualNode::COUNT_FOR_TEST + 1),
2332 &last_vnode
2333 ));
2334 }
2335}