Skip to main content

mito2/read/
series_reader.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
15//! Reads selected metric series from partition ranges.
16
17use std::collections::HashSet;
18use std::sync::Arc;
19use std::time::Instant;
20
21use async_stream::try_stream;
22use futures::TryStreamExt;
23use mito_codec::row_converter::{PrimaryKeyFilter, SparsePrimaryKeyCodec};
24use snafu::ResultExt;
25use store_api::region_engine::PartitionRange;
26use tokio::sync::Semaphore;
27
28#[cfg(feature = "enterprise")]
29use crate::error::InvalidRequestSnafu;
30use crate::error::{JoinSnafu, Result, UnexpectedSnafu};
31use crate::read::BoxedRecordBatchStream;
32use crate::read::pruner::PartitionPruner;
33use crate::read::range_cache::{
34    build_series_range_cache_key, cache_flat_range_stream, cached_flat_range_stream,
35};
36use crate::read::scan_region::StreamContext;
37use crate::read::scan_util::{
38    PartitionMetrics, SplitRecordBatchStream, compute_average_batch_size,
39    compute_parallel_channel_size, new_filter_metrics, scan_flat_mem_ranges,
40    should_split_flat_batches_for_merge,
41};
42use crate::read::seq_scan::SeqScan;
43use crate::read::series_candidate::validate_metric_metadata;
44use crate::series_index::MetricSeriesId;
45use crate::sst::parquet::DEFAULT_READ_BATCH_SIZE;
46use crate::sst::parquet::flat_format::primary_key_column_index;
47use crate::sst::parquet::prefilter::prefilter_flat_batch_by_primary_key;
48use crate::sst::parquet::reader::ReaderMetrics;
49use crate::sst::parquet::row_group::ParquetFetchMetrics;
50
51const TSID_DOMAIN_END: u128 = 1u128 << u64::BITS;
52
53/// A stable partition of the TSID integer domain.
54#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
55pub(crate) struct SeriesRange {
56    start: u128,
57    end: u128,
58}
59
60impl SeriesRange {
61    pub(crate) fn new(partition: usize, partitions: usize) -> Option<Self> {
62        if partitions == 0 || partition >= partitions {
63            return None;
64        }
65
66        let partitions = partitions as u128;
67        let partition = partition as u128;
68        let boundary = |partition: u128| {
69            let numerator = partition * TSID_DOMAIN_END;
70            numerator.div_ceil(partitions)
71        };
72        Some(Self {
73            start: boundary(partition),
74            end: boundary(partition + 1),
75        })
76    }
77
78    fn partition_for(tsid: u64, partitions: usize) -> usize {
79        ((tsid as u128 * partitions as u128) >> u64::BITS) as usize
80    }
81}
82
83/// All series assigned to one data-reader partition.
84#[derive(Debug)]
85pub(crate) struct AssignedSeriesBatch {
86    range: SeriesRange,
87    series: Vec<MetricSeriesId>,
88    enable_range_cache: bool,
89}
90
91impl AssignedSeriesBatch {
92    fn new(range: SeriesRange, series: Vec<MetricSeriesId>, enable_range_cache: bool) -> Self {
93        Self {
94            range,
95            series,
96            enable_range_cache,
97        }
98    }
99
100    pub(crate) fn range(&self) -> SeriesRange {
101        self.range
102    }
103
104    pub(crate) fn series(&self) -> &[MetricSeriesId] {
105        &self.series
106    }
107
108    pub(crate) fn enable_range_cache(&self) -> bool {
109        self.enable_range_cache
110    }
111}
112
113/// Collects candidate batches and assigns every TSID by its integer range.
114pub(crate) struct SeriesBatchCollector {
115    assignments: Vec<Vec<MetricSeriesId>>,
116    num_series: usize,
117}
118
119impl SeriesBatchCollector {
120    pub(crate) fn new(partitions: usize) -> Option<Self> {
121        (partitions > 0).then(|| Self {
122            assignments: (0..partitions).map(|_| Vec::new()).collect(),
123            num_series: 0,
124        })
125    }
126
127    pub(crate) fn push(&mut self, batch: Vec<MetricSeriesId>) {
128        self.num_series += batch.len();
129        let partitions = self.assignments.len();
130        for series in batch {
131            let partition = SeriesRange::partition_for(series.tsid, partitions);
132            self.assignments[partition].push(series);
133        }
134    }
135
136    pub(crate) fn len(&self) -> usize {
137        self.num_series
138    }
139
140    pub(crate) fn finish(self, enable_range_cache: bool) -> Vec<AssignedSeriesBatch> {
141        let partitions = self.assignments.len();
142        self.assignments
143            .into_iter()
144            .enumerate()
145            .map(|(partition, series)| {
146                AssignedSeriesBatch::new(
147                    SeriesRange::new(partition, partitions).unwrap(),
148                    series,
149                    enable_range_cache,
150                )
151            })
152            .collect()
153    }
154}
155
156/// Immutable allow-list for one partition's metric series.
157#[derive(Clone, Debug)]
158struct MetricSeriesFilter {
159    range: SeriesRange,
160    series: Arc<HashSet<MetricSeriesId>>,
161    enable_range_cache: bool,
162}
163
164impl MetricSeriesFilter {
165    fn new(assigned: &AssignedSeriesBatch) -> Self {
166        let series = assigned.series().iter().copied().collect();
167        Self {
168            range: assigned.range(),
169            series: Arc::new(series),
170            enable_range_cache: assigned.enable_range_cache(),
171        }
172    }
173
174    fn primary_key_filter(&self, codec: SparsePrimaryKeyCodec) -> Box<dyn PrimaryKeyFilter> {
175        Box::new(MetricSeriesPrimaryKeyFilter {
176            codec,
177            series: self.series.clone(),
178            last_primary_key: Vec::new(),
179            last_match: None,
180        })
181    }
182
183    fn overlaps_encoded_bounds(
184        &self,
185        codec: &SparsePrimaryKeyCodec,
186        encoded_min: &[u8],
187        encoded_max: &[u8],
188    ) -> Option<bool> {
189        let (min_table_id, min_tsid) = codec.decode_ids(encoded_min).ok()?;
190        let (max_table_id, max_tsid) = codec.decode_ids(encoded_max).ok()?;
191        let min = MetricSeriesId {
192            table_id: min_table_id,
193            tsid: min_tsid,
194        };
195        let max = MetricSeriesId {
196            table_id: max_table_id,
197            tsid: max_tsid,
198        };
199        if min > max {
200            return None;
201        }
202        if min_table_id != max_table_id {
203            return Some(true);
204        }
205
206        Some(u128::from(min_tsid) < self.range.end && u128::from(max_tsid) >= self.range.start)
207    }
208}
209
210struct MetricSeriesPrimaryKeyFilter {
211    codec: SparsePrimaryKeyCodec,
212    series: Arc<HashSet<MetricSeriesId>>,
213    last_primary_key: Vec<u8>,
214    last_match: Option<bool>,
215}
216
217impl PrimaryKeyFilter for MetricSeriesPrimaryKeyFilter {
218    fn matches(&mut self, primary_key: &[u8]) -> mito_codec::error::Result<bool> {
219        if let Some(last_match) = self.last_match
220            && self.last_primary_key == primary_key
221        {
222            return Ok(last_match);
223        }
224
225        let (table_id, tsid) = self.codec.decode_ids(primary_key)?;
226        let matched = self.series.contains(&MetricSeriesId { table_id, tsid });
227        self.last_primary_key.clear();
228        self.last_primary_key.extend_from_slice(primary_key);
229        self.last_match = Some(matched);
230        Ok(matched)
231    }
232}
233
234fn filter_flat_stream_by_series(
235    mut input: BoxedRecordBatchStream,
236    codec: SparsePrimaryKeyCodec,
237    filter: MetricSeriesFilter,
238) -> BoxedRecordBatchStream {
239    Box::pin(try_stream! {
240        let mut primary_key_filter = filter.primary_key_filter(codec);
241        while let Some(batch) = input.try_next().await? {
242            let pk_idx = primary_key_column_index(batch.num_columns());
243            if let Some(batch) = prefilter_flat_batch_by_primary_key(
244                batch,
245                pk_idx,
246                primary_key_filter.as_mut(),
247            )? {
248                yield batch;
249            }
250        }
251    })
252}
253
254/// Reads all collected metric series assigned to one partition.
255pub(crate) struct SeriesReader {
256    stream_ctx: Arc<StreamContext>,
257    partition_ranges: Vec<PartitionRange>,
258    range: SeriesRange,
259    filter: MetricSeriesFilter,
260    codec: SparsePrimaryKeyCodec,
261    partition_pruner: Arc<PartitionPruner>,
262    range_semaphore: Arc<Semaphore>,
263    part_metrics: PartitionMetrics,
264}
265
266impl SeriesReader {
267    /// Creates a reader for the series assigned to one data partition.
268    ///
269    /// `partition_pruner` must come from the candidate scanner's pruner, which
270    /// has predicate prefiltering disabled. The file read path applies precise
271    /// filters with tag filtering skipped because candidate discovery has already
272    /// enforced tag predicates through the exact assigned-series set.
273    ///
274    /// A single `range_semaphore` covers both the range-build phase and the final
275    /// merge.
276    pub(crate) fn try_new(
277        stream_ctx: Arc<StreamContext>,
278        partition_ranges: Vec<PartitionRange>,
279        assigned_series: AssignedSeriesBatch,
280        partition_pruner: Arc<PartitionPruner>,
281        range_semaphore: Arc<Semaphore>,
282        part_metrics: PartitionMetrics,
283    ) -> Result<Self> {
284        validate_metric_metadata(&stream_ctx)?;
285        #[cfg(feature = "enterprise")]
286        snafu::ensure!(
287            stream_ctx.input.extension_ranges().is_empty(),
288            InvalidRequestSnafu {
289                region_id: stream_ctx.input.region_metadata().region_id,
290                reason: "series reader does not support extension ranges",
291            }
292        );
293
294        let range = assigned_series.range();
295        let filter = MetricSeriesFilter::new(&assigned_series);
296        let codec = SparsePrimaryKeyCodec::new(stream_ctx.input.region_metadata());
297        Ok(Self {
298            stream_ctx,
299            partition_ranges,
300            range,
301            filter,
302            codec,
303            partition_pruner,
304            range_semaphore,
305            part_metrics,
306        })
307    }
308
309    pub(crate) async fn build_stream(&self) -> Result<BoxedRecordBatchStream> {
310        if self.partition_ranges.is_empty() || self.filter.series.is_empty() {
311            return Ok(Box::pin(futures::stream::empty()));
312        }
313
314        let mut tasks = Vec::with_capacity(self.partition_ranges.len());
315        for part_range in self.partition_ranges.iter().copied() {
316            let stream_ctx = self.stream_ctx.clone();
317            let filter = self.filter.clone();
318            let codec = self.codec.clone();
319            let partition_pruner = self.partition_pruner.clone();
320            let range_semaphore = self.range_semaphore.clone();
321            let part_metrics = self.part_metrics.clone();
322            let range = self.range;
323            tasks.push(common_runtime::spawn_query(async move {
324                let _permit = range_semaphore.acquire().await.map_err(|error| {
325                    UnexpectedSnafu {
326                        reason: format!("failed to acquire series range permit: {error}"),
327                    }
328                    .build()
329                })?;
330                build_series_partition_range(
331                    stream_ctx,
332                    part_range,
333                    range,
334                    filter,
335                    codec,
336                    partition_pruner,
337                    part_metrics,
338                )
339                .await
340            }));
341        }
342
343        let mut range_streams = Vec::with_capacity(tasks.len());
344        let mut estimated_batch_sizes = Vec::with_capacity(tasks.len());
345        for task in tasks {
346            let (stream, estimated_batch_size) = task.await.context(JoinSnafu)??;
347            range_streams.push(stream);
348            estimated_batch_sizes.push(estimated_batch_size);
349        }
350
351        // Every range task above has finished, so all build permits are released
352        // and the final merge can reuse the same semaphore.
353        let estimated_batch_size = compute_average_batch_size(estimated_batch_sizes);
354        SeqScan::build_flat_reader_from_sources(
355            &self.stream_ctx,
356            range_streams,
357            Some(self.range_semaphore.clone()),
358            Some(&self.part_metrics),
359            true,
360            compute_parallel_channel_size(estimated_batch_size),
361        )
362        .await
363    }
364}
365
366async fn build_series_partition_range(
367    stream_ctx: Arc<StreamContext>,
368    part_range: PartitionRange,
369    range: SeriesRange,
370    filter: MetricSeriesFilter,
371    codec: SparsePrimaryKeyCodec,
372    partition_pruner: Arc<PartitionPruner>,
373    part_metrics: PartitionMetrics,
374) -> Result<(BoxedRecordBatchStream, usize)> {
375    let cache_key = filter
376        .enable_range_cache
377        .then(|| build_series_range_cache_key(&stream_ctx, &part_range, range))
378        .flatten();
379    if let Some(key) = cache_key.as_ref() {
380        if let Some(value) = stream_ctx.input.cache_strategy.get_range_result(key) {
381            part_metrics.inc_range_cache_hit();
382            return Ok((cached_flat_range_stream(value), DEFAULT_READ_BATCH_SIZE));
383        }
384        part_metrics.inc_range_cache_miss();
385    }
386
387    let range_meta = &stream_ctx.ranges[part_range.identifier];
388    let split_batch_size = should_split_flat_batches_for_merge(&stream_ctx, range_meta);
389    let mut sources = Vec::with_capacity(range_meta.row_group_indices.len());
390
391    for index in range_meta.row_group_indices.iter().copied() {
392        if stream_ctx.is_mem_range_index(index) {
393            let stream = Box::pin(scan_flat_mem_ranges(
394                stream_ctx.clone(),
395                part_metrics.clone(),
396                index,
397                range_meta.time_range,
398            ));
399            sources.push(filter_flat_stream_by_series(
400                stream,
401                codec.clone(),
402                filter.clone(),
403            ));
404            continue;
405        }
406
407        if stream_ctx.is_file_range_index(index) {
408            let file = stream_ctx.input.file_from_index(index);
409            if matches!(
410                file.primary_key_range()
411                    .and_then(|(min, max)| { filter.overlaps_encoded_bounds(&codec, &min, &max) }),
412                Some(false)
413            ) {
414                continue;
415            }
416
417            if partition_pruner.try_skip_manifest_pruned_file_range(index, &part_metrics) {
418                continue;
419            }
420            let mut reader_metrics = ReaderMetrics {
421                filter_metrics: new_filter_metrics(part_metrics.explain_verbose()),
422                ..Default::default()
423            };
424            let file_ranges = partition_pruner
425                .build_file_ranges(index, &part_metrics, &mut reader_metrics)
426                .await?;
427            part_metrics.inc_num_file_ranges(file_ranges.len());
428            part_metrics.merge_reader_metrics(&reader_metrics, None);
429            let ranges = file_ranges
430                .iter()
431                .filter(|file_range| {
432                    !matches!(
433                        file_range.primary_key_range().and_then(|(min, max)| {
434                            filter.overlaps_encoded_bounds(&codec, min, max)
435                        }),
436                        Some(false)
437                    )
438                })
439                .cloned()
440                .collect::<smallvec::SmallVec<[_; 2]>>();
441            if ranges.is_empty() {
442                continue;
443            }
444
445            let stream = scan_series_file_ranges(
446                part_metrics.clone(),
447                ranges,
448                filter.clone(),
449                codec.clone(),
450            );
451            sources.push(Box::pin(stream) as BoxedRecordBatchStream);
452            continue;
453        }
454
455        return UnexpectedSnafu {
456            reason: format!(
457                "series reader received unsupported range index {}",
458                index.index
459            ),
460        }
461        .fail();
462    }
463
464    if split_batch_size.is_some() {
465        sources = sources
466            .into_iter()
467            .map(|stream| Box::pin(SplitRecordBatchStream::new(stream)) as BoxedRecordBatchStream)
468            .collect();
469    }
470    let estimated_batch_size = split_batch_size.unwrap_or(DEFAULT_READ_BATCH_SIZE);
471    let stream = SeqScan::build_flat_reader_from_sources(
472        &stream_ctx,
473        sources,
474        None,
475        Some(&part_metrics),
476        false,
477        compute_parallel_channel_size(estimated_batch_size),
478    )
479    .await?;
480    let stream = match cache_key {
481        Some(key) => cache_flat_range_stream(
482            stream,
483            stream_ctx.input.cache_strategy.clone(),
484            key,
485            part_metrics,
486        ),
487        None => stream,
488    };
489    Ok((stream, estimated_batch_size))
490}
491
492fn scan_series_file_ranges(
493    part_metrics: PartitionMetrics,
494    ranges: smallvec::SmallVec<[crate::sst::parquet::file_range::FileRange; 2]>,
495    filter: MetricSeriesFilter,
496    codec: SparsePrimaryKeyCodec,
497) -> impl futures::Stream<Item = Result<datatypes::arrow::record_batch::RecordBatch>> {
498    try_stream! {
499        let fetch_metrics = part_metrics
500            .explain_verbose()
501            .then(|| Arc::new(ParquetFetchMetrics::default()));
502        let mut reader_metrics = ReaderMetrics {
503            fetch_metrics: fetch_metrics.clone(),
504            ..Default::default()
505        };
506        let mut primary_key_filter = filter.primary_key_filter(codec);
507
508        for range in ranges {
509            let build_start = Instant::now();
510            let Some(mut reader) = range
511                .reader_by_primary_key(
512                    primary_key_filter.as_mut(),
513                    fetch_metrics.as_deref(),
514                )
515                .await?
516            else {
517                continue;
518            };
519            let build_cost = build_start.elapsed();
520            reader_metrics.build_cost += build_cost;
521            part_metrics.inc_build_reader_cost(build_cost);
522
523            let scan_start = Instant::now();
524            while let Some(record_batch) = reader.next_batch().await? {
525                reader_metrics.num_record_batches += 1;
526                reader_metrics.num_batches += 1;
527                reader_metrics.num_rows += record_batch.num_rows();
528
529                let num_rows_before_filter = record_batch.num_rows();
530                let Some(record_batch) = range.precise_filter_flat(
531                    record_batch,
532                    range.pre_filter_mode().skip_fields(),
533                    true,
534                )? else {
535                    reader_metrics.filter_metrics.rows_precise_filtered +=
536                        num_rows_before_filter;
537                    continue;
538                };
539                reader_metrics.filter_metrics.rows_precise_filtered +=
540                    num_rows_before_filter - record_batch.num_rows();
541
542                let record_batch = if let Some(mapper) = range.compaction_projection_mapper() {
543                    mapper.project(record_batch)?
544                } else {
545                    record_batch
546                };
547                if let Some(compat) = range.compat_batch() {
548                    yield compat.compat(record_batch)?;
549                } else {
550                    yield record_batch;
551                }
552            }
553            reader_metrics.scan_cost += scan_start.elapsed();
554        }
555
556        reader_metrics.observe_rows("series_data");
557        reader_metrics.filter_metrics.observe();
558        part_metrics.merge_reader_metrics(&reader_metrics, None);
559    }
560}
561
562#[cfg(test)]
563mod tests {
564    use store_api::codec::PrimaryKeyEncoding;
565
566    use super::*;
567    use crate::error::DecodeSnafu;
568    use crate::test_util::sst_util::sst_region_metadata_with_encoding;
569
570    fn series(table_id: u32, tsid: u64) -> MetricSeriesId {
571        MetricSeriesId { table_id, tsid }
572    }
573
574    fn assigned_batch(
575        partitions: usize,
576        partition: usize,
577        series: Vec<MetricSeriesId>,
578    ) -> AssignedSeriesBatch {
579        let mut collector = SeriesBatchCollector::new(partitions).unwrap();
580        collector.push(series);
581        collector.finish(true).remove(partition)
582    }
583
584    #[test]
585    fn series_range_assignment_is_stable_across_batch_boundaries() {
586        let input = vec![
587            series(1, 0),
588            series(2, 1u64 << 62),
589            series(1, 1u64 << 63),
590            series(2, 3u64 << 62),
591            series(1, u64::MAX),
592            series(3, 0),
593        ];
594        let mut first = SeriesBatchCollector::new(4).unwrap();
595        first.push(input.clone());
596        let first = first.finish(true);
597
598        let mut second = SeriesBatchCollector::new(4).unwrap();
599        for chunk in input.chunks(2) {
600            second.push(chunk.to_vec());
601        }
602        let second = second.finish(true);
603
604        assert_eq!(
605            first
606                .iter()
607                .map(AssignedSeriesBatch::series)
608                .collect::<Vec<_>>(),
609            second
610                .iter()
611                .map(AssignedSeriesBatch::series)
612                .collect::<Vec<_>>()
613        );
614        assert_eq!(&[series(1, 0), series(3, 0)], first[0].series());
615        assert_eq!(&[series(2, 1u64 << 62)], first[1].series());
616        assert_eq!(
617            input.len(),
618            first
619                .iter()
620                .map(|batch| batch.series().len())
621                .sum::<usize>()
622        );
623    }
624
625    #[test]
626    fn collector_tracks_size_and_cache_policy() {
627        let mut collector = SeriesBatchCollector::new(2).unwrap();
628        collector.push(vec![series(1, 1), series(1, u64::MAX)]);
629        collector.push(vec![series(2, 2)]);
630        assert_eq!(3, collector.len());
631
632        let assignments = collector.finish(false);
633        assert_eq!(
634            3,
635            assignments
636                .iter()
637                .map(|batch| batch.series().len())
638                .sum::<usize>()
639        );
640        assert!(assignments.iter().all(|batch| !batch.enable_range_cache()));
641    }
642
643    #[test]
644    fn series_ranges_cover_the_tsid_domain() {
645        let ranges = (0..3)
646            .map(|partition| SeriesRange::new(partition, 3).unwrap())
647            .collect::<Vec<_>>();
648        assert_eq!(0, ranges[0].start);
649        assert_eq!(TSID_DOMAIN_END, ranges[2].end);
650        assert_eq!(ranges[0].end, ranges[1].start);
651        assert_eq!(ranges[1].end, ranges[2].start);
652        assert_eq!(0, SeriesRange::partition_for(0, 3));
653        assert_eq!(2, SeriesRange::partition_for(u64::MAX, 3));
654
655        for (partition, range) in ranges.iter().enumerate() {
656            assert_eq!(partition, SeriesRange::partition_for(range.start as u64, 3));
657            if range.end < TSID_DOMAIN_END {
658                assert_eq!(
659                    partition + 1,
660                    SeriesRange::partition_for(range.end as u64, 3)
661                );
662            }
663        }
664    }
665
666    #[test]
667    fn metric_series_filter_matches_encoded_primary_key() {
668        let metadata = Arc::new(sst_region_metadata_with_encoding(
669            PrimaryKeyEncoding::Sparse,
670        ));
671        let assigned = assigned_batch(1, 0, vec![series(1, 10), series(2, 20)]);
672        let filter = MetricSeriesFilter::new(&assigned);
673        let codec = SparsePrimaryKeyCodec::new(&metadata);
674        let mut primary_key_filter = filter.primary_key_filter(codec);
675
676        for selected in [series(1, 10), series(2, 20)] {
677            assert!(
678                primary_key_filter
679                    .matches(&encode_series(&metadata, selected.table_id, selected.tsid))
680                    .context(DecodeSnafu)
681                    .unwrap()
682            );
683        }
684        for unselected in [series(2, 10), series(1, 20)] {
685            assert!(
686                !primary_key_filter
687                    .matches(&encode_series(
688                        &metadata,
689                        unselected.table_id,
690                        unselected.tsid,
691                    ))
692                    .context(DecodeSnafu)
693                    .unwrap()
694            );
695        }
696    }
697
698    fn encode_series(
699        metadata: &store_api::metadata::RegionMetadataRef,
700        table_id: u32,
701        tsid: u64,
702    ) -> Vec<u8> {
703        let codec = SparsePrimaryKeyCodec::new(metadata);
704        let mut primary_key = Vec::new();
705        codec
706            .encode_internal(table_id, tsid, &mut primary_key)
707            .unwrap();
708        primary_key
709    }
710
711    #[test]
712    fn series_range_overlaps_single_table_primary_key_bounds() {
713        let metadata = Arc::new(sst_region_metadata_with_encoding(
714            PrimaryKeyEncoding::Sparse,
715        ));
716        let assigned = assigned_batch(2, 0, vec![series(1, 10), series(2, 20)]);
717        let filter = MetricSeriesFilter::new(&assigned);
718        let codec = SparsePrimaryKeyCodec::new(&metadata);
719
720        // The full assigned range overlaps even though the actual candidate bounds do not.
721        assert_eq!(
722            Some(true),
723            filter.overlaps_encoded_bounds(
724                &codec,
725                &encode_series(&metadata, 1, 100),
726                &encode_series(&metadata, 1, 200),
727            )
728        );
729        assert_eq!(
730            Some(false),
731            filter.overlaps_encoded_bounds(
732                &codec,
733                &encode_series(&metadata, 1, 1u64 << 63),
734                &encode_series(&metadata, 1, u64::MAX),
735            )
736        );
737        assert_eq!(
738            Some(true),
739            filter.overlaps_encoded_bounds(
740                &codec,
741                &encode_series(&metadata, 3, 10),
742                &encode_series(&metadata, 3, 20),
743            )
744        );
745    }
746
747    #[test]
748    fn series_range_keeps_multiple_table_primary_key_bounds() {
749        let metadata = Arc::new(sst_region_metadata_with_encoding(
750            PrimaryKeyEncoding::Sparse,
751        ));
752        let assigned = assigned_batch(2, 0, vec![series(1, 10), series(2, 20)]);
753        let filter = MetricSeriesFilter::new(&assigned);
754        let codec = SparsePrimaryKeyCodec::new(&metadata);
755
756        assert_eq!(
757            Some(true),
758            filter.overlaps_encoded_bounds(
759                &codec,
760                &encode_series(&metadata, 1, 1u64 << 63),
761                &encode_series(&metadata, 2, u64::MAX),
762            )
763        );
764    }
765
766    #[test]
767    fn series_statistics_keep_invalid_or_inverted_bounds() {
768        let metadata = Arc::new(sst_region_metadata_with_encoding(
769            PrimaryKeyEncoding::Sparse,
770        ));
771        let assigned = assigned_batch(2, 0, vec![series(1, 10)]);
772        let filter = MetricSeriesFilter::new(&assigned);
773        let codec = SparsePrimaryKeyCodec::new(&metadata);
774
775        assert_eq!(
776            None,
777            filter.overlaps_encoded_bounds(&codec, b"invalid", b"bounds")
778        );
779        assert_eq!(
780            None,
781            filter.overlaps_encoded_bounds(
782                &codec,
783                &encode_series(&metadata, 2, 0),
784                &encode_series(&metadata, 1, 0),
785            )
786        );
787    }
788}