Skip to main content

mito2/cache/
index.rs

1// Copyright 2023 Greptime Team
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
15pub mod bloom_filter_index;
16pub mod inverted_index;
17pub mod result_cache;
18
19use std::collections::{HashMap, HashSet};
20use std::future::Future;
21use std::hash::Hash;
22use std::ops::Range;
23use std::sync::Arc;
24
25use bytes::Bytes;
26use object_store::Buffer;
27
28use crate::metrics::{CACHE_BYTES, CACHE_HIT, CACHE_MISS};
29
30/// Metrics for index metadata.
31const INDEX_METADATA_TYPE: &str = "index_metadata";
32/// Metrics for index content.
33const INDEX_CONTENT_TYPE: &str = "index_content";
34
35/// Metrics collected from IndexCache operations.
36#[derive(Debug, Default, Clone)]
37pub struct IndexCacheMetrics {
38    /// Number of cache hits.
39    pub cache_hit: usize,
40    /// Number of cache misses.
41    pub cache_miss: usize,
42    /// Number of pages accessed.
43    pub num_pages: usize,
44    /// Total bytes from pages.
45    pub page_bytes: u64,
46}
47
48#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
49pub struct PageKey {
50    page_id: u64,
51}
52
53impl PageKey {
54    /// Converts an offset to a page ID based on the page size.
55    fn calculate_page_id(offset: u64, page_size: u64) -> u64 {
56        offset / page_size
57    }
58
59    /// Calculates the total number of pages that a given size spans, starting from a specific offset.
60    fn calculate_page_count(offset: u64, size: u32, page_size: u64) -> u32 {
61        let start_page = Self::calculate_page_id(offset, page_size);
62        let end_page = Self::calculate_page_id(offset + (size as u64) - 1, page_size);
63        (end_page + 1 - start_page) as u32
64    }
65
66    /// Calculates the byte range for data retrieval based on the specified offset and size.
67    ///
68    /// This function determines the starting and ending byte positions required for reading data.
69    /// For example, with an offset of 5000 and a size of 5000, using a PAGE_SIZE of 4096,
70    /// the resulting byte range will be 904..5904. This indicates that:
71    /// - The reader will first access fixed-size pages [4096, 8192) and [8192, 12288).
72    /// - To read the range [5000..10000), it only needs to fetch bytes within the range [904, 5904) across two pages.
73    fn calculate_range(offset: u64, size: u32, page_size: u64) -> Range<usize> {
74        let start = (offset % page_size) as usize;
75        let end = start + size as usize;
76        start..end
77    }
78
79    /// Generates a iterator of `IndexKey` for the pages that a given offset and size span.
80    fn generate_page_keys(offset: u64, size: u32, page_size: u64) -> impl Iterator<Item = Self> {
81        let start_page = Self::calculate_page_id(offset, page_size);
82        let total_pages = Self::calculate_page_count(offset, size, page_size);
83        (0..total_pages).map(move |i| Self {
84            page_id: start_page + i as u64,
85        })
86    }
87}
88
89/// Cache for index metadata and content.
90pub struct IndexCache<K, M> {
91    /// Cache for index metadata
92    index_metadata: moka::sync::Cache<K, Arc<M>>,
93    /// Cache for index content.
94    index: moka::sync::Cache<(K, PageKey), Bytes>,
95    // Page size for index content.
96    page_size: u64,
97
98    /// Weighter for metadata.
99    weight_of_metadata: fn(&K, &Arc<M>) -> u32,
100    /// Weighter for content.
101    weight_of_content: fn(&(K, PageKey), &Bytes) -> u32,
102}
103
104impl<K, M> IndexCache<K, M>
105where
106    K: Hash + Eq + Send + Sync + 'static,
107    M: Send + Sync + 'static,
108{
109    pub fn new_with_weighter(
110        index_metadata_cap: u64,
111        index_content_cap: u64,
112        page_size: u64,
113        index_type: &'static str,
114        weight_of_metadata: fn(&K, &Arc<M>) -> u32,
115        weight_of_content: fn(&(K, PageKey), &Bytes) -> u32,
116    ) -> Self {
117        common_telemetry::debug!(
118            "Building IndexCache with metadata size: {index_metadata_cap}, content size: {index_content_cap}, page size: {page_size}, index type: {index_type}"
119        );
120        let index_metadata = moka::sync::CacheBuilder::new(index_metadata_cap)
121            .name(&format!("index_metadata_{}", index_type))
122            .weigher(weight_of_metadata)
123            .eviction_listener(move |k, v, _cause| {
124                let size = weight_of_metadata(&k, &v);
125                CACHE_BYTES
126                    .with_label_values(&[INDEX_METADATA_TYPE])
127                    .sub(size.into());
128            })
129            .support_invalidation_closures()
130            .build();
131        let index_cache = moka::sync::CacheBuilder::new(index_content_cap)
132            .name(&format!("index_content_{}", index_type))
133            .weigher(weight_of_content)
134            .eviction_listener(move |k, v, _cause| {
135                let size = weight_of_content(&k, &v);
136                CACHE_BYTES
137                    .with_label_values(&[INDEX_CONTENT_TYPE])
138                    .sub(size.into());
139            })
140            .support_invalidation_closures()
141            .build();
142        Self {
143            index_metadata,
144            index: index_cache,
145            page_size,
146            weight_of_content,
147            weight_of_metadata,
148        }
149    }
150}
151
152impl<K, M> IndexCache<K, M>
153where
154    K: Hash + Eq + Clone + Copy + Send + Sync + 'static,
155    M: Send + Sync + 'static,
156{
157    pub fn get_metadata(&self, key: K) -> Option<Arc<M>> {
158        self.index_metadata.get(&key)
159    }
160
161    pub fn put_metadata(&self, key: K, metadata: Arc<M>) {
162        CACHE_BYTES
163            .with_label_values(&[INDEX_METADATA_TYPE])
164            .add((self.weight_of_metadata)(&key, &metadata).into());
165        self.index_metadata.insert(key, metadata)
166    }
167
168    /// Gets given range of index data from cache, and loads from source if the file
169    /// is not already cached.
170    async fn get_or_load<F, Fut, E>(
171        &self,
172        key: K,
173        file_size: u64,
174        offset: u64,
175        size: u32,
176        load: F,
177    ) -> Result<(Vec<u8>, IndexCacheMetrics), E>
178    where
179        F: Fn(Vec<Range<u64>>) -> Fut,
180        Fut: Future<Output = Result<Vec<Bytes>, E>>,
181        E: std::error::Error,
182    {
183        let mut metrics = IndexCacheMetrics::default();
184        let page_keys =
185            PageKey::generate_page_keys(offset, size, self.page_size).collect::<Vec<_>>();
186        // Size is 0, return empty data.
187        if page_keys.is_empty() {
188            return Ok((Vec::new(), metrics));
189        }
190        metrics.num_pages = page_keys.len();
191        let mut data = Vec::with_capacity(page_keys.len());
192        data.resize(page_keys.len(), Bytes::new());
193        let mut cache_miss_range = vec![];
194        let mut cache_miss_idx = vec![];
195        let last_index = page_keys.len() - 1;
196        // TODO: Avoid copy as much as possible.
197        for (i, page_key) in page_keys.iter().enumerate() {
198            match self.get_page(key, *page_key) {
199                Some(page) => {
200                    CACHE_HIT.with_label_values(&[INDEX_CONTENT_TYPE]).inc();
201                    metrics.cache_hit += 1;
202                    metrics.page_bytes += page.len() as u64;
203                    data[i] = page;
204                }
205                None => {
206                    CACHE_MISS.with_label_values(&[INDEX_CONTENT_TYPE]).inc();
207                    metrics.cache_miss += 1;
208                    let base_offset = page_key.page_id * self.page_size;
209                    let pruned_size = if i == last_index {
210                        prune_size(page_keys.iter(), file_size, self.page_size)
211                    } else {
212                        self.page_size
213                    };
214                    cache_miss_range.push(base_offset..base_offset + pruned_size);
215                    cache_miss_idx.push(i);
216                }
217            }
218        }
219        if !cache_miss_range.is_empty() {
220            let pages = load(cache_miss_range).await?;
221            for (i, page) in cache_miss_idx.into_iter().zip(pages) {
222                let page_key = page_keys[i];
223                metrics.page_bytes += page.len() as u64;
224                data[i] = page.clone();
225                self.put_page(key, page_key, page.clone());
226            }
227        }
228        let buffer = Buffer::from_iter(data);
229        Ok((
230            buffer
231                .slice(PageKey::calculate_range(offset, size, self.page_size))
232                .to_vec(),
233            metrics,
234        ))
235    }
236
237    /// Like [`IndexCache::get_or_load`] for several ranges, but loads the missing pages of
238    /// all ranges with a single `load` call so a cold read costs one round trip instead of
239    /// one per range.
240    async fn get_or_load_vec<F, Fut, E>(
241        &self,
242        key: K,
243        file_size: u64,
244        ranges: &[Range<u64>],
245        load: F,
246    ) -> Result<(Vec<Bytes>, IndexCacheMetrics), E>
247    where
248        F: FnOnce(Vec<Range<u64>>) -> Fut,
249        Fut: Future<Output = Result<Vec<Bytes>, E>>,
250        E: std::error::Error,
251    {
252        let mut metrics = IndexCacheMetrics::default();
253        let mut pages: HashMap<u64, Bytes> = HashMap::new();
254        let mut missing: Vec<u64> = Vec::new();
255        let mut seen = HashSet::new();
256        for range in ranges.iter().filter(|r| r.end > r.start) {
257            let size = (range.end - range.start) as u32;
258            for page_key in PageKey::generate_page_keys(range.start, size, self.page_size) {
259                if !seen.insert(page_key.page_id) {
260                    continue;
261                }
262                metrics.num_pages += 1;
263                match self.get_page(key, page_key) {
264                    Some(page) => {
265                        CACHE_HIT.with_label_values(&[INDEX_CONTENT_TYPE]).inc();
266                        metrics.cache_hit += 1;
267                        metrics.page_bytes += page.len() as u64;
268                        pages.insert(page_key.page_id, page);
269                    }
270                    None => {
271                        CACHE_MISS.with_label_values(&[INDEX_CONTENT_TYPE]).inc();
272                        metrics.cache_miss += 1;
273                        missing.push(page_key.page_id);
274                    }
275                }
276            }
277        }
278
279        if !missing.is_empty() {
280            missing.sort_unstable();
281            let load_ranges = missing
282                .iter()
283                .map(|page_id| {
284                    let start = page_id * self.page_size;
285                    start..(start + self.page_size).min(file_size)
286                })
287                .collect();
288            let loaded = load(load_ranges).await?;
289            debug_assert_eq!(loaded.len(), missing.len());
290            for (page_id, page) in missing.into_iter().zip(loaded) {
291                metrics.page_bytes += page.len() as u64;
292                self.put_page(key, PageKey { page_id }, page.clone());
293                pages.insert(page_id, page);
294            }
295        }
296
297        let result = ranges
298            .iter()
299            .map(|range| {
300                let size = (range.end - range.start) as u32;
301                if size == 0 {
302                    return Bytes::new();
303                }
304                let buffer = Buffer::from_iter(
305                    PageKey::generate_page_keys(range.start, size, self.page_size)
306                        .map(|p| pages[&p.page_id].clone()),
307                );
308                buffer
309                    .slice(PageKey::calculate_range(range.start, size, self.page_size))
310                    .to_bytes()
311            })
312            .collect();
313        Ok((result, metrics))
314    }
315
316    fn get_page(&self, key: K, page_key: PageKey) -> Option<Bytes> {
317        self.index.get(&(key, page_key))
318    }
319
320    fn put_page(&self, key: K, page_key: PageKey, value: Bytes) {
321        // Clones the value to ensure it doesn't reference a larger buffer.
322        let value = Bytes::from(value.to_vec());
323        CACHE_BYTES
324            .with_label_values(&[INDEX_CONTENT_TYPE])
325            .add((self.weight_of_content)(&(key, page_key), &value).into());
326        self.index.insert((key, page_key), value);
327    }
328
329    /// Invalidates all cache entries whose keys satisfy `predicate`.
330    pub fn invalidate_if<F>(&self, predicate: F)
331    where
332        F: Fn(&K) -> bool + Send + Sync + 'static,
333    {
334        let predicate = Arc::new(predicate);
335        let metadata_predicate = Arc::clone(&predicate);
336
337        self.index_metadata
338            .invalidate_entries_if(move |key, _| metadata_predicate(key))
339            .expect("cache should support invalidation closures");
340
341        self.index
342            .invalidate_entries_if(move |(key, _), _| predicate(key))
343            .expect("cache should support invalidation closures");
344    }
345}
346
347/// Prunes the size of the last page based on the indexes.
348/// We have following cases:
349/// 1. The rest file size is less than the page size, read to the end of the file.
350/// 2. Otherwise, read the page size.
351fn prune_size<'a>(
352    indexes: impl Iterator<Item = &'a PageKey>,
353    file_size: u64,
354    page_size: u64,
355) -> u64 {
356    let last_page_start = indexes.last().map(|i| i.page_id * page_size).unwrap_or(0);
357    page_size.min(file_size - last_page_start)
358}