1use std::convert::TryFrom;
16use std::sync::Arc;
17
18use common_base::range_read::RangeReader;
19use common_telemetry::warn;
20use greptime_proto::v1::index::{BloomFilterMeta, InvertedIndexMeta, InvertedIndexMetas};
21use index::bitmap::BitmapType;
22use index::bloom_filter::reader::{BloomFilterReader, BloomFilterReaderImpl};
23use index::fulltext_index::Config as FulltextConfig;
24use index::inverted_index::format::reader::{InvertedIndexBlobReader, InvertedIndexReader};
25use index::target::IndexTarget;
26use puffin::blob_metadata::BlobMetadata;
27use puffin::puffin_manager::{PuffinManager, PuffinReader};
28use serde_json::{Map, Value, json};
29use store_api::sst_entry::{
30 PUFFIN_INDEX_TYPE_BLOOM_FILTER, PUFFIN_INDEX_TYPE_FULLTEXT_BLOOM,
31 PUFFIN_INDEX_TYPE_FULLTEXT_TANTIVY, PUFFIN_INDEX_TYPE_INVERTED, PuffinIndexMetaEntry,
32};
33use store_api::storage::{ColumnId, RegionGroup, RegionId, RegionNumber, RegionSeq, TableId};
34
35use crate::cache::index::bloom_filter_index::{
36 BloomFilterIndexCacheRef, CachedBloomFilterIndexBlobReader, Tag,
37};
38use crate::cache::index::inverted_index::{CachedInvertedIndexBlobReader, InvertedIndexCacheRef};
39use crate::sst::file::RegionIndexId;
40use crate::sst::index::bloom_filter::INDEX_BLOB_TYPE as BLOOM_BLOB_TYPE;
41use crate::sst::index::fulltext_index::{
42 INDEX_BLOB_TYPE_BLOOM as FULLTEXT_BLOOM_BLOB_TYPE,
43 INDEX_BLOB_TYPE_TANTIVY as FULLTEXT_TANTIVY_BLOB_TYPE,
44};
45use crate::sst::index::inverted_index::INDEX_BLOB_TYPE as INVERTED_BLOB_TYPE;
46use crate::sst::index::puffin_manager::{SstPuffinManager, SstPuffinReader};
47
48const TARGET_TYPE_UNKNOWN: &str = "unknown";
49
50const TARGET_TYPE_COLUMN: &str = "column";
51
52pub(crate) struct IndexEntryContext<'a> {
53 pub(crate) table_dir: &'a str,
54 pub(crate) index_file_path: &'a str,
55 pub(crate) region_id: RegionId,
56 pub(crate) table_id: TableId,
57 pub(crate) region_number: RegionNumber,
58 pub(crate) region_group: RegionGroup,
59 pub(crate) region_sequence: RegionSeq,
60 pub(crate) file_id: &'a str,
61 pub(crate) index_file_size: Option<u64>,
62 pub(crate) node_id: Option<u64>,
63}
64
65pub(crate) async fn collect_index_entries_from_puffin(
67 manager: SstPuffinManager,
68 region_index_id: RegionIndexId,
69 context: IndexEntryContext<'_>,
70 bloom_filter_cache: Option<BloomFilterIndexCacheRef>,
71 inverted_index_cache: Option<InvertedIndexCacheRef>,
72) -> Vec<PuffinIndexMetaEntry> {
73 let mut entries = Vec::new();
74
75 let reader = match manager.reader(®ion_index_id).await {
76 Ok(reader) => reader,
77 Err(err) => {
78 warn!(
79 err;
80 "Failed to open puffin index file, table_dir: {}, file_id: {}",
81 context.table_dir,
82 context.file_id
83 );
84 return entries;
85 }
86 };
87
88 let file_metadata = match reader.metadata().await {
89 Ok(metadata) => metadata,
90 Err(err) => {
91 warn!(
92 err;
93 "Failed to read puffin file metadata, table_dir: {}, file_id: {}",
94 context.table_dir,
95 context.file_id
96 );
97 return entries;
98 }
99 };
100
101 for blob in &file_metadata.blobs {
102 match BlobIndexTypeTargetKey::from_blob_type(&blob.blob_type) {
103 Some(BlobIndexTypeTargetKey::BloomFilter(target_key)) => {
104 let bloom_meta = try_read_bloom_meta(
105 &reader,
106 region_index_id,
107 blob.blob_type.as_str(),
108 target_key,
109 bloom_filter_cache.as_ref(),
110 Tag::Skipping,
111 &context,
112 )
113 .await;
114
115 let bloom_value = bloom_meta.as_deref().map(bloom_meta_value);
116 let (target_type, target_json) = decode_target_info(target_key);
117 let meta_json = build_meta_json(bloom_value, None, None);
118 let entry = build_index_entry(
119 &context,
120 PUFFIN_INDEX_TYPE_BLOOM_FILTER,
121 target_type,
122 target_key.to_string(),
123 target_json,
124 blob.length as u64,
125 meta_json,
126 );
127 entries.push(entry);
128 }
129 Some(BlobIndexTypeTargetKey::FulltextBloom(target_key)) => {
130 let bloom_meta = try_read_bloom_meta(
131 &reader,
132 region_index_id,
133 blob.blob_type.as_str(),
134 target_key,
135 bloom_filter_cache.as_ref(),
136 Tag::Fulltext,
137 &context,
138 )
139 .await;
140
141 let bloom_value = bloom_meta.as_deref().map(bloom_meta_value);
142 let fulltext_value = Some(fulltext_meta_value(blob));
143 let (target_type, target_json) = decode_target_info(target_key);
144 let meta_json = build_meta_json(bloom_value, fulltext_value, None);
145 let entry = build_index_entry(
146 &context,
147 PUFFIN_INDEX_TYPE_FULLTEXT_BLOOM,
148 target_type,
149 target_key.to_string(),
150 target_json,
151 blob.length as u64,
152 meta_json,
153 );
154 entries.push(entry);
155 }
156 Some(BlobIndexTypeTargetKey::FulltextTantivy(target_key)) => {
157 let fulltext_value = Some(fulltext_meta_value(blob));
158 let (target_type, target_json) = decode_target_info(target_key);
159 let meta_json = build_meta_json(None, fulltext_value, None);
160 let entry = build_index_entry(
161 &context,
162 PUFFIN_INDEX_TYPE_FULLTEXT_TANTIVY,
163 target_type,
164 target_key.to_string(),
165 target_json,
166 blob.length as u64,
167 meta_json,
168 );
169 entries.push(entry);
170 }
171 Some(BlobIndexTypeTargetKey::Inverted) => {
172 let mut inverted_entries = collect_inverted_entries(
173 &reader,
174 region_index_id,
175 inverted_index_cache.as_ref(),
176 &context,
177 )
178 .await;
179 entries.append(&mut inverted_entries);
180 }
181 None => {}
182 }
183 }
184
185 entries
186}
187
188async fn collect_inverted_entries(
189 reader: &SstPuffinReader,
190 region_index_id: RegionIndexId,
191 cache: Option<&InvertedIndexCacheRef>,
192 context: &IndexEntryContext<'_>,
193) -> Vec<PuffinIndexMetaEntry> {
194 let file_id = region_index_id.file_id();
196
197 let guard = match reader.blob(INVERTED_BLOB_TYPE).await {
198 Ok(guard) => guard,
199 Err(err) => {
200 warn!(
201 err;
202 "Failed to open inverted index blob, table_dir: {}, file_id: {}",
203 context.table_dir,
204 context.file_id
205 );
206 return Vec::new();
207 }
208 };
209
210 let blob_reader = match guard.reader().await {
211 Ok(reader) => reader,
212 Err(err) => {
213 warn!(
214 err;
215 "Failed to build inverted index blob reader, table_dir: {}, file_id: {}",
216 context.table_dir,
217 context.file_id
218 );
219 return Vec::new();
220 }
221 };
222
223 let blob_size = blob_reader
224 .metadata()
225 .await
226 .ok()
227 .map(|meta| meta.content_length);
228 let metas = if let (Some(cache), Some(blob_size)) = (cache, blob_size) {
229 let reader = CachedInvertedIndexBlobReader::new(
230 file_id,
231 region_index_id.version,
232 blob_size,
233 InvertedIndexBlobReader::new(blob_reader),
234 cache.clone(),
235 );
236 match reader.metadata(None).await {
237 Ok(metas) => metas,
238 Err(err) => {
239 warn!(
240 err;
241 "Failed to read inverted index metadata, table_dir: {}, file_id: {}",
242 context.table_dir,
243 context.file_id
244 );
245 return Vec::new();
246 }
247 }
248 } else {
249 let reader = InvertedIndexBlobReader::new(blob_reader);
250 match reader.metadata(None).await {
251 Ok(metas) => metas,
252 Err(err) => {
253 warn!(
254 err;
255 "Failed to read inverted index metadata, table_dir: {}, file_id: {}",
256 context.table_dir,
257 context.file_id
258 );
259 return Vec::new();
260 }
261 }
262 };
263
264 build_inverted_entries(context, metas.as_ref())
265}
266
267fn build_inverted_entries(
268 context: &IndexEntryContext<'_>,
269 metas: &InvertedIndexMetas,
270) -> Vec<PuffinIndexMetaEntry> {
271 let mut entries = Vec::new();
272 for (name, meta) in &metas.metas {
273 let (target_type, target_json) = decode_target_info(name);
274 let inverted_value = inverted_meta_value(meta, metas);
275 let meta_json = build_meta_json(None, None, Some(inverted_value));
276 let entry = build_index_entry(
277 context,
278 PUFFIN_INDEX_TYPE_INVERTED,
279 target_type,
280 name.clone(),
281 target_json,
282 meta.inverted_index_size,
283 meta_json,
284 );
285 entries.push(entry);
286 }
287 entries
288}
289
290async fn try_read_bloom_meta(
291 reader: &SstPuffinReader,
292 region_index_id: RegionIndexId,
293 blob_type: &str,
294 target_key: &str,
295 cache: Option<&BloomFilterIndexCacheRef>,
296 tag: Tag,
297 context: &IndexEntryContext<'_>,
298) -> Option<Arc<BloomFilterMeta>> {
299 let column_id = decode_column_id(target_key);
300
301 match reader.blob(blob_type).await {
303 Ok(guard) => match guard.reader().await {
304 Ok(blob_reader) => {
305 let blob_size = blob_reader
306 .metadata()
307 .await
308 .ok()
309 .map(|meta| meta.content_length);
310 let bloom_reader = BloomFilterReaderImpl::new(blob_reader);
311 let result = match (cache, column_id, blob_size) {
312 (Some(cache), Some(column_id), Some(blob_size)) => {
313 CachedBloomFilterIndexBlobReader::new(
314 region_index_id.file_id(),
315 region_index_id.version,
316 column_id,
317 tag,
318 blob_size,
319 bloom_reader,
320 cache.clone(),
321 )
322 .metadata(None)
323 .await
324 }
325 _ => bloom_reader.metadata(None).await,
326 };
327
328 match result {
329 Ok(meta) => Some(meta),
330 Err(err) => {
331 warn!(
332 err;
333 "Failed to read index metadata, table_dir: {}, file_id: {}, blob: {}",
334 context.table_dir,
335 context.file_id,
336 blob_type
337 );
338 None
339 }
340 }
341 }
342 Err(err) => {
343 warn!(
344 err;
345 "Failed to open index blob reader, table_dir: {}, file_id: {}, blob: {}",
346 context.table_dir,
347 context.file_id,
348 blob_type
349 );
350 None
351 }
352 },
353 Err(err) => {
354 warn!(
355 err;
356 "Failed to open index blob, table_dir: {}, file_id: {}, blob: {}",
357 context.table_dir,
358 context.file_id,
359 blob_type
360 );
361 None
362 }
363 }
364}
365
366fn decode_target_info(target_key: &str) -> (String, String) {
367 match IndexTarget::decode(target_key) {
368 Ok(IndexTarget::ColumnId(id)) => (
369 TARGET_TYPE_COLUMN.to_string(),
370 json!({ "column": id }).to_string(),
371 ),
372 _ => (
373 TARGET_TYPE_UNKNOWN.to_string(),
374 json!({ "error": "failed_to_decode" }).to_string(),
375 ),
376 }
377}
378
379fn decode_column_id(target_key: &str) -> Option<ColumnId> {
380 match IndexTarget::decode(target_key) {
381 Ok(IndexTarget::ColumnId(id)) => Some(id),
382 _ => None,
383 }
384}
385
386fn bloom_meta_value(meta: &BloomFilterMeta) -> Value {
387 json!({
388 "rows_per_segment": meta.rows_per_segment,
389 "segment_count": meta.segment_count,
390 "row_count": meta.row_count,
391 "bloom_filter_size": meta.bloom_filter_size,
392 })
393}
394
395fn fulltext_meta_value(blob: &BlobMetadata) -> Value {
396 let config = FulltextConfig::from_blob_metadata(blob).unwrap_or_default();
397 json!({
398 "analyzer": config.analyzer.to_str(),
399 "case_sensitive": config.case_sensitive,
400 })
401}
402
403fn inverted_meta_value(meta: &InvertedIndexMeta, metas: &InvertedIndexMetas) -> Value {
404 let bitmap_type = BitmapType::try_from(meta.bitmap_type)
405 .map(|bt| format!("{:?}", bt))
406 .unwrap_or_else(|_| meta.bitmap_type.to_string());
407 json!({
408 "bitmap_type": bitmap_type,
409 "base_offset": meta.base_offset,
410 "inverted_index_size": meta.inverted_index_size,
411 "relative_fst_offset": meta.relative_fst_offset,
412 "fst_size": meta.fst_size,
413 "relative_null_bitmap_offset": meta.relative_null_bitmap_offset,
414 "null_bitmap_size": meta.null_bitmap_size,
415 "segment_row_count": metas.segment_row_count,
416 "total_row_count": metas.total_row_count,
417 })
418}
419
420fn build_meta_json(
421 bloom: Option<Value>,
422 fulltext: Option<Value>,
423 inverted: Option<Value>,
424) -> Option<String> {
425 let mut map = Map::new();
426 if let Some(value) = bloom {
427 map.insert("bloom".to_string(), value);
428 }
429 if let Some(value) = fulltext {
430 map.insert("fulltext".to_string(), value);
431 }
432 if let Some(value) = inverted {
433 map.insert("inverted".to_string(), value);
434 }
435 if map.is_empty() {
436 None
437 } else {
438 Some(Value::Object(map).to_string())
439 }
440}
441
442enum BlobIndexTypeTargetKey<'a> {
443 BloomFilter(&'a str),
444 FulltextBloom(&'a str),
445 FulltextTantivy(&'a str),
446 Inverted,
447}
448
449impl<'a> BlobIndexTypeTargetKey<'a> {
450 fn from_blob_type(blob_type: &'a str) -> Option<Self> {
451 if let Some(target_key) = Self::target_key_from_blob(blob_type, BLOOM_BLOB_TYPE) {
452 Some(BlobIndexTypeTargetKey::BloomFilter(target_key))
453 } else if let Some(target_key) =
454 Self::target_key_from_blob(blob_type, FULLTEXT_BLOOM_BLOB_TYPE)
455 {
456 Some(BlobIndexTypeTargetKey::FulltextBloom(target_key))
457 } else if let Some(target_key) =
458 Self::target_key_from_blob(blob_type, FULLTEXT_TANTIVY_BLOB_TYPE)
459 {
460 Some(BlobIndexTypeTargetKey::FulltextTantivy(target_key))
461 } else if blob_type == INVERTED_BLOB_TYPE {
462 Some(BlobIndexTypeTargetKey::Inverted)
463 } else {
464 None
465 }
466 }
467
468 fn target_key_from_blob(blob_type: &'a str, prefix: &str) -> Option<&'a str> {
469 blob_type
471 .strip_prefix(prefix)
472 .and_then(|suffix| suffix.strip_prefix('-'))
473 }
474}
475
476fn build_index_entry(
477 context: &IndexEntryContext<'_>,
478 index_type: &str,
479 target_type: String,
480 target_key: String,
481 target_json: String,
482 blob_size: u64,
483 meta_json: Option<String>,
484) -> PuffinIndexMetaEntry {
485 PuffinIndexMetaEntry {
486 table_dir: context.table_dir.to_string(),
487 index_file_path: context.index_file_path.to_string(),
488 region_id: context.region_id,
489 table_id: context.table_id,
490 region_number: context.region_number,
491 region_group: context.region_group,
492 region_sequence: context.region_sequence,
493 file_id: context.file_id.to_string(),
494 index_file_size: context.index_file_size,
495 index_type: index_type.to_string(),
496 target_type,
497 target_key,
498 target_json,
499 blob_size,
500 meta_json,
501 node_id: context.node_id,
502 }
503}