Skip to main content

mito2/engine/
puffin_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
15use 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
65/// Collect index metadata entries present in the SST puffin file.
66pub(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(&region_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    // Read the inverted index blob and surface its per-column metadata entries.
195    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    // Failures are logged but do not abort the overall metadata collection.
302    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 types encode their target as "<prefix>-<target>".
470        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}