1use core::ops::Range;
16use std::sync::Arc;
17use std::time::Instant;
18
19use api::v1::index::InvertedIndexMetas;
20use async_trait::async_trait;
21use bytes::Bytes;
22use index::inverted_index::error::Result;
23use index::inverted_index::format::reader::{InvertedIndexReadMetrics, InvertedIndexReader};
24use prost::Message;
25use store_api::storage::{FileId, IndexVersion};
26
27use crate::cache::index::{INDEX_METADATA_TYPE, IndexCache, PageKey};
28use crate::metrics::{CACHE_HIT, CACHE_MISS};
29
30const INDEX_TYPE_INVERTED_INDEX: &str = "inverted_index";
31
32pub type InvertedIndexCache = IndexCache<(FileId, IndexVersion), InvertedIndexMetas>;
34pub type InvertedIndexCacheRef = Arc<InvertedIndexCache>;
35
36impl InvertedIndexCache {
37 pub fn new(index_metadata_cap: u64, index_content_cap: u64, page_size: u64) -> Self {
39 Self::new_with_weighter(
40 index_metadata_cap,
41 index_content_cap,
42 page_size,
43 INDEX_TYPE_INVERTED_INDEX,
44 inverted_index_metadata_weight,
45 inverted_index_content_weight,
46 )
47 }
48
49 pub fn invalidate_file(&self, file_id: FileId) {
51 self.invalidate_if(move |key| key.0 == file_id);
52 }
53}
54
55fn inverted_index_metadata_weight(k: &(FileId, IndexVersion), v: &Arc<InvertedIndexMetas>) -> u32 {
57 (k.0.as_bytes().len() + size_of::<IndexVersion>() + v.encoded_len()) as u32
58}
59
60fn inverted_index_content_weight((k, _): &((FileId, IndexVersion), PageKey), v: &Bytes) -> u32 {
62 (k.0.as_bytes().len() + size_of::<IndexVersion>() + v.len()) as u32
63}
64
65pub struct CachedInvertedIndexBlobReader<R> {
67 file_id: FileId,
68 index_version: IndexVersion,
69 blob_size: u64,
70 inner: R,
71 cache: InvertedIndexCacheRef,
72}
73
74impl<R> CachedInvertedIndexBlobReader<R> {
75 pub fn new(
77 file_id: FileId,
78 index_version: IndexVersion,
79 blob_size: u64,
80 inner: R,
81 cache: InvertedIndexCacheRef,
82 ) -> Self {
83 Self {
84 file_id,
85 index_version,
86 blob_size,
87 inner,
88 cache,
89 }
90 }
91}
92
93#[async_trait]
94impl<R: InvertedIndexReader> InvertedIndexReader for CachedInvertedIndexBlobReader<R> {
95 async fn range_read<'a>(
96 &self,
97 offset: u64,
98 size: u32,
99 metrics: Option<&'a mut InvertedIndexReadMetrics>,
100 ) -> Result<Vec<u8>> {
101 let start = metrics.as_ref().map(|_| Instant::now());
102
103 let inner = &self.inner;
104 let (result, cache_metrics) = self
105 .cache
106 .get_or_load(
107 (self.file_id, self.index_version),
108 self.blob_size,
109 offset,
110 size,
111 move |ranges| async move { inner.read_vec(&ranges, None).await },
112 )
113 .await?;
114
115 if let Some(m) = metrics {
116 m.total_bytes += cache_metrics.page_bytes;
117 m.total_ranges += cache_metrics.num_pages;
118 m.cache_hit += cache_metrics.cache_hit;
119 m.cache_miss += cache_metrics.cache_miss;
120 m.fetch_elapsed += start.unwrap().elapsed();
121 }
122
123 Ok(result)
124 }
125
126 async fn read_vec<'a>(
127 &self,
128 ranges: &[Range<u64>],
129 metrics: Option<&'a mut InvertedIndexReadMetrics>,
130 ) -> Result<Vec<Bytes>> {
131 let start = metrics.as_ref().map(|_| Instant::now());
132
133 let inner = &self.inner;
134 let (pages, total_cache_metrics) = self
135 .cache
136 .get_or_load_vec(
137 (self.file_id, self.index_version),
138 self.blob_size,
139 ranges,
140 move |ranges| async move { inner.read_vec(&ranges, None).await },
141 )
142 .await?;
143 if let Some(m) = metrics {
144 m.total_bytes += total_cache_metrics.page_bytes;
145 m.total_ranges += total_cache_metrics.num_pages;
146 m.cache_hit += total_cache_metrics.cache_hit;
147 m.cache_miss += total_cache_metrics.cache_miss;
148 m.fetch_elapsed += start.unwrap().elapsed();
149 }
150
151 Ok(pages)
152 }
153
154 async fn metadata<'a>(
155 &self,
156 metrics: Option<&'a mut InvertedIndexReadMetrics>,
157 ) -> Result<Arc<InvertedIndexMetas>> {
158 if let Some(cached) = self.cache.get_metadata((self.file_id, self.index_version)) {
159 CACHE_HIT.with_label_values(&[INDEX_METADATA_TYPE]).inc();
160 if let Some(m) = metrics {
161 m.cache_hit += 1;
162 }
163 Ok(cached)
164 } else {
165 let meta = self.inner.metadata(metrics).await?;
166 self.cache
167 .put_metadata((self.file_id, self.index_version), meta.clone());
168 CACHE_MISS.with_label_values(&[INDEX_METADATA_TYPE]).inc();
169 Ok(meta)
170 }
171 }
172}
173
174#[cfg(test)]
175mod test {
176 use std::num::NonZeroUsize;
177
178 use futures::stream;
179 use index::Bytes;
180 use index::bitmap::{Bitmap, BitmapType};
181 use index::inverted_index::format::reader::{InvertedIndexBlobReader, InvertedIndexReader};
182 use index::inverted_index::format::writer::{InvertedIndexBlobWriter, InvertedIndexWriter};
183 use prometheus::register_int_counter_vec;
184 use rand::{Rng, RngCore};
185
186 use super::*;
187 use crate::sst::index::store::InstrumentedStore;
188 use crate::test_util::TestEnv;
189
190 const FUZZ_REPEAT_TIMES: usize = 100;
192
193 #[test]
195 fn fuzz_index_calculation() {
196 let mut rng = rand::rng();
198 let mut data = vec![0u8; 1024 * 1024];
199 rng.fill_bytes(&mut data);
200
201 for _ in 0..FUZZ_REPEAT_TIMES {
202 let offset = rng.random_range(0..data.len() as u64);
203 let size = rng.random_range(0..data.len() as u32 - offset as u32);
204 let page_size: usize = rng.random_range(1..1024);
205
206 let indexes =
207 PageKey::generate_page_keys(offset, size, page_size as u64).collect::<Vec<_>>();
208 let page_num = indexes.len();
209 let mut read = Vec::with_capacity(size as usize);
210 for key in indexes.into_iter() {
211 let start = key.page_id as usize * page_size;
212 let page = if start + page_size < data.len() {
213 &data[start..start + page_size]
214 } else {
215 &data[start..]
216 };
217 read.extend_from_slice(page);
218 }
219 let expected_range = offset as usize..(offset + size as u64 as u64) as usize;
220 let read = read[PageKey::calculate_range(offset, size, page_size as u64)].to_vec();
221 if read != data.get(expected_range).unwrap() {
222 panic!(
223 "fuzz_read_index failed, offset: {}, size: {}, page_size: {}\nread len: {}, expected len: {}\nrange: {:?}, page num: {}",
224 offset,
225 size,
226 page_size,
227 read.len(),
228 size as usize,
229 PageKey::calculate_range(offset, size, page_size as u64),
230 page_num
231 );
232 }
233 }
234 }
235
236 fn unpack(fst_value: u64) -> [u32; 2] {
237 bytemuck::cast::<u64, [u32; 2]>(fst_value)
238 }
239
240 async fn create_inverted_index_blob() -> Vec<u8> {
241 let mut blob = Vec::new();
242 let mut writer = InvertedIndexBlobWriter::new(&mut blob);
243 writer
244 .add_index(
245 "tag0".to_string(),
246 Bitmap::from_lsb0_bytes(&[0b0000_0001, 0b0000_0000], BitmapType::Roaring),
247 Box::new(stream::iter(vec![
248 Ok((
249 Bytes::from("a"),
250 Bitmap::from_lsb0_bytes(&[0b0000_0001], BitmapType::Roaring),
251 )),
252 Ok((
253 Bytes::from("b"),
254 Bitmap::from_lsb0_bytes(&[0b0010_0000], BitmapType::Roaring),
255 )),
256 Ok((
257 Bytes::from("c"),
258 Bitmap::from_lsb0_bytes(&[0b0000_0001], BitmapType::Roaring),
259 )),
260 ])),
261 index::bitmap::BitmapType::Roaring,
262 )
263 .await
264 .unwrap();
265 writer
266 .add_index(
267 "tag1".to_string(),
268 Bitmap::from_lsb0_bytes(&[0b0000_0001, 0b0000_0000], BitmapType::Roaring),
269 Box::new(stream::iter(vec![
270 Ok((
271 Bytes::from("x"),
272 Bitmap::from_lsb0_bytes(&[0b0000_0001], BitmapType::Roaring),
273 )),
274 Ok((
275 Bytes::from("y"),
276 Bitmap::from_lsb0_bytes(&[0b0010_0000], BitmapType::Roaring),
277 )),
278 Ok((
279 Bytes::from("z"),
280 Bitmap::from_lsb0_bytes(&[0b0000_0001], BitmapType::Roaring),
281 )),
282 ])),
283 index::bitmap::BitmapType::Roaring,
284 )
285 .await
286 .unwrap();
287 writer
288 .finish(8, NonZeroUsize::new(1).unwrap())
289 .await
290 .unwrap();
291
292 blob
293 }
294
295 #[tokio::test]
296 async fn test_inverted_index_cache() {
297 let blob = create_inverted_index_blob().await;
298
299 let mut env = TestEnv::new().await;
301 let file_size = blob.len() as u64;
302 let index_version = 0;
303 let store = env.init_object_store_manager();
304 let temp_path = "data";
305 store.write(temp_path, blob).await.unwrap();
306 let store = InstrumentedStore::new(store);
307 let metric =
308 register_int_counter_vec!("test_bytes", "a counter for test", &["test"]).unwrap();
309 let counter = metric.with_label_values(&["test"]);
310 let range_reader = store
311 .range_reader("data", &counter, &counter)
312 .await
313 .unwrap();
314
315 let reader = InvertedIndexBlobReader::new(range_reader);
316 let cached_reader = CachedInvertedIndexBlobReader::new(
317 FileId::random(),
318 index_version,
319 file_size,
320 reader,
321 Arc::new(InvertedIndexCache::new(8192, 8192, 50)),
322 );
323 let metadata = cached_reader.metadata(None).await.unwrap();
324 assert_eq!(metadata.total_row_count, 8);
325 assert_eq!(metadata.segment_row_count, 1);
326 assert_eq!(metadata.metas.len(), 2);
327 let tag0 = metadata.metas.get("tag0").unwrap();
329 let stats0 = tag0.stats.as_ref().unwrap();
330 assert_eq!(stats0.distinct_count, 3);
331 assert_eq!(stats0.null_count, 1);
332 assert_eq!(stats0.min_value, Bytes::from("a"));
333 assert_eq!(stats0.max_value, Bytes::from("c"));
334 let fst0 = cached_reader
335 .fst(
336 tag0.base_offset + tag0.relative_fst_offset as u64,
337 tag0.fst_size,
338 None,
339 )
340 .await
341 .unwrap();
342 assert_eq!(fst0.len(), 3);
343 let [offset, size] = unpack(fst0.get(b"a").unwrap());
344 let bitmap = cached_reader
345 .bitmap(
346 tag0.base_offset + offset as u64,
347 size,
348 BitmapType::Roaring,
349 None,
350 )
351 .await
352 .unwrap();
353 assert_eq!(
354 bitmap,
355 Bitmap::from_lsb0_bytes(&[0b0000_0001], BitmapType::Roaring)
356 );
357 let [offset, size] = unpack(fst0.get(b"b").unwrap());
358 let bitmap = cached_reader
359 .bitmap(
360 tag0.base_offset + offset as u64,
361 size,
362 BitmapType::Roaring,
363 None,
364 )
365 .await
366 .unwrap();
367 assert_eq!(
368 bitmap,
369 Bitmap::from_lsb0_bytes(&[0b0010_0000], BitmapType::Roaring)
370 );
371 let [offset, size] = unpack(fst0.get(b"c").unwrap());
372 let bitmap = cached_reader
373 .bitmap(
374 tag0.base_offset + offset as u64,
375 size,
376 BitmapType::Roaring,
377 None,
378 )
379 .await
380 .unwrap();
381 assert_eq!(
382 bitmap,
383 Bitmap::from_lsb0_bytes(&[0b0000_0001], BitmapType::Roaring)
384 );
385
386 let tag1 = metadata.metas.get("tag1").unwrap();
388 let stats1 = tag1.stats.as_ref().unwrap();
389 assert_eq!(stats1.distinct_count, 3);
390 assert_eq!(stats1.null_count, 1);
391 assert_eq!(stats1.min_value, Bytes::from("x"));
392 assert_eq!(stats1.max_value, Bytes::from("z"));
393 let fst1 = cached_reader
394 .fst(
395 tag1.base_offset + tag1.relative_fst_offset as u64,
396 tag1.fst_size,
397 None,
398 )
399 .await
400 .unwrap();
401 assert_eq!(fst1.len(), 3);
402 let [offset, size] = unpack(fst1.get(b"x").unwrap());
403 let bitmap = cached_reader
404 .bitmap(
405 tag1.base_offset + offset as u64,
406 size,
407 BitmapType::Roaring,
408 None,
409 )
410 .await
411 .unwrap();
412 assert_eq!(
413 bitmap,
414 Bitmap::from_lsb0_bytes(&[0b0000_0001], BitmapType::Roaring)
415 );
416 let [offset, size] = unpack(fst1.get(b"y").unwrap());
417 let bitmap = cached_reader
418 .bitmap(
419 tag1.base_offset + offset as u64,
420 size,
421 BitmapType::Roaring,
422 None,
423 )
424 .await
425 .unwrap();
426 assert_eq!(
427 bitmap,
428 Bitmap::from_lsb0_bytes(&[0b0010_0000], BitmapType::Roaring)
429 );
430 let [offset, size] = unpack(fst1.get(b"z").unwrap());
431 let bitmap = cached_reader
432 .bitmap(
433 tag1.base_offset + offset as u64,
434 size,
435 BitmapType::Roaring,
436 None,
437 )
438 .await
439 .unwrap();
440 assert_eq!(
441 bitmap,
442 Bitmap::from_lsb0_bytes(&[0b0000_0001], BitmapType::Roaring)
443 );
444
445 let mut rng = rand::rng();
447 for _ in 0..FUZZ_REPEAT_TIMES {
448 let offset = rng.random_range(0..file_size);
449 let size = rng.random_range(0..file_size as u32 - offset as u32);
450 let expected = cached_reader.range_read(offset, size, None).await.unwrap();
451 let inner = &cached_reader.inner;
452 let (read, _cache_metrics) = cached_reader
453 .cache
454 .get_or_load(
455 (cached_reader.file_id, cached_reader.index_version),
456 file_size,
457 offset,
458 size,
459 |ranges| async move { inner.read_vec(&ranges, None).await },
460 )
461 .await
462 .unwrap();
463 assert_eq!(read, expected);
464 }
465 }
466
467 #[tokio::test]
468 async fn test_get_or_load_vec_loads_missing_pages_once() {
469 use std::sync::atomic::{AtomicUsize, Ordering};
470
471 let mut rng = rand::rng();
472 let mut data = vec![0u8; 64 * 1024];
473 rng.fill_bytes(&mut data);
474 let file_size = data.len() as u64;
475 let cache = InvertedIndexCache::new(1024 * 1024, 1024 * 1024, 1000);
476 let key = (FileId::random(), 0);
477
478 for _ in 0..FUZZ_REPEAT_TIMES {
479 cache.invalidate_file(key.0);
480 let ranges = (0..rng.random_range(1..20))
481 .map(|_| {
482 let start = rng.random_range(0..file_size);
483 start..rng.random_range(start..=file_size.min(start + 5000))
484 })
485 .collect::<Vec<_>>();
486 let expected = ranges
487 .iter()
488 .map(|r| bytes::Bytes::copy_from_slice(&data[r.start as usize..r.end as usize]))
489 .collect::<Vec<_>>();
490
491 let loads = AtomicUsize::new(0);
492 let load = |pages: Vec<Range<u64>>| {
493 loads.fetch_add(1, Ordering::Relaxed);
494 assert!(pages.windows(2).all(|w| w[0].end <= w[1].start));
496 let data = &data;
497 async move {
498 Ok::<_, std::io::Error>(
499 pages
500 .iter()
501 .map(|r| {
502 bytes::Bytes::copy_from_slice(
503 &data[r.start as usize..r.end as usize],
504 )
505 })
506 .collect(),
507 )
508 }
509 };
510 let (cold, cold_metrics) = cache
511 .get_or_load_vec(key, file_size, &ranges, load)
512 .await
513 .unwrap();
514 assert_eq!(cold, expected);
515 let cold_loads = usize::from(cold_metrics.cache_miss > 0);
516 assert_eq!(loads.load(Ordering::Relaxed), cold_loads);
517
518 let (warm, metrics) = cache
519 .get_or_load_vec(key, file_size, &ranges, load)
520 .await
521 .unwrap();
522 assert_eq!(warm, expected);
523 assert_eq!(metrics.cache_miss, 0);
524 assert_eq!(loads.load(Ordering::Relaxed), cold_loads);
525 }
526 }
527}