Skip to main content

mito2/series_index/
bucket.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//! Event-time bucket planning and source coverage.
16
17use 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/// Index files sharing a non-overlapping, half-open time bucket.
30#[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    /// Zero means that merged indexes use incompatible window widths.
36    pub(crate) compaction_window_secs: i64,
37    /// Indexed SST summaries keyed by aligned start.
38    /// See [`SeriesIndexEntry::window_sequences`] for the layout and sequence assumption.
39    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                // Do not compare maps with different window boundaries, including
64                // a previous incompatible merge, against current SST coverage.
65                self.compaction_window_secs = 0;
66                self.window_sequences.clear();
67            }
68        }
69        self.index_ids.append(&mut other.index_ids);
70    }
71
72    /// Inserts a bucket, consuming overlapping entries and expanding their interval.
73    /// Each map key equals the stored bucket start; adjacent intervals remain separate.
74    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        // Expanding the end may expose more overlaps, so look up the next entry again.
81        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/// SSTs grouped into a half-open time interval for an aggregate series-index build.
91///
92/// Reconciliation may expand the interval through existing index coverage. A changed
93/// bucket is rebuilt from all its SSTs. Unknown source sequences prevent builds,
94/// since their coverage cannot be determined safely.
95#[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    /// Maximum known source sequence, or zero when all sequences are unknown.
102    pub(crate) max_file_sequence: u64,
103    pub(crate) compaction_window_secs: i64,
104    /// SST summaries keyed by compaction-window-aligned start.
105    /// Each SST contributes one entry; equal starts merge by maximum end and sequence.
106    /// Ranges may overlap. See [`WindowSequence`] for the sequence assumption.
107    /// Missing SST sequences use zero as a placeholder and set `has_unknown_sequence`,
108    /// preventing index builds and reuse.
109    pub(crate) window_sequences: BTreeMap<i64, WindowSequence>,
110}
111
112/// Next bucket coverage and the work required to publish it, computed without I/O.
113pub(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    /// Retire only after replacement builds and catalog publication succeed.
118    pub(crate) superseded_index_ids: Vec<FileId>,
119    pub(crate) computed_buckets: usize,
120    pub(crate) skipped_buckets: usize,
121}
122
123/// Plans whole-bucket replacements when source SST summaries change. The returned bucket
124/// map is publishable only after every planned build and the catalog writes succeed.
125pub(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    // Reconciliation changes geometry, not established coverage. In particular,
132    // a deferred bridge must not change the indexed snapshot.
133    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
202/// Groups SSTs across levels into sorted, disjoint time buckets without planning builds.
203///
204/// Both widths must be positive, and `width_secs` must be a multiple of the compaction
205/// window width. Inclusive SST ranges are rounded outward to aligned,
206/// half-open intervals in seconds. Overlapping intervals merge; adjacent ones stay
207/// separate. Each merged bucket tracks its maximum sequence and any unknown sequence.
208pub(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
255/// Expands SST buckets through existing index coverage before grouping build inputs.
256///
257/// `buckets` must be sorted by start and disjoint. `index_buckets` must contain
258/// disjoint intervals keyed by their starts. Expansion preserves start ordering,
259/// but may introduce overlaps, which are merged in the returned disjoint buckets.
260/// Established index coverage is borrowed so deferred builds cannot modify it.
261fn 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        // Disjoint index intervals make the last overlapping interval's end the
273        // furthest boundary. Expanding to it cannot expose another index interval.
274        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
281/// Merges overlapping spans sorted by nondecreasing start; equal starts are allowed.
282/// Returns sorted, disjoint buckets. Adjacent half-open intervals remain separate.
283/// All spans must use the same compaction-window width for their SST summaries.
284fn 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
302/// Merges summaries sharing a start, relying on [`WindowSequence`]'s sequence assumption.
303fn 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    /// Creates entry metadata with a fresh UUID and sorted source file IDs.
320    /// Returns `None` for unknown sequences or too few files.
321    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        // Include sequence 15 and both files at 32, but exclude the adjacent bucket.
510        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        // New file IDs with already indexed sequences do not invalidate the index.
548        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        // TTL retires the replacement; the previous index is already superseded.
598        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        // Splitting changes the summaries and conservatively triggers one rebuild.
658        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        // The maximum end survives even when the highest sequence is in a shorter SST.
704        assert_eq!(
705            coverage(&[(0, 30, 30), (40, 50, 20)]),
706            initial.builds[0].1.window_sequences
707        );
708        // New data has a sequence greater than those in the indexed snapshot.
709        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}