Skip to main content

risingwave_storage/hummock/
sstable_store.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::clone::Clone;
16use std::collections::VecDeque;
17use std::ops::Deref;
18use std::pin::Pin;
19use std::sync::Arc;
20use std::sync::atomic::{AtomicUsize, Ordering};
21
22use await_tree::{InstrumentAwait, SpanExt};
23use bytes::Bytes;
24use fail::fail_point;
25use foyer::{
26    Cache, CacheBuilder, CacheEntry, EventListener, Hint, HybridCache, HybridCacheBuilder,
27    HybridCacheEntry, HybridCacheProperties,
28};
29use futures::{FutureExt, StreamExt, future};
30use prost::Message;
31use risingwave_hummock_sdk::sstable_info::SstableInfo;
32use risingwave_hummock_sdk::vector_index::{HnswGraphFileInfo, VectorFileInfo};
33use risingwave_hummock_sdk::{
34    HummockHnswGraphFileId, HummockObjectId, HummockRawObjectId, HummockSstableObjectId,
35    HummockVectorFileId, SST_OBJECT_SUFFIX,
36};
37use risingwave_hummock_trace::TracedCachePolicy;
38use risingwave_object_store::object::{
39    ObjectError, ObjectMetadataIter, ObjectResult, ObjectStoreRef, ObjectStreamingUploader,
40};
41use risingwave_pb::hummock::PbHnswGraph;
42use serde::{Deserialize, Serialize};
43use thiserror_ext::AsReport;
44use tokio::time::Instant;
45
46use super::{
47    BatchUploadWriter, Block, BlockMeta, BlockResponse, RecentFilter, Sstable, SstableMeta,
48    SstableWriterOptions,
49};
50use crate::hummock::block_stream::{BlockDataStream, MemoryUsageTracker, PrefetchBlockStream};
51use crate::hummock::none::NoneRecentFilter;
52use crate::hummock::vector::file::{VectorBlock, VectorBlockMeta, VectorFileMeta};
53use crate::hummock::vector::monitor::VectorStoreCacheStats;
54use crate::hummock::{BlockEntry, BlockHolder, HummockError, HummockResult, RecentFilterTrait};
55use crate::monitor::{HummockStateStoreMetrics, StoreLocalStatistic};
56
57macro_rules! impl_vector_index_meta_file {
58    ($($type_name:ident),+) => {
59        pub enum HummockVectorIndexMetaFile {
60            $(
61                $type_name(Pin<Box<$type_name>>),
62            )+
63        }
64
65        $(
66            impl From<$type_name> for HummockVectorIndexMetaFile {
67                fn from(v: $type_name) -> Self {
68                    Self::$type_name(Box::pin(v))
69                }
70            }
71
72            unsafe impl Send for VectorMetaFileHolder<$type_name> {}
73
74            impl VectorMetaFileHolder<$type_name> {
75                fn try_from_entry(
76                    entry: CacheEntry<HummockRawObjectId, HummockVectorIndexMetaFile>,
77                    object_id: HummockRawObjectId
78                ) -> HummockResult<Self> {
79                    let HummockVectorIndexMetaFile::$type_name(file_meta) = &*entry else {
80                        return Err(HummockError::decode_error(format!(
81                            "expect {} for object {}",
82                            stringify!($type_name),
83                            object_id
84                        )));
85                    };
86                    let ptr = file_meta.as_ref().get_ref() as *const _;
87                    Ok(VectorMetaFileHolder {
88                        _cache_entry: entry,
89                        ptr,
90                    })
91                }
92            }
93        )+
94    };
95}
96
97impl_vector_index_meta_file!(VectorFileMeta, PbHnswGraph);
98
99pub struct VectorMetaFileHolder<T> {
100    _cache_entry: CacheEntry<HummockRawObjectId, HummockVectorIndexMetaFile>,
101    ptr: *const T,
102}
103
104impl<T> Deref for VectorMetaFileHolder<T> {
105    type Target = T;
106
107    fn deref(&self) -> &Self::Target {
108        // SAFETY: VectorFileHolder is exposed only as immutable, and `VectorFileMeta` is pinned via box
109        unsafe { &*self.ptr }
110    }
111}
112
113pub type TableHolder = HybridCacheEntry<HummockSstableObjectId, Box<Sstable>>;
114
115pub type VectorBlockHolder = CacheEntry<(HummockVectorFileId, usize), Box<VectorBlock>>;
116
117pub type VectorFileHolder = VectorMetaFileHolder<VectorFileMeta>;
118pub type HnswGraphFileHolder = VectorMetaFileHolder<PbHnswGraph>;
119
120#[derive(Debug, Clone, Copy, PartialEq, PartialOrd, Eq, Ord, Hash, Serialize, Deserialize)]
121pub struct SstableBlockIndex {
122    pub sst_id: HummockSstableObjectId,
123    pub block_idx: u64,
124}
125
126pub struct BlockCacheEventListener {
127    metrics: Arc<HummockStateStoreMetrics>,
128}
129
130impl BlockCacheEventListener {
131    pub fn new(metrics: Arc<HummockStateStoreMetrics>) -> Self {
132        Self { metrics }
133    }
134}
135
136impl EventListener for BlockCacheEventListener {
137    type Key = SstableBlockIndex;
138    type Value = Box<Block>;
139
140    fn on_leave(&self, _reason: foyer::Event, _key: &Self::Key, value: &Self::Value)
141    where
142        Self::Key: foyer::Key,
143        Self::Value: foyer::Value,
144    {
145        self.metrics
146            .block_efficiency_histogram
147            .observe(value.efficiency());
148    }
149}
150
151// TODO: Define policy based on use cases (read / compaction / ...).
152#[derive(Clone, Copy, Eq, PartialEq)]
153pub enum CachePolicy {
154    /// Disable read cache and not fill the cache afterwards.
155    Disable,
156    /// Try reading the cache and fill the cache afterwards.
157    Fill(Hint),
158    /// Read the cache but not fill the cache afterwards.
159    NotFill,
160}
161
162impl Default for CachePolicy {
163    fn default() -> Self {
164        CachePolicy::Fill(Hint::Normal)
165    }
166}
167
168impl From<TracedCachePolicy> for CachePolicy {
169    fn from(policy: TracedCachePolicy) -> Self {
170        match policy {
171            TracedCachePolicy::Disable => Self::Disable,
172            TracedCachePolicy::Fill(priority) => Self::Fill(priority.into()),
173            TracedCachePolicy::NotFill => Self::NotFill,
174        }
175    }
176}
177
178impl From<CachePolicy> for TracedCachePolicy {
179    fn from(policy: CachePolicy) -> Self {
180        match policy {
181            CachePolicy::Disable => Self::Disable,
182            CachePolicy::Fill(priority) => Self::Fill(priority.into()),
183            CachePolicy::NotFill => Self::NotFill,
184        }
185    }
186}
187
188pub struct SstableStoreConfig {
189    pub store: ObjectStoreRef,
190    pub path: String,
191
192    pub prefetch_buffer_capacity: usize,
193    pub max_prefetch_block_number: usize,
194    pub recent_filter: Arc<RecentFilter<(HummockSstableObjectId, usize)>>,
195    pub state_store_metrics: Arc<HummockStateStoreMetrics>,
196    pub use_new_object_prefix_strategy: bool,
197    pub skip_bloom_filter_in_serde: bool,
198
199    pub meta_cache: HybridCache<HummockSstableObjectId, Box<Sstable>>,
200    pub block_cache: HybridCache<SstableBlockIndex, Box<Block>>,
201
202    pub vector_meta_cache: Cache<HummockRawObjectId, HummockVectorIndexMetaFile>,
203    pub vector_block_cache: Cache<(HummockVectorFileId, usize), Box<VectorBlock>>,
204}
205
206pub struct SstableStore {
207    path: String,
208    store: ObjectStoreRef,
209
210    meta_cache: HybridCache<HummockSstableObjectId, Box<Sstable>>,
211    block_cache: HybridCache<SstableBlockIndex, Box<Block>>,
212    pub vector_meta_cache: Cache<HummockRawObjectId, HummockVectorIndexMetaFile>,
213    pub vector_block_cache: Cache<(HummockVectorFileId, usize), Box<VectorBlock>>,
214
215    /// Recent filter for `(sst_obj_id, blk_idx)`.
216    ///
217    /// `blk_idx == USIZE::MAX` stands for `sst_obj_id` only entry.
218    recent_filter: Arc<RecentFilter<(HummockSstableObjectId, usize)>>,
219    prefetch_buffer_usage: Arc<AtomicUsize>,
220    prefetch_buffer_capacity: usize,
221    max_prefetch_block_number: usize,
222    /// Whether the object store is divided into prefixes depends on two factors:
223    ///   1. The specific object store type.
224    ///   2. Whether the existing cluster is a new cluster.
225    ///
226    /// The value of `use_new_object_prefix_strategy` is determined by the `use_new_object_prefix_strategy` field in the system parameters.
227    /// For a new cluster, `use_new_object_prefix_strategy` is set to True.
228    /// For an old cluster, `use_new_object_prefix_strategy` is set to False.
229    /// The final decision of whether to divide prefixes is based on this field and the specific object store type, this approach is implemented to ensure backward compatibility.
230    use_new_object_prefix_strategy: bool,
231
232    /// sst serde happens when a sst meta is written to meta disk cache.
233    /// Excluding the SST filter from serde can reduce the meta disk cache entry size
234    /// and reduce disk IO throughput at the cost of making the SST filter useless.
235    skip_bloom_filter_in_serde: bool,
236}
237
238impl SstableStore {
239    pub fn new(config: SstableStoreConfig) -> Self {
240        // TODO: We should validate path early. Otherwise object store won't report invalid path
241        // error until first write attempt.
242
243        Self {
244            path: config.path,
245            store: config.store,
246
247            meta_cache: config.meta_cache,
248            block_cache: config.block_cache,
249            vector_meta_cache: config.vector_meta_cache,
250            vector_block_cache: config.vector_block_cache,
251
252            recent_filter: config.recent_filter,
253            prefetch_buffer_usage: Arc::new(AtomicUsize::new(0)),
254            prefetch_buffer_capacity: config.prefetch_buffer_capacity,
255            max_prefetch_block_number: config.max_prefetch_block_number,
256            use_new_object_prefix_strategy: config.use_new_object_prefix_strategy,
257            skip_bloom_filter_in_serde: config.skip_bloom_filter_in_serde,
258        }
259    }
260
261    /// For compactor, we do not need a high concurrency load for cache. Instead, we need the cache
262    ///  can be evict more effective.
263    #[expect(clippy::borrowed_box)]
264    pub async fn for_compactor(
265        store: ObjectStoreRef,
266        path: String,
267        block_cache_capacity: usize,
268        meta_cache_capacity: usize,
269        use_new_object_prefix_strategy: bool,
270    ) -> HummockResult<Self> {
271        let meta_cache = HybridCacheBuilder::new()
272            .memory(meta_cache_capacity)
273            .with_shards(1)
274            .with_weighter(|_: &HummockSstableObjectId, value: &Box<Sstable>| {
275                std::mem::size_of::<HummockSstableObjectId>()
276                    + value.estimated_meta_cache_memory_weight()
277            })
278            .storage()
279            .build()
280            .await
281            .map_err(HummockError::foyer_error)?;
282
283        let block_cache = HybridCacheBuilder::new()
284            .memory(block_cache_capacity)
285            .with_shards(1)
286            .with_weighter(|_: &SstableBlockIndex, value: &Box<Block>| {
287                std::mem::size_of::<SstableBlockIndex>() + value.estimated_memory_weight()
288            })
289            .storage()
290            .build()
291            .await
292            .map_err(HummockError::foyer_error)?;
293
294        Ok(Self {
295            path,
296            store,
297
298            prefetch_buffer_usage: Arc::new(AtomicUsize::new(0)),
299            prefetch_buffer_capacity: block_cache_capacity,
300            max_prefetch_block_number: 16, /* compactor won't use this parameter, so just assign a default value. */
301            recent_filter: Arc::new(NoneRecentFilter::default().into()),
302            use_new_object_prefix_strategy,
303            skip_bloom_filter_in_serde: false,
304
305            meta_cache,
306            block_cache,
307            vector_meta_cache: CacheBuilder::new(1 << 10).build(),
308            vector_block_cache: CacheBuilder::new(1 << 10).build(),
309        })
310    }
311
312    pub async fn delete(&self, object_id: HummockSstableObjectId) -> HummockResult<()> {
313        self.store
314            .delete(self.get_sst_data_path(object_id).as_str())
315            .await?;
316        self.meta_cache.remove(&object_id);
317        // TODO(MrCroxx): support group remove in foyer.
318        Ok(())
319    }
320
321    pub fn delete_cache(&self, object_id: HummockSstableObjectId) -> HummockResult<()> {
322        self.meta_cache.remove(&object_id);
323        Ok(())
324    }
325
326    pub(crate) async fn put_sst_data(
327        &self,
328        object_id: HummockSstableObjectId,
329        data: Bytes,
330    ) -> HummockResult<()> {
331        let data_path = self.get_sst_data_path(object_id);
332        self.store
333            .upload(&data_path, data)
334            .await
335            .map_err(Into::into)
336    }
337
338    /// Buffers consecutive decoded blocks starting at `block_index`, up to `end_index` (exclusive).
339    /// Cache hits, the memory budget, and the prefetch limit may shorten the returned sequence.
340    /// All I/O and decoding errors are returned here, before the buffered blocks are consumed.
341    pub async fn prefetch_blocks(
342        &self,
343        sst: &Sstable,
344        block_index: usize,
345        end_index: usize,
346        policy: CachePolicy,
347        stats: &mut StoreLocalStatistic,
348    ) -> HummockResult<Box<PrefetchBlockStream>> {
349        let object_id = sst.id;
350        if self.prefetch_buffer_usage.load(Ordering::Acquire) > self.prefetch_buffer_capacity {
351            let block = self.get(sst, block_index, policy, stats).await?;
352            return Ok(Box::new(PrefetchBlockStream::new(
353                VecDeque::from([block]),
354                block_index,
355                None,
356            )));
357        }
358        if policy != CachePolicy::Disable
359            && let Some(entry) = self
360                .block_cache
361                .get(&SstableBlockIndex {
362                    sst_id: object_id,
363                    block_idx: block_index as _,
364                })
365                .await
366                .map_err(HummockError::foyer_error)?
367        {
368            stats.cache_data_block_total += 1;
369            if entry.source() == foyer::Source::Outer {
370                stats.cache_data_block_miss += 1;
371            }
372            let block = BlockHolder::from_hybrid_cache_entry(entry);
373            return Ok(Box::new(PrefetchBlockStream::new(
374                VecDeque::from([block]),
375                block_index,
376                None,
377            )));
378        }
379        let end_index = std::cmp::min(
380            end_index,
381            block_index.saturating_add(self.max_prefetch_block_number),
382        );
383        let mut end_index = std::cmp::min(end_index, sst.meta.block_metas.len());
384        let start_offset = sst.meta.block_metas[block_index].offset as usize;
385        if policy != CachePolicy::Disable {
386            let mut min_hit_index = end_index;
387            let mut hit_count = 0;
388            for idx in block_index..end_index {
389                if self.block_cache.contains(&SstableBlockIndex {
390                    sst_id: object_id,
391                    block_idx: idx as _,
392                }) {
393                    if min_hit_index > idx && idx > block_index {
394                        min_hit_index = idx;
395                    }
396                    hit_count += 1;
397                }
398            }
399
400            if hit_count * 3 >= (end_index - block_index)
401                || min_hit_index * 2 > block_index + end_index
402            {
403                end_index = min_hit_index;
404            }
405        }
406        stats.cache_data_prefetch_count += 1;
407        stats.cache_data_prefetch_block_count += (end_index - block_index) as u64;
408        let end_offset = start_offset
409            + sst.meta.block_metas[block_index..end_index]
410                .iter()
411                .map(|meta| meta.len as usize)
412                .sum::<usize>();
413        let data_path = self.get_sst_data_path(object_id);
414        let memory_usage = end_offset - start_offset;
415        let tracker = MemoryUsageTracker::new(self.prefetch_buffer_usage.clone(), memory_usage);
416        let span = await_tree::span!("Prefetch SST-{}", object_id).verbose();
417        let store = self.store.clone();
418        let join_handle = tokio::spawn(async move {
419            store
420                .read(&data_path, start_offset..end_offset)
421                .instrument_await(span)
422                .await
423        });
424        let buf = match join_handle.await {
425            Ok(Ok(data)) => data,
426            Ok(Err(e)) => {
427                tracing::error!(
428                    "prefetch meet error when read {}..{} from sst-{} ({})",
429                    start_offset,
430                    end_offset,
431                    object_id,
432                    sst.meta.estimated_size,
433                );
434                return Err(e.into());
435            }
436            Err(_) => {
437                return Err(HummockError::other("cancel by other thread"));
438            }
439        };
440        let mut offset = 0;
441        let mut blocks = VecDeque::default();
442        for idx in block_index..end_index {
443            let end = offset + sst.meta.block_metas[idx].len as usize;
444            if end > buf.len() {
445                return Err(ObjectError::internal("read unexpected EOF").into());
446            }
447            // copy again to avoid holding a large data in memory.
448            let block = Block::decode_with_copy(
449                buf.slice(offset..end),
450                sst.meta.block_metas[idx].uncompressed_size as usize,
451                true,
452            )?;
453            let holder = if let CachePolicy::Fill(hint) = policy {
454                let hint = if idx == block_index { hint } else { Hint::Low };
455                let entry = self.block_cache.insert_with_properties(
456                    SstableBlockIndex {
457                        sst_id: object_id,
458                        block_idx: idx as _,
459                    },
460                    Box::new(block),
461                    HybridCacheProperties::default().with_hint(hint),
462                );
463                BlockHolder::from_hybrid_cache_entry(entry)
464            } else {
465                BlockHolder::from_owned_block(Box::new(block))
466            };
467
468            blocks.push_back(holder);
469            offset = end;
470        }
471        Ok(Box::new(PrefetchBlockStream::new(
472            blocks,
473            block_index,
474            Some(tracker),
475        )))
476    }
477
478    pub async fn get_block_response(
479        &self,
480        sst: &Sstable,
481        block_index: usize,
482        policy: CachePolicy,
483    ) -> HummockResult<BlockResponse> {
484        let object_id = sst.id;
485        let (range, uncompressed_capacity) = sst.calculate_block_info(block_index);
486        let store = self.store.clone();
487
488        let file_size = sst.meta.estimated_size;
489        let data_path = Arc::new(self.get_sst_data_path(object_id));
490
491        let disable_cache: fn() -> bool = || {
492            fail_point!("disable_block_cache", |_| true);
493            false
494        };
495
496        let policy = if disable_cache() {
497            CachePolicy::Disable
498        } else {
499            policy
500        };
501
502        let idx = SstableBlockIndex {
503            sst_id: object_id,
504            block_idx: block_index as _,
505        };
506
507        self.recent_filter
508            .extend([(object_id, usize::MAX), (object_id, block_index)]);
509
510        // future: fetch block if hybrid cache miss
511        let fetch_block = async move {
512            let block_data = match store
513                .read(&data_path, range.clone())
514                .instrument_await("get_block_response".verbose())
515                .await
516            {
517                Ok(data) => data,
518                Err(e) => {
519                    tracing::error!(
520                        "get_block_response meet error when read {:?} from sst-{}, total length: {}",
521                        range,
522                        object_id,
523                        file_size
524                    );
525                    return Err(HummockError::from(e));
526                }
527            };
528            let block = Box::new(Block::decode(block_data, uncompressed_capacity)?);
529            Ok(block)
530        };
531
532        match policy {
533            CachePolicy::Fill(hint) => {
534                let properties = HybridCacheProperties::default().with_hint(hint);
535                let fetch = self.block_cache.get_or_fetch(&idx, || {
536                    fetch_block.map(|res| res.map(|block| (block, properties)))
537                });
538                Ok(BlockResponse::Fetch(fetch))
539            }
540            CachePolicy::NotFill => {
541                match self
542                    .block_cache
543                    .get(&idx)
544                    .await
545                    .map_err(HummockError::foyer_error)?
546                {
547                    Some(entry) => Ok(BlockResponse::Block(BlockHolder::from_hybrid_cache_entry(
548                        entry,
549                    ))),
550                    _ => {
551                        let block = fetch_block.await?;
552                        Ok(BlockResponse::Block(BlockHolder::from_owned_block(block)))
553                    }
554                }
555            }
556            CachePolicy::Disable => {
557                let block = fetch_block.await?;
558                Ok(BlockResponse::Block(BlockHolder::from_owned_block(block)))
559            }
560        }
561    }
562
563    pub async fn get(
564        &self,
565        sst: &Sstable,
566        block_index: usize,
567        policy: CachePolicy,
568        stats: &mut StoreLocalStatistic,
569    ) -> HummockResult<BlockHolder> {
570        let block_response = self.get_block_response(sst, block_index, policy).await?;
571        let block_holder = block_response.wait().await?;
572        stats.cache_data_block_total += 1;
573        if let BlockEntry::HybridCache(entry) = block_holder.entry()
574            && entry.source() == foyer::Source::Outer
575        {
576            stats.cache_data_block_miss += 1;
577        }
578        Ok(block_holder)
579    }
580
581    pub async fn get_vector_file_meta(
582        &self,
583        vector_file: &VectorFileInfo,
584        stats: &mut VectorStoreCacheStats,
585    ) -> HummockResult<VectorFileHolder> {
586        let store = self.store.clone();
587        let path = self.get_object_data_path(HummockObjectId::VectorFile(vector_file.object_id));
588        let meta_offset = vector_file.meta_offset;
589        let entry = self
590            .vector_meta_cache
591            .get_or_fetch(&vector_file.object_id.as_raw(), || async move {
592                let encoded_footer = store.read(&path, meta_offset..).await?;
593                let meta = VectorFileMeta::decode_footer(&encoded_footer)?;
594                Ok::<_, anyhow::Error>(HummockVectorIndexMetaFile::from(meta))
595            })
596            .await?;
597        stats.file_meta_total += 1;
598        if entry.source() == foyer::Source::Outer {
599            stats.file_meta_miss += 1;
600        }
601        VectorFileHolder::try_from_entry(entry, vector_file.object_id.as_raw())
602    }
603
604    pub async fn get_vector_block(
605        &self,
606        vector_file: &VectorFileInfo,
607        block_idx: usize,
608        block_meta: &VectorBlockMeta,
609        stats: &mut VectorStoreCacheStats,
610    ) -> HummockResult<VectorBlockHolder> {
611        let store = self.store.clone();
612        let path = self.get_object_data_path(HummockObjectId::VectorFile(vector_file.object_id));
613        let start_offset = block_meta.offset;
614        let end_offset = start_offset + block_meta.block_size;
615        let entry = self
616            .vector_block_cache
617            .get_or_fetch(&(vector_file.object_id, block_idx), || async move {
618                let encoded_block = store.read(&path, start_offset..end_offset).await?;
619                let block = VectorBlock::decode(&encoded_block)?;
620                Ok::<_, anyhow::Error>(Box::new(block))
621            })
622            .await
623            .map_err(HummockError::foyer_error)?;
624
625        stats.file_block_total += 1;
626        if entry.source() == foyer::Source::Outer {
627            stats.file_block_miss += 1;
628        }
629        Ok(entry)
630    }
631
632    pub fn insert_vector_cache(
633        &self,
634        object_id: HummockVectorFileId,
635        meta: VectorFileMeta,
636        blocks: Vec<VectorBlock>,
637    ) {
638        self.vector_meta_cache
639            .insert(object_id.as_raw(), meta.into());
640        for (idx, block) in blocks.into_iter().enumerate() {
641            self.vector_block_cache
642                .insert((object_id, idx), Box::new(block));
643        }
644    }
645
646    pub fn insert_hnsw_graph_cache(&self, object_id: HummockHnswGraphFileId, graph: PbHnswGraph) {
647        self.vector_meta_cache
648            .insert(object_id.as_raw(), graph.into());
649    }
650
651    pub async fn get_hnsw_graph(
652        &self,
653        graph_file: &HnswGraphFileInfo,
654        stats: &mut VectorStoreCacheStats,
655    ) -> HummockResult<HnswGraphFileHolder> {
656        let store = self.store.clone();
657        let graph_file_path =
658            self.get_object_data_path(HummockObjectId::HnswGraphFile(graph_file.object_id));
659        let entry = self
660            .vector_meta_cache
661            .get_or_fetch(&graph_file.object_id.as_raw(), || async move {
662                let encoded_graph = store.read(&graph_file_path, ..).await?;
663                let graph = PbHnswGraph::decode(encoded_graph.as_ref())?;
664                Ok::<_, anyhow::Error>(HummockVectorIndexMetaFile::from(graph))
665            })
666            .await
667            .map_err(HummockError::foyer_error)?;
668        stats.hnsw_graph_total += 1;
669        if entry.source() == foyer::Source::Outer {
670            stats.hnsw_graph_miss += 1;
671        }
672        HnswGraphFileHolder::try_from_entry(entry, graph_file.object_id.as_raw())
673    }
674
675    pub fn get_sst_data_path(&self, object_id: impl Into<HummockSstableObjectId>) -> String {
676        self.get_object_data_path(HummockObjectId::Sstable(object_id.into()))
677    }
678
679    pub fn get_object_data_path(&self, object_id: HummockObjectId) -> String {
680        let obj_prefix = self.store.get_object_prefix(
681            object_id.as_raw().as_raw_id(),
682            self.use_new_object_prefix_strategy,
683        );
684        risingwave_hummock_sdk::get_object_data_path(&obj_prefix, &self.path, object_id)
685    }
686
687    pub fn get_object_id_from_path(path: &str) -> HummockObjectId {
688        risingwave_hummock_sdk::get_object_id_from_path(path)
689    }
690
691    pub fn store(&self) -> ObjectStoreRef {
692        self.store.clone()
693    }
694
695    #[cfg(any(test, feature = "test"))]
696    pub async fn clear_block_cache(&self) -> HummockResult<()> {
697        self.block_cache
698            .clear()
699            .await
700            .map_err(HummockError::foyer_error)
701    }
702
703    #[cfg(any(test, feature = "test"))]
704    pub async fn clear_meta_cache(&self) -> HummockResult<()> {
705        self.meta_cache
706            .clear()
707            .await
708            .map_err(HummockError::foyer_error)
709    }
710
711    pub async fn sstable_cached(
712        &self,
713        sst_obj_id: HummockSstableObjectId,
714    ) -> HummockResult<Option<HybridCacheEntry<HummockSstableObjectId, Box<Sstable>>>> {
715        self.meta_cache
716            .get(&sst_obj_id)
717            .await
718            .map_err(HummockError::foyer_error)
719    }
720
721    /// Returns `table_holder`
722    pub async fn sstable(
723        &self,
724        sstable_info_ref: &SstableInfo,
725        stats: &mut StoreLocalStatistic,
726    ) -> HummockResult<TableHolder> {
727        let object_id = sstable_info_ref.object_id;
728        let store = self.store.clone();
729        let meta_path = self.get_sst_data_path(object_id);
730        let stats_ptr = stats.remote_io_time.clone();
731        let range = sstable_info_ref.meta_offset as usize..;
732        let skip_bloom_filter_in_serde = self.skip_bloom_filter_in_serde;
733
734        let fetch = self.meta_cache.get_or_fetch(&object_id, || async move {
735            let now = Instant::now();
736            let buf = store
737                .read(&meta_path, range)
738                .instrument_await("get_meta_response".verbose())
739                .await?;
740            let meta = SstableMeta::decode(&buf[..])?;
741
742            let sst = Sstable::new(object_id, meta, skip_bloom_filter_in_serde);
743            let add = (now.elapsed().as_secs_f64() * 1000.0).ceil();
744            stats_ptr.fetch_add(add as u64, Ordering::Relaxed);
745            Ok::<_, anyhow::Error>(Box::new(sst))
746        });
747
748        stats.cache_meta_block_total += 1;
749        let entry = fetch
750            .instrument_await("fetch_meta".verbose())
751            .await
752            .map_err(HummockError::foyer_error);
753        if let Ok(ref entry) = entry
754            && entry.source() == foyer::Source::Outer
755        {
756            stats.cache_meta_block_miss += 1;
757        }
758        entry
759    }
760
761    pub async fn list_sst_object_metadata_from_object_store(
762        &self,
763        prefix: Option<String>,
764        start_after: Option<String>,
765        limit: Option<usize>,
766    ) -> HummockResult<ObjectMetadataIter> {
767        let list_path = format!("{}/{}", self.path, prefix.unwrap_or("".into()));
768        let raw_iter = self.store.list(&list_path, start_after, limit).await?;
769        let iter = raw_iter.filter(|r| match r {
770            Ok(i) => future::ready(i.key.ends_with(&format!(".{}", SST_OBJECT_SUFFIX))),
771            Err(_) => future::ready(true),
772        });
773        Ok(Box::pin(iter))
774    }
775
776    pub fn create_sst_writer(
777        self: Arc<Self>,
778        object_id: impl Into<HummockSstableObjectId>,
779        options: SstableWriterOptions,
780    ) -> BatchUploadWriter {
781        BatchUploadWriter::new(object_id, self, options)
782    }
783
784    pub fn insert_meta_cache(&self, object_id: HummockSstableObjectId, meta: SstableMeta) {
785        let sst = Sstable::new(object_id, meta, self.skip_bloom_filter_in_serde);
786        self.meta_cache.insert(object_id, Box::new(sst));
787    }
788
789    pub fn insert_block_cache(
790        &self,
791        object_id: HummockSstableObjectId,
792        block_index: u64,
793        block: Box<Block>,
794    ) {
795        self.block_cache.insert(
796            SstableBlockIndex {
797                sst_id: object_id,
798                block_idx: block_index,
799            },
800            block,
801        );
802    }
803
804    pub fn get_prefetch_memory_usage(&self) -> usize {
805        self.prefetch_buffer_usage.load(Ordering::Acquire)
806    }
807
808    pub async fn get_stream_for_blocks(
809        &self,
810        object_id: HummockSstableObjectId,
811        metas: &[BlockMeta],
812    ) -> HummockResult<BlockDataStream> {
813        fail_point!("get_stream_err");
814        let data_path = self.get_sst_data_path(object_id);
815        let store = self.store();
816        let block_meta = &metas[0];
817        let start_pos = block_meta.offset as usize;
818        let end_pos = metas.iter().map(|meta| meta.len as usize).sum::<usize>() + start_pos;
819        let range = start_pos..end_pos;
820        // spawn to tokio pool because the object-storage sdk may not be safe to cancel.
821        let ret = tokio::spawn(async move { store.streaming_read(&data_path, range).await }).await;
822
823        let reader = match ret {
824            Ok(Ok(reader)) => reader,
825            Ok(Err(e)) => return Err(HummockError::from(e)),
826            Err(e) => {
827                return Err(HummockError::other(format!(
828                    "failed to get result, this read request may be canceled: {}",
829                    e.as_report()
830                )));
831            }
832        };
833        Ok(BlockDataStream::new(reader, metas))
834    }
835
836    pub fn meta_cache(&self) -> &HybridCache<HummockSstableObjectId, Box<Sstable>> {
837        &self.meta_cache
838    }
839
840    pub fn block_cache(&self) -> &HybridCache<SstableBlockIndex, Box<Block>> {
841        &self.block_cache
842    }
843
844    pub fn recent_filter(&self) -> &Arc<RecentFilter<(HummockSstableObjectId, usize)>> {
845        &self.recent_filter
846    }
847
848    pub async fn create_streaming_uploader(
849        &self,
850        path: &str,
851    ) -> ObjectResult<ObjectStreamingUploader> {
852        self.store.streaming_upload(path).await
853    }
854}
855
856pub type SstableStoreRef = Arc<SstableStore>;
857#[cfg(test)]
858mod tests {
859    use std::ops::Range;
860    use std::sync::Arc;
861
862    use risingwave_hummock_sdk::HummockObjectId;
863    use risingwave_hummock_sdk::sstable_info::SstableInfo;
864
865    use super::{SstableStoreRef, SstableWriterOptions};
866    use crate::hummock::iterator::HummockIterator;
867    use crate::hummock::iterator::test_utils::{iterator_test_key_of, mock_sstable_store};
868    use crate::hummock::sstable::SstableIteratorReadOptions;
869    use crate::hummock::test_utils::{
870        default_builder_opt_for_test, gen_default_test_sstable, gen_test_sstable_data, put_sst,
871        test_key_of,
872    };
873    use crate::hummock::value::HummockValue;
874    use crate::hummock::{CachePolicy, SstableIterator, SstableMeta, SstableStore};
875    use crate::monitor::StoreLocalStatistic;
876
877    const SST_ID: u64 = 1;
878
879    fn get_hummock_value(x: usize) -> HummockValue<Vec<u8>> {
880        HummockValue::put(format!("overlapped_new_{}", x).as_bytes().to_vec())
881    }
882
883    async fn validate_sst(
884        sstable_store: SstableStoreRef,
885        info: &SstableInfo,
886        mut meta: SstableMeta,
887        x_range: Range<usize>,
888    ) {
889        let mut stats = StoreLocalStatistic::default();
890        let holder = sstable_store.sstable(info, &mut stats).await.unwrap();
891        std::mem::take(&mut meta.bloom_filter);
892        assert_eq!(holder.meta, meta);
893        let holder = sstable_store.sstable(info, &mut stats).await.unwrap();
894        assert_eq!(holder.meta, meta);
895        let mut iter = SstableIterator::new(
896            holder,
897            sstable_store,
898            Arc::new(SstableIteratorReadOptions::default()),
899            info,
900        );
901        iter.rewind().await.unwrap();
902        for i in x_range {
903            let key = iter.key();
904            let value = iter.value();
905            assert_eq!(key, iterator_test_key_of(i).to_ref());
906            assert_eq!(value, get_hummock_value(i).as_slice());
907            iter.next().await.unwrap();
908        }
909    }
910
911    #[tokio::test]
912    async fn test_batch_upload() {
913        let sstable_store = mock_sstable_store().await;
914        let x_range = 0..100;
915        let (data, meta) = gen_test_sstable_data(
916            default_builder_opt_for_test(),
917            x_range
918                .clone()
919                .map(|x| (iterator_test_key_of(x), get_hummock_value(x))),
920        )
921        .await;
922        let writer_opts = SstableWriterOptions {
923            capacity_hint: None,
924            tracker: None,
925            policy: CachePolicy::Disable,
926        };
927        let info = put_sst(
928            SST_ID,
929            data.clone(),
930            meta.clone(),
931            sstable_store.clone(),
932            writer_opts,
933            vec![0],
934        )
935        .await
936        .unwrap();
937
938        validate_sst(sstable_store, &info, meta, x_range).await;
939    }
940
941    #[tokio::test]
942    async fn test_streaming_upload() {
943        // Generate test data.
944        let sstable_store = mock_sstable_store().await;
945        let x_range = 0..100;
946        let (data, meta) = gen_test_sstable_data(
947            default_builder_opt_for_test(),
948            x_range
949                .clone()
950                .map(|x| (iterator_test_key_of(x), get_hummock_value(x))),
951        )
952        .await;
953        let writer_opts = SstableWriterOptions {
954            capacity_hint: None,
955            tracker: None,
956            policy: CachePolicy::Disable,
957        };
958        let info = put_sst(
959            SST_ID,
960            data.clone(),
961            meta.clone(),
962            sstable_store.clone(),
963            writer_opts,
964            vec![0],
965        )
966        .await
967        .unwrap();
968
969        validate_sst(sstable_store, &info, meta, x_range).await;
970    }
971
972    #[tokio::test]
973    async fn test_empty_prefetch_falls_back_to_block_get() {
974        let mut sstable_store = mock_sstable_store().await;
975        // Direct SstableStoreConfig construction can bypass the serde nonzero check.
976        Arc::get_mut(&mut sstable_store)
977            .unwrap()
978            .max_prefetch_block_number = 0;
979        let (sstable, info) =
980            gen_default_test_sstable(default_builder_opt_for_test(), 0, sstable_store.clone())
981                .await;
982        sstable_store.clear_block_cache().await.unwrap();
983        let mut iter = SstableIterator::new(
984            sstable,
985            sstable_store.clone(),
986            Arc::new(SstableIteratorReadOptions {
987                prefetch: true,
988                cache_policy: CachePolicy::Disable,
989                ..Default::default()
990            }),
991            &info,
992        );
993        tokio::time::timeout(std::time::Duration::from_secs(5), iter.rewind())
994            .await
995            .expect("empty prefetch must not cause an unbounded refill loop")
996            .unwrap();
997        assert_eq!(iter.key(), test_key_of(0).to_ref());
998        let mut stats = StoreLocalStatistic::default();
999        iter.collect_local_statistic(&mut stats);
1000        assert_eq!(stats.cache_data_prefetch_count, 1);
1001        assert_eq!(stats.cache_data_block_total, 1);
1002        assert_eq!(sstable_store.get_prefetch_memory_usage(), 0);
1003    }
1004
1005    #[tokio::test]
1006    async fn test_prefetch_failure_falls_back_and_releases_tracker() {
1007        for truncate in [true, false] {
1008            let sstable_store = mock_sstable_store().await;
1009            let (sstable, info) =
1010                gen_default_test_sstable(default_builder_opt_for_test(), 0, sstable_store.clone())
1011                    .await;
1012            assert!(sstable.meta.block_metas.len() > 1);
1013            sstable_store.clear_block_cache().await.unwrap();
1014            let path = sstable_store.get_sst_data_path(sstable.id);
1015            let mut data = sstable_store.store.read(&path, ..).await.unwrap().to_vec();
1016            let second_offset = sstable.meta.block_metas[1].offset as usize;
1017            if truncate {
1018                // A multi-block read fails, but reading block 0 still succeeds.
1019                data.truncate(second_offset);
1020            } else {
1021                // Decoding block 1 fails its checksum after block 0 was already decoded.
1022                data[second_offset] ^= 1;
1023            }
1024            sstable_store
1025                .store
1026                .upload(&path, data.into())
1027                .await
1028                .unwrap();
1029            let mut iter = SstableIterator::new(
1030                sstable.clone(),
1031                sstable_store.clone(),
1032                Arc::new(SstableIteratorReadOptions {
1033                    prefetch: true,
1034                    cache_policy: CachePolicy::Disable,
1035                    ..Default::default()
1036                }),
1037                &info,
1038            );
1039            iter.rewind().await.unwrap();
1040            assert_eq!(iter.key(), test_key_of(0).to_ref());
1041            assert_eq!(sstable_store.get_prefetch_memory_usage(), 0);
1042            let mut stats = StoreLocalStatistic::default();
1043            iter.collect_local_statistic(&mut stats);
1044            assert_eq!(stats.cache_data_prefetch_count, 1);
1045            assert_eq!(stats.cache_data_block_total, 1);
1046            // If the single-block read also fails, the error must still reach the caller.
1047            let second_key = risingwave_hummock_sdk::key::FullKey::decode(
1048                &sstable.meta.block_metas[1].smallest_key,
1049            );
1050            assert!(iter.seek(second_key).await.is_err());
1051            assert_eq!(sstable_store.get_prefetch_memory_usage(), 0);
1052        }
1053    }
1054
1055    #[tokio::test]
1056    async fn test_prefetch_releases_old_budget_before_refill() {
1057        let mut sstable_store = mock_sstable_store().await;
1058        // The first prefetch may start at zero usage. A retained old tracker would force
1059        // the next prefetch into the single-block path instead of reading another batch.
1060        Arc::get_mut(&mut sstable_store)
1061            .unwrap()
1062            .prefetch_buffer_capacity = 0;
1063        let (sstable, info) =
1064            gen_default_test_sstable(default_builder_opt_for_test(), 0, sstable_store.clone())
1065                .await;
1066        sstable_store.clear_block_cache().await.unwrap();
1067        let next_batch = sstable_store.max_prefetch_block_number;
1068        assert!(sstable.meta.block_metas.len() > next_batch + 1);
1069        let mut iter = SstableIterator::new(
1070            sstable.clone(),
1071            sstable_store.clone(),
1072            Arc::new(SstableIteratorReadOptions {
1073                prefetch: true,
1074                cache_policy: CachePolicy::Disable,
1075                ..Default::default()
1076            }),
1077            &info,
1078        );
1079        for idx in [0, next_batch] {
1080            let key = risingwave_hummock_sdk::key::FullKey::decode(
1081                &sstable.meta.block_metas[idx].smallest_key,
1082            );
1083            iter.seek(key).await.unwrap();
1084            assert_eq!(iter.key(), key);
1085            assert!(sstable_store.get_prefetch_memory_usage() > 0);
1086        }
1087        let mut stats = StoreLocalStatistic::default();
1088        iter.collect_local_statistic(&mut stats);
1089        assert_eq!(stats.cache_data_prefetch_count, 2);
1090        assert_eq!(stats.cache_data_block_total, 0);
1091        drop(iter);
1092        assert_eq!(sstable_store.get_prefetch_memory_usage(), 0);
1093    }
1094
1095    #[tokio::test]
1096    async fn test_prefetch_producer_large_limit() {
1097        let mut sstable_store = mock_sstable_store().await;
1098        Arc::get_mut(&mut sstable_store)
1099            .unwrap()
1100            .max_prefetch_block_number = usize::MAX;
1101        let (sstable, _) =
1102            gen_default_test_sstable(default_builder_opt_for_test(), 0, sstable_store.clone())
1103                .await;
1104        sstable_store.clear_block_cache().await.unwrap();
1105        let mut stats = StoreLocalStatistic::default();
1106        let mut stream = sstable_store
1107            .prefetch_blocks(&sstable, 1, 3, CachePolicy::Disable, &mut stats)
1108            .await
1109            .unwrap();
1110        assert_eq!(stats.cache_data_prefetch_block_count, 2);
1111        for idx in [1, 2] {
1112            assert!(matches!(
1113                stream.take_block(idx),
1114                crate::hummock::block_stream::PrefetchLookup::Hit(_)
1115            ));
1116        }
1117        assert!(matches!(
1118            stream.take_block(3),
1119            crate::hummock::block_stream::PrefetchLookup::Exhausted
1120        ));
1121        drop(stream);
1122        assert_eq!(sstable_store.get_prefetch_memory_usage(), 0);
1123    }
1124
1125    #[tokio::test]
1126    async fn test_prefetch_producer_respects_disabled_cache() {
1127        let sstable_store = mock_sstable_store().await;
1128        let (sstable, _) =
1129            gen_default_test_sstable(default_builder_opt_for_test(), 0, sstable_store.clone())
1130                .await;
1131        sstable_store.clear_block_cache().await.unwrap();
1132        // A cached block in the requested range must not shorten a cache-disabled prefetch.
1133        sstable_store
1134            .get(
1135                &sstable,
1136                1,
1137                CachePolicy::Fill(foyer::Hint::Normal),
1138                &mut StoreLocalStatistic::default(),
1139            )
1140            .await
1141            .unwrap();
1142        let mut stats = StoreLocalStatistic::default();
1143        let stream = sstable_store
1144            .prefetch_blocks(&sstable, 0, 3, CachePolicy::Disable, &mut stats)
1145            .await
1146            .unwrap();
1147        assert_eq!(stats.cache_data_prefetch_block_count, 3);
1148        drop(stream);
1149        assert!(
1150            !sstable_store
1151                .block_cache
1152                .contains(&super::SstableBlockIndex {
1153                    sst_id: sstable.id,
1154                    block_idx: 0,
1155                })
1156        );
1157
1158        // With block 0 cached and the object removed, NotFill may succeed but Disable must fail.
1159        sstable_store
1160            .get(
1161                &sstable,
1162                0,
1163                CachePolicy::Fill(foyer::Hint::Normal),
1164                &mut StoreLocalStatistic::default(),
1165            )
1166            .await
1167            .unwrap();
1168        sstable_store.delete(sstable.id).await.unwrap();
1169        assert!(
1170            sstable_store
1171                .prefetch_blocks(
1172                    &sstable,
1173                    0,
1174                    3,
1175                    CachePolicy::NotFill,
1176                    &mut StoreLocalStatistic::default()
1177                )
1178                .await
1179                .is_ok()
1180        );
1181        assert!(
1182            sstable_store
1183                .prefetch_blocks(
1184                    &sstable,
1185                    0,
1186                    3,
1187                    CachePolicy::Disable,
1188                    &mut StoreLocalStatistic::default()
1189                )
1190                .await
1191                .is_err()
1192        );
1193        assert_eq!(sstable_store.get_prefetch_memory_usage(), 0);
1194    }
1195
1196    #[tokio::test]
1197    async fn test_basic() {
1198        let sstable_store = mock_sstable_store().await;
1199        let object_id = 123;
1200        let data_path = sstable_store.get_sst_data_path(object_id);
1201        assert_eq!(data_path, "test/123.data");
1202        assert_eq!(
1203            SstableStore::get_object_id_from_path(&data_path),
1204            HummockObjectId::Sstable(object_id.into())
1205        );
1206    }
1207
1208    #[tokio::test]
1209    async fn test_clear_file_cache_and_refetch() {
1210        let sstable_store = mock_sstable_store().await;
1211        let x_range = 0..10;
1212        let (data, meta) = gen_test_sstable_data(
1213            default_builder_opt_for_test(),
1214            x_range.map(|x| (iterator_test_key_of(x), get_hummock_value(x))),
1215        )
1216        .await;
1217        let info = put_sst(
1218            SST_ID,
1219            data,
1220            meta,
1221            sstable_store.clone(),
1222            SstableWriterOptions {
1223                capacity_hint: None,
1224                tracker: None,
1225                policy: CachePolicy::Fill(foyer::Hint::Normal),
1226            },
1227            vec![0],
1228        )
1229        .await
1230        .unwrap();
1231
1232        let mut stats = StoreLocalStatistic::default();
1233        assert!(
1234            sstable_store
1235                .sstable_cached(info.object_id)
1236                .await
1237                .unwrap()
1238                .is_some()
1239        );
1240        let sstable = sstable_store.sstable(&info, &mut stats).await.unwrap();
1241        sstable_store
1242            .get(
1243                &sstable,
1244                0,
1245                CachePolicy::Fill(foyer::Hint::Normal),
1246                &mut stats,
1247            )
1248            .await
1249            .unwrap();
1250        assert!(
1251            sstable_store
1252                .block_cache()
1253                .get(&super::SstableBlockIndex {
1254                    sst_id: info.object_id,
1255                    block_idx: 0,
1256                })
1257                .await
1258                .unwrap()
1259                .is_some()
1260        );
1261
1262        sstable_store.clear_meta_cache().await.unwrap();
1263        assert!(
1264            sstable_store
1265                .sstable_cached(info.object_id)
1266                .await
1267                .unwrap()
1268                .is_none()
1269        );
1270        sstable_store.sstable(&info, &mut stats).await.unwrap();
1271        assert!(
1272            sstable_store
1273                .sstable_cached(info.object_id)
1274                .await
1275                .unwrap()
1276                .is_some()
1277        );
1278
1279        sstable_store.clear_block_cache().await.unwrap();
1280        assert!(
1281            sstable_store
1282                .block_cache()
1283                .get(&super::SstableBlockIndex {
1284                    sst_id: info.object_id,
1285                    block_idx: 0,
1286                })
1287                .await
1288                .unwrap()
1289                .is_none()
1290        );
1291        sstable_store
1292            .get(
1293                &sstable,
1294                0,
1295                CachePolicy::Fill(foyer::Hint::Normal),
1296                &mut stats,
1297            )
1298            .await
1299            .unwrap();
1300        assert!(
1301            sstable_store
1302                .block_cache()
1303                .get(&super::SstableBlockIndex {
1304                    sst_id: info.object_id,
1305                    block_idx: 0,
1306                })
1307                .await
1308                .unwrap()
1309                .is_some()
1310        );
1311    }
1312}