Skip to main content

mito2/series_index/
searcher.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 async_stream::try_stream;
16use common_recordbatch::filter::SimpleFilterEvaluator;
17use common_time::range::TimestampRange;
18use common_time::timestamp::TimeUnit;
19use datafusion_expr::{Expr, col, lit};
20use datatypes::arrow::array::{ArrayRef, UInt32Array, UInt64Array};
21use datatypes::arrow::buffer::BooleanBuffer;
22use datatypes::arrow::datatypes::{DataType, SchemaRef};
23use datatypes::arrow::record_batch::RecordBatch;
24use datatypes::value::timestamp_to_scalar_value;
25use futures::TryStreamExt;
26use object_store::ObjectStore;
27use snafu::{OptionExt, ResultExt, ensure};
28use store_api::metadata::RegionMetadataRef;
29use table::predicate::Predicate;
30
31use crate::error::{InvalidRecordBatchSnafu, RecordBatchSnafu, Result, UnexpectedSnafu};
32use crate::series_index::{
33    MAX_TS_COLUMN, METRIC_SERIES_ID_BATCH_SIZE, MIN_TS_COLUMN, MetricSeriesId,
34    MetricSeriesIdStream, ROW_COUNT_COLUMN, SeriesIndexFileHandle, TABLE_ID_COLUMN, TSID_COLUMN,
35    series_index_path, series_index_schema,
36};
37use crate::sst::parquet::index_reader::ParquetIndexReader;
38use crate::sst::parquet::prefilter::simple_tag_filters;
39
40/// Searches a series-index file for metric series matching query predicates.
41pub struct SeriesIndexSearcher {
42    /// Pins the index file until this searcher and all of its streams are dropped.
43    file_handle: SeriesIndexFileHandle,
44    /// The file to search, or `None` when an empty time range rules out
45    /// every series without opening the file.
46    reader: Option<ParquetIndexReader>,
47    /// Row-group pruning predicate built from all of the file's filters.
48    pruning_predicate: Predicate,
49    filters: Vec<SimpleFilterEvaluator>,
50}
51
52impl SeriesIndexSearcher {
53    /// Creates a searcher for the series-index file protected by `file_handle`.
54    /// Predicates are built from the unit recorded in the file's schema, so a
55    /// file written before a time index unit widening keeps being interpreted
56    /// in its own unit.
57    pub(crate) async fn try_new(
58        metadata: RegionMetadataRef,
59        object_store: ObjectStore,
60        file_handle: SeriesIndexFileHandle,
61        predicate: Option<&Predicate>,
62        time_range: Option<TimestampRange>,
63    ) -> Result<Self> {
64        // Keep search-time metadata validation identical to the writer.
65        series_index_schema(&metadata)?;
66        if time_range.as_ref().is_some_and(TimestampRange::is_empty) {
67            return Ok(Self {
68                file_handle,
69                reader: None,
70                pruning_predicate: Predicate::new(Vec::new()),
71                filters: Vec::new(),
72            });
73        }
74
75        let file_id = file_handle.file_id();
76        let path = series_index_path(file_id.region_id(), file_id.file_id());
77        let reader = ParquetIndexReader::open(object_store, &path).await?;
78        let unit = validate_index_schema(reader.schema())?;
79
80        let mut filters = simple_tag_filters(&metadata, None, predicate);
81        for expr in time_range_exprs(unit, time_range.as_ref()) {
82            let filter = SimpleFilterEvaluator::try_new(&expr).context(UnexpectedSnafu {
83                reason: "failed to build an internal series-index time filter",
84            })?;
85            filters.push((expr, filter));
86        }
87        // An older index file may not contain tags added by schema evolution.
88        // Ignore filters on those tags to preserve a conservative candidate set.
89        let (pruning_predicate, filters) = filters_for_schema(reader.schema(), &filters);
90
91        Ok(Self {
92            file_handle,
93            reader: Some(reader),
94            pruning_predicate,
95            filters,
96        })
97    }
98
99    /// Searches the index file and returns sorted batches of matching
100    /// metric-series IDs.
101    pub fn search(&self) -> Result<MetricSeriesIdStream> {
102        let Some(reader) = self.reader.as_ref() else {
103            return Ok(Box::pin(futures::stream::empty()));
104        };
105        let mut projection_columns = Vec::with_capacity(self.filters.len() + 2);
106        projection_columns.extend([TABLE_ID_COLUMN, TSID_COLUMN]);
107        projection_columns.extend(self.filters.iter().map(SimpleFilterEvaluator::column_name));
108        let mut batches = reader.read(&self.pruning_predicate, &projection_columns)?;
109        let filters = self.filters.clone();
110        let file_handle = self.file_handle.clone();
111
112        Ok(Box::pin(try_stream! {
113            let mut last_series = None;
114            let mut output = Vec::with_capacity(METRIC_SERIES_ID_BATCH_SIZE);
115            while let Some(batch) = batches.try_next().await? {
116                let mut mask = BooleanBuffer::new_set(batch.num_rows());
117                for filter in &filters {
118                    let column = column(&batch, filter.column_name())?;
119                    let evaluated = filter.evaluate_array(column).context(RecordBatchSnafu)?;
120                    mask = &mask & &evaluated;
121                }
122
123                let table_ids = column(&batch, TABLE_ID_COLUMN)?
124                    .as_any()
125                    .downcast_ref::<UInt32Array>()
126                    .context(InvalidRecordBatchSnafu {
127                        reason: "series index __table_id is not UInt32",
128                    })?;
129                let tsids = column(&batch, TSID_COLUMN)?
130                    .as_any()
131                    .downcast_ref::<UInt64Array>()
132                    .context(InvalidRecordBatchSnafu {
133                        reason: "series index __tsid is not UInt64",
134                    })?;
135
136                for (row, matched) in mask.iter().enumerate() {
137                    if !matched {
138                        continue;
139                    }
140                    let series = MetricSeriesId {
141                        table_id: table_ids.value(row),
142                        tsid: tsids.value(row),
143                    };
144                    if last_series == Some(series) {
145                        continue;
146                    }
147                    last_series = Some(series);
148                    output.push(series);
149                    if output.len() == METRIC_SERIES_ID_BATCH_SIZE {
150                        yield std::mem::replace(
151                            &mut output,
152                            Vec::with_capacity(METRIC_SERIES_ID_BATCH_SIZE),
153                        );
154                    }
155                }
156            }
157            if !output.is_empty() {
158                yield output;
159            }
160            drop(file_handle);
161        }))
162    }
163}
164
165fn filters_for_schema(
166    schema: &SchemaRef,
167    filters: &[(Expr, SimpleFilterEvaluator)],
168) -> (Predicate, Vec<SimpleFilterEvaluator>) {
169    let (exprs, filters): (Vec<_>, Vec<_>) = filters
170        .iter()
171        .filter(|(_, filter)| schema.field_with_name(filter.column_name()).is_ok())
172        .cloned()
173        .unzip();
174    (Predicate::new(exprs), filters)
175}
176
177// Builds `__series_min_ts`/`__series_max_ts` predicates in `unit`, the unit
178// recorded in the file being searched: its raw i64 bounds were written in that
179// unit even if the region's time index has since been widened. The searcher
180// handles empty ranges itself, so `time_range`, if present, is non-empty.
181fn time_range_exprs(unit: TimeUnit, time_range: Option<&TimestampRange>) -> Vec<Expr> {
182    let Some(time_range) = time_range else {
183        return Vec::new();
184    };
185    let mut exprs = Vec::with_capacity(2);
186    // A series overlaps [start, end) only if its maximum is at least start.
187    // Round start up so a series ending before an unaligned start is pruned.
188    if let Some(start) = time_range
189        .start()
190        .and_then(|start| start.convert_to_ceil(unit))
191    {
192        exprs.push(
193            col(MAX_TS_COLUMN).gt_eq(lit(timestamp_to_scalar_value(unit, Some(start.value())))),
194        );
195    }
196    // A series overlaps [start, end) only if its minimum is less than end.
197    // Round the exclusive end up to avoid pruning the containing unit interval.
198    if let Some(end) = time_range.end().and_then(|end| end.convert_to_ceil(unit)) {
199        exprs.push(col(MIN_TS_COLUMN).lt(lit(timestamp_to_scalar_value(unit, Some(end.value())))));
200    }
201    exprs
202}
203
204/// Validates an index file's schema and returns the time unit of its
205/// `__series_min_ts`/`__series_max_ts` columns. They are native
206/// `Timestamp(unit)` columns, so the unit is part of the datatype and each
207/// file is interpreted in the unit it was written with.
208fn validate_index_schema(schema: &SchemaRef) -> Result<TimeUnit> {
209    for (name, data_type) in [
210        (ROW_COUNT_COLUMN, DataType::UInt64),
211        (TABLE_ID_COLUMN, DataType::UInt32),
212        (TSID_COLUMN, DataType::UInt64),
213    ] {
214        let field = schema
215            .field_with_name(name)
216            .ok()
217            .with_context(|| InvalidRecordBatchSnafu {
218                reason: format!("series index is missing internal column {name}"),
219            })?;
220        ensure!(
221            field.data_type() == &data_type && !field.is_nullable(),
222            InvalidRecordBatchSnafu {
223                reason: format!(
224                    "series index internal column {name} must be non-nullable {data_type:?}, got {:?}",
225                    field.data_type()
226                ),
227            }
228        );
229    }
230    let unit = |name: &str| {
231        let field = schema
232            .field_with_name(name)
233            .ok()
234            .context(InvalidRecordBatchSnafu {
235                reason: format!("series index is missing internal column {name}"),
236            })?;
237        ensure!(
238            !field.is_nullable(),
239            InvalidRecordBatchSnafu {
240                reason: format!("series index column {name} must be non-nullable"),
241            }
242        );
243        match field.data_type() {
244            DataType::Timestamp(unit, _) => Ok(unit.into()),
245            data_type => InvalidRecordBatchSnafu {
246                reason: format!(
247                    "series index column {name} must be a Timestamp, got {data_type:?}"
248                ),
249            }
250            .fail(),
251        }
252    };
253    let min_unit = unit(MIN_TS_COLUMN)?;
254    let max_unit = unit(MAX_TS_COLUMN)?;
255    ensure!(
256        min_unit == max_unit,
257        InvalidRecordBatchSnafu {
258            reason: format!(
259                "series index columns {MIN_TS_COLUMN} and {MAX_TS_COLUMN} have different time units"
260            ),
261        }
262    );
263    Ok(min_unit)
264}
265
266fn column<'a>(batch: &'a RecordBatch, name: &str) -> Result<&'a ArrayRef> {
267    let index = batch
268        .schema()
269        .index_of(name)
270        .ok()
271        .with_context(|| InvalidRecordBatchSnafu {
272            reason: format!("series index batch is missing column {name}"),
273        })?;
274    Ok(batch.column(index))
275}
276
277#[cfg(test)]
278mod tests {
279    use std::sync::Arc;
280
281    use api::v1::SemanticType;
282    use datafusion_expr::{col, lit};
283    use datatypes::arrow::array::{
284        BinaryArray, TimestampMicrosecondArray, TimestampMillisecondArray,
285        TimestampNanosecondArray, TimestampSecondArray, UInt8Array,
286    };
287    use datatypes::arrow::datatypes::{Field, Schema, TimeUnit as ArrowTimeUnit};
288    use datatypes::arrow::record_batch::RecordBatch;
289    use datatypes::prelude::ConcreteDataType;
290    use datatypes::schema::ColumnSchema;
291    use futures::TryStreamExt;
292    use object_store::services::Memory;
293    use store_api::codec::PrimaryKeyEncoding;
294    use store_api::metadata::{ColumnMetadata, RegionMetadataBuilder};
295    use store_api::storage::FileId;
296    use tokio::sync::mpsc::UnboundedReceiver;
297
298    use super::*;
299    use crate::series_index::purger::PurgeRequest;
300    use crate::series_index::{
301        SeriesIndexEntry, SeriesIndexWriter, SeriesIndexWriterOptions, series_index_channel,
302    };
303    use crate::test_util::sst_util::{new_sparse_primary_key, sst_region_metadata_with_encoding};
304
305    fn object_store() -> ObjectStore {
306        ObjectStore::new(Memory::default()).unwrap()
307    }
308
309    fn index_handle(
310        metadata: &RegionMetadataRef,
311        store: &ObjectStore,
312    ) -> (SeriesIndexFileHandle, UnboundedReceiver<PurgeRequest>) {
313        let (purger, receiver) = series_index_channel(store.clone());
314        let entry = SeriesIndexEntry {
315            file_size: 0,
316            index_uuid: FileId::random(),
317            bucket_start: common_time::Timestamp::new_second(0),
318            bucket_end: common_time::Timestamp::new_second(60),
319            source_file_ids: Vec::new(),
320            min_file_sequence: 0,
321            max_file_sequence: 0,
322            compaction_window_secs: 60,
323            window_sequences: Default::default(),
324        };
325        (
326            SeriesIndexFileHandle::new(metadata.region_id, entry, purger),
327            receiver,
328        )
329    }
330
331    fn flat_batch_with_time_unit(
332        primary_keys: &[Vec<u8>],
333        timestamps: &[i64],
334        unit: ArrowTimeUnit,
335    ) -> RecordBatch {
336        let ts_column = match unit {
337            ArrowTimeUnit::Second => {
338                Arc::new(TimestampSecondArray::from(timestamps.to_vec())) as ArrayRef
339            }
340            ArrowTimeUnit::Millisecond => {
341                Arc::new(TimestampMillisecondArray::from(timestamps.to_vec())) as ArrayRef
342            }
343            ArrowTimeUnit::Microsecond => {
344                Arc::new(TimestampMicrosecondArray::from(timestamps.to_vec())) as ArrayRef
345            }
346            ArrowTimeUnit::Nanosecond => {
347                Arc::new(TimestampNanosecondArray::from(timestamps.to_vec())) as ArrayRef
348            }
349        };
350        let schema = Arc::new(Schema::new(vec![
351            Field::new("ts", DataType::Timestamp(unit, None), false),
352            Field::new("__primary_key", DataType::Binary, false),
353            Field::new("__sequence", DataType::UInt64, false),
354            Field::new("__op_type", DataType::UInt8, false),
355        ]));
356        RecordBatch::try_new(
357            schema,
358            vec![
359                ts_column,
360                Arc::new(BinaryArray::from_iter_values(
361                    primary_keys.iter().map(Vec::as_slice),
362                )),
363                Arc::new(UInt64Array::from(vec![1; timestamps.len()])),
364                Arc::new(UInt8Array::from(vec![0; timestamps.len()])),
365            ],
366        )
367        .unwrap()
368    }
369
370    async fn write_index(
371        metadata: RegionMetadataRef,
372        object_store: ObjectStore,
373        rows: &[(u32, u64, &str, &str, i64)],
374        row_group_size: usize,
375    ) -> (SeriesIndexFileHandle, UnboundedReceiver<PurgeRequest>) {
376        write_index_with_time_unit(
377            metadata,
378            object_store,
379            rows,
380            row_group_size,
381            ArrowTimeUnit::Millisecond,
382        )
383        .await
384    }
385
386    async fn write_index_with_time_unit(
387        metadata: RegionMetadataRef,
388        object_store: ObjectStore,
389        rows: &[(u32, u64, &str, &str, i64)],
390        row_group_size: usize,
391        unit: ArrowTimeUnit,
392    ) -> (SeriesIndexFileHandle, UnboundedReceiver<PurgeRequest>) {
393        let (file_handle, receiver) = index_handle(&metadata, &object_store);
394        let file_id = file_handle.file_id();
395        let path = series_index_path(file_id.region_id(), file_id.file_id());
396        let primary_keys = rows
397            .iter()
398            .map(|(table_id, tsid, tag_0, tag_1, _)| {
399                new_sparse_primary_key(&[*tag_0, *tag_1], &metadata, *table_id, *tsid)
400            })
401            .collect::<Vec<_>>();
402        let timestamps = rows.iter().map(|row| row.4).collect::<Vec<_>>();
403        let mut writer = SeriesIndexWriter::try_new(
404            metadata,
405            object_store,
406            &path,
407            SeriesIndexWriterOptions { row_group_size },
408            None,
409        )
410        .await
411        .unwrap();
412        writer
413            .write(&flat_batch_with_time_unit(&primary_keys, &timestamps, unit))
414            .await
415            .unwrap();
416        writer.finish().await.unwrap();
417        (file_handle, receiver)
418    }
419
420    async fn collect_ids(stream: MetricSeriesIdStream) -> Vec<MetricSeriesId> {
421        stream
422            .try_collect::<Vec<_>>()
423            .await
424            .unwrap()
425            .into_iter()
426            .flatten()
427            .collect()
428    }
429
430    #[tokio::test]
431    async fn search_applies_candidate_tag_filters_and_time_overlap() {
432        let metadata = Arc::new(sst_region_metadata_with_encoding(
433            PrimaryKeyEncoding::Sparse,
434        ));
435        let object_store = object_store();
436        let (index, _receiver) = write_index(
437            metadata.clone(),
438            object_store.clone(),
439            &[
440                (1, 10, "a", "x", 10),
441                (1, 20, "b", "x", 20),
442                (1, 30, "a", "y", 30),
443            ],
444            2,
445        )
446        .await;
447
448        // The field filter is not available in the series index and is ignored,
449        // matching candidate-primary-key filter behavior.
450        let predicate = Predicate::new(vec![
451            col("tag_0").eq(lit("a")),
452            col("field_0").gt(lit(0_u64)),
453        ]);
454        let time_range = TimestampRange::new(
455            common_time::Timestamp::new_millisecond(20),
456            common_time::Timestamp::new_millisecond(31),
457        )
458        .unwrap();
459        let searcher = SeriesIndexSearcher::try_new(
460            metadata.clone(),
461            object_store.clone(),
462            index.clone(),
463            Some(&predicate),
464            Some(time_range),
465        )
466        .await
467        .unwrap();
468        let ids = collect_ids(searcher.search().unwrap()).await;
469        assert_eq!(
470            ids,
471            vec![MetricSeriesId {
472                table_id: 1,
473                tsid: 30
474            }]
475        );
476
477        // Both bounds fall between millisecond ticks. Rounding the inclusive
478        // start and exclusive end upward leaves only the 30 ms series.
479        let time_range = TimestampRange::new(
480            common_time::Timestamp::new_microsecond(20_001),
481            common_time::Timestamp::new_microsecond(30_001),
482        )
483        .unwrap();
484        let searcher = SeriesIndexSearcher::try_new(
485            metadata.clone(),
486            object_store.clone(),
487            index.clone(),
488            None,
489            Some(time_range),
490        )
491        .await
492        .unwrap();
493        let ids = collect_ids(searcher.search().unwrap()).await;
494        assert_eq!(
495            ids,
496            vec![MetricSeriesId {
497                table_id: 1,
498                tsid: 30
499            }]
500        );
501
502        // The stored maximum equal to the query start intersects, while the
503        // stored minimum equal to the exclusive query end does not.
504        let time_range = TimestampRange::new(
505            common_time::Timestamp::new_millisecond(20),
506            common_time::Timestamp::new_millisecond(30),
507        )
508        .unwrap();
509        let searcher =
510            SeriesIndexSearcher::try_new(metadata, object_store, index, None, Some(time_range))
511                .await
512                .unwrap();
513        let ids = collect_ids(searcher.search().unwrap()).await;
514        assert_eq!(
515            ids,
516            vec![MetricSeriesId {
517                table_id: 1,
518                tsid: 20
519            }]
520        );
521    }
522
523    #[tokio::test]
524    async fn search_skips_filters_for_columns_missing_from_older_index() {
525        let old_metadata = Arc::new(sst_region_metadata_with_encoding(
526            PrimaryKeyEncoding::Sparse,
527        ));
528        let object_store = object_store();
529        let (index, _receiver) = write_index(
530            old_metadata.clone(),
531            object_store.clone(),
532            &[
533                (1, 10, "a", "x", 10),
534                (1, 20, "b", "x", 20),
535                (1, 30, "a", "y", 30),
536            ],
537            2,
538        )
539        .await;
540
541        let mut builder = RegionMetadataBuilder::from_existing(old_metadata.as_ref().clone());
542        builder.push_column_metadata(ColumnMetadata {
543            column_schema: ColumnSchema::new("tag_2", ConcreteDataType::string_datatype(), true),
544            semantic_type: SemanticType::Tag,
545            column_id: 4,
546        });
547        let mut primary_key = old_metadata.primary_key.clone();
548        primary_key.push(4);
549        builder.primary_key(primary_key);
550        let current_metadata = Arc::new(builder.build().unwrap());
551
552        let predicate =
553            Predicate::new(vec![col("tag_0").eq(lit("a")), col("tag_2").eq(lit("new"))]);
554        let searcher = SeriesIndexSearcher::try_new(
555            current_metadata,
556            object_store,
557            index,
558            Some(&predicate),
559            None,
560        )
561        .await
562        .unwrap();
563        let ids = collect_ids(searcher.search().unwrap()).await;
564        assert_eq!(
565            ids,
566            vec![
567                MetricSeriesId {
568                    table_id: 1,
569                    tsid: 10,
570                },
571                MetricSeriesId {
572                    table_id: 1,
573                    tsid: 30,
574                },
575            ]
576        );
577    }
578
579    #[tokio::test]
580    async fn search_uses_file_time_unit_after_time_index_widen() {
581        // The index is written while the region's time index is milliseconds.
582        let metadata = Arc::new(sst_region_metadata_with_encoding(
583            PrimaryKeyEncoding::Sparse,
584        ));
585        let object_store = object_store();
586        let (index, _receiver) = write_index(
587            metadata.clone(),
588            object_store.clone(),
589            &[
590                (1, 10, "a", "x", 10),
591                (1, 20, "b", "x", 20),
592                (1, 30, "c", "x", 30),
593            ],
594            2,
595        )
596        .await;
597
598        // The region's time index is then widened to microseconds.
599        let mut widened = (*metadata).clone();
600        for column in &mut widened.column_metadatas {
601            if column.column_schema.name == "ts" {
602                column.column_schema.data_type = ConcreteDataType::timestamp_microsecond_datatype();
603            }
604        }
605        let widened = Arc::new(widened);
606
607        // A query range of [10.001ms, 25ms) expressed in microseconds must be
608        // compared against the file's millisecond bounds (ceil: [11ms,
609        // 25ms)): only the 20ms series intersects. Interpreting the file in
610        // microseconds instead would match nothing.
611        let time_range = TimestampRange::new(
612            common_time::Timestamp::new_microsecond(10_001),
613            common_time::Timestamp::new_microsecond(25_000),
614        )
615        .unwrap();
616        let searcher = SeriesIndexSearcher::try_new(
617            widened,
618            object_store.clone(),
619            index.clone(),
620            None,
621            Some(time_range),
622        )
623        .await
624        .unwrap();
625        let ids = collect_ids(searcher.search().unwrap()).await;
626        assert_eq!(
627            ids,
628            vec![MetricSeriesId {
629                table_id: 1,
630                tsid: 20
631            }]
632        );
633
634        // A millisecond-unit range behaves identically to the pre-widen
635        // searcher: [20ms, 30ms) keeps only the 20ms series.
636        let time_range = TimestampRange::new(
637            common_time::Timestamp::new_millisecond(20),
638            common_time::Timestamp::new_millisecond(30),
639        )
640        .unwrap();
641        let searcher =
642            SeriesIndexSearcher::try_new(metadata, object_store, index, None, Some(time_range))
643                .await
644                .unwrap();
645        let ids = collect_ids(searcher.search().unwrap()).await;
646        assert_eq!(
647            ids,
648            vec![MetricSeriesId {
649                table_id: 1,
650                tsid: 20
651            }]
652        );
653    }
654
655    #[tokio::test]
656    async fn search_reads_each_file_in_its_recorded_unit() {
657        // The first file is written while the region's time index is
658        // milliseconds; its series sit at 10ms and 20ms.
659        let metadata = Arc::new(sst_region_metadata_with_encoding(
660            PrimaryKeyEncoding::Sparse,
661        ));
662        let object_store = object_store();
663        let (index, _receiver) = write_index(
664            metadata.clone(),
665            object_store.clone(),
666            &[(1, 10, "a", "x", 10), (1, 20, "b", "x", 20)],
667            2,
668        )
669        .await;
670
671        // The region's time index is then widened to microseconds and a
672        // second file is written; its series sit at 15_000µs and 15µs.
673        let mut widened = (*metadata).clone();
674        for column in &mut widened.column_metadatas {
675            if column.column_schema.name == "ts" {
676                column.column_schema.data_type = ConcreteDataType::timestamp_microsecond_datatype();
677            }
678        }
679        let widened = Arc::new(widened);
680        let (new_index, _new_receiver) = write_index_with_time_unit(
681            widened.clone(),
682            object_store.clone(),
683            &[(1, 30, "c", "x", 15_000), (1, 40, "d", "x", 15)],
684            2,
685            ArrowTimeUnit::Microsecond,
686        )
687        .await;
688
689        // Each file is searched by its own searcher, which interprets it in
690        // the unit recorded in its schema. For [10.5ms, 25ms):
691        // - the millisecond file (ceil: [11ms, 25ms)) keeps only the 20ms
692        //   series; reading it in microseconds would match nothing;
693        // - the microsecond file keeps only the 15_000µs series; reading it
694        //   in milliseconds would also keep the 15µs series.
695        let time_range = TimestampRange::new(
696            common_time::Timestamp::new_microsecond(10_500),
697            common_time::Timestamp::new_microsecond(25_000),
698        )
699        .unwrap();
700        let searcher = SeriesIndexSearcher::try_new(
701            widened.clone(),
702            object_store.clone(),
703            index,
704            None,
705            Some(time_range),
706        )
707        .await
708        .unwrap();
709        let ids = collect_ids(searcher.search().unwrap()).await;
710        assert_eq!(
711            ids,
712            vec![MetricSeriesId {
713                table_id: 1,
714                tsid: 20
715            }]
716        );
717        let searcher =
718            SeriesIndexSearcher::try_new(widened, object_store, new_index, None, Some(time_range))
719                .await
720                .unwrap();
721        let ids = collect_ids(searcher.search().unwrap()).await;
722        assert_eq!(
723            ids,
724            vec![MetricSeriesId {
725                table_id: 1,
726                tsid: 30
727            }]
728        );
729    }
730
731    fn ts_field(name: &str, unit: Option<ArrowTimeUnit>) -> Field {
732        let data_type = match unit {
733            Some(unit) => DataType::Timestamp(unit, None),
734            None => DataType::Int64,
735        };
736        Field::new(name, data_type, false)
737    }
738
739    fn index_file_schema(
740        min_unit: Option<ArrowTimeUnit>,
741        max_unit: Option<ArrowTimeUnit>,
742    ) -> SchemaRef {
743        Arc::new(Schema::new(vec![
744            ts_field(MIN_TS_COLUMN, min_unit),
745            ts_field(MAX_TS_COLUMN, max_unit),
746            Field::new(ROW_COUNT_COLUMN, DataType::UInt64, false),
747            Field::new(TABLE_ID_COLUMN, DataType::UInt32, false),
748            Field::new(TSID_COLUMN, DataType::UInt64, false),
749        ]))
750    }
751
752    #[test]
753    fn validate_index_schema_rejects_unusable_columns() {
754        // Min/max columns that are not Timestamps are rejected.
755        let err = validate_index_schema(&index_file_schema(None, None))
756            .unwrap_err()
757            .to_string();
758        assert!(err.contains("must be a Timestamp, got Int64"), "{err}");
759
760        // The min and max columns must agree on the unit.
761        let err = validate_index_schema(&index_file_schema(
762            Some(ArrowTimeUnit::Millisecond),
763            Some(ArrowTimeUnit::Microsecond),
764        ))
765        .unwrap_err()
766        .to_string();
767        assert!(err.contains("have different time units"), "{err}");
768
769        assert_eq!(
770            TimeUnit::Nanosecond,
771            validate_index_schema(&index_file_schema(
772                Some(ArrowTimeUnit::Nanosecond),
773                Some(ArrowTimeUnit::Nanosecond),
774            ))
775            .unwrap()
776        );
777    }
778
779    #[tokio::test]
780    async fn search_streams_fixed_size_batches() {
781        let metadata = Arc::new(sst_region_metadata_with_encoding(
782            PrimaryKeyEncoding::Sparse,
783        ));
784        let object_store = object_store();
785        let rows = (0..501_u64)
786            .map(|tsid| (1, tsid, "a", "x", tsid as i64))
787            .collect::<Vec<_>>();
788        let (index, _receiver) =
789            write_index(metadata.clone(), object_store.clone(), &rows, 100).await;
790
791        let searcher = SeriesIndexSearcher::try_new(metadata, object_store, index, None, None)
792            .await
793            .unwrap();
794        let batches = searcher
795            .search()
796            .unwrap()
797            .try_collect::<Vec<_>>()
798            .await
799            .unwrap();
800        assert_eq!(batches.iter().map(Vec::len).collect::<Vec<_>>(), [500, 1]);
801        assert_eq!(
802            batches[0][0],
803            MetricSeriesId {
804                table_id: 1,
805                tsid: 0
806            }
807        );
808        assert_eq!(
809            batches[1][0],
810            MetricSeriesId {
811                table_id: 1,
812                tsid: 500
813            }
814        );
815    }
816
817    #[tokio::test]
818    async fn search_streams_pin_deleted_index_until_released() {
819        let metadata = Arc::new(sst_region_metadata_with_encoding(
820            PrimaryKeyEncoding::Sparse,
821        ));
822        let store = object_store();
823        let rows = (0..501_u64)
824            .map(|tsid| (1, tsid, "a", "x", tsid as i64))
825            .collect::<Vec<_>>();
826        let (index, mut receiver) = write_index(metadata.clone(), store.clone(), &rows, 100).await;
827        let file_id = index.file_id();
828        let searcher = SeriesIndexSearcher::try_new(metadata, store, index.clone(), None, None)
829            .await
830            .unwrap();
831        index.mark_deleted();
832        drop(index);
833        assert!(receiver.try_recv().is_err());
834
835        let mut stream = searcher.search().unwrap();
836        let cancelled = searcher.search().unwrap();
837        drop(searcher);
838        // Even a stream that has not been polled must pin its file.
839        assert!(receiver.try_recv().is_err());
840        assert_eq!(stream.try_next().await.unwrap().unwrap().len(), 500);
841        assert!(receiver.try_recv().is_err());
842        assert_eq!(
843            collect_ids(stream).await,
844            vec![MetricSeriesId {
845                table_id: 1,
846                tsid: 500
847            }]
848        );
849        assert!(receiver.try_recv().is_err());
850
851        // The unpolled stream is now the final owner; cancellation releases it.
852        drop(cancelled);
853        assert_eq!(receiver.try_recv().unwrap().file_id, file_id);
854    }
855
856    #[tokio::test]
857    async fn search_prunes_row_groups_and_empty_ranges() {
858        let metadata = Arc::new(sst_region_metadata_with_encoding(
859            PrimaryKeyEncoding::Sparse,
860        ));
861        let object_store = object_store();
862        let (index, _receiver) = write_index(
863            metadata.clone(),
864            object_store.clone(),
865            &[
866                (1, 0, "a", "x", 0),
867                (1, 1, "b", "x", 1),
868                (1, 2, "m", "x", 2),
869                (1, 3, "m", "x", 3),
870                (1, 4, "y", "x", 4),
871                (1, 5, "z", "x", 5),
872            ],
873            2,
874        )
875        .await;
876
877        let file_id = index.file_id();
878        let path = series_index_path(file_id.region_id(), file_id.file_id());
879        let reader = ParquetIndexReader::open(object_store.clone(), &path)
880            .await
881            .unwrap();
882        let unit = validate_index_schema(reader.schema()).unwrap();
883
884        // Tag predicates prune row groups through the tag-column statistics.
885        let predicate = Predicate::new(vec![col("tag_0").eq(lit("m"))]);
886        let mut filters = simple_tag_filters(&metadata, None, Some(&predicate));
887        let (pruning_predicate, _) = filters_for_schema(reader.schema(), &filters);
888        assert_eq!(reader.row_groups_to_read(&pruning_predicate), vec![1]);
889
890        // Time-range predicates prune row groups in the file's own unit:
891        // [2ms, 4ms) keeps only the row group holding the 2ms and 3ms series.
892        let time_range = TimestampRange::new(
893            common_time::Timestamp::new_millisecond(2),
894            common_time::Timestamp::new_millisecond(4),
895        )
896        .unwrap();
897        for expr in time_range_exprs(unit, Some(&time_range)) {
898            let filter = SimpleFilterEvaluator::try_new(&expr)
899                .context(UnexpectedSnafu {
900                    reason: "failed to build an internal series-index time filter",
901                })
902                .unwrap();
903            filters.push((expr, filter));
904        }
905        let (pruning_predicate, _) = filters_for_schema(reader.schema(), &filters);
906        assert_eq!(reader.row_groups_to_read(&pruning_predicate), vec![1]);
907
908        // An empty time range yields no series without opening the file.
909        let (missing_index, _receiver) = index_handle(&metadata, &object_store);
910        let empty = SeriesIndexSearcher::try_new(
911            metadata,
912            object_store,
913            missing_index,
914            None,
915            Some(TimestampRange::empty()),
916        )
917        .await
918        .unwrap();
919        assert!(
920            empty
921                .search()
922                .unwrap()
923                .try_collect::<Vec<_>>()
924                .await
925                .unwrap()
926                .is_empty()
927        );
928    }
929}