1pub 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
30const INDEX_METADATA_TYPE: &str = "index_metadata";
32const INDEX_CONTENT_TYPE: &str = "index_content";
34
35#[derive(Debug, Default, Clone)]
37pub struct IndexCacheMetrics {
38 pub cache_hit: usize,
40 pub cache_miss: usize,
42 pub num_pages: usize,
44 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 fn calculate_page_id(offset: u64, page_size: u64) -> u64 {
56 offset / page_size
57 }
58
59 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 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 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
89pub struct IndexCache<K, M> {
91 index_metadata: moka::sync::Cache<K, Arc<M>>,
93 index: moka::sync::Cache<(K, PageKey), Bytes>,
95 page_size: u64,
97
98 weight_of_metadata: fn(&K, &Arc<M>) -> u32,
100 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 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 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 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 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 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 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
347fn 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}