Skip to main content

risingwave_storage/hummock/store/
version.rs

1// Copyright 2022 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::cmp::Ordering;
16use std::collections::HashMap;
17use std::collections::vec_deque::VecDeque;
18use std::ops::Bound::{self};
19use std::sync::Arc;
20use std::time::Instant;
21
22use bytes::Bytes;
23use futures::future::try_join_all;
24use itertools::Itertools;
25use parking_lot::RwLock;
26use risingwave_common::array::VectorRef;
27use risingwave_common::bitmap::Bitmap;
28use risingwave_common::catalog::{TableId, TableOption};
29use risingwave_common::hash::VirtualNode;
30use risingwave_common::util::epoch::MAX_SPILL_TIMES;
31use risingwave_hummock_sdk::key::{
32    FullKey, TableKey, TableKeyRange, UserKey, bound_table_key_range,
33};
34use risingwave_hummock_sdk::key_range::KeyRangeCommon;
35use risingwave_hummock_sdk::sstable_info::SstableInfo;
36use risingwave_hummock_sdk::table_watermark::{
37    TableWatermarksIndex, VnodeWatermark, WatermarkDirection, WatermarkSerdeType,
38};
39use risingwave_hummock_sdk::vector_index::VectorIndexImpl;
40use risingwave_hummock_sdk::{EpochWithGap, HummockEpoch, LocalSstableInfo};
41use risingwave_pb::hummock::LevelType;
42use sync_point::sync_point;
43use tracing::warn;
44
45use crate::error::StorageResult;
46use crate::hummock::event_handler::LocalInstanceId;
47use crate::hummock::iterator::change_log::ChangeLogIterator;
48use crate::hummock::iterator::{
49    BackwardUserIterator, HummockIterator, IteratorFactory, MergeIterator, UserIterator,
50};
51use crate::hummock::local_version::pinned_version::PinnedVersion;
52use crate::hummock::sstable::{SstableIteratorReadOptions, SstableIteratorType};
53use crate::hummock::sstable_store::SstableStoreRef;
54use crate::hummock::table_change_log_manager::TableChangeLogManager;
55use crate::hummock::utils::{
56    MemoryTracker, filter_single_sst, prune_nonoverlapping_ssts, prune_overlapping_ssts,
57    range_overlap, search_sst_idx,
58};
59use crate::hummock::vector::file::{FileVectorStore, FileVectorStoreCtx};
60use crate::hummock::vector::monitor::{VectorStoreCacheStats, report_hnsw_stat};
61use crate::hummock::{
62    BackwardIteratorFactory, ForwardIteratorFactory, HummockError, HummockResult,
63    HummockStorageIterator, HummockStorageIteratorInner, HummockStorageRevIteratorInner,
64    ReadVersionTuple, Sstable, SstableIterator, get_from_batch, get_from_sstable_info,
65    hit_sstable_filter,
66};
67use crate::mem_table::{
68    ImmId, ImmutableMemtable, MemTableHummockIterator, MemTableHummockRevIterator,
69};
70use crate::monitor::{
71    GetLocalMetricsGuard, HummockStateStoreMetrics, IterLocalMetricsGuard, StoreLocalStatistic,
72};
73use crate::store::{
74    OnNearestItemFn, ReadLogOptions, ReadOptions, VectorNearestOptions, gen_min_epoch,
75};
76use crate::vector::hnsw::nearest;
77use crate::vector::{MeasureDistanceBuilder, NearestBuilder};
78
79pub type CommittedVersion = PinnedVersion;
80
81/// Data not committed to Hummock. There are two types of staging data:
82/// - Immutable memtable: data that has been written into local state store but not persisted.
83/// - Uncommitted SST: data that has been uploaded to persistent storage but not committed to
84///   hummock version.
85
86#[derive(Clone, Debug, PartialEq)]
87pub struct StagingSstableInfo {
88    // newer data comes first
89    sstable_infos: Vec<LocalSstableInfo>,
90    old_value_sstable_infos: Vec<LocalSstableInfo>,
91    /// Epochs whose data are included in the Sstable. The newer epoch comes first.
92    /// The field must not be empty.
93    epochs: Vec<HummockEpoch>,
94    // newer data at the front
95    imm_ids: HashMap<LocalInstanceId, Vec<ImmId>>,
96    imm_size: usize,
97}
98
99impl StagingSstableInfo {
100    pub fn new(
101        sstable_infos: Vec<LocalSstableInfo>,
102        old_value_sstable_infos: Vec<LocalSstableInfo>,
103        epochs: Vec<HummockEpoch>,
104        imm_ids: HashMap<LocalInstanceId, Vec<ImmId>>,
105        imm_size: usize,
106    ) -> Self {
107        // the epochs are sorted from higher epoch to lower epoch
108        assert!(epochs.is_sorted_by(|epoch1, epoch2| epoch2 <= epoch1));
109        Self {
110            sstable_infos,
111            old_value_sstable_infos,
112            epochs,
113            imm_ids,
114            imm_size,
115        }
116    }
117
118    pub fn sstable_infos(&self) -> &Vec<LocalSstableInfo> {
119        &self.sstable_infos
120    }
121
122    pub fn old_value_sstable_infos(&self) -> &Vec<LocalSstableInfo> {
123        &self.old_value_sstable_infos
124    }
125
126    pub fn imm_size(&self) -> usize {
127        self.imm_size
128    }
129
130    pub fn epochs(&self) -> &Vec<HummockEpoch> {
131        &self.epochs
132    }
133
134    pub fn imm_ids(&self) -> &HashMap<LocalInstanceId, Vec<ImmId>> {
135        &self.imm_ids
136    }
137}
138
139pub enum VersionUpdate {
140    Sst(Arc<StagingSstableInfo>),
141    CommittedSnapshot(CommittedVersion),
142    NewTableWatermark {
143        direction: WatermarkDirection,
144        epoch: HummockEpoch,
145        vnode_watermarks: Vec<VnodeWatermark>,
146        watermark_type: WatermarkSerdeType,
147    },
148}
149
150pub struct StagingVersion {
151    pending_imm_size: usize,
152    /// It contains the imms added but not sent to the uploader of hummock event handler.
153    /// It is non-empty only when `upload_on_flush` is false.
154    ///
155    /// It will be sent to the uploader when `pending_imm_size` exceed threshold or on `seal_current_epoch`.
156    ///
157    /// newer data comes last
158    pub pending_imms: Vec<(ImmutableMemtable, MemoryTracker)>,
159    /// It contains the imms already sent to uploader of hummock event handler.
160    /// Note: Currently, building imm and writing to staging version is not atomic, and therefore
161    /// imm of smaller batch id may be added later than one with greater batch id
162    ///
163    /// Newer data comes first.
164    pub uploading_imms: VecDeque<ImmutableMemtable>,
165
166    // newer data comes first
167    pub sst: VecDeque<Arc<StagingSstableInfo>>,
168}
169
170impl StagingVersion {
171    /// Get the overlapping `imm`s and `sst`s that overlap respectively with `table_key_range` and
172    /// the user key range derived from `table_id`, `epoch` and `table_key_range`.
173    pub fn prune_overlap<'a>(
174        &'a self,
175        max_epoch_inclusive: HummockEpoch,
176        table_id: TableId,
177        table_key_range: &'a TableKeyRange,
178    ) -> (
179        impl Iterator<Item = &'a ImmutableMemtable> + 'a,
180        impl Iterator<Item = &'a SstableInfo> + 'a,
181    ) {
182        let (left, right) = table_key_range;
183        let left = left.as_ref().map(|key| TableKey(key.0.as_ref()));
184        let right = right.as_ref().map(|key| TableKey(key.0.as_ref()));
185        let overlapped_imms = self
186            .pending_imms
187            .iter()
188            .map(|(imm, _)| imm)
189            .rev() // rev to let newer imm come first
190            .chain(self.uploading_imms.iter())
191            .filter(move |imm| {
192                // retain imm which is overlapped with (min_epoch_exclusive, max_epoch_inclusive]
193                imm.epoch() <= max_epoch_inclusive
194                    && imm.table_id == table_id
195                    && range_overlap(
196                        &(left, right),
197                        &imm.start_table_key(),
198                        Bound::Included(&imm.end_table_key()),
199                    )
200            });
201
202        // TODO: Remove duplicate sst based on sst id
203        let overlapped_ssts = self
204            .sst
205            .iter()
206            .filter(move |staging_sst| {
207                let sst_max_epoch = *staging_sst.epochs.last().expect("epochs not empty");
208                sst_max_epoch <= max_epoch_inclusive
209            })
210            .flat_map(move |staging_sst| {
211                // TODO: sstable info should be concat-able after each streaming table owns a read
212                // version. May use concat sstable iter instead in some cases.
213                staging_sst
214                    .sstable_infos
215                    .iter()
216                    .map(|sstable| &sstable.sst_info)
217                    .filter(move |sstable: &&SstableInfo| {
218                        filter_single_sst(sstable, table_id, table_key_range)
219                    })
220            });
221        (overlapped_imms, overlapped_ssts)
222    }
223
224    pub fn is_empty(&self) -> bool {
225        self.pending_imms.is_empty() && self.uploading_imms.is_empty() && self.sst.is_empty()
226    }
227}
228
229/// A container of information required for reading from hummock.
230pub struct HummockReadVersion {
231    table_id: TableId,
232    instance_id: LocalInstanceId,
233
234    is_initialized: bool,
235
236    /// Local version for staging data.
237    staging: StagingVersion,
238
239    /// Remote version for committed data.
240    committed: CommittedVersion,
241
242    /// Indicate if this is replicated. If it is, we should ignore it during
243    /// global state store read, to avoid duplicated results.
244    /// Otherwise for local state store, it is fine, see we will see the
245    /// `ReadVersion` just for that local state store.
246    is_replicated: bool,
247
248    table_watermarks: Option<TableWatermarksIndex>,
249
250    // Vnode bitmap corresponding to the read version
251    // It will be initialized after local state store init
252    vnodes: Arc<Bitmap>,
253}
254
255impl HummockReadVersion {
256    pub fn new_with_replication_option(
257        table_id: TableId,
258        instance_id: LocalInstanceId,
259        committed_version: CommittedVersion,
260        is_replicated: bool,
261        vnodes: Arc<Bitmap>,
262    ) -> Self {
263        // before build `HummockReadVersion`, we need to get the a initial version which obtained
264        // from meta. want this initialization after version is initialized (now with
265        // notification), so add a assert condition to guarantee correct initialization order
266        assert!(committed_version.is_valid());
267        Self {
268            table_id,
269            instance_id,
270            table_watermarks: {
271                match committed_version.table_watermarks.get(&table_id) {
272                    Some(table_watermarks) => Some(TableWatermarksIndex::new_committed(
273                        table_watermarks.clone(),
274                        committed_version
275                            .state_table_info
276                            .info()
277                            .get(&table_id)
278                            .expect("should exist")
279                            .committed_epoch,
280                        table_watermarks.watermark_type,
281                    )),
282                    None => None,
283                }
284            },
285            staging: StagingVersion {
286                pending_imm_size: 0,
287                pending_imms: Vec::default(),
288                uploading_imms: VecDeque::default(),
289                sst: VecDeque::default(),
290            },
291
292            committed: committed_version,
293
294            is_replicated,
295            vnodes,
296            is_initialized: false,
297        }
298    }
299
300    pub fn new(
301        table_id: TableId,
302        instance_id: LocalInstanceId,
303        committed_version: CommittedVersion,
304        vnodes: Arc<Bitmap>,
305    ) -> Self {
306        Self::new_with_replication_option(table_id, instance_id, committed_version, false, vnodes)
307    }
308
309    pub fn table_id(&self) -> TableId {
310        self.table_id
311    }
312
313    pub fn init(&mut self) {
314        assert!(!self.is_initialized);
315        self.is_initialized = true;
316    }
317
318    pub fn add_pending_imm(&mut self, imm: ImmutableMemtable, tracker: MemoryTracker) {
319        assert!(self.is_initialized);
320        assert!(!self.is_replicated);
321        if let Some(item) = self
322            .staging
323            .pending_imms
324            .last()
325            .map(|(imm, _)| imm)
326            .or_else(|| self.staging.uploading_imms.front())
327        {
328            // check batch_id order from newest to old
329            assert!(item.batch_id() < imm.batch_id());
330        }
331
332        self.staging.pending_imm_size += imm.size();
333        self.staging.pending_imms.push((imm, tracker));
334    }
335
336    pub fn add_replicated_imm(&mut self, imm: ImmutableMemtable) {
337        assert!(self.is_initialized);
338        assert!(self.is_replicated);
339        assert!(self.staging.pending_imms.is_empty());
340        if let Some(item) = self.staging.uploading_imms.front() {
341            // check batch_id order from newest to old
342            assert!(item.batch_id() < imm.batch_id());
343        }
344        self.staging.uploading_imms.push_front(imm);
345    }
346
347    pub fn pending_imm_size(&self) -> usize {
348        self.staging.pending_imm_size
349    }
350
351    pub fn start_upload_pending_imms(&mut self) -> Vec<(ImmutableMemtable, MemoryTracker)> {
352        assert!(self.is_initialized);
353        assert!(!self.is_replicated);
354        let pending_imms = std::mem::take(&mut self.staging.pending_imms);
355        for (imm, _) in &pending_imms {
356            self.staging.uploading_imms.push_front(imm.clone());
357        }
358        self.staging.pending_imm_size = 0;
359        pending_imms
360    }
361
362    /// Updates the read version with `VersionUpdate`.
363    /// There will be three data types to be processed
364    /// `VersionUpdate::Staging`
365    ///     - `StagingData::ImmMem` -> Insert into memory's `staging_imm`
366    ///     - `StagingData::Sst` -> Update the sst to memory's `staging_sst` and remove the
367    ///       corresponding `staging_imms` according to the `batch_id`
368    /// `VersionUpdate::CommittedDelta` -> Unimplemented yet
369    /// `VersionUpdate::CommittedSnapshot` -> Update `committed_version` , and clean up related
370    /// `staging_sst` and `staging_imm` in memory according to epoch
371    pub fn update(&mut self, info: VersionUpdate) {
372        match info {
373            VersionUpdate::Sst(staging_sst_ref) => {
374                {
375                    assert!(!self.is_replicated);
376                    let Some(imms) = staging_sst_ref.imm_ids.get(&self.instance_id) else {
377                        warn!(
378                            instance_id = self.instance_id,
379                            "no related imm in sst input"
380                        );
381                        return;
382                    };
383
384                    // old data comes first
385                    for imm_id in imms.iter().rev() {
386                        let check_err = match self.staging.uploading_imms.pop_back() {
387                            None => Some("empty".to_owned()),
388                            Some(prev_imm_id) => {
389                                if prev_imm_id.batch_id() == *imm_id {
390                                    None
391                                } else {
392                                    Some(format!(
393                                        "miss match id {} {}",
394                                        prev_imm_id.batch_id(),
395                                        *imm_id
396                                    ))
397                                }
398                            }
399                        };
400                        assert!(
401                            check_err.is_none(),
402                            "should be valid staging_sst.size {},
403                                    staging_sst.imm_ids {:?},
404                                    staging_sst.epochs {:?},
405                                    local_pending_imm_ids {:?},
406                                    local_uploading_imm_ids {:?},
407                                    instance_id {}
408                                    check_err {:?}",
409                            staging_sst_ref.imm_size,
410                            staging_sst_ref.imm_ids,
411                            staging_sst_ref.epochs,
412                            self.staging
413                                .pending_imms
414                                .iter()
415                                .map(|(imm, _)| imm.batch_id())
416                                .collect_vec(),
417                            self.staging
418                                .uploading_imms
419                                .iter()
420                                .map(|imm| imm.batch_id())
421                                .collect_vec(),
422                            self.instance_id,
423                            check_err
424                        );
425                    }
426
427                    self.staging.sst.push_front(staging_sst_ref);
428                }
429            }
430
431            VersionUpdate::CommittedSnapshot(committed_version) => {
432                if let Some(info) = committed_version
433                    .state_table_info
434                    .info()
435                    .get(&self.table_id)
436                {
437                    let committed_epoch = info.committed_epoch;
438                    if self.is_replicated {
439                        self.staging
440                            .uploading_imms
441                            .retain(|imm| imm.epoch() > committed_epoch);
442                        self.staging
443                            .pending_imms
444                            .retain(|(imm, _)| imm.epoch() > committed_epoch);
445                    } else {
446                        self.staging
447                            .pending_imms
448                            .iter()
449                            .map(|(imm, _)| imm)
450                            .chain(self.staging.uploading_imms.iter())
451                            .for_each(|imm| {
452                                assert!(
453                                    imm.epoch() > committed_epoch,
454                                    "imm of table {} min_epoch {} should be greater than committed_epoch {}",
455                                    imm.table_id,
456                                    imm.epoch(),
457                                    committed_epoch
458                                )
459                            });
460                    }
461
462                    self.staging.sst.retain(|sst| {
463                        sst.epochs.first().expect("epochs not empty") > &committed_epoch
464                    });
465
466                    // check epochs.last() > MCE
467                    assert!(self.staging.sst.iter().all(|sst| {
468                        sst.epochs.last().expect("epochs not empty") > &committed_epoch
469                    }));
470
471                    if let Some(committed_watermarks) =
472                        committed_version.table_watermarks.get(&self.table_id)
473                    {
474                        if let Some(watermark_index) = &mut self.table_watermarks {
475                            watermark_index.apply_committed_watermarks(
476                                committed_watermarks.clone(),
477                                committed_epoch,
478                            );
479                        } else {
480                            self.table_watermarks = Some(TableWatermarksIndex::new_committed(
481                                committed_watermarks.clone(),
482                                committed_epoch,
483                                committed_watermarks.watermark_type,
484                            ));
485                        }
486                    }
487                }
488
489                self.committed = committed_version;
490            }
491            VersionUpdate::NewTableWatermark {
492                direction,
493                epoch,
494                vnode_watermarks,
495                watermark_type,
496            } => {
497                if let Some(watermark_index) = &mut self.table_watermarks {
498                    watermark_index.add_epoch_watermark(
499                        epoch,
500                        Arc::from(vnode_watermarks),
501                        direction,
502                    );
503                } else {
504                    self.table_watermarks = Some(TableWatermarksIndex::new(
505                        direction,
506                        epoch,
507                        vnode_watermarks,
508                        self.committed.table_committed_epoch(self.table_id),
509                        watermark_type,
510                    ));
511                }
512            }
513        }
514    }
515
516    pub fn staging(&self) -> &StagingVersion {
517        &self.staging
518    }
519
520    pub fn committed(&self) -> &CommittedVersion {
521        &self.committed
522    }
523
524    /// We have assumption that the watermark is increasing monotonically. Therefore,
525    /// here if the upper layer usage has passed an regressed watermark, we should
526    /// filter out the regressed watermark. Currently the kv log store may write
527    /// regressed watermark
528    pub fn filter_regress_watermarks(&self, watermarks: &mut Vec<VnodeWatermark>) {
529        if let Some(watermark_index) = &self.table_watermarks {
530            watermark_index.filter_regress_watermarks(watermarks)
531        }
532    }
533
534    pub fn latest_watermark(&self, vnode: VirtualNode) -> Option<Bytes> {
535        self.table_watermarks
536            .as_ref()
537            .and_then(|watermark_index| watermark_index.latest_watermark(vnode))
538    }
539
540    pub fn is_initialized(&self) -> bool {
541        self.is_initialized
542    }
543
544    pub fn is_replicated(&self) -> bool {
545        self.is_replicated
546    }
547
548    pub fn update_vnode_bitmap(&mut self, vnodes: Arc<Bitmap>) -> Arc<Bitmap> {
549        std::mem::replace(&mut self.vnodes, vnodes)
550    }
551
552    pub fn contains(&self, vnode: VirtualNode) -> bool {
553        self.vnodes.is_set(vnode.to_index())
554    }
555
556    pub fn vnodes(&self) -> Arc<Bitmap> {
557        self.vnodes.clone()
558    }
559}
560
561pub fn read_filter_for_version(
562    epoch: HummockEpoch,
563    table_id: TableId,
564    mut table_key_range: TableKeyRange,
565    read_version: &RwLock<HummockReadVersion>,
566) -> StorageResult<(TableKeyRange, ReadVersionTuple)> {
567    let read_version_guard = read_version.read();
568
569    let committed_version = read_version_guard.committed().clone();
570
571    if let Some(watermark) = read_version_guard.table_watermarks.as_ref() {
572        watermark.rewrite_range_with_table_watermark(epoch, &mut table_key_range)
573    }
574
575    let (imm_iter, sst_iter) =
576        read_version_guard
577            .staging()
578            .prune_overlap(epoch, table_id, &table_key_range);
579
580    let imms = imm_iter.cloned().collect();
581    let ssts = sst_iter.cloned().collect();
582
583    Ok((table_key_range, (imms, ssts, committed_version)))
584}
585
586#[derive(Clone)]
587pub struct HummockVersionReader {
588    sstable_store: SstableStoreRef,
589
590    /// Statistics
591    state_store_metrics: Arc<HummockStateStoreMetrics>,
592    preload_retry_times: usize,
593}
594
595/// use `HummockVersionReader` to reuse `get` and `iter` implement for both `batch_query` and
596/// `streaming_query`
597impl HummockVersionReader {
598    pub fn new(
599        sstable_store: SstableStoreRef,
600        state_store_metrics: Arc<HummockStateStoreMetrics>,
601        preload_retry_times: usize,
602    ) -> Self {
603        Self {
604            sstable_store,
605            state_store_metrics,
606            preload_retry_times,
607        }
608    }
609
610    pub fn stats(&self) -> &Arc<HummockStateStoreMetrics> {
611        &self.state_store_metrics
612    }
613}
614
615const SLOW_ITER_FETCH_META_DURATION_SECOND: f64 = 5.0;
616
617impl HummockVersionReader {
618    fn skip_get_by_vnode_user_key_range(
619        sstable_info: &SstableInfo,
620        vnode: VirtualNode,
621        user_key: UserKey<&[u8]>,
622        local_stats: &mut StoreLocalStatistic,
623    ) -> bool {
624        if let Some(vnode_statistics) = &sstable_info.vnode_statistics {
625            // Only skip if vnode-statistics exists and key is out of range.
626            // If vnode-statistics not found, it may be due to incomplete stats (reached limit).
627            if let Some((vnode_min, vnode_max)) = vnode_statistics.get_vnode_user_key_range(vnode) {
628                local_stats.vnode_checked_get_count += 1;
629                if user_key < vnode_min.as_ref() || user_key > vnode_max.as_ref() {
630                    local_stats.vnode_pruned_get_count += 1;
631                    return true;
632                }
633            }
634        }
635        false
636    }
637
638    pub async fn get<'a, O>(
639        &'a self,
640        table_key: TableKey<Bytes>,
641        epoch: u64,
642        table_id: TableId,
643        table_option: TableOption,
644        read_options: ReadOptions,
645        read_version_tuple: ReadVersionTuple,
646        on_key_value_fn: impl crate::store::KeyValueFn<'a, O>,
647    ) -> StorageResult<Option<O>> {
648        let (imms, uncommitted_ssts, committed_version) = read_version_tuple;
649
650        let min_epoch = gen_min_epoch(epoch, table_option.retention_seconds);
651        let mut stats_guard = GetLocalMetricsGuard::new(self.state_store_metrics.clone(), table_id);
652        let local_stats = &mut stats_guard.local_stats;
653        local_stats.found_key = true;
654
655        // 1. read staging data
656        for imm in &imms {
657            // skip imm that only holding out-of-date data
658            if imm.epoch() < min_epoch {
659                continue;
660            }
661
662            local_stats.staging_imm_get_count += 1;
663
664            if let Some((data, data_epoch)) = get_from_batch(
665                imm,
666                TableKey(table_key.as_ref()),
667                epoch,
668                &read_options,
669                local_stats,
670            ) {
671                return Ok(if data_epoch.pure_epoch() < min_epoch {
672                    None
673                } else {
674                    data.into_user_value()
675                        .map(|v| {
676                            on_key_value_fn(
677                                FullKey::new_with_gap_epoch(
678                                    table_id,
679                                    table_key.to_ref(),
680                                    data_epoch,
681                                ),
682                                v.as_ref(),
683                            )
684                        })
685                        .transpose()?
686                });
687            }
688        }
689
690        // 2. order guarantee: imm -> sst
691        let dist_key_hash = read_options
692            .prefix_hint
693            .as_ref()
694            .map(|dist_key| Sstable::hash_for_filter(dist_key.as_ref(), table_id.as_raw_id()));
695
696        // Here epoch passed in is pure epoch, and we will seek the constructed `full_key` later.
697        // Therefore, it is necessary to construct the `full_key` with `MAX_SPILL_TIMES`, otherwise, the iterator might skip keys with spill offset greater than 0.
698        let full_key = FullKey::new_with_gap_epoch(
699            table_id,
700            TableKey(table_key.clone()),
701            EpochWithGap::new(epoch, MAX_SPILL_TIMES),
702        );
703        let single_table_key_range = table_key.clone()..=table_key.clone();
704
705        // prune uncommitted ssts with the keyrange
706        let pruned_uncommitted_ssts =
707            prune_overlapping_ssts(&uncommitted_ssts, table_id, &single_table_key_range);
708        for local_sst in pruned_uncommitted_ssts {
709            local_stats.staging_sst_get_count += 1;
710            if let Some(iter) = get_from_sstable_info(
711                self.sstable_store.clone(),
712                local_sst,
713                full_key.to_ref(),
714                &read_options,
715                dist_key_hash,
716                local_stats,
717            )
718            .await?
719            {
720                debug_assert!(iter.is_valid());
721                let data_epoch = iter.key().epoch_with_gap;
722                return Ok(if data_epoch.pure_epoch() < min_epoch {
723                    None
724                } else {
725                    iter.value()
726                        .into_user_value()
727                        .map(|v| {
728                            on_key_value_fn(
729                                FullKey::new_with_gap_epoch(
730                                    table_id,
731                                    table_key.to_ref(),
732                                    data_epoch,
733                                ),
734                                v,
735                            )
736                        })
737                        .transpose()?
738                });
739            }
740        }
741        // 3. read from committed_version sst file
742        // Because SST meta records encoded key range,
743        // the filter key needs to be encoded as well.
744        assert!(committed_version.is_valid());
745        for level in committed_version.levels(table_id) {
746            if level.table_infos.is_empty() {
747                continue;
748            }
749
750            match level.level_type {
751                LevelType::Overlapping | LevelType::Unspecified => {
752                    let sstable_infos = prune_overlapping_ssts(
753                        &level.table_infos,
754                        table_id,
755                        &single_table_key_range,
756                    );
757                    for sstable_info in sstable_infos {
758                        // filter vnode-key range that is definitely not containing the key
759                        if Self::skip_get_by_vnode_user_key_range(
760                            sstable_info,
761                            VirtualNode::from_index(full_key.user_key.get_vnode_id()),
762                            full_key.user_key.as_ref(),
763                            local_stats,
764                        ) {
765                            continue;
766                        }
767
768                        local_stats.overlapping_get_count += 1;
769                        if let Some(iter) = get_from_sstable_info(
770                            self.sstable_store.clone(),
771                            sstable_info,
772                            full_key.to_ref(),
773                            &read_options,
774                            dist_key_hash,
775                            local_stats,
776                        )
777                        .await?
778                        {
779                            debug_assert!(iter.is_valid());
780                            let data_epoch = iter.key().epoch_with_gap;
781                            return Ok(if data_epoch.pure_epoch() < min_epoch {
782                                None
783                            } else {
784                                iter.value()
785                                    .into_user_value()
786                                    .map(|v| {
787                                        on_key_value_fn(
788                                            FullKey::new_with_gap_epoch(
789                                                table_id,
790                                                table_key.to_ref(),
791                                                data_epoch,
792                                            ),
793                                            v,
794                                        )
795                                    })
796                                    .transpose()?
797                            });
798                        }
799                    }
800                }
801                LevelType::Nonoverlapping => {
802                    let mut table_info_idx =
803                        search_sst_idx(&level.table_infos, full_key.user_key.as_ref());
804                    if table_info_idx == 0 {
805                        continue;
806                    }
807                    table_info_idx = table_info_idx.saturating_sub(1);
808                    let sstable_info = &level.table_infos[table_info_idx];
809
810                    if sstable_info.table_ids.binary_search(&table_id).is_err() {
811                        continue;
812                    }
813
814                    // Filter SSTs that definitely cannot contain the key.
815                    let ord = sstable_info
816                        .key_range
817                        .compare_right_with_user_key(full_key.user_key.as_ref());
818                    // the case that the key falls into the gap between two ssts
819                    if ord == Ordering::Less {
820                        sync_point!("HUMMOCK_V2::GET::SKIP_BY_NO_FILE");
821                        continue;
822                    }
823
824                    // Filter vnode-key range that is definitely not containing the key.
825                    if Self::skip_get_by_vnode_user_key_range(
826                        sstable_info,
827                        VirtualNode::from_index(full_key.user_key.get_vnode_id()),
828                        full_key.user_key.as_ref(),
829                        local_stats,
830                    ) {
831                        continue;
832                    }
833
834                    local_stats.non_overlapping_get_count += 1;
835                    if let Some(iter) = get_from_sstable_info(
836                        self.sstable_store.clone(),
837                        sstable_info,
838                        full_key.to_ref(),
839                        &read_options,
840                        dist_key_hash,
841                        local_stats,
842                    )
843                    .await?
844                    {
845                        debug_assert!(iter.is_valid());
846                        let data_epoch = iter.key().epoch_with_gap;
847                        return Ok(if data_epoch.pure_epoch() < min_epoch {
848                            None
849                        } else {
850                            iter.value()
851                                .into_user_value()
852                                .map(|v| {
853                                    on_key_value_fn(
854                                        FullKey::new_with_gap_epoch(
855                                            table_id,
856                                            table_key.to_ref(),
857                                            data_epoch,
858                                        ),
859                                        v,
860                                    )
861                                })
862                                .transpose()?
863                        });
864                    }
865                }
866            }
867        }
868        stats_guard.local_stats.found_key = false;
869        Ok(None)
870    }
871
872    pub async fn iter(
873        &self,
874        table_key_range: TableKeyRange,
875        epoch: u64,
876        table_id: TableId,
877        table_option: TableOption,
878        read_options: ReadOptions,
879        read_version_tuple: (Vec<ImmutableMemtable>, Vec<SstableInfo>, CommittedVersion),
880    ) -> StorageResult<HummockStorageIterator> {
881        self.iter_with_memtable(
882            table_key_range,
883            epoch,
884            table_id,
885            table_option,
886            read_options,
887            read_version_tuple,
888            None,
889        )
890        .await
891    }
892
893    pub async fn iter_with_memtable<'b>(
894        &self,
895        table_key_range: TableKeyRange,
896        epoch: u64,
897        table_id: TableId,
898        table_option: TableOption,
899        read_options: ReadOptions,
900        read_version_tuple: (Vec<ImmutableMemtable>, Vec<SstableInfo>, CommittedVersion),
901        memtable_iter: Option<MemTableHummockIterator<'b>>,
902    ) -> StorageResult<HummockStorageIteratorInner<'b>> {
903        let user_key_range_ref = bound_table_key_range(table_id, &table_key_range);
904        let user_key_range = (
905            user_key_range_ref.0.map(|key| key.cloned()),
906            user_key_range_ref.1.map(|key| key.cloned()),
907        );
908        let mut factory = ForwardIteratorFactory::default();
909        let mut local_stats = StoreLocalStatistic::default();
910        let (imms, uncommitted_ssts, committed) = read_version_tuple;
911        let min_epoch = gen_min_epoch(epoch, table_option.retention_seconds);
912        self.iter_inner(
913            table_key_range,
914            epoch,
915            table_id,
916            read_options,
917            imms,
918            uncommitted_ssts,
919            &committed,
920            &mut local_stats,
921            &mut factory,
922        )
923        .await?;
924        let merge_iter = factory.build(memtable_iter);
925        // the epoch_range left bound for iterator read
926        let mut user_iter = UserIterator::new(
927            merge_iter,
928            user_key_range,
929            epoch,
930            min_epoch,
931            Some(committed),
932        );
933        user_iter.rewind().await?;
934        Ok(HummockStorageIteratorInner::new(
935            user_iter,
936            self.state_store_metrics.clone(),
937            table_id,
938            local_stats,
939        ))
940    }
941
942    pub async fn rev_iter<'b>(
943        &self,
944        table_key_range: TableKeyRange,
945        epoch: u64,
946        table_id: TableId,
947        table_option: TableOption,
948        read_options: ReadOptions,
949        read_version_tuple: (Vec<ImmutableMemtable>, Vec<SstableInfo>, CommittedVersion),
950        memtable_iter: Option<MemTableHummockRevIterator<'b>>,
951    ) -> StorageResult<HummockStorageRevIteratorInner<'b>> {
952        let user_key_range_ref = bound_table_key_range(table_id, &table_key_range);
953        let user_key_range = (
954            user_key_range_ref.0.map(|key| key.cloned()),
955            user_key_range_ref.1.map(|key| key.cloned()),
956        );
957        let mut factory = BackwardIteratorFactory::default();
958        let mut local_stats = StoreLocalStatistic::default();
959        let (imms, uncommitted_ssts, committed) = read_version_tuple;
960        let min_epoch = gen_min_epoch(epoch, table_option.retention_seconds);
961        self.iter_inner(
962            table_key_range,
963            epoch,
964            table_id,
965            read_options,
966            imms,
967            uncommitted_ssts,
968            &committed,
969            &mut local_stats,
970            &mut factory,
971        )
972        .await?;
973        let merge_iter = factory.build(memtable_iter);
974        // the epoch_range left bound for iterator read
975        let mut user_iter = BackwardUserIterator::new(
976            merge_iter,
977            user_key_range,
978            epoch,
979            min_epoch,
980            Some(committed),
981        );
982        user_iter.rewind().await?;
983        Ok(HummockStorageRevIteratorInner::new(
984            user_iter,
985            self.state_store_metrics.clone(),
986            table_id,
987            local_stats,
988        ))
989    }
990
991    async fn iter_inner<F: IteratorFactory>(
992        &self,
993        table_key_range: TableKeyRange,
994        epoch: u64,
995        table_id: TableId,
996        read_options: ReadOptions,
997        imms: Vec<ImmutableMemtable>,
998        uncommitted_ssts: Vec<SstableInfo>,
999        committed: &CommittedVersion,
1000        local_stats: &mut StoreLocalStatistic,
1001        factory: &mut F,
1002    ) -> StorageResult<()> {
1003        {
1004            fn bound_inner<T>(bound: &Bound<T>) -> Option<&T> {
1005                match bound {
1006                    Bound::Included(bound) | Bound::Excluded(bound) => Some(bound),
1007                    Bound::Unbounded => None,
1008                }
1009            }
1010            let (left, right) = &table_key_range;
1011            if let (Some(left), Some(right)) = (bound_inner(left), bound_inner(right))
1012                && right < left
1013            {
1014                if cfg!(debug_assertions) {
1015                    panic!("invalid iter key range: {table_id} {left:?} {right:?}")
1016                } else {
1017                    return Err(HummockError::other(format!(
1018                        "invalid iter key range: {table_id} {left:?} {right:?}"
1019                    ))
1020                    .into());
1021                }
1022            }
1023        }
1024
1025        local_stats.staging_imm_iter_count = imms.len() as u64;
1026        for imm in imms {
1027            factory.add_batch_iter(imm);
1028        }
1029
1030        // 2. build iterator from committed
1031        // Because SST meta records encoded key range,
1032        // the filter key range needs to be encoded as well.
1033        let user_key_range = bound_table_key_range(table_id, &table_key_range);
1034        let user_key_range_ref = (
1035            user_key_range.0.as_ref().map(UserKey::as_ref),
1036            user_key_range.1.as_ref().map(UserKey::as_ref),
1037        );
1038        let mut staging_sst_iter_count = 0;
1039        // encode once
1040        let filter_prefix_hash = read_options
1041            .prefix_hint
1042            .as_ref()
1043            .map(|hint| Sstable::hash_for_filter(hint, table_id.as_raw_id()));
1044        let mut sst_read_options = SstableIteratorReadOptions::from_read_options(&read_options);
1045        sst_read_options.read_table_id = Some(table_id);
1046        sst_read_options.scan_end_user_key = Some(user_key_range.1.map(|key| key.cloned()));
1047        sst_read_options.prefetch = read_options.prefetch_options.prefetch;
1048        if sst_read_options.prefetch {
1049            sst_read_options.max_preload_retry_times = self.preload_retry_times;
1050        }
1051        let sst_read_options = Arc::new(sst_read_options);
1052        for sstable_info in &uncommitted_ssts {
1053            let table_holder = self
1054                .sstable_store
1055                .sstable(sstable_info, local_stats)
1056                .await?;
1057
1058            if let Some(prefix_hash) = filter_prefix_hash.as_ref()
1059                && !hit_sstable_filter(
1060                    &table_holder,
1061                    &user_key_range_ref,
1062                    *prefix_hash,
1063                    local_stats,
1064                )
1065            {
1066                continue;
1067            }
1068
1069            staging_sst_iter_count += 1;
1070            factory.add_staging_sst_iter(F::SstableIteratorType::create(
1071                table_holder,
1072                self.sstable_store.clone(),
1073                sst_read_options.clone(),
1074                sstable_info,
1075            ));
1076        }
1077        local_stats.staging_sst_iter_count = staging_sst_iter_count;
1078
1079        let timer = Instant::now();
1080
1081        for level in committed.levels(table_id) {
1082            if level.table_infos.is_empty() {
1083                continue;
1084            }
1085
1086            if level.level_type == LevelType::Nonoverlapping {
1087                let mut table_infos =
1088                    prune_nonoverlapping_ssts(&level.table_infos, user_key_range_ref, table_id)
1089                        .peekable();
1090
1091                if table_infos.peek().is_none() {
1092                    continue;
1093                }
1094                let sstable_infos = table_infos.cloned().collect_vec();
1095                if sstable_infos.len() > 1 {
1096                    factory.add_concat_sst_iter(
1097                        sstable_infos,
1098                        self.sstable_store.clone(),
1099                        sst_read_options.clone(),
1100                    );
1101                    local_stats.non_overlapping_iter_count += 1;
1102                } else {
1103                    let sstable_info = &sstable_infos[0];
1104
1105                    let sstable = self
1106                        .sstable_store
1107                        .sstable(sstable_info, local_stats)
1108                        .await?;
1109
1110                    if let Some(dist_hash) = filter_prefix_hash.as_ref()
1111                        && !hit_sstable_filter(
1112                            &sstable,
1113                            &user_key_range_ref,
1114                            *dist_hash,
1115                            local_stats,
1116                        )
1117                    {
1118                        continue;
1119                    }
1120                    // Since there is only one sst to be included for the current non-overlapping
1121                    // level, there is no need to create a ConcatIterator on it.
1122                    // We put the SstableIterator in `overlapping_iters` just for convenience since
1123                    // it overlaps with SSTs in other levels. In metrics reporting, we still count
1124                    // it in `non_overlapping_iter_count`.
1125                    factory.add_overlapping_sst_iter(F::SstableIteratorType::create(
1126                        sstable,
1127                        self.sstable_store.clone(),
1128                        sst_read_options.clone(),
1129                        sstable_info,
1130                    ));
1131                    local_stats.non_overlapping_iter_count += 1;
1132                }
1133            } else {
1134                let table_infos =
1135                    prune_overlapping_ssts(&level.table_infos, table_id, &table_key_range);
1136                // Overlapping
1137                let fetch_meta_req = table_infos.rev().collect_vec();
1138                if fetch_meta_req.is_empty() {
1139                    continue;
1140                }
1141                for sstable_info in fetch_meta_req {
1142                    let sstable = self
1143                        .sstable_store
1144                        .sstable(sstable_info, local_stats)
1145                        .await?;
1146                    assert_eq!(sstable_info.object_id, sstable.id);
1147                    if let Some(dist_hash) = filter_prefix_hash.as_ref()
1148                        && !hit_sstable_filter(
1149                            &sstable,
1150                            &user_key_range_ref,
1151                            *dist_hash,
1152                            local_stats,
1153                        )
1154                    {
1155                        continue;
1156                    }
1157                    factory.add_overlapping_sst_iter(F::SstableIteratorType::create(
1158                        sstable,
1159                        self.sstable_store.clone(),
1160                        sst_read_options.clone(),
1161                        sstable_info,
1162                    ));
1163                    local_stats.overlapping_iter_count += 1;
1164                }
1165            }
1166        }
1167        let fetch_meta_duration_sec = timer.elapsed().as_secs_f64();
1168        if fetch_meta_duration_sec > SLOW_ITER_FETCH_META_DURATION_SECOND {
1169            let table_id_string = table_id.to_string();
1170            tracing::warn!(
1171                "Fetching meta while creating an iter to read table_id {:?} at epoch {:?} is slow: duration = {:?}s, cache unhits = {:?}.",
1172                table_id_string,
1173                epoch,
1174                fetch_meta_duration_sec,
1175                local_stats.cache_meta_block_miss
1176            );
1177            self.state_store_metrics
1178                .iter_slow_fetch_meta_cache_unhits
1179                .set(local_stats.cache_meta_block_miss as i64);
1180        }
1181        Ok(())
1182    }
1183
1184    pub async fn iter_log(
1185        &self,
1186        epoch_range: (u64, u64),
1187        key_range: TableKeyRange,
1188        options: ReadLogOptions,
1189        table_change_log_manager: Arc<TableChangeLogManager>,
1190    ) -> HummockResult<ChangeLogIterator> {
1191        // The end value of `epoch_range` is not greater than max committed epoch, guaranteed by the caller `BatchTableInnerIterLogInner`.
1192        let change_log: Vec<_> = {
1193            let table_change_logs = table_change_log_manager
1194                .fetch_table_change_logs(options.table_id, epoch_range, false, None)
1195                .await?;
1196            if let Some(change_log) = table_change_logs.get(&options.table_id) {
1197                change_log.filter_epoch(epoch_range).cloned().collect_vec()
1198            } else {
1199                Vec::new()
1200            }
1201        };
1202
1203        if let Some(max_epoch_change_log) = change_log.last() {
1204            let (_, max_epoch) = epoch_range;
1205            if !max_epoch_change_log.epochs().contains(&max_epoch) {
1206                warn!(
1207                    max_epoch,
1208                    change_log_epochs = ?change_log.iter().flat_map(|epoch_log| epoch_log.epochs()).collect_vec(),
1209                    table_id = %options.table_id,
1210                    "max_epoch does not exist"
1211                );
1212            }
1213        }
1214        let read_options = Arc::new(SstableIteratorReadOptions {
1215            cache_policy: Default::default(),
1216            read_table_id: Some(options.table_id),
1217            scan_end_user_key: None,
1218            prefetch: false,
1219            max_preload_retry_times: 0,
1220        });
1221
1222        async fn make_iter(
1223            sstable_infos: impl Iterator<Item = &SstableInfo>,
1224            sstable_store: &SstableStoreRef,
1225            read_options: Arc<SstableIteratorReadOptions>,
1226            local_stat: &mut StoreLocalStatistic,
1227        ) -> HummockResult<MergeIterator<SstableIterator>> {
1228            let iters = try_join_all(sstable_infos.map(|sstable_info| {
1229                let sstable_store = sstable_store.clone();
1230                let read_options = read_options.clone();
1231                async move {
1232                    let mut local_stat = StoreLocalStatistic::default();
1233                    let table_holder = sstable_store.sstable(sstable_info, &mut local_stat).await?;
1234                    Ok::<_, HummockError>((
1235                        SstableIterator::new(
1236                            table_holder,
1237                            sstable_store,
1238                            read_options,
1239                            sstable_info,
1240                        ),
1241                        local_stat,
1242                    ))
1243                }
1244            }))
1245            .await?;
1246            Ok::<_, HummockError>(MergeIterator::new(iters.into_iter().map(
1247                |(iter, stats)| {
1248                    local_stat.add(&stats);
1249                    iter
1250                },
1251            )))
1252        }
1253
1254        let mut local_stat = StoreLocalStatistic::default();
1255
1256        let new_value_iter = make_iter(
1257            change_log
1258                .iter()
1259                .flat_map(|log| log.new_value.iter())
1260                .filter(|sst| filter_single_sst(sst, options.table_id, &key_range)),
1261            &self.sstable_store,
1262            read_options.clone(),
1263            &mut local_stat,
1264        )
1265        .await?;
1266        let old_value_iter = make_iter(
1267            change_log
1268                .iter()
1269                .flat_map(|log| log.old_value.iter())
1270                .filter(|sst| filter_single_sst(sst, options.table_id, &key_range)),
1271            &self.sstable_store,
1272            read_options.clone(),
1273            &mut local_stat,
1274        )
1275        .await?;
1276        ChangeLogIterator::new(
1277            epoch_range,
1278            key_range,
1279            new_value_iter,
1280            old_value_iter,
1281            options.table_id,
1282            IterLocalMetricsGuard::new(
1283                self.state_store_metrics.clone(),
1284                options.table_id,
1285                local_stat,
1286            ),
1287        )
1288        .await
1289    }
1290
1291    pub async fn nearest<'a, M: MeasureDistanceBuilder, O: Send>(
1292        &'a self,
1293        version: PinnedVersion,
1294        table_id: TableId,
1295        target: VectorRef<'a>,
1296        options: VectorNearestOptions,
1297        on_nearest_item_fn: impl OnNearestItemFn<'a, O>,
1298    ) -> HummockResult<Vec<O>> {
1299        let Some(index) = version.vector_indexes.get(&table_id) else {
1300            return Ok(vec![]);
1301        };
1302        if target.dimension() != index.dimension {
1303            return Err(HummockError::other(format!(
1304                "target dimension {} not match index dimension {}",
1305                target.dimension(),
1306                index.dimension
1307            )));
1308        }
1309        match &index.inner {
1310            VectorIndexImpl::Flat(flat) => {
1311                let mut builder = NearestBuilder::<'_, O, M>::new(target, options.top_n);
1312                let mut cache_stat = VectorStoreCacheStats::default();
1313                for vector_file in &flat.vector_store_info.vector_files {
1314                    let meta = self
1315                        .sstable_store
1316                        .get_vector_file_meta(vector_file, &mut cache_stat)
1317                        .await?;
1318                    for (i, block_meta) in meta.block_metas.iter().enumerate() {
1319                        let block = self
1320                            .sstable_store
1321                            .get_vector_block(vector_file, i, block_meta, &mut cache_stat)
1322                            .await?;
1323                        builder.add(&**block, &on_nearest_item_fn);
1324                    }
1325                }
1326                cache_stat.report(table_id, "flat", self.stats());
1327                Ok(builder.finish())
1328            }
1329            VectorIndexImpl::HnswFlat(hnsw_flat) => {
1330                let Some(graph_file) = &hnsw_flat.graph_file else {
1331                    return Ok(vec![]);
1332                };
1333
1334                let mut ctx = FileVectorStoreCtx::default();
1335
1336                let graph = self
1337                    .sstable_store
1338                    .get_hnsw_graph(graph_file, &mut ctx.stats)
1339                    .await?;
1340
1341                let vector_store =
1342                    FileVectorStore::new_for_reader(hnsw_flat, self.sstable_store.clone());
1343                let (items, stats) = nearest::<O, M, _>(
1344                    &vector_store,
1345                    &mut ctx,
1346                    &*graph,
1347                    target,
1348                    on_nearest_item_fn,
1349                    options.hnsw_ef_search,
1350                    options.top_n,
1351                )
1352                .await?;
1353                ctx.stats.report(table_id, "hnsw_read", self.stats());
1354                report_hnsw_stat(
1355                    self.stats(),
1356                    table_id,
1357                    "hnsw_read",
1358                    options.top_n,
1359                    options.hnsw_ef_search,
1360                    [stats],
1361                );
1362                Ok(items)
1363            }
1364        }
1365    }
1366}
1367
1368#[cfg(test)]
1369mod tests {
1370    use std::collections::{BTreeMap, HashMap, HashSet};
1371    use std::sync::Arc;
1372
1373    use bytes::Bytes;
1374    use prometheus::Registry;
1375    use risingwave_common::catalog::{TableId, TableOption};
1376    use risingwave_common::config::MetricLevel;
1377    use risingwave_common::hash::VirtualNode;
1378    use risingwave_common::util::epoch::test_epoch;
1379    use risingwave_hummock_sdk::compaction_group::StaticCompactionGroupId;
1380    use risingwave_hummock_sdk::key::{FullKey, TableKey, UserKey, gen_key_from_bytes};
1381    use risingwave_hummock_sdk::key_range::KeyRange;
1382    use risingwave_hummock_sdk::level::{Level, Levels, OverlappingLevel};
1383    use risingwave_hummock_sdk::sstable_info::{SstableInfo, SstableInfoInner, VnodeStatistics};
1384    use risingwave_hummock_sdk::version::HummockVersion;
1385    use risingwave_hummock_sdk::{EpochWithGap, HummockSstableObjectId};
1386    use risingwave_pb::hummock::hummock_version::PbLevels;
1387    use risingwave_pb::hummock::{
1388        LevelType as PbLevelType, PbHummockVersion, PbLevel, PbOverlappingLevel,
1389        PbSstableFilterLayout, PbSstableFilterType, PbStateTableInfo, StateTableInfoDelta,
1390    };
1391    use tokio::sync::mpsc::unbounded_channel;
1392
1393    use crate::hummock::HummockValue;
1394    use crate::hummock::iterator::test_utils::mock_sstable_store;
1395    use crate::hummock::local_version::pinned_version::{PinVersionAction, PinnedVersion};
1396    use crate::hummock::store::version::{CommittedVersion, HummockVersionReader};
1397    use crate::hummock::test_utils::{
1398        default_builder_opt_for_test, gen_test_sstable_with_table_ids,
1399    };
1400    use crate::monitor::{HummockStateStoreMetrics, flush_local_metrics_for_test};
1401    use crate::store::ReadOptions;
1402
1403    /// In a nonoverlapping level, `search_sst_idx` may locate an SST whose key range covers
1404    /// the query user key but whose `table_ids` does not contain the queried table. The get
1405    /// path must skip such SSTs instead of reading them.
1406    #[tokio::test]
1407    async fn test_get_skips_sst_by_table_id_filter() {
1408        let query_table_id = TableId::new(100);
1409        let epoch: u64 = (31 * 1000) << 16;
1410        let compaction_group_id = StaticCompactionGroupId::StateDefault;
1411
1412        // SST key range: table_id 50..150, but table_ids = [50, 150] (no 100).
1413        let sst_info = SstableInfoInner {
1414            sst_id: 1.into(),
1415            object_id: 1.into(),
1416            key_range: KeyRange {
1417                left: Bytes::from(
1418                    FullKey::for_test(TableId::new(50), b"aaa".to_vec(), epoch).encode(),
1419                ),
1420                right: Bytes::from(
1421                    FullKey::for_test(TableId::new(150), b"zzz".to_vec(), epoch).encode(),
1422                ),
1423                right_exclusive: false,
1424            },
1425            table_ids: vec![TableId::new(50), TableId::new(150)],
1426            file_size: 1024,
1427            ..Default::default()
1428        }
1429        .into();
1430
1431        let level = Level {
1432            level_idx: 1,
1433            level_type: PbLevelType::Nonoverlapping,
1434            table_infos: vec![sst_info],
1435            total_file_size: 0,
1436            sub_level_id: 0,
1437            uncompressed_file_size: 0,
1438            vnode_partition_count: 0,
1439        };
1440
1441        #[allow(deprecated)]
1442        let levels = Levels {
1443            levels: vec![level],
1444            l0: OverlappingLevel::default(),
1445            group_id: compaction_group_id,
1446            parent_group_id: compaction_group_id,
1447            member_table_ids: vec![],
1448            compaction_group_version_id: 0,
1449        };
1450
1451        let mut version = HummockVersion::from_persisted_protobuf_owned(PbHummockVersion {
1452            id: 1u64.into(),
1453            ..Default::default()
1454        });
1455        version.levels.insert(compaction_group_id, levels);
1456        version.state_table_info.apply_delta(
1457            &HashMap::from([(
1458                query_table_id,
1459                StateTableInfoDelta {
1460                    committed_epoch: epoch,
1461                    compaction_group_id,
1462                },
1463            )]),
1464            &HashSet::new(),
1465        );
1466
1467        let pinned_version = PinnedVersion::new(version, unbounded_channel().0);
1468        let reader = HummockVersionReader::new(
1469            mock_sstable_store().await,
1470            Arc::new(HummockStateStoreMetrics::unused()),
1471            0,
1472        );
1473
1474        let result = reader
1475            .get(
1476                TableKey(Bytes::from("test_key")),
1477                epoch,
1478                query_table_id,
1479                TableOption::default(),
1480                ReadOptions::default(),
1481                (vec![], vec![], pinned_version),
1482                |_key, _value| Ok(()),
1483            )
1484            .await
1485            .unwrap();
1486
1487        assert!(result.is_none());
1488    }
1489
1490    /// Build a committed version containing a single SST with custom vnode stats.
1491    #[allow(deprecated)]
1492    fn build_version_with_vnode_stats(
1493        table_id: TableId,
1494        vnode_stats: VnodeStatistics,
1495        key_range: (Vec<u8>, Vec<u8>),
1496        level_type: PbLevelType,
1497    ) -> (SstableInfo, CommittedVersion) {
1498        let object_id = HummockSstableObjectId::new(1);
1499        let left_full_key = FullKey::new_with_gap_epoch(
1500            table_id,
1501            TableKey(Bytes::from(key_range.0)),
1502            EpochWithGap::new_from_epoch(test_epoch(0)),
1503        )
1504        .encode();
1505        let right_full_key = FullKey::new_with_gap_epoch(
1506            table_id,
1507            TableKey(Bytes::from(key_range.1)),
1508            EpochWithGap::new_from_epoch(test_epoch(0)),
1509        )
1510        .encode();
1511
1512        let sstable_info: SstableInfo = SstableInfoInner {
1513            object_id,
1514            sst_id: object_id.as_raw_id().into(),
1515            key_range: KeyRange {
1516                left: Bytes::from(left_full_key),
1517                right: Bytes::from(right_full_key),
1518                right_exclusive: false,
1519            },
1520            file_size: 1,
1521            table_ids: vec![table_id],
1522            meta_offset: 0,
1523            stale_key_count: 0,
1524            total_key_count: 0,
1525            min_epoch: 0,
1526            max_epoch: 0,
1527            uncompressed_file_size: 0,
1528            range_tombstone_count: 0,
1529            filter_type: PbSstableFilterType::SstableFilterXor16,
1530            filter_layout: PbSstableFilterLayout::Plain,
1531            sst_size: 1,
1532            vnode_statistics: Some(vnode_stats),
1533        }
1534        .into();
1535        let pb_level = PbLevel {
1536            level_idx: if level_type == PbLevelType::Overlapping {
1537                0
1538            } else {
1539                1
1540            },
1541            level_type: level_type as i32,
1542            table_infos: vec![sstable_info.clone().into()],
1543            total_file_size: 1,
1544            sub_level_id: 0,
1545            uncompressed_file_size: 1,
1546            vnode_partition_count: 0,
1547        };
1548
1549        let (levels, l0) = if level_type == PbLevelType::Overlapping {
1550            (
1551                vec![],
1552                Some(PbOverlappingLevel {
1553                    sub_levels: vec![pb_level],
1554                    total_file_size: 1,
1555                    uncompressed_file_size: 1,
1556                }),
1557            )
1558        } else {
1559            (vec![pb_level], Some(PbOverlappingLevel::default()))
1560        };
1561
1562        let pb_levels = PbLevels {
1563            levels,
1564            l0,
1565            group_id: StaticCompactionGroupId::NewCompactionGroup,
1566            parent_group_id: 0.into(),
1567            member_table_ids: vec![],
1568            compaction_group_version_id: 0,
1569        };
1570
1571        let pb_version = PbHummockVersion {
1572            id: 1.into(),
1573            levels: HashMap::from_iter([(StaticCompactionGroupId::NewCompactionGroup, pb_levels)]),
1574            max_committed_epoch: 0,
1575            table_watermarks: HashMap::new(),
1576            table_change_logs: HashMap::new(),
1577            state_table_info: HashMap::from_iter([(
1578                table_id,
1579                PbStateTableInfo {
1580                    committed_epoch: 0,
1581                    compaction_group_id: StaticCompactionGroupId::NewCompactionGroup,
1582                },
1583            )]),
1584            vector_indexes: HashMap::new(),
1585        };
1586
1587        let version = HummockVersion::from(&pb_version);
1588        let (tx, _rx) = unbounded_channel::<PinVersionAction>();
1589        let pinned = PinnedVersion::new(version, tx);
1590        (sstable_info, pinned)
1591    }
1592
1593    /// Build a committed version from an existing SST (with real object in the store).
1594    #[allow(deprecated)]
1595    fn build_version_from_sstables(
1596        table_id: TableId,
1597        sstable_infos: Vec<SstableInfo>,
1598        level_type: PbLevelType,
1599    ) -> CommittedVersion {
1600        let total_file_size = sstable_infos.iter().map(|sst| sst.file_size).sum::<u64>();
1601        let uncompressed_file_size = sstable_infos
1602            .iter()
1603            .map(|sst| sst.uncompressed_file_size)
1604            .sum::<u64>();
1605        let pb_level = PbLevel {
1606            level_idx: if level_type == PbLevelType::Overlapping {
1607                0
1608            } else {
1609                1
1610            },
1611            level_type: level_type as i32,
1612            table_infos: sstable_infos.into_iter().map(Into::into).collect(),
1613            total_file_size,
1614            sub_level_id: 0,
1615            uncompressed_file_size,
1616            vnode_partition_count: 0,
1617        };
1618
1619        let (levels, l0) = if level_type == PbLevelType::Overlapping {
1620            (
1621                vec![],
1622                Some(PbOverlappingLevel {
1623                    sub_levels: vec![pb_level],
1624                    total_file_size,
1625                    uncompressed_file_size,
1626                }),
1627            )
1628        } else {
1629            (vec![pb_level], Some(PbOverlappingLevel::default()))
1630        };
1631
1632        let pb_levels = PbLevels {
1633            levels,
1634            l0,
1635            group_id: StaticCompactionGroupId::NewCompactionGroup,
1636            parent_group_id: 0.into(),
1637            member_table_ids: vec![],
1638            compaction_group_version_id: 0,
1639        };
1640
1641        let pb_version = PbHummockVersion {
1642            id: 1.into(),
1643            levels: HashMap::from_iter([(StaticCompactionGroupId::NewCompactionGroup, pb_levels)]),
1644            max_committed_epoch: 0,
1645            table_watermarks: HashMap::new(),
1646            table_change_logs: HashMap::new(),
1647            state_table_info: HashMap::from_iter([(
1648                table_id,
1649                PbStateTableInfo {
1650                    committed_epoch: 0,
1651                    compaction_group_id: StaticCompactionGroupId::NewCompactionGroup,
1652                },
1653            )]),
1654            vector_indexes: HashMap::new(),
1655        };
1656
1657        let version = HummockVersion::from(&pb_version);
1658        let (tx, _rx) = unbounded_channel::<PinVersionAction>();
1659        PinnedVersion::new(version, tx)
1660    }
1661
1662    /// Build a committed version from one existing non-overlapping SST.
1663    #[allow(deprecated)]
1664    fn build_version_from_sstable(
1665        table_id: TableId,
1666        sstable_info: SstableInfo,
1667    ) -> CommittedVersion {
1668        build_version_from_sstables(table_id, vec![sstable_info], PbLevelType::Nonoverlapping)
1669    }
1670
1671    fn vnode_prune_counts(
1672        metrics: &HummockStateStoreMetrics,
1673        table_id: TableId,
1674        operation: &str,
1675    ) -> (u64, u64) {
1676        let table_label = table_id.to_string();
1677        let checked = metrics
1678            .vnode_pruning_counts
1679            .with_guarded_label_values(&[
1680                table_label.clone(),
1681                operation.to_owned(),
1682                "checked".to_owned(),
1683            ])
1684            .get();
1685        let pruned = metrics
1686            .vnode_pruning_counts
1687            .with_guarded_label_values(&[table_label, operation.to_owned(), "pruned".to_owned()])
1688            .get();
1689        (checked, pruned)
1690    }
1691
1692    async fn assert_vnode_prune_get_skips_out_of_range_key(
1693        table_id: TableId,
1694        epoch: u64,
1695        level_type: PbLevelType,
1696    ) {
1697        let sstable_store = mock_sstable_store().await;
1698        let registry = Registry::new();
1699        let metrics = Arc::new(HummockStateStoreMetrics::new(&registry, MetricLevel::Debug));
1700        let reader = HummockVersionReader::new(sstable_store, metrics.clone(), 0);
1701        let (checked_before, pruned_before) = vnode_prune_counts(&metrics, table_id, "get");
1702
1703        let make_user_key = |vnode: VirtualNode, suffix: &str| {
1704            let mut raw = vnode.to_be_bytes().to_vec();
1705            raw.extend_from_slice(suffix.as_bytes());
1706            UserKey::new(table_id, TableKey(raw.into()))
1707        };
1708
1709        // Stats cover vnode 1 only up to "bb".
1710        let vnode_stats = VnodeStatistics::from_map(BTreeMap::from_iter([(
1711            VirtualNode::from_index(1),
1712            (
1713                make_user_key(VirtualNode::from_index(1), "aa"),
1714                make_user_key(VirtualNode::from_index(1), "bb"),
1715            ),
1716        )]));
1717
1718        // SST key range is wide enough to include the queried key, but vnode stats should prune it.
1719        let key_range = {
1720            let mut left = VirtualNode::from_index(0).to_be_bytes().to_vec();
1721            left.extend_from_slice(b"aa");
1722            let mut right = VirtualNode::from_index(1).to_be_bytes().to_vec();
1723            right.extend_from_slice(b"zzzz");
1724            (left, right)
1725        };
1726
1727        let (_sst, committed) =
1728            build_version_with_vnode_stats(table_id, vnode_stats, key_range, level_type);
1729
1730        // Query vnode 1 but with suffix beyond the recorded max -> should be pruned.
1731        let mut raw = VirtualNode::from_index(1).to_be_bytes().to_vec();
1732        raw.extend_from_slice(b"zz");
1733        let table_key = TableKey(Bytes::from(raw.clone()));
1734
1735        let result = reader
1736            .get(
1737                table_key,
1738                epoch,
1739                table_id,
1740                TableOption::default(),
1741                ReadOptions::default(),
1742                (vec![], vec![], committed),
1743                |_k, v| Ok(Bytes::copy_from_slice(v)),
1744            )
1745            .await
1746            .unwrap();
1747        flush_local_metrics_for_test();
1748
1749        assert!(
1750            result.is_none(),
1751            "vnode pruning should skip SST without reading data"
1752        );
1753        let (checked_after, pruned_after) = vnode_prune_counts(&metrics, table_id, "get");
1754        assert_eq!(checked_before + 1, checked_after);
1755        assert_eq!(pruned_before + 1, pruned_after);
1756    }
1757
1758    async fn assert_vnode_prune_get_not_pruned_nonoverlapping() {
1759        let table_id = TableId::new(42);
1760        let epoch = test_epoch(3);
1761        let sstable_store = mock_sstable_store().await;
1762        let registry = Registry::new();
1763        let metrics = Arc::new(HummockStateStoreMetrics::new(&registry, MetricLevel::Debug));
1764        let reader = HummockVersionReader::new(sstable_store.clone(), metrics.clone(), 0);
1765        let (checked_before, pruned_before) = vnode_prune_counts(&metrics, table_id, "get");
1766
1767        let mut opts = default_builder_opt_for_test();
1768        opts.max_vnode_key_range_bytes = None;
1769        let mut kvs = vec![
1770            (
1771                FullKey::new_with_gap_epoch(
1772                    table_id,
1773                    gen_key_from_bytes(VirtualNode::from_index(1), b"aa"),
1774                    EpochWithGap::new_from_epoch(epoch),
1775                ),
1776                HummockValue::put(Bytes::from_static(b"v1")),
1777            ),
1778            (
1779                FullKey::new_with_gap_epoch(
1780                    table_id,
1781                    gen_key_from_bytes(VirtualNode::ZERO, b"cc"),
1782                    EpochWithGap::new_from_epoch(epoch),
1783                ),
1784                HummockValue::put(Bytes::from_static(b"v0")),
1785            ),
1786        ];
1787        kvs.sort_by(|(k1, _), (k2, _)| k1.cmp(k2));
1788        let (_, mut sstable_info): (crate::hummock::sstable_store::TableHolder, SstableInfo) =
1789            gen_test_sstable_with_table_ids(
1790                opts,
1791                10,
1792                kvs.into_iter(),
1793                sstable_store.clone(),
1794                vec![table_id.as_raw_id()],
1795            )
1796            .await;
1797        // Override vnode stats to ensure the queried key falls inside the recorded range.
1798        let mut inner = sstable_info.get_inner();
1799        inner.vnode_statistics = Some(VnodeStatistics::from_map(BTreeMap::from_iter([(
1800            VirtualNode::from_index(1),
1801            (
1802                UserKey::new(
1803                    table_id,
1804                    gen_key_from_bytes(VirtualNode::from_index(1), b"aa"),
1805                ),
1806                UserKey::new(
1807                    table_id,
1808                    gen_key_from_bytes(VirtualNode::from_index(1), b"zz"),
1809                ),
1810            ),
1811        )])));
1812        sstable_info = inner.into();
1813        let committed = build_version_from_sstable(table_id, sstable_info.clone());
1814
1815        // Key is within vnode range, should not be pruned.
1816        let mut raw = VirtualNode::from_index(1).to_be_bytes().to_vec();
1817        raw.extend_from_slice(b"aa");
1818        let table_key = TableKey(Bytes::from(raw.clone()));
1819
1820        let result = reader
1821            .get(
1822                table_key,
1823                epoch,
1824                table_id,
1825                TableOption::default(),
1826                ReadOptions::default(),
1827                (vec![], vec![], committed),
1828                |_k, v| Ok(Bytes::copy_from_slice(v)),
1829            )
1830            .await
1831            .unwrap();
1832        flush_local_metrics_for_test();
1833        assert!(result.is_some(), "key should be read when not pruned");
1834        let (checked_after, pruned_after) = vnode_prune_counts(&metrics, table_id, "get");
1835        assert_eq!(checked_before + 1, checked_after);
1836        assert_eq!(pruned_before, pruned_after);
1837    }
1838
1839    #[tokio::test]
1840    async fn test_vnode_prune_get_single_sst_cases() {
1841        assert_vnode_prune_get_skips_out_of_range_key(
1842            TableId::default(),
1843            test_epoch(1),
1844            PbLevelType::Nonoverlapping,
1845        )
1846        .await;
1847        assert_vnode_prune_get_skips_out_of_range_key(
1848            TableId::new(7),
1849            test_epoch(2),
1850            PbLevelType::Overlapping,
1851        )
1852        .await;
1853        assert_vnode_prune_get_not_pruned_nonoverlapping().await;
1854    }
1855
1856    #[tokio::test]
1857    async fn test_vnode_prune_get_overlapping_distribution_prunes_only_out_of_range_sst() {
1858        let table_id = TableId::new(77);
1859        let epoch = test_epoch(4);
1860        let vnode = VirtualNode::from_index(1);
1861        let sstable_store = mock_sstable_store().await;
1862        let registry = Registry::new();
1863        let metrics = Arc::new(HummockStateStoreMetrics::new(&registry, MetricLevel::Debug));
1864        let reader = HummockVersionReader::new(sstable_store.clone(), metrics.clone(), 0);
1865        let (checked_before, pruned_before) = vnode_prune_counts(&metrics, table_id, "get");
1866
1867        let mut opts = default_builder_opt_for_test();
1868        opts.max_vnode_key_range_bytes = None;
1869
1870        let mut kvs1 = vec![
1871            (
1872                FullKey::new_with_gap_epoch(
1873                    table_id,
1874                    gen_key_from_bytes(vnode, b"aa"),
1875                    EpochWithGap::new_from_epoch(epoch),
1876                ),
1877                HummockValue::put(Bytes::from_static(b"s1_aa")),
1878            ),
1879            (
1880                FullKey::new_with_gap_epoch(
1881                    table_id,
1882                    gen_key_from_bytes(vnode, b"zz"),
1883                    EpochWithGap::new_from_epoch(epoch),
1884                ),
1885                HummockValue::put(Bytes::from_static(b"s1_zz")),
1886            ),
1887        ];
1888        kvs1.sort_by(|(k1, _), (k2, _)| k1.cmp(k2));
1889        let (_, mut sst1): (crate::hummock::sstable_store::TableHolder, SstableInfo) =
1890            gen_test_sstable_with_table_ids(
1891                opts.clone(),
1892                11,
1893                kvs1.into_iter(),
1894                sstable_store.clone(),
1895                vec![table_id.as_raw_id()],
1896            )
1897            .await;
1898        let mut sst1_inner = sst1.get_inner();
1899        sst1_inner.vnode_statistics = Some(VnodeStatistics::from_map(BTreeMap::from_iter([(
1900            vnode,
1901            (
1902                UserKey::new(table_id, gen_key_from_bytes(vnode, b"aa")),
1903                UserKey::new(table_id, gen_key_from_bytes(vnode, b"bb")),
1904            ),
1905        )])));
1906        sst1 = sst1_inner.into();
1907
1908        let mut kvs2 = vec![(
1909            FullKey::new_with_gap_epoch(
1910                table_id,
1911                gen_key_from_bytes(vnode, b"mm"),
1912                EpochWithGap::new_from_epoch(epoch),
1913            ),
1914            HummockValue::put(Bytes::from_static(b"hit")),
1915        )];
1916        kvs2.sort_by(|(k1, _), (k2, _)| k1.cmp(k2));
1917        let (_, mut sst2): (crate::hummock::sstable_store::TableHolder, SstableInfo) =
1918            gen_test_sstable_with_table_ids(
1919                opts,
1920                12,
1921                kvs2.into_iter(),
1922                sstable_store.clone(),
1923                vec![table_id.as_raw_id()],
1924            )
1925            .await;
1926        let mut sst2_inner = sst2.get_inner();
1927        sst2_inner.vnode_statistics = Some(VnodeStatistics::from_map(BTreeMap::from_iter([(
1928            vnode,
1929            (
1930                UserKey::new(table_id, gen_key_from_bytes(vnode, b"aa")),
1931                UserKey::new(table_id, gen_key_from_bytes(vnode, b"zz")),
1932            ),
1933        )])));
1934        sst2 = sst2_inner.into();
1935
1936        let committed =
1937            build_version_from_sstables(table_id, vec![sst1, sst2], PbLevelType::Overlapping);
1938
1939        let table_key = gen_key_from_bytes(vnode, b"mm");
1940        let result = reader
1941            .get(
1942                table_key,
1943                epoch,
1944                table_id,
1945                TableOption::default(),
1946                ReadOptions::default(),
1947                (vec![], vec![], committed),
1948                |_k, v| Ok(Bytes::copy_from_slice(v)),
1949            )
1950            .await
1951            .unwrap();
1952        flush_local_metrics_for_test();
1953
1954        assert_eq!(result, Some(Bytes::from_static(b"hit")));
1955        let (checked_after, pruned_after) = vnode_prune_counts(&metrics, table_id, "get");
1956        // First SST is pruned by vnode stats; second SST is checked and read.
1957        assert_eq!(checked_before + 2, checked_after);
1958        assert_eq!(pruned_before + 1, pruned_after);
1959    }
1960}