1use std::collections::BTreeMap;
18use std::time::Duration;
19
20use common_time::{TimeToLive, Timestamp};
21use smallvec::{SmallVec, smallvec};
22use store_api::storage::FileId;
23
24use crate::series_index::catalog::{SeriesIndexEntry, WindowSequence};
25use crate::sst::file::FileHandle;
26
27const SERIES_INDEX_TRIGGER_FILES: usize = 4;
28
29#[derive(Debug, Clone, PartialEq, Eq)]
31pub(crate) struct IndexBucket {
32 pub(crate) start: Timestamp,
33 pub(crate) end: Timestamp,
34 pub(crate) index_ids: SmallVec<[FileId; 2]>,
35 pub(crate) compaction_window_secs: i64,
37 pub(crate) window_sequences: BTreeMap<i64, WindowSequence>,
40}
41
42impl IndexBucket {
43 pub(crate) fn from_entry(entry: &SeriesIndexEntry) -> Self {
44 Self {
45 start: entry.bucket_start,
46 end: entry.bucket_end,
47 index_ids: smallvec![entry.index_uuid],
48 compaction_window_secs: entry.compaction_window_secs,
49 window_sequences: entry.window_sequences.clone(),
50 }
51 }
52
53 fn merge(&mut self, mut other: Self) {
54 self.start = self.start.min(other.start);
55 self.end = self.end.max(other.end);
56 if self.index_ids.is_empty() {
57 self.compaction_window_secs = other.compaction_window_secs;
58 self.window_sequences = std::mem::take(&mut other.window_sequences);
59 } else if !other.index_ids.is_empty() {
60 if self.compaction_window_secs == other.compaction_window_secs {
61 merge_window_sequences(&mut self.window_sequences, other.window_sequences);
62 } else {
63 self.compaction_window_secs = 0;
66 self.window_sequences.clear();
67 }
68 }
69 self.index_ids.append(&mut other.index_ids);
70 }
71
72 pub(crate) fn insert_into(mut self, buckets: &mut BTreeMap<Timestamp, Self>) {
75 if let Some((&previous_start, previous)) = buckets.range(..self.start).next_back()
76 && previous.end > self.start
77 {
78 self.start = previous_start;
79 }
80 while let Some((&next_start, _)) = buckets.range(self.start..self.end).next() {
82 if let Some(other) = buckets.remove(&next_start) {
83 self.merge(other);
84 }
85 }
86 buckets.insert(self.start, self);
87 }
88}
89
90#[derive(Debug, Clone)]
96pub(crate) struct SeriesBucket {
97 pub(crate) start: Timestamp,
98 pub(crate) end: Timestamp,
99 pub(crate) files: SmallVec<[FileHandle; 2]>,
100 pub(crate) has_unknown_sequence: bool,
101 pub(crate) max_file_sequence: u64,
103 pub(crate) compaction_window_secs: i64,
104 pub(crate) window_sequences: BTreeMap<i64, WindowSequence>,
110}
111
112pub(crate) struct SeriesIndexPlan {
114 pub(crate) index_buckets: BTreeMap<Timestamp, IndexBucket>,
115 pub(crate) builds: Vec<(SeriesBucket, SeriesIndexEntry)>,
116 pub(crate) expired_index_ids: Vec<FileId>,
117 pub(crate) superseded_index_ids: Vec<FileId>,
119 pub(crate) computed_buckets: usize,
120 pub(crate) skipped_buckets: usize,
121}
122
123pub(crate) fn plan_series_indexes(
126 buckets: Vec<SeriesBucket>,
127 mut index_buckets: BTreeMap<Timestamp, IndexBucket>,
128 ttl: Option<TimeToLive>,
129 now_ms: i64,
130) -> SeriesIndexPlan {
131 let buckets = reconcile_series_buckets(buckets, &index_buckets);
134 let computed_buckets = buckets.len();
135 let expired = |end| {
136 ttl.is_some_and(|ttl| {
137 ttl.is_expired(&end, &Timestamp::new_millisecond(now_ms))
138 .unwrap_or(false)
139 })
140 };
141 let mut expired_index_ids = Vec::new();
142 index_buckets.retain(|_, bucket| {
143 if expired(bucket.end) {
144 expired_index_ids.extend_from_slice(&bucket.index_ids);
145 false
146 } else {
147 true
148 }
149 });
150 let mut builds = Vec::new();
151 let mut superseded_index_ids = Vec::new();
152 for bucket in buckets {
153 if expired(bucket.end) {
154 continue;
155 }
156 if bucket.has_unknown_sequence {
157 continue;
158 }
159 if index_buckets.get(&bucket.start).is_some_and(|indexed| {
160 indexed.end == bucket.end
161 && indexed.compaction_window_secs == bucket.compaction_window_secs
162 && indexed.window_sequences == bucket.window_sequences
163 }) {
164 continue;
165 }
166 if let Some(entry) = bucket.to_series_entry() {
167 index_buckets.retain(|_, indexed| {
168 if indexed.start < bucket.end && bucket.start < indexed.end {
169 superseded_index_ids.extend_from_slice(&indexed.index_ids);
170 false
171 } else {
172 true
173 }
174 });
175 IndexBucket::from_entry(&entry).insert_into(&mut index_buckets);
176 builds.push((bucket, entry));
177 }
178 }
179 index_buckets.retain(|_, bucket| !bucket.index_ids.is_empty());
180 SeriesIndexPlan {
181 index_buckets,
182 skipped_buckets: computed_buckets - builds.len(),
183 builds,
184 expired_index_ids,
185 superseded_index_ids,
186 computed_buckets,
187 }
188}
189
190pub(crate) fn rounded_bucket_width(
191 requested: Duration,
192 compaction_window: Duration,
193) -> Option<i64> {
194 let window_secs = i64::try_from(compaction_window.as_secs()).ok()?.max(1);
195 let requested_secs = i64::try_from(requested.as_secs())
196 .unwrap_or(i64::MAX)
197 .max(1);
198 let multiples = requested_secs / window_secs + i64::from(requested_secs % window_secs != 0);
199 multiples.checked_mul(window_secs)
200}
201
202pub(crate) fn group_files_into_series_buckets(
209 files: &[FileHandle],
210 width_secs: i64,
211 compaction_window_secs: i64,
212) -> Vec<SeriesBucket> {
213 let mut spans = files
214 .iter()
215 .map(|file| {
216 let start = file.time_range().0.split().0;
217 let end = file.time_range().1.split().0;
218 let sequence = file
219 .meta_ref()
220 .sequence
221 .map_or(0, |sequence| sequence.get());
222 let first_window = start.div_euclid(compaction_window_secs);
223 let last_window = end.div_euclid(compaction_window_secs);
224 let window_sequences = BTreeMap::from([(
225 first_window.saturating_mul(compaction_window_secs),
226 WindowSequence {
227 start: first_window.saturating_mul(compaction_window_secs),
228 end: last_window
229 .saturating_add(1)
230 .saturating_mul(compaction_window_secs),
231 max_sequence: sequence,
232 },
233 )]);
234 SeriesBucket {
235 start: Timestamp::new_second(
236 start.div_euclid(width_secs).saturating_mul(width_secs),
237 ),
238 end: Timestamp::new_second(
239 end.div_euclid(width_secs)
240 .saturating_add(1)
241 .saturating_mul(width_secs),
242 ),
243 files: smallvec![file.clone()],
244 has_unknown_sequence: file.meta_ref().sequence.is_none(),
245 max_file_sequence: sequence,
246 compaction_window_secs,
247 window_sequences,
248 }
249 })
250 .collect::<Vec<_>>();
251 spans.sort_unstable_by(|a, b| a.start.cmp(&b.start).then_with(|| b.end.cmp(&a.end)));
252 group_series_buckets(spans)
253}
254
255fn reconcile_series_buckets(
262 mut buckets: Vec<SeriesBucket>,
263 index_buckets: &BTreeMap<Timestamp, IndexBucket>,
264) -> Vec<SeriesBucket> {
265 for bucket in &mut buckets {
266 if let Some((&start, previous)) = index_buckets.range(..=bucket.start).next_back()
267 && previous.end > bucket.start
268 {
269 bucket.start = start;
270 bucket.end = bucket.end.max(previous.end);
271 }
272 if let Some((_, last)) = index_buckets.range(bucket.start..bucket.end).next_back() {
275 bucket.end = bucket.end.max(last.end);
276 }
277 }
278 group_series_buckets(buckets)
279}
280
281fn group_series_buckets(spans: Vec<SeriesBucket>) -> Vec<SeriesBucket> {
285 let mut buckets: Vec<SeriesBucket> = Vec::new();
286 for mut span in spans {
287 if let Some(last) = buckets.last_mut()
288 && span.start < last.end
289 {
290 last.end = last.end.max(span.end);
291 last.files.append(&mut span.files);
292 last.has_unknown_sequence |= span.has_unknown_sequence;
293 last.max_file_sequence = last.max_file_sequence.max(span.max_file_sequence);
294 merge_window_sequences(&mut last.window_sequences, span.window_sequences);
295 } else {
296 buckets.push(span);
297 }
298 }
299 buckets
300}
301
302fn merge_window_sequences(
304 target: &mut BTreeMap<i64, WindowSequence>,
305 source: BTreeMap<i64, WindowSequence>,
306) {
307 for (start, summary) in source {
308 target
309 .entry(start)
310 .and_modify(|current| {
311 current.end = current.end.max(summary.end);
312 current.max_sequence = current.max_sequence.max(summary.max_sequence);
313 })
314 .or_insert(summary);
315 }
316}
317
318impl SeriesBucket {
319 fn to_series_entry(&self) -> Option<SeriesIndexEntry> {
322 if self.has_unknown_sequence || self.files.len() < SERIES_INDEX_TRIGGER_FILES {
323 return None;
324 }
325 let mut source_file_ids = self
326 .files
327 .iter()
328 .map(|file| file.file_id().file_id())
329 .collect::<Vec<_>>();
330 source_file_ids.sort_unstable_by(|left, right| left.as_bytes().cmp(right.as_bytes()));
331 let min_file_sequence = self
332 .files
333 .iter()
334 .filter_map(|file| file.meta_ref().sequence.map(|sequence| sequence.get()))
335 .min()?;
336 Some(SeriesIndexEntry {
337 file_size: 0,
338 index_uuid: FileId::random(),
339 bucket_start: self.start,
340 bucket_end: self.end,
341 source_file_ids,
342 min_file_sequence,
343 max_file_sequence: self.max_file_sequence,
344 compaction_window_secs: self.compaction_window_secs,
345 window_sequences: self.window_sequences.clone(),
346 })
347 }
348}
349
350#[cfg(test)]
351mod tests {
352 use std::collections::HashSet;
353 use std::num::NonZeroU64;
354
355 use common_time::timestamp::TimeUnit;
356
357 use super::*;
358 use crate::sst::file::FileMeta;
359 use crate::test_util::new_noop_file_purger;
360
361 fn coverage(intervals: &[(i64, i64, u64)]) -> BTreeMap<i64, WindowSequence> {
362 intervals
363 .iter()
364 .map(|&(start, end, max_sequence)| {
365 (
366 start,
367 WindowSequence {
368 start,
369 end,
370 max_sequence,
371 },
372 )
373 })
374 .collect()
375 }
376
377 fn file(sequence: Option<u64>, level: u8, start: Timestamp, end: Timestamp) -> FileHandle {
378 FileHandle::new(
379 FileMeta {
380 file_id: FileId::random(),
381 sequence: sequence.and_then(NonZeroU64::new),
382 level,
383 time_range: (start, end),
384 ..Default::default()
385 },
386 new_noop_file_purger(),
387 )
388 }
389
390 #[test]
391 fn test_second_resolution_buckets_merge_spans_across_levels() {
392 let width = rounded_bucket_width(Duration::from_secs(11), Duration::from_secs(10)).unwrap();
393 let files = [
394 file(
395 Some(1),
396 0,
397 Timestamp::new_millisecond(-1),
398 Timestamp::new_millisecond(19999),
399 ),
400 file(
401 Some(2),
402 1,
403 Timestamp::new_microsecond(19000000),
404 Timestamp::new_microsecond(39000000),
405 ),
406 file(
407 Some(3),
408 2,
409 Timestamp::new_nanosecond(39000000000),
410 Timestamp::new_nanosecond(40000000000),
411 ),
412 file(
413 Some(4),
414 1,
415 Timestamp::new_second(60),
416 Timestamp::new_second(61),
417 ),
418 ];
419 let buckets = group_files_into_series_buckets(&files, width, 10);
420 let spans = buckets
421 .iter()
422 .map(|b| (b.start, b.end, b.files.len()))
423 .collect::<Vec<_>>();
424 assert_eq!(
425 vec![
426 (Timestamp::new_second(-20), Timestamp::new_second(60), 3),
427 (Timestamp::new_second(60), Timestamp::new_second(80), 1),
428 ],
429 spans
430 );
431 assert_eq!(3, buckets[0].max_file_sequence);
432 assert_eq!(4, buckets[1].max_file_sequence);
433 assert!(buckets[0].to_series_entry().is_none());
434 assert!(buckets[1].to_series_entry().is_none());
435 let mut files = files.to_vec();
436 files.push(file(
437 None,
438 0,
439 Timestamp::new_second(0),
440 Timestamp::new_second(1),
441 ));
442 assert!(
443 group_files_into_series_buckets(&files, width, 10)[0]
444 .to_series_entry()
445 .is_none()
446 );
447 }
448
449 #[test]
450 fn test_seconds_do_not_require_millisecond_conversion() {
451 let start = Timestamp::new_second(i64::MAX / 1000 + 100);
452 assert!(start.convert_to(TimeUnit::Millisecond).is_none());
453 let buckets = group_files_into_series_buckets(&[file(Some(1), 0, start, start)], 1, 1);
454 assert_eq!(start, buckets[0].start);
455 assert_eq!(Timestamp::new_second(start.value() + 1), buckets[0].end);
456
457 let width = rounded_bucket_width(Duration::ZERO, Duration::from_millis(100)).unwrap();
458 let buckets = group_files_into_series_buckets(
459 &[file(
460 Some(1),
461 0,
462 Timestamp::new_nanosecond(-1),
463 Timestamp::new_microsecond(1),
464 )],
465 width,
466 1,
467 );
468 assert_eq!(
469 (Timestamp::new_second(-1), Timestamp::new_second(1)),
470 (buckets[0].start, buckets[0].end)
471 );
472 }
473
474 #[test]
475 fn test_reconcile_bridge_expands_through_index_and_sst_buckets() {
476 let ts = Timestamp::new_second;
477 let mut indexes = BTreeMap::new();
478 let ids = [FileId::random(), FileId::random(), FileId::random()];
479 for (start, end, max_sequence, id) in [
480 (0, 20, 10, ids[0]),
481 (30, 60, 30, ids[1]),
482 (70, 100, 20, ids[2]),
483 ] {
484 IndexBucket {
485 start: ts(start),
486 end: ts(end),
487 index_ids: smallvec![id],
488 compaction_window_secs: 10,
489 window_sequences: coverage(&[(start, end, max_sequence)]),
490 }
491 .insert_into(&mut indexes);
492 }
493 let files = [
494 file(Some(15), 1, ts(10), ts(30)),
495 file(Some(31), 0, ts(50), ts(70)),
496 file(Some(32), 0, ts(90), ts(95)),
497 file(Some(32), 1, ts(90), ts(95)),
498 file(Some(33), 0, ts(90), ts(95)),
499 file(Some(34), 0, ts(100), ts(105)),
500 ];
501 let planned = group_files_into_series_buckets(&files, 10, 10);
502 assert_eq!(4, planned.len());
503 let plan = plan_series_indexes(planned, indexes, None, 0);
504 assert_eq!((2, 1), (plan.computed_buckets, plan.skipped_buckets));
505 let [(bucket, entry)] = plan.builds.as_slice() else {
506 panic!("expected one replacement build");
507 };
508 assert_eq!((ts(0), ts(100)), (bucket.start, bucket.end));
509 let expected = files[..5]
511 .iter()
512 .map(|file| file.file_id().file_id())
513 .collect::<HashSet<_>>();
514 assert_eq!(
515 expected,
516 bucket
517 .files
518 .iter()
519 .map(|file| file.file_id().file_id())
520 .collect()
521 );
522 assert_eq!(expected, entry.source_file_ids.iter().copied().collect());
523 assert_eq!((15, 33), (entry.min_file_sequence, entry.max_file_sequence));
524 assert!(plan.expired_index_ids.is_empty());
525 assert_eq!(1, plan.index_buckets.len());
526 let merged = &plan.index_buckets[&ts(0)];
527 assert_eq!((ts(0), ts(100)), (merged.start, merged.end));
528 assert_eq!(33, bucket.max_file_sequence);
529 assert_eq!(bucket.window_sequences, merged.window_sequences);
530 assert_eq!(ids.as_slice(), plan.superseded_index_ids.as_slice());
531 assert_eq!([entry.index_uuid].as_slice(), merged.index_ids.as_slice());
532 }
533
534 #[test]
535 fn test_plan_reuses_replaced_sources_and_rebuilds_changed_bucket() {
536 let ts = Timestamp::new_second;
537 let make_file = |sequence| file(Some(sequence), 0, ts(1), ts(2));
538 let original = (1..=4).map(make_file).collect::<Vec<_>>();
539 let initial = plan_series_indexes(
540 group_files_into_series_buckets(&original, 10, 10),
541 BTreeMap::new(),
542 None,
543 0,
544 );
545 let first_id = initial.builds[0].1.index_uuid;
546
547 let mut files = (1..=4).map(make_file).collect::<Vec<_>>();
549 let replaced = plan_series_indexes(
550 group_files_into_series_buckets(&files, 10, 10),
551 initial.index_buckets.clone(),
552 None,
553 0,
554 );
555 assert!(replaced.builds.is_empty());
556 assert!(replaced.expired_index_ids.is_empty());
557 assert_eq!(initial.index_buckets, replaced.index_buckets);
558
559 files.push(make_file(5));
560 let ready = plan_series_indexes(
561 group_files_into_series_buckets(&files, 10, 10),
562 replaced.index_buckets,
563 None,
564 0,
565 );
566 let [(bucket, entry)] = ready.builds.as_slice() else {
567 panic!("a changed window must rebuild the whole bucket");
568 };
569 let expected = files[..]
570 .iter()
571 .map(|file| file.file_id().file_id())
572 .collect::<HashSet<_>>();
573 assert_eq!(
574 expected,
575 bucket
576 .files
577 .iter()
578 .map(|file| file.file_id().file_id())
579 .collect()
580 );
581 assert_eq!(expected, entry.source_file_ids.iter().copied().collect());
582 assert_eq!((1, 5), (entry.min_file_sequence, entry.max_file_sequence));
583 assert_eq!(
584 [entry.index_uuid].as_slice(),
585 ready.index_buckets[&ts(0)].index_ids.as_slice()
586 );
587 assert_eq!(vec![first_id], ready.superseded_index_ids);
588 let repeated = plan_series_indexes(
589 group_files_into_series_buckets(&files, 10, 10),
590 ready.index_buckets.clone(),
591 None,
592 0,
593 );
594 assert!(repeated.builds.is_empty());
595 assert_eq!(ready.index_buckets, repeated.index_buckets);
596
597 let expired = plan_series_indexes(
599 Vec::new(),
600 ready.index_buckets,
601 Some(TimeToLive::Duration(Duration::from_secs(10))),
602 21_000,
603 );
604 assert!(expired.builds.is_empty());
605 assert!(expired.index_buckets.is_empty());
606 assert_eq!(
607 [entry.index_uuid].as_slice(),
608 expired.expired_index_ids.as_slice()
609 );
610 }
611
612 #[test]
613 fn test_sst_summaries_round_ranges_outward() {
614 let ts = Timestamp::new_millisecond;
615 let files = [
616 file(Some(10), 0, ts(-1), ts(19_999)),
617 file(Some(30), 1, ts(10_000), ts(20_000)),
618 ];
619 let buckets = group_files_into_series_buckets(&files, 100, 10);
620 assert_eq!(1, buckets.len());
621 assert_eq!(
622 coverage(&[(-10, 20, 10), (10, 30, 30)]),
623 buckets[0].window_sequences
624 );
625 }
626
627 #[rstest::rstest]
628 #[case(32)]
629 #[case(33)]
630 #[case(100_000)]
631 fn test_wide_sst_reuses_unchanged_inputs_and_rebuilds_after_splitting(#[case] windows: i64) {
632 let ts = Timestamp::new_second;
633 let end = windows * 10;
634 let files = (1..=4)
635 .map(|seq| file(Some(seq), 0, ts(0), ts(end - 1)))
636 .collect::<Vec<_>>();
637 let initial = plan_series_indexes(
638 group_files_into_series_buckets(&files, end, 10),
639 BTreeMap::new(),
640 None,
641 0,
642 );
643 assert_eq!(1, initial.builds.len());
644 assert_eq!(
645 coverage(&[(0, end, 4)]),
646 initial.builds[0].1.window_sequences
647 );
648 let repeated = plan_series_indexes(
649 group_files_into_series_buckets(&files, end, 10),
650 initial.index_buckets.clone(),
651 None,
652 0,
653 );
654 assert!(repeated.builds.is_empty());
655 assert_eq!(initial.index_buckets, repeated.index_buckets);
656
657 let split = (1..=4)
659 .rev()
660 .flat_map(|seq| {
661 [
662 file(Some(seq), 1, ts(end / 2), ts(end - 1)),
663 file(Some(seq), 1, ts(0), ts(end / 2 - 1)),
664 ]
665 })
666 .collect::<Vec<_>>();
667 let replaced = plan_series_indexes(
668 group_files_into_series_buckets(&split, end, 10),
669 repeated.index_buckets,
670 None,
671 0,
672 );
673 assert_eq!(1, replaced.builds.len());
674 assert_eq!(1, replaced.superseded_index_ids.len());
675 let repeated = plan_series_indexes(
676 group_files_into_series_buckets(&split, end, 10),
677 replaced.index_buckets.clone(),
678 None,
679 0,
680 );
681 assert!(repeated.builds.is_empty());
682 assert_eq!(replaced.index_buckets, repeated.index_buckets);
683 }
684
685 #[rstest::rstest]
686 #[case::existing_start(0)]
687 #[case::new_start(20)]
688 fn test_new_data_changes_sst_summary(#[case] start: i64) {
689 let ts = Timestamp::new_second;
690 let mut files = vec![
691 file(Some(1), 0, ts(1), ts(29)),
692 file(Some(10), 0, ts(1), ts(29)),
693 file(Some(30), 0, ts(1), ts(9)),
694 file(Some(20), 0, ts(41), ts(49)),
695 ];
696 let initial = plan_series_indexes(
697 group_files_into_series_buckets(&files, 100, 10),
698 BTreeMap::new(),
699 None,
700 0,
701 );
702 assert_eq!(1, initial.builds.len());
703 assert_eq!(
705 coverage(&[(0, 30, 30), (40, 50, 20)]),
706 initial.builds[0].1.window_sequences
707 );
708 files.push(file(Some(31), 0, ts(start + 1), ts(start + 2)));
710 let plan = plan_series_indexes(
711 group_files_into_series_buckets(&files, 100, 10),
712 initial.index_buckets,
713 None,
714 0,
715 );
716 let [(bucket, entry)] = plan.builds.as_slice() else {
717 panic!("new data must rebuild the bucket");
718 };
719 assert_eq!(files.len(), bucket.files.len());
720 assert_eq!(31, entry.max_file_sequence);
721 let expected = if start == 0 {
722 coverage(&[(0, 30, 31), (40, 50, 20)])
723 } else {
724 coverage(&[(0, 30, 30), (20, 30, 31), (40, 50, 20)])
725 };
726 assert_eq!(expected, entry.window_sequences);
727 let repeated = plan_series_indexes(
728 group_files_into_series_buckets(&files, 100, 10),
729 plan.index_buckets,
730 None,
731 0,
732 );
733 assert!(repeated.builds.is_empty());
734 }
735
736 #[test]
737 fn test_deferred_bridge_preserves_established_coverage() {
738 let ts = Timestamp::new_second;
739 let mut indexes = BTreeMap::new();
740 for (start, seq) in [(0, 10), (20, 30)] {
741 IndexBucket {
742 start: ts(start),
743 end: ts(start + 10),
744 index_ids: smallvec![FileId::random()],
745 compaction_window_secs: 10,
746 window_sequences: coverage(&[(start, start + 10, seq)]),
747 }
748 .insert_into(&mut indexes);
749 }
750 let mut files = vec![
751 file(Some(13), 0, ts(1), ts(21)),
752 file(Some(30), 0, ts(21), ts(22)),
753 ];
754 let deferred = plan_series_indexes(
755 group_files_into_series_buckets(&files, 10, 10),
756 indexes.clone(),
757 None,
758 0,
759 );
760 assert!(deferred.builds.is_empty());
761 assert!(deferred.superseded_index_ids.is_empty());
762 assert_eq!(indexes, deferred.index_buckets);
763 files.extend((11..=12).map(|seq| file(Some(seq), 0, ts(1), ts(2))));
764 let plan = plan_series_indexes(
765 group_files_into_series_buckets(&files, 10, 10),
766 deferred.index_buckets,
767 None,
768 0,
769 );
770 assert_eq!(1, plan.builds.len());
771 assert_eq!(4, plan.builds[0].0.files.len());
772 assert_eq!(2, plan.superseded_index_ids.len());
773 assert_eq!(
774 coverage(&[(0, 30, 13), (20, 30, 30)]),
775 plan.builds[0].1.window_sequences
776 );
777
778 files.push(file(None, 0, ts(1), ts(2)));
779 let unknown = plan_series_indexes(
780 group_files_into_series_buckets(&files, 10, 10),
781 indexes.clone(),
782 None,
783 0,
784 );
785 assert!(unknown.builds.is_empty());
786 assert_eq!(indexes, unknown.index_buckets);
787 }
788
789 #[rstest::rstest]
790 #[case::removed_window(100, 10, false)]
791 #[case::added_window(100, 10, true)]
792 #[case::window_width(100, 20, false)]
793 #[case::bucket_width(200, 10, false)]
794 fn test_coverage_shape_changes_rebuild(
795 #[case] bucket_width: i64,
796 #[case] window_width: i64,
797 #[case] add_window: bool,
798 ) {
799 let ts = Timestamp::new_second;
800 let mut files = (27..=30)
801 .map(|seq| file(Some(seq), 0, ts(11), ts(12)))
802 .collect::<Vec<_>>();
803 let mut original = files.clone();
804 original.push(file(Some(10), 0, ts(1), ts(2)));
805 let initial = plan_series_indexes(
806 group_files_into_series_buckets(&original, 100, 10),
807 BTreeMap::new(),
808 None,
809 0,
810 );
811 if add_window {
812 files = original;
813 files.push(file(Some(15), 0, ts(21), ts(22)));
814 } else if bucket_width != 100 || window_width != 10 {
815 files = original;
816 }
817 let plan = plan_series_indexes(
818 group_files_into_series_buckets(&files, bucket_width, window_width),
819 initial.index_buckets,
820 None,
821 0,
822 );
823 assert_eq!(1, plan.builds.len());
824 assert_eq!(1, plan.superseded_index_ids.len());
825 let repeated = plan_series_indexes(
826 group_files_into_series_buckets(&files, bucket_width, window_width),
827 plan.index_buckets,
828 None,
829 0,
830 );
831 assert!(repeated.builds.is_empty());
832 }
833
834 #[test]
835 fn test_index_map_merge_is_order_independent() {
836 let ts = Timestamp::new_second;
837 let make_index = |start, end, width, windows: &[(i64, u64)]| IndexBucket {
838 start: ts(start),
839 end: ts(end),
840 index_ids: smallvec![FileId::random()],
841 compaction_window_secs: width,
842 window_sequences: windows
843 .iter()
844 .map(|&(start, max_sequence)| {
845 (
846 start,
847 WindowSequence {
848 start,
849 end: start + width,
850 max_sequence,
851 },
852 )
853 })
854 .collect(),
855 };
856 let indexes = [
857 make_index(0, 20, 10, &[(0, 10), (10, 20)]),
858 make_index(10, 30, 10, &[(10, 30), (20, 15)]),
859 make_index(20, 40, 10, &[(20, 25), (30, 5)]),
860 ];
861 for order in [[0, 1, 2], [2, 1, 0], [1, 0, 2], [0, 2, 1]] {
862 let mut map = BTreeMap::new();
863 for i in order {
864 indexes[i].clone().insert_into(&mut map);
865 }
866 assert_eq!(1, map.len());
867 assert_eq!(10, map[&ts(0)].compaction_window_secs);
868 assert_eq!(
869 coverage(&[(0, 10, 10), (10, 20, 30), (20, 30, 25), (30, 40, 5)]),
870 map[&ts(0)].window_sequences
871 );
872 make_index(10, 30, 20, &[(0, 30)]).insert_into(&mut map);
873 indexes[0].clone().insert_into(&mut map);
874 assert_eq!(0, map[&ts(0)].compaction_window_secs);
875 assert!(map[&ts(0)].window_sequences.is_empty());
876 }
877 }
878}