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