Skip to main content

risingwave_storage/hummock/event_handler/
refiller.rs

1// Copyright 2023 RisingWave Labs
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use 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    /// Cache refill timeout.
195    pub timeout: Duration,
196
197    /// Data file cache refill levels.
198    pub data_refill_levels: HashSet<u32>,
199
200    /// Meta file cache refill concurrency.
201    pub meta_refill_concurrency: usize,
202
203    /// Data file cache refill concurrency.
204    pub concurrency: usize,
205
206    /// Data file cache refill unit (blocks).
207    pub unit: usize,
208
209    /// Data file cache reill unit threshold.
210    ///
211    /// Only units whose admit rate > threshold will be refilled.
212    pub threshold: f64,
213
214    /// Skip recent filter.
215    pub skip_recent_filter: bool,
216
217    /// Skip inheritance filter.
218    pub skip_inheritance_filter: bool,
219
220    /// Default table cache refill policy.
221    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    // first current version, second new version
260    dyn Fn(Vec<SstDeltaInfo>, CacheRefillContext, PinnedVersion, PinnedVersion) -> JoinHandle<()>
261        + Send
262        + Sync
263        + 'static,
264>;
265
266pub type TableCacheRefillContextMap = HashMap<TableId, TableCacheRefillContext>;
267
268/// Per-table metadata captured for a refill task to decide whether an sstable block should be
269/// refilled. Mutable runtime state used to build this snapshot is owned by `CacheRefiller`.
270#[derive(Clone)]
271pub struct TableCacheRefillContext {
272    /// Vnodes covered by local streaming read versions on this compute node.
273    pub streaming_vnode_bitmap: Option<Bitmap>,
274    /// Vnodes served by this compute node according to the serving vnode mapping.
275    pub serving_vnode_bitmap: Option<Bitmap>,
276    /// Effective refill policy after applying the default policy and per-table overrides.
277    pub policy: CacheRefillPolicy,
278}
279
280/// Read-only data cloned from `CacheRefiller` for monitor/debugging APIs.
281///
282/// Streaming vnode mapping is the table-level union maintained by the refiller.
283#[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        // Insert-only deltas have no delete-side evidence for recent/inheritance filters.
319        // Only serving-owned blocks need refill, because streaming writers already populated
320        // their local cache.
321        (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        // A table switch always starts a new block. The next table's smallest key has an
345        // unrelated vnode, so use the current table's terminal range instead.
346        Some(next_block_meta) if next_block_meta.table_id() != block_meta.table_id() => {
347            Bound::Unbounded
348        }
349        // Full-key versions of the same table key may span adjacent blocks. After projecting
350        // away the epoch, the boundary vnode therefore remains part of the current block.
351        Some(next_block_meta) => Bound::Included(
352            FullKey::decode(&next_block_meta.smallest_key)
353                .user_key
354                .table_key,
355        ),
356        // `SstableMeta::largest_key` is the actual last key, unlike the next block's smallest
357        // key above. Keep it inclusive, especially for singleton tables whose key contains only
358        // the vnode prefix.
359        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    // Block-meta separators may shorten the table key below the vnode prefix. They are valid
371    // full-key search boundaries but cannot identify a vnode, so fail open instead of panicking
372    // or dropping a block that may belong to this worker.
373    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
385/// A cache refiller for hummock data.
386pub(crate) struct CacheRefiller {
387    /// order: old => new
388    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            // Writer-appended L0 SSTs are already warm on the streaming side. Their data refill
450            // may therefore only be needed by serving workers.
451            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            // This is deliberately a whole-SST admission check before Meta load. The serving
460            // mapping key set identifies result tables, but bitmap contents remain for the exact
461            // block/vnode decision in DataCacheRefillTaskGenerator after Meta load. An SST stays
462            // when any contained table matches either the serving or streaming refill lane.
463            delta.insert_sst_infos.retain(|sst| {
464                sst.table_ids.iter().any(|table_id| {
465                    // A missing entry means there is no table override, not that the table is
466                    // absent. Preserve the configured legacy/default policy in that case.
467                    let policy = self
468                        .table_cache_refill_policies
469                        .get(table_id)
470                        .copied()
471                        .unwrap_or(self.default_policy);
472
473                    // Enabled preserves legacy full refill. For scoped policies, mapping keys are
474                    // only a whole-SST coarse gate; bitmap bits still filter blocks post-Meta.
475                    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    /// Replaces the complete policy snapshot applicable to this worker.
532    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    /// Replaces the complete serving vnode mapping snapshot.
540    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                // `Enabled` normally does not use bitmap filtering. The only exception is L0
584                // insert-only refill, where serving workers still need serving-locality evidence.
585                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        // Skip data cache refill if data disk cache is not enabled.
669        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        // CN-written SSTs are appended to L0 without replacing parent SSTs. Other inserted SSTs
679        // need delete-side evidence for recent and inheritance filtering.
680        debug_assert!(has_parent_ssts || self.delta.insert_sst_level == 0);
681
682        // Return if the target level is not in the refill levels
683        if !self
684            .context
685            .config
686            .data_refill_levels
687            .contains(&self.delta.insert_sst_level)
688        {
689            return tasks;
690        }
691
692        // Cache refill units must not cross a table boundary. A logical SST projection still
693        // decides whether to admit each single-table unit.
694        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                // SstableBuilder ends a block before the table ID changes, so block metadata
703                // defines the exact physical boundary. `table_ids` below only admits logical
704                // projections and must not make a unit span another table.
705                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        // Policy/vnode ownership defines refill responsibility first, but it does not bypass
740        // recent admission for normal insert+delete refill.
741        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        // Skipping the recent filter selects full refill. Inheritance filtering only applies to
759        // non-L0 normal refill after real recent-filter admission.
760        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    // Return if recent filter is required and no deleted sst ids are in the recent filter.
773    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        // Get parent sst metas from cache.
789        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        // assert units in asc order
815        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                    // `largest_key` can be included or excluded, both are treated as included here
831                    &psst.meta.largest_key
832                } else {
833                    &psst.meta.block_metas[pblk + 1].smallest_key
834                };
835
836                // partition point: unit.right < pblk.left
837                let uleft = originals.partition_point(|task| {
838                    KeyComparator::compare_encoded_full_key(task.largest_key(), pleft)
839                        == std::cmp::Ordering::Less
840                });
841                // partition point: unit.left <= pblk.right
842                let uright = originals.partition_point(|task| {
843                    KeyComparator::compare_encoded_full_key(task.smallest_key(), pright)
844                        != std::cmp::Ordering::Greater
845                });
846
847                // overlapping: uleft..uright
848                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                    // Main counts after recent admission but before inheritance.
931                    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            // update filter for sst id only
985            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                            // The entry should always be `Some(..)`, use if here for compatible.
1044                            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            // A synthetic delete marks this as a normal delta; recent and inheritance filters
2064            // are disabled below, so the test does not rely on a matching parent SST.
2065            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}