Skip to main content

mito2/series_index/
maintenance.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//! Worker-owned series-index reconciliation and publication.
16
17use std::collections::HashSet;
18use std::sync::Arc;
19use std::time::{Duration, Instant};
20
21use common_telemetry::{debug, info};
22use object_store::ObjectStore;
23use store_api::storage::{FileId, RegionId};
24
25use crate::error::Result;
26use crate::metrics::{SERIES_INDEX_RECONCILE_ELAPSED, SERIES_INDEX_RECONCILE_TOTAL};
27use crate::read::series_candidate::is_sparse_metric_metadata;
28use crate::region::version::VersionRef;
29use crate::region::{MitoRegionRef, RegionLeaderState, RegionRoleState};
30use crate::series_index::bucket::{
31    group_files_into_series_buckets, plan_series_indexes, rounded_bucket_width,
32};
33use crate::series_index::builder::{build_range_index, build_series_index};
34use crate::series_index::catalog::{
35    RangeIndexCatalog, SeriesIndexCatalog, delete_catalogs, range_catalog_path,
36    series_catalog_path, store_catalog,
37};
38use crate::series_index::purger::IndexFilePurger;
39use crate::series_index::version::{SeriesIndexFileHandle, SeriesIndexVersion};
40
41/// Retires newly completed series files unless their snapshot is published.
42#[derive(Default)]
43struct UnpublishedSeriesFiles(Vec<SeriesIndexFileHandle>);
44
45impl UnpublishedSeriesFiles {
46    fn disarm(&mut self) {
47        self.0.clear();
48    }
49}
50
51impl Drop for UnpublishedSeriesFiles {
52    fn drop(&mut self) {
53        for handle in &self.0 {
54            handle.mark_deleted();
55        }
56    }
57}
58
59#[derive(Debug, Default)]
60pub(crate) struct ReconcileStats {
61    pub(crate) source_files: usize,
62    pub(crate) built_range: usize,
63    pub(crate) built_series: usize,
64    pub(crate) removed_range: usize,
65    pub(crate) removed_series: usize,
66    pub(crate) computed_buckets: usize,
67    pub(crate) skipped_buckets: usize,
68}
69
70impl ReconcileStats {
71    fn changed(&self) -> bool {
72        self.built_range + self.built_series + self.removed_range + self.removed_series > 0
73    }
74}
75
76/// Reconciles indexes for one region snapshot, persists catalogs, then atomically publishes it.
77#[allow(clippy::too_many_arguments)]
78pub(crate) async fn reconcile_series_indexes(
79    worker_id: u32,
80    store: ObjectStore,
81    region: MitoRegionRef,
82    requested_bucket_width: Duration,
83    now_ms: i64,
84    purger: IndexFilePurger,
85    enable_range_index: bool,
86    allow_builds: bool,
87) -> Result<ReconcileStats> {
88    let total_start = Instant::now();
89    // Use this snapshot throughout reconciliation, even if the region version advances.
90    let version = region.version_control.current().version;
91    if !is_sparse_metric_metadata(&version.metadata) {
92        SERIES_INDEX_RECONCILE_TOTAL
93            .with_label_values(&["noop"])
94            .inc();
95        return Ok(ReconcileStats::default());
96    }
97    let build_start = Instant::now();
98    let mut unpublished = UnpublishedSeriesFiles::default();
99    let (next, mut stats) = build_index_version(
100        worker_id,
101        &store,
102        &region,
103        &version,
104        requested_bucket_width,
105        now_ms,
106        &purger,
107        &mut unpublished,
108        enable_range_index,
109        allow_builds,
110    )
111    .await?;
112    SERIES_INDEX_RECONCILE_ELAPSED
113        .with_label_values(&["build"])
114        .observe(build_start.elapsed().as_secs_f64());
115    // Persist changed catalogs before making the new snapshot visible to readers.
116    let publish_result: Result<()> = async {
117        if let Some(next) = next {
118            persist_index_catalogs(&store, region.region_id, &next, &stats).await?;
119            publish_index_version(&region, Arc::new(next));
120            unpublished.disarm();
121        }
122        Ok(())
123    }
124    .await;
125    // Drop may have cleaned up while catalogs were being written, even if a write failed.
126    // If dropping starts after this check, the normal drop path cleans up our publication.
127    // This assumes the region ID is not reopened or replaced during cleanup.
128    if region.state() == RegionRoleState::Leader(RegionLeaderState::Dropping) {
129        delete_catalogs(&store, region.region_id).await;
130        region.series_index_version_control.mark_dropped();
131        stats = ReconcileStats::default();
132    }
133    publish_result?;
134    let result = if stats.changed() { "changed" } else { "noop" };
135    SERIES_INDEX_RECONCILE_TOTAL
136        .with_label_values(&[result])
137        .inc();
138    SERIES_INDEX_RECONCILE_ELAPSED
139        .with_label_values(&["total"])
140        .observe(total_start.elapsed().as_secs_f64());
141    if stats.changed() {
142        info!(
143            "Reconciled series-index snapshot, worker: {worker_id}, region: {}, elapsed: {:?}, stats: {:?}",
144            region.region_id,
145            total_start.elapsed(),
146            stats
147        );
148    } else {
149        debug!(
150            "Series-index reconciliation made no changes, worker: {worker_id}, region: {}",
151            region.region_id
152        );
153    }
154    Ok(stats)
155}
156
157/// Builds a changed snapshot, retaining reusable indexes and removing obsolete coverage.
158/// Returns `None` when no range or series indexes were added or removed.
159#[allow(clippy::too_many_arguments)]
160async fn build_index_version(
161    worker_id: u32,
162    store: &ObjectStore,
163    region: &MitoRegionRef,
164    version: &VersionRef,
165    requested_bucket_width: Duration,
166    now_ms: i64,
167    purger: &IndexFilePurger,
168    unpublished: &mut UnpublishedSeriesFiles,
169    enable_range_index: bool,
170    allow_builds: bool,
171) -> Result<(Option<SeriesIndexVersion>, ReconcileStats)> {
172    let mut stats = ReconcileStats::default();
173    let files = version
174        .ssts
175        .levels()
176        .iter()
177        .flat_map(|level| level.files())
178        .cloned()
179        .collect::<Vec<_>>();
180    stats.source_files = files.len();
181    let visible = files
182        .iter()
183        .map(|file| file.file_id().file_id())
184        .collect::<HashSet<_>>();
185    let current = region.series_index_version();
186    // Empty build inputs still let the planner expire established bucket coverage.
187    let buckets = match (allow_builds, version.compaction_time_window) {
188        (true, Some(window)) => rounded_bucket_width(requested_bucket_width, window)
189            // Successful rounding guarantees the window fits in i64; subsecond windows
190            // use the same one-second minimum as rounded_bucket_width.
191            .map(|width| {
192                group_files_into_series_buckets(&files, width, (window.as_secs() as i64).max(1))
193            })
194            .unwrap_or_default(),
195        (true, None) => {
196            debug!(
197                "Deferring series indexes without compaction window, worker: {worker_id}, region: {}",
198                region.region_id
199            );
200            Vec::new()
201        }
202        (false, _) => Vec::new(),
203    };
204    let plan = plan_series_indexes(
205        buckets,
206        current.index_buckets.clone(),
207        version.options.ttl,
208        now_ms,
209    );
210    stats.computed_buckets = plan.computed_buckets;
211    stats.skipped_buckets = plan.skipped_buckets;
212    let mut next = prune_index_version(&current, &visible, &plan.expired_index_ids, &mut stats);
213    if !allow_builds {
214        if let Some(next) = &mut next {
215            next.index_buckets = plan.index_buckets;
216        }
217        return Ok((next, stats));
218    }
219    if plan.builds.is_empty()
220        && next.is_none()
221        && (!enable_range_index || current.range_indexes.len() == visible.len())
222    {
223        return Ok((None, stats));
224    }
225    let SeriesIndexVersion {
226        mut range_indexes,
227        mut series_indexes,
228        ..
229    } = next.unwrap_or_else(|| SeriesIndexVersion {
230        range_indexes: current.range_indexes.clone(),
231        series_indexes: current.series_indexes.clone(),
232        index_buckets: Default::default(),
233    });
234    let index_buckets = plan.index_buckets;
235    for id in &plan.superseded_index_ids {
236        series_indexes.remove(id);
237    }
238    for (bucket, expected) in plan.builds {
239        // Complete companion indexes independently so a failed series build preserves them.
240        for file in &bucket.files {
241            if enable_range_index
242                && !range_indexes.contains_key(&file.file_id().file_id())
243                && let Some(entry) = build_range_index(store, region, version, file.clone()).await?
244            {
245                stats.built_range += 1;
246                range_indexes.insert(entry.file_id, entry);
247            }
248        }
249        let series_handle =
250            build_series_index(store, region, version, &bucket, &expected, purger).await?;
251        unpublished.0.push(series_handle.clone());
252        stats.built_series += 1;
253        series_indexes.insert(expected.index_uuid, series_handle);
254    }
255    // Cover SSTs outside planned aggregate builds, including skipped buckets.
256    if enable_range_index {
257        for file in files {
258            let file_id = file.file_id().file_id();
259            if range_indexes.contains_key(&file_id) {
260                continue;
261            }
262            if let Some(entry) = build_range_index(store, region, version, file).await? {
263                stats.built_range += 1;
264                range_indexes.insert(entry.file_id, entry);
265            }
266        }
267    }
268    stats.removed_series = current
269        .series_indexes
270        .keys()
271        .filter(|id| !series_indexes.contains_key(id))
272        .count();
273    // Bucket reconciliation is speculative until indexes change. Recompute its coverage
274    // next time rather than replacing the published snapshot on a no-op pass.
275    if !stats.changed() {
276        return Ok((None, stats));
277    }
278    let next = SeriesIndexVersion {
279        range_indexes,
280        series_indexes,
281        index_buckets,
282    };
283    Ok((Some(next), stats))
284}
285
286/// Prunes obsolete metadata without building indexes or retiring published handles.
287/// Returns `None` without cloning the current version when no entries need removal.
288/// Physical range deletion belongs to the SST purger; series handles are retired only
289/// after the cleaned catalogs and snapshot have been published successfully.
290fn prune_index_version(
291    current: &SeriesIndexVersion,
292    visible: &HashSet<FileId>,
293    expired_index_ids: &[FileId],
294    stats: &mut ReconcileStats,
295) -> Option<SeriesIndexVersion> {
296    if current.range_indexes.keys().all(|id| visible.contains(id))
297        && expired_index_ids
298            .iter()
299            .all(|id| !current.series_indexes.contains_key(id))
300    {
301        return None;
302    }
303    let mut next = SeriesIndexVersion {
304        range_indexes: current.range_indexes.clone(),
305        series_indexes: current.series_indexes.clone(),
306        index_buckets: current.index_buckets.clone(),
307    };
308    next.range_indexes
309        .retain(|file_id, _| visible.contains(file_id));
310    for id in expired_index_ids {
311        next.series_indexes.remove(id);
312    }
313    stats.removed_range = current.range_indexes.len() - next.range_indexes.len();
314    stats.removed_series = current.series_indexes.len() - next.series_indexes.len();
315    Some(next)
316}
317
318/// Writes changed catalogs in a stable order; the two writes are not atomic together.
319async fn persist_index_catalogs(
320    store: &ObjectStore,
321    region_id: RegionId,
322    next: &SeriesIndexVersion,
323    stats: &ReconcileStats,
324) -> Result<()> {
325    if stats.built_range + stats.removed_range > 0 {
326        let mut range_entries = next.range_indexes.values().copied().collect::<Vec<_>>();
327        range_entries
328            .sort_unstable_by(|left, right| left.file_id.as_bytes().cmp(right.file_id.as_bytes()));
329        store_catalog(
330            store,
331            &range_catalog_path(region_id),
332            &RangeIndexCatalog {
333                indexes: range_entries,
334            },
335        )
336        .await?;
337    }
338    if stats.built_series + stats.removed_series > 0 {
339        let mut series_entries = next
340            .series_indexes
341            .values()
342            .map(|handle| handle.entry().clone())
343            .collect::<Vec<_>>();
344        series_entries.sort_unstable_by_key(|entry| {
345            (
346                entry.bucket_start,
347                entry.bucket_end,
348                entry.min_file_sequence,
349                entry.max_file_sequence,
350            )
351        });
352        store_catalog(
353            store,
354            &series_catalog_path(region_id),
355            &SeriesIndexCatalog {
356                indexes: series_entries,
357            },
358        )
359        .await?;
360    }
361    Ok(())
362}
363
364/// Publishes the snapshot and retires series files absent from the new version.
365fn publish_index_version(region: &MitoRegionRef, next: Arc<SeriesIndexVersion>) {
366    let previous = region.series_index_version_control.publish(next.clone());
367    for (id, handle) in &previous.series_indexes {
368        if !next.series_indexes.contains_key(id) {
369            // Purge only after readers release their retained handles.
370            handle.mark_deleted();
371        }
372    }
373}
374
375#[cfg(test)]
376mod tests {
377    use std::collections::HashMap;
378
379    use common_time::Timestamp;
380    use object_store::services::Memory;
381
382    use super::*;
383    use crate::series_index::catalog::{RangeIndexEntry, SeriesIndexEntry};
384    use crate::series_index::purger::series_index_channel;
385
386    #[rstest::rstest]
387    fn test_prune_index_version(
388        #[values(false, true)] remove_range: bool,
389        #[values(false, true)] remove_series: bool,
390    ) {
391        let range_id = FileId::random();
392        let series_id = FileId::random();
393        let store = ObjectStore::new(Memory::default()).unwrap();
394        let (purger, mut receiver) = series_index_channel(store);
395        let current = SeriesIndexVersion::new(
396            HashMap::from([(
397                range_id,
398                RangeIndexEntry {
399                    file_id: range_id,
400                    file_size: 10,
401                },
402            )]),
403            HashMap::from([(
404                series_id,
405                SeriesIndexFileHandle::new(
406                    RegionId::new(1, 1),
407                    SeriesIndexEntry {
408                        index_uuid: series_id,
409                        file_size: 20,
410                        bucket_start: Timestamp::new_second(0),
411                        bucket_end: Timestamp::new_second(100),
412                        source_file_ids: vec![range_id],
413                        min_file_sequence: 1,
414                        max_file_sequence: 1,
415                        compaction_window_secs: 100,
416                        window_sequences: Default::default(),
417                    },
418                    purger,
419                ),
420            )]),
421        );
422        let visible = if remove_range {
423            HashSet::new()
424        } else {
425            HashSet::from([range_id])
426        };
427        // Unknown and repeated IDs must not inflate removal counts.
428        let mut expired = vec![FileId::random()];
429        if remove_series {
430            expired.extend([series_id, series_id]);
431        }
432        let mut stats = ReconcileStats::default();
433        let next = prune_index_version(&current, &visible, &expired, &mut stats);
434        assert_eq!(remove_range || remove_series, next.is_some());
435        assert_eq!(usize::from(remove_range), stats.removed_range);
436        assert_eq!(usize::from(remove_series), stats.removed_series);
437        if let Some(next) = next {
438            assert_eq!(!remove_range, next.range_indexes.contains_key(&range_id));
439            assert_eq!(!remove_series, next.series_indexes.contains_key(&series_id));
440        }
441        assert_eq!(1, current.range_indexes.len());
442        assert_eq!(1, current.series_indexes.len());
443        drop(current);
444        // Pruning alone must never retire published files.
445        assert!(receiver.try_recv().is_err());
446    }
447}