Skip to main content

mito2/sst/index/inverted_index/
creator.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::collections::HashSet;
16use std::num::NonZeroUsize;
17use std::sync::Arc;
18use std::sync::atomic::AtomicUsize;
19
20use api::v1::SemanticType;
21use common_telemetry::{debug, warn};
22use datatypes::arrow::record_batch::RecordBatch;
23use datatypes::vectors::Helper;
24use index::inverted_index::create::InvertedIndexCreator;
25use index::inverted_index::create::sort::external_sort::ExternalSorter;
26use index::inverted_index::create::sort_create::SortIndexCreator;
27use index::inverted_index::format::writer::InvertedIndexBlobWriter;
28use index::target::IndexTarget;
29use mito_codec::index::{IndexValueCodec, IndexValuesCodec};
30use mito_codec::row_converter::sparse::SparsePrimaryKeyView;
31use mito_codec::row_converter::{SortField, SparseOffsetsCache};
32use puffin::puffin_manager::{PuffinWriter, PutOptions};
33use smallvec::SmallVec;
34use snafu::{ResultExt, ensure};
35use store_api::codec::PrimaryKeyEncoding;
36use store_api::metadata::RegionMetadataRef;
37use store_api::storage::{ColumnId, FileId};
38use tokio::io::duplex;
39use tokio_util::compat::{TokioAsyncReadCompatExt, TokioAsyncWriteCompatExt};
40
41use crate::error::{
42    BiErrorsSnafu, DecodeSnafu, EncodeSnafu, IndexFinishSnafu, OperateAbortedIndexSnafu,
43    PuffinAddBlobSnafu, PushIndexValueSnafu, Result,
44};
45use crate::read::Batch;
46use crate::sst::index::TYPE_INVERTED_INDEX;
47use crate::sst::index::column::column_index_rows;
48use crate::sst::index::intermediate::{
49    IntermediateLocation, IntermediateManager, TempFileProvider,
50};
51use crate::sst::index::inverted_index::INDEX_BLOB_TYPE;
52use crate::sst::index::primary_key::PrimaryKeyRuns;
53use crate::sst::index::puffin_manager::SstPuffinWriter;
54use crate::sst::index::statistics::{ByteCount, RowCount, Statistics};
55
56/// The minimum memory usage threshold for one column.
57const MIN_MEMORY_USAGE_THRESHOLD_PER_COLUMN: usize = 1024 * 1024; // 1MB
58
59/// The buffer size for the pipe used to send index data to the puffin blob.
60const PIPE_BUFFER_SIZE_FOR_SENDING_BLOB: usize = 8192;
61
62/// `InvertedIndexer` creates inverted index for SST files.
63pub struct InvertedIndexer {
64    /// The index creator.
65    index_creator: Box<dyn InvertedIndexCreator>,
66    /// The provider of intermediate files.
67    temp_file_provider: Arc<TempFileProvider>,
68
69    /// Codec for decoding primary keys.
70    codec: IndexValuesCodec,
71    /// Reusable buffer for encoding index values.
72    value_buf: Vec<u8>,
73    /// Scratch offsets shared by indexed tags of one sparse primary key.
74    pk_offsets: SparseOffsetsCache,
75
76    /// Statistics of index creation.
77    stats: Statistics,
78    /// Whether the index creation is aborted.
79    aborted: bool,
80
81    /// The memory usage of the index creator.
82    memory_usage: Arc<AtomicUsize>,
83
84    /// Ids of indexed columns and their encoded target keys.
85    indexed_column_ids: Vec<(ColumnId, String)>,
86
87    /// Region metadata for column lookups.
88    metadata: RegionMetadataRef,
89}
90
91impl InvertedIndexer {
92    /// Creates a new `InvertedIndexer`.
93    /// Should ensure that the number of tag columns is greater than 0.
94    pub fn new(
95        sst_file_id: FileId,
96        metadata: &RegionMetadataRef,
97        intermediate_manager: IntermediateManager,
98        memory_usage_threshold: Option<usize>,
99        segment_row_count: NonZeroUsize,
100        indexed_column_ids: HashSet<ColumnId>,
101    ) -> Self {
102        let temp_file_provider = Arc::new(TempFileProvider::new(
103            IntermediateLocation::new(&metadata.region_id, &sst_file_id),
104            intermediate_manager,
105        ));
106
107        let memory_usage = Arc::new(AtomicUsize::new(0));
108
109        let sorter = ExternalSorter::factory(
110            temp_file_provider.clone() as _,
111            Some(MIN_MEMORY_USAGE_THRESHOLD_PER_COLUMN),
112            memory_usage.clone(),
113            memory_usage_threshold,
114        );
115        let index_creator = Box::new(SortIndexCreator::new(sorter, segment_row_count));
116
117        let codec = IndexValuesCodec::from_tag_columns(
118            metadata.primary_key_encoding,
119            metadata.primary_key_columns(),
120        );
121        let indexed_column_ids = indexed_column_ids
122            .into_iter()
123            .map(|col_id| {
124                let target_key = format!("{}", IndexTarget::ColumnId(col_id));
125                (col_id, target_key)
126            })
127            .collect();
128        Self {
129            codec,
130            index_creator,
131            temp_file_provider,
132            value_buf: vec![],
133            pk_offsets: SparseOffsetsCache::new(),
134            stats: Statistics::new(TYPE_INVERTED_INDEX),
135            aborted: false,
136            memory_usage,
137            indexed_column_ids,
138            metadata: metadata.clone(),
139        }
140    }
141
142    /// Updates index with a batch of rows.
143    /// Garbage will be cleaned up if failed to update.
144    pub async fn update(&mut self, batch: &mut Batch) -> Result<()> {
145        ensure!(!self.aborted, OperateAbortedIndexSnafu);
146
147        if batch.is_empty() {
148            return Ok(());
149        }
150
151        if let Err(update_err) = self.do_update(batch).await {
152            // clean up garbage if failed to update
153            if let Err(err) = self.do_cleanup().await {
154                if cfg!(any(test, feature = "test")) {
155                    panic!("Failed to clean up index creator, err: {err}",);
156                } else {
157                    warn!(err; "Failed to clean up index creator");
158                }
159            }
160            return Err(update_err);
161        }
162
163        Ok(())
164    }
165
166    /// Updates the inverted index with the given flat format RecordBatch.
167    pub async fn update_flat(&mut self, batch: &RecordBatch) -> Result<()> {
168        ensure!(!self.aborted, OperateAbortedIndexSnafu);
169
170        if batch.num_rows() == 0 {
171            return Ok(());
172        }
173
174        self.do_update_flat(batch).await
175    }
176
177    async fn do_update_flat(&mut self, batch: &RecordBatch) -> Result<()> {
178        let mut guard = self.stats.record_update();
179
180        guard.inc_row_count(batch.num_rows());
181
182        let is_sparse = self.metadata.primary_key_encoding == PrimaryKeyEncoding::Sparse;
183        let mut sparse_columns: SmallVec<[(ColumnId, &str); 8]> = SmallVec::new();
184
185        for (col_id, target_key) in &self.indexed_column_ids {
186            let Some(column_meta) = self.metadata.column_by_id(*col_id) else {
187                debug!(
188                    "Column {} not found in the metadata during building inverted index",
189                    col_id
190                );
191                continue;
192            };
193            let column_name = &column_meta.column_schema.name;
194            if let Some(column_array) = batch.column_by_name(column_name) {
195                // Convert Arrow array to VectorRef using Helper
196                let vector = Helper::try_into_vector(column_array.clone())
197                    .context(crate::error::ConvertVectorSnafu)?;
198                let sort_field = SortField::new(vector.data_type());
199
200                for (row, count) in column_index_rows(batch, column_meta.semantic_type) {
201                    let elem = IndexValueCodec::encode_value(
202                        vector.get_ref(row),
203                        &sort_field,
204                        &mut self.value_buf,
205                    )
206                    .context(EncodeSnafu)?;
207                    if self.index_creator.push_with_name_n(target_key, elem, count) {
208                        self.index_creator
209                            .spill()
210                            .await
211                            .context(PushIndexValueSnafu)?;
212                    }
213                }
214            } else if is_sparse && column_meta.semantic_type == SemanticType::Tag {
215                if self.codec.pk_col_info(*col_id).is_some() {
216                    sparse_columns.push((*col_id, target_key));
217                }
218            } else {
219                debug!(
220                    "Column {} not found in the batch during building inverted index",
221                    col_id
222                );
223            }
224        }
225
226        if !sparse_columns.is_empty() {
227            for (pk, count) in PrimaryKeyRuns::try_new(batch)? {
228                let mut view =
229                    SparsePrimaryKeyView::new(pk, &mut self.pk_offsets).context(DecodeSnafu)?;
230                // Visit all needed tags before moving to the next PK so offset discovery is shared.
231                for &(col_id, target_key) in &sparse_columns {
232                    let value = IndexValueCodec::encode_sparse_value(
233                        &mut view,
234                        col_id,
235                        &mut self.value_buf,
236                    )
237                    .context(DecodeSnafu)?;
238                    if self
239                        .index_creator
240                        .push_with_name_n(target_key, value, count)
241                    {
242                        self.index_creator
243                            .spill()
244                            .await
245                            .context(PushIndexValueSnafu)?;
246                    }
247                }
248            }
249        }
250
251        Ok(())
252    }
253
254    /// Finishes index creation and cleans up garbage.
255    /// Returns the number of rows and bytes written.
256    pub(crate) async fn finish(
257        &mut self,
258        puffin_writer: &mut SstPuffinWriter,
259    ) -> Result<(RowCount, ByteCount)> {
260        ensure!(!self.aborted, OperateAbortedIndexSnafu);
261
262        if self.stats.row_count() == 0 {
263            // no IO is performed, no garbage to clean up, just return
264            return Ok((0, 0));
265        }
266
267        let finish_res = self.do_finish(puffin_writer).await;
268        // clean up garbage no matter finish successfully or not
269        if let Err(err) = self.do_cleanup().await {
270            if cfg!(any(test, feature = "test")) {
271                panic!("Failed to clean up index creator, err: {err}",);
272            } else {
273                warn!(err; "Failed to clean up index creator");
274            }
275        }
276
277        finish_res.map(|_| (self.stats.row_count(), self.stats.byte_count()))
278    }
279
280    /// Aborts index creation and clean up garbage.
281    pub async fn abort(&mut self) -> Result<()> {
282        if self.aborted {
283            return Ok(());
284        }
285        self.aborted = true;
286
287        self.do_cleanup().await
288    }
289
290    async fn do_update(&mut self, batch: &mut Batch) -> Result<()> {
291        let mut guard = self.stats.record_update();
292
293        let n = batch.num_rows();
294        guard.inc_row_count(n);
295
296        for (col_id, target_key) in &self.indexed_column_ids {
297            match self.codec.pk_col_info(*col_id) {
298                // pk
299                Some(col_info) => {
300                    let pk_idx = col_info.idx;
301                    let field = &col_info.field;
302                    let value = batch
303                        .pk_col_value(self.codec.decoder(), pk_idx, *col_id)?
304                        .filter(|v| !v.is_null())
305                        .map(|v| {
306                            self.value_buf.clear();
307                            IndexValueCodec::encode_nonnull_value(
308                                v.as_value_ref(),
309                                field,
310                                &mut self.value_buf,
311                            )
312                            .context(EncodeSnafu)?;
313                            Ok(self.value_buf.as_slice())
314                        })
315                        .transpose()?;
316
317                    if self.index_creator.push_with_name_n(target_key, value, n) {
318                        self.index_creator
319                            .spill()
320                            .await
321                            .context(PushIndexValueSnafu)?;
322                    }
323                }
324                // fields
325                None => {
326                    let Some(values) = batch.field_col_value(*col_id) else {
327                        debug!(
328                            "Column {} not found in the batch during building inverted index",
329                            col_id
330                        );
331                        continue;
332                    };
333                    let sort_field = SortField::new(values.data.data_type());
334                    for i in 0..n {
335                        self.value_buf.clear();
336                        let value = values.data.get_ref(i);
337                        if value.is_null() {
338                            if self.index_creator.push_with_name(target_key, None) {
339                                self.index_creator
340                                    .spill()
341                                    .await
342                                    .context(PushIndexValueSnafu)?;
343                            }
344                        } else {
345                            IndexValueCodec::encode_nonnull_value(
346                                value,
347                                &sort_field,
348                                &mut self.value_buf,
349                            )
350                            .context(EncodeSnafu)?;
351                            if self
352                                .index_creator
353                                .push_with_name(target_key, Some(&self.value_buf))
354                            {
355                                self.index_creator
356                                    .spill()
357                                    .await
358                                    .context(PushIndexValueSnafu)?;
359                            }
360                        }
361                    }
362                }
363            }
364        }
365
366        Ok(())
367    }
368
369    /// Data flow of finishing index:
370    ///
371    /// ```text
372    ///                               (In Memory Buffer)
373    ///                                    ┌──────┐
374    ///  ┌─────────────┐                   │ PIPE │
375    ///  │             │ write index data  │      │
376    ///  │ IndexWriter ├──────────────────►│ tx   │
377    ///  │             │                   │      │
378    ///  └─────────────┘                   │      │
379    ///                  ┌─────────────────┤ rx   │
380    ///  ┌─────────────┐ │ read as blob    └──────┘
381    ///  │             │ │
382    ///  │ PuffinWriter├─┤
383    ///  │             │ │ copy to file    ┌──────┐
384    ///  └─────────────┘ └────────────────►│ File │
385    ///                                    └──────┘
386    /// ```
387    async fn do_finish(&mut self, puffin_writer: &mut SstPuffinWriter) -> Result<()> {
388        let mut guard = self.stats.record_finish();
389
390        let (tx, rx) = duplex(PIPE_BUFFER_SIZE_FOR_SENDING_BLOB);
391        let mut index_writer = InvertedIndexBlobWriter::new(tx.compat_write());
392
393        let (index_finish, puffin_add_blob) = futures::join!(
394            // TODO(zhongzc): config bitmap type
395            self.index_creator
396                .finish(&mut index_writer, index::bitmap::BitmapType::Roaring),
397            puffin_writer.put_blob(
398                INDEX_BLOB_TYPE,
399                rx.compat(),
400                PutOptions::default(),
401                Default::default(),
402            )
403        );
404
405        match (
406            puffin_add_blob.context(PuffinAddBlobSnafu),
407            index_finish.context(IndexFinishSnafu),
408        ) {
409            (Err(e1), Err(e2)) => BiErrorsSnafu {
410                first: Box::new(e1),
411                second: Box::new(e2),
412            }
413            .fail()?,
414
415            (Ok(_), e @ Err(_)) => e?,
416            (e @ Err(_), Ok(_)) => e.map(|_| ())?,
417            (Ok(written_bytes), Ok(_)) => {
418                guard.inc_byte_count(written_bytes);
419            }
420        }
421
422        Ok(())
423    }
424
425    async fn do_cleanup(&mut self) -> Result<()> {
426        let _guard = self.stats.record_cleanup();
427
428        self.temp_file_provider.cleanup().await
429    }
430
431    pub fn column_ids(&self) -> impl Iterator<Item = ColumnId> + '_ {
432        self.indexed_column_ids.iter().map(|(col_id, _)| *col_id)
433    }
434
435    pub fn memory_usage(&self) -> usize {
436        self.memory_usage.load(std::sync::atomic::Ordering::Relaxed)
437    }
438}
439
440#[cfg(test)]
441mod tests {
442    use std::collections::BTreeSet;
443
444    use api::v1::SemanticType;
445    use datafusion_expr::{Expr as DfExpr, Operator, binary_expr, col, lit};
446    use datatypes::data_type::ConcreteDataType;
447    use datatypes::schema::ColumnSchema;
448    use datatypes::value::ValueRef;
449    use datatypes::vectors::{UInt8Vector, UInt64Vector};
450    use futures::future::BoxFuture;
451    use mito_codec::row_converter::{DensePrimaryKeyCodec, PrimaryKeyCodecExt};
452    use object_store::ObjectStore;
453    use object_store::services::Memory;
454    use puffin::puffin_manager::PuffinManager;
455    use puffin::puffin_manager::cache::PuffinMetadataCache;
456    use store_api::metadata::{ColumnMetadata, RegionMetadataBuilder};
457    use store_api::region_request::PathType;
458    use store_api::storage::RegionId;
459
460    use super::*;
461    use crate::access_layer::RegionFilePathFactory;
462    use crate::cache::index::inverted_index::InvertedIndexCache;
463    use crate::metrics::CACHE_BYTES;
464    use crate::read::BatchColumn;
465    use crate::sst::file::{RegionFileId, RegionIndexId};
466    use crate::sst::index::inverted_index::applier::builder::InvertedIndexApplierBuilder;
467    use crate::sst::index::puffin_manager::PuffinManagerFactory;
468
469    fn mock_object_store() -> ObjectStore {
470        ObjectStore::new(Memory::default()).unwrap()
471    }
472
473    async fn new_intm_mgr(path: impl AsRef<str>) -> IntermediateManager {
474        IntermediateManager::init_fs(path).await.unwrap()
475    }
476
477    fn mock_region_metadata() -> RegionMetadataRef {
478        let mut builder = RegionMetadataBuilder::new(RegionId::new(1, 2));
479        builder
480            .push_column_metadata(ColumnMetadata {
481                column_schema: ColumnSchema::new(
482                    "tag_str",
483                    ConcreteDataType::string_datatype(),
484                    false,
485                ),
486                semantic_type: SemanticType::Tag,
487                column_id: 1,
488            })
489            .push_column_metadata(ColumnMetadata {
490                column_schema: ColumnSchema::new(
491                    "tag_i32",
492                    ConcreteDataType::int32_datatype(),
493                    false,
494                ),
495                semantic_type: SemanticType::Tag,
496                column_id: 2,
497            })
498            .push_column_metadata(ColumnMetadata {
499                column_schema: ColumnSchema::new(
500                    "ts",
501                    ConcreteDataType::timestamp_millisecond_datatype(),
502                    false,
503                ),
504                semantic_type: SemanticType::Timestamp,
505                column_id: 3,
506            })
507            .push_column_metadata(ColumnMetadata {
508                column_schema: ColumnSchema::new(
509                    "field_u64",
510                    ConcreteDataType::uint64_datatype(),
511                    false,
512                ),
513                semantic_type: SemanticType::Field,
514                column_id: 4,
515            })
516            .primary_key(vec![1, 2]);
517
518        Arc::new(builder.build().unwrap())
519    }
520
521    fn new_batch(
522        str_tag: impl AsRef<str>,
523        i32_tag: impl Into<i32>,
524        u64_field: impl IntoIterator<Item = u64>,
525    ) -> Batch {
526        let fields = vec![
527            (0, SortField::new(ConcreteDataType::string_datatype())),
528            (1, SortField::new(ConcreteDataType::int32_datatype())),
529        ];
530        let codec = DensePrimaryKeyCodec::with_fields(fields);
531        let row: [ValueRef; 2] = [str_tag.as_ref().into(), i32_tag.into().into()];
532        let primary_key = codec.encode(row.into_iter()).unwrap();
533
534        let u64_field = BatchColumn {
535            column_id: 4,
536            data: Arc::new(UInt64Vector::from_iter_values(u64_field)),
537        };
538        let num_rows = u64_field.data.len();
539
540        Batch::new(
541            primary_key,
542            Arc::new(UInt64Vector::from_iter_values(std::iter::repeat_n(
543                0, num_rows,
544            ))),
545            Arc::new(UInt64Vector::from_iter_values(std::iter::repeat_n(
546                0, num_rows,
547            ))),
548            Arc::new(UInt8Vector::from_iter_values(std::iter::repeat_n(
549                1, num_rows,
550            ))),
551            vec![u64_field],
552        )
553        .unwrap()
554    }
555
556    async fn build_applier_factory(
557        prefix: &str,
558        rows: BTreeSet<(&'static str, i32, [u64; 2])>,
559    ) -> impl Fn(DfExpr) -> BoxFuture<'static, Vec<usize>> {
560        let (d, factory) = PuffinManagerFactory::new_for_test_async(prefix).await;
561        let table_dir = "table0".to_string();
562        let sst_file_id = FileId::random();
563        let object_store = mock_object_store();
564        let region_metadata = mock_region_metadata();
565        let intm_mgr = new_intm_mgr(d.path().to_string_lossy()).await;
566        let memory_threshold = None;
567        let segment_row_count = 2;
568        let indexed_column_ids = HashSet::from_iter([1, 2, 4]);
569
570        let mut creator = InvertedIndexer::new(
571            sst_file_id,
572            &region_metadata,
573            intm_mgr,
574            memory_threshold,
575            NonZeroUsize::new(segment_row_count).unwrap(),
576            indexed_column_ids.clone(),
577        );
578
579        for (str_tag, i32_tag, u64_field) in &rows {
580            let mut batch = new_batch(str_tag, *i32_tag, u64_field.iter().copied());
581            creator.update(&mut batch).await.unwrap();
582        }
583
584        let puffin_manager = factory.build(
585            object_store.clone(),
586            RegionFilePathFactory::new(table_dir.clone(), PathType::Bare),
587        );
588
589        let sst_file_id = RegionFileId::new(region_metadata.region_id, sst_file_id);
590        let index_id = RegionIndexId::new(sst_file_id, 0);
591        let mut writer = puffin_manager.writer(&index_id).await.unwrap();
592        let (row_count, _) = creator.finish(&mut writer).await.unwrap();
593        assert_eq!(row_count, rows.len() * segment_row_count);
594        writer.finish().await.unwrap();
595
596        move |expr| {
597            let _d = &d;
598            let cache = Arc::new(InvertedIndexCache::new(10, 10, 100));
599            let puffin_metadata_cache = Arc::new(PuffinMetadataCache::new(10, &CACHE_BYTES));
600            let applier = InvertedIndexApplierBuilder::new(
601                table_dir.clone(),
602                PathType::Bare,
603                object_store.clone(),
604                &region_metadata,
605                indexed_column_ids.clone(),
606                factory.clone(),
607            )
608            .with_inverted_index_cache(Some(cache))
609            .with_puffin_metadata_cache(Some(puffin_metadata_cache))
610            .build(&[expr])
611            .unwrap()
612            .unwrap();
613            let sst_metadata = Arc::new(region_metadata.clone());
614            let plan = applier.plan_for_sst(&sst_metadata).unwrap().unwrap();
615            Box::pin(async move {
616                applier
617                    .apply(index_id, None, &plan.index_applier, None)
618                    .await
619                    .unwrap()
620                    .matched_segment_ids
621                    .iter_ones()
622                    .collect()
623            })
624        }
625    }
626
627    #[tokio::test]
628    async fn test_create_and_query_get_key() {
629        let rows = BTreeSet::from_iter([
630            ("aaa", 1, [1, 2]),
631            ("aaa", 2, [2, 3]),
632            ("aaa", 3, [3, 4]),
633            ("aab", 1, [4, 5]),
634            ("aab", 2, [5, 6]),
635            ("aab", 3, [6, 7]),
636            ("abc", 1, [7, 8]),
637            ("abc", 2, [8, 9]),
638            ("abc", 3, [9, 10]),
639        ]);
640
641        let applier_factory = build_applier_factory("test_create_and_query_get_key_", rows).await;
642
643        let expr = col("tag_str").eq(lit("aaa"));
644        let res = applier_factory(expr).await;
645        assert_eq!(res, vec![0, 1, 2]);
646
647        let expr = col("tag_i32").eq(lit(2));
648        let res = applier_factory(expr).await;
649        assert_eq!(res, vec![1, 4, 7]);
650
651        let expr = col("tag_str").eq(lit("aaa")).and(col("tag_i32").eq(lit(2)));
652        let res = applier_factory(expr).await;
653        assert_eq!(res, vec![1]);
654
655        let expr = col("tag_str")
656            .eq(lit("aaa"))
657            .or(col("tag_str").eq(lit("abc")));
658        let res = applier_factory(expr).await;
659        assert_eq!(res, vec![0, 1, 2, 6, 7, 8]);
660
661        let expr = col("tag_str").in_list(vec![lit("aaa"), lit("abc")], false);
662        let res = applier_factory(expr).await;
663        assert_eq!(res, vec![0, 1, 2, 6, 7, 8]);
664
665        let expr = col("field_u64").eq(lit(2u64));
666        let res = applier_factory(expr).await;
667        assert_eq!(res, vec![0, 1]);
668    }
669
670    #[tokio::test]
671    async fn test_create_and_query_range() {
672        let rows = BTreeSet::from_iter([
673            ("aaa", 1, [1, 2]),
674            ("aaa", 2, [2, 3]),
675            ("aaa", 3, [3, 4]),
676            ("aab", 1, [4, 5]),
677            ("aab", 2, [5, 6]),
678            ("aab", 3, [6, 7]),
679            ("abc", 1, [7, 8]),
680            ("abc", 2, [8, 9]),
681            ("abc", 3, [9, 10]),
682        ]);
683
684        let applier_factory = build_applier_factory("test_create_and_query_range_", rows).await;
685
686        let expr = col("tag_str").between(lit("aaa"), lit("aab"));
687        let res = applier_factory(expr).await;
688        assert_eq!(res, vec![0, 1, 2, 3, 4, 5]);
689
690        let expr = col("tag_i32").between(lit(2), lit(3));
691        let res = applier_factory(expr).await;
692        assert_eq!(res, vec![1, 2, 4, 5, 7, 8]);
693
694        let expr = col("tag_str").between(lit("aaa"), lit("aaa"));
695        let res = applier_factory(expr).await;
696        assert_eq!(res, vec![0, 1, 2]);
697
698        let expr = col("tag_i32").between(lit(2), lit(2));
699        let res = applier_factory(expr).await;
700        assert_eq!(res, vec![1, 4, 7]);
701
702        let expr = col("field_u64").between(lit(2u64), lit(5u64));
703        let res = applier_factory(expr).await;
704        assert_eq!(res, vec![0, 1, 2, 3, 4]);
705    }
706
707    #[tokio::test]
708    async fn test_create_and_query_comparison() {
709        let rows = BTreeSet::from_iter([
710            ("aaa", 1, [1, 2]),
711            ("aaa", 2, [2, 3]),
712            ("aaa", 3, [3, 4]),
713            ("aab", 1, [4, 5]),
714            ("aab", 2, [5, 6]),
715            ("aab", 3, [6, 7]),
716            ("abc", 1, [7, 8]),
717            ("abc", 2, [8, 9]),
718            ("abc", 3, [9, 10]),
719        ]);
720
721        let applier_factory =
722            build_applier_factory("test_create_and_query_comparison_", rows).await;
723
724        let expr = col("tag_str").lt(lit("aab"));
725        let res = applier_factory(expr).await;
726        assert_eq!(res, vec![0, 1, 2]);
727
728        let expr = col("tag_i32").lt(lit(2));
729        let res = applier_factory(expr).await;
730        assert_eq!(res, vec![0, 3, 6]);
731
732        let expr = col("field_u64").lt(lit(2u64));
733        let res = applier_factory(expr).await;
734        assert_eq!(res, vec![0]);
735
736        let expr = col("tag_str").gt(lit("aab"));
737        let res = applier_factory(expr).await;
738        assert_eq!(res, vec![6, 7, 8]);
739
740        let expr = col("tag_i32").gt(lit(2));
741        let res = applier_factory(expr).await;
742        assert_eq!(res, vec![2, 5, 8]);
743
744        let expr = col("field_u64").gt(lit(8u64));
745        let res = applier_factory(expr).await;
746        assert_eq!(res, vec![7, 8]);
747
748        let expr = col("tag_str").lt_eq(lit("aab"));
749        let res = applier_factory(expr).await;
750        assert_eq!(res, vec![0, 1, 2, 3, 4, 5]);
751
752        let expr = col("tag_i32").lt_eq(lit(2));
753        let res = applier_factory(expr).await;
754        assert_eq!(res, vec![0, 1, 3, 4, 6, 7]);
755
756        let expr = col("field_u64").lt_eq(lit(2u64));
757        let res = applier_factory(expr).await;
758        assert_eq!(res, vec![0, 1]);
759
760        let expr = col("tag_str").gt_eq(lit("aab"));
761        let res = applier_factory(expr).await;
762        assert_eq!(res, vec![3, 4, 5, 6, 7, 8]);
763
764        let expr = col("tag_i32").gt_eq(lit(2));
765        let res = applier_factory(expr).await;
766        assert_eq!(res, vec![1, 2, 4, 5, 7, 8]);
767
768        let expr = col("field_u64").gt_eq(lit(8u64));
769        let res = applier_factory(expr).await;
770        assert_eq!(res, vec![6, 7, 8]);
771
772        let expr = col("tag_str")
773            .gt(lit("aaa"))
774            .and(col("tag_str").lt(lit("abc")));
775        let res = applier_factory(expr).await;
776        assert_eq!(res, vec![3, 4, 5]);
777
778        let expr = col("tag_i32").gt(lit(1)).and(col("tag_i32").lt(lit(3)));
779        let res = applier_factory(expr).await;
780        assert_eq!(res, vec![1, 4, 7]);
781
782        let expr = col("field_u64")
783            .gt(lit(2u64))
784            .and(col("field_u64").lt(lit(9u64)));
785        let res = applier_factory(expr).await;
786        assert_eq!(res, vec![1, 2, 3, 4, 5, 6, 7]);
787    }
788
789    #[tokio::test]
790    async fn test_create_and_query_regex() {
791        let rows = BTreeSet::from_iter([
792            ("aaa", 1, [1, 2]),
793            ("aaa", 2, [2, 3]),
794            ("aaa", 3, [3, 4]),
795            ("aab", 1, [4, 5]),
796            ("aab", 2, [5, 6]),
797            ("aab", 3, [6, 7]),
798            ("abc", 1, [7, 8]),
799            ("abc", 2, [8, 9]),
800            ("abc", 3, [9, 10]),
801        ]);
802
803        let applier_factory = build_applier_factory("test_create_and_query_regex_", rows).await;
804
805        let expr = binary_expr(col("tag_str"), Operator::RegexMatch, lit(".*"));
806        let res = applier_factory(expr).await;
807        assert_eq!(res, vec![0, 1, 2, 3, 4, 5, 6, 7, 8]);
808
809        let expr = binary_expr(col("tag_str"), Operator::RegexMatch, lit("a.*c"));
810        let res = applier_factory(expr).await;
811        assert_eq!(res, vec![6, 7, 8]);
812
813        let expr = binary_expr(col("tag_str"), Operator::RegexMatch, lit("a.*b$"));
814        let res = applier_factory(expr).await;
815        assert_eq!(res, vec![3, 4, 5]);
816
817        let expr = binary_expr(col("tag_str"), Operator::RegexMatch, lit("\\w"));
818        let res = applier_factory(expr).await;
819        assert_eq!(res, vec![0, 1, 2, 3, 4, 5, 6, 7, 8]);
820
821        let expr = binary_expr(col("tag_str"), Operator::RegexMatch, lit("\\d"));
822        let res = applier_factory(expr).await;
823        assert!(res.is_empty());
824
825        let expr = binary_expr(col("tag_str"), Operator::RegexMatch, lit("^aaa$"));
826        let res = applier_factory(expr).await;
827        assert_eq!(res, vec![0, 1, 2]);
828    }
829}