diff --git a/config/config.md b/config/config.md
index 2d85ddd717..5347406112 100644
--- a/config/config.md
+++ b/config/config.md
@@ -165,6 +165,9 @@
| `region_engine.mito.worker_request_batch_size` | Integer | `64` | Max batch size for a worker to handle requests. |
| `region_engine.mito.manifest_checkpoint_distance` | Integer | `10` | Number of meta action updated to trigger a new checkpoint for the manifest. |
| `region_engine.mito.compress_manifest` | Bool | `false` | Whether to compress manifest and checkpoint file by gzip (default false). |
+| `region_engine.mito.experimental_series_index_root` | String | `""` | Under development; do not enable. Root directory for local series indexes.
Empty disables the feature. Relative paths resolve under `data_home`. |
+| `region_engine.mito.experimental_series_index_maintenance_interval` | String | `5m` | Interval between series-index maintenance runs. Zero uses the default of 5 min. |
+| `region_engine.mito.experimental_series_index_bucket_width` | String | `5days` | Requested minimum series-index bucket width (default: 5 days), rounded up to
an exact multiple of each region's compaction time window. |
| `region_engine.mito.max_background_flushes` | Integer | Auto | Max number of running background flush jobs (default: 1/2 of cpu cores). |
| `region_engine.mito.max_background_compactions` | Integer | Auto | Max number of running background compaction jobs (default: 1/4 of cpu cores). |
| `region_engine.mito.max_background_purges` | Integer | Auto | Max number of running background purge jobs (default: number of cpu cores). |
@@ -612,6 +615,9 @@
| `region_engine.mito.experimental_manifest_keep_removed_file_count` | Integer | `256` | Number of removed files to keep in manifest's `removed_files` field before also
remove them from `removed_files`. Mostly for debugging purpose.
If set to 0, it will only use `keep_removed_file_ttl` to decide when to remove files
from `removed_files` field. |
| `region_engine.mito.experimental_manifest_keep_removed_file_ttl` | String | `1h` | How long to keep removed files in the `removed_files` field of manifest
after they are removed from manifest.
files will only be removed from `removed_files` field
if both `keep_removed_file_count` and `keep_removed_file_ttl` is reached. |
| `region_engine.mito.compress_manifest` | Bool | `false` | Whether to compress manifest and checkpoint file by gzip (default false). |
+| `region_engine.mito.experimental_series_index_root` | String | `""` | Under development; do not enable. Root directory for local series indexes.
Empty disables the feature. Relative paths resolve under `data_home`. |
+| `region_engine.mito.experimental_series_index_maintenance_interval` | String | `5m` | Interval between series-index maintenance runs. Zero uses the default of 5 min. |
+| `region_engine.mito.experimental_series_index_bucket_width` | String | `5days` | Requested minimum series-index bucket width (default: 5 days), rounded up to
an exact multiple of each region's compaction time window. |
| `region_engine.mito.max_background_flushes` | Integer | Auto | Max number of running background flush jobs (default: 1/2 of cpu cores). |
| `region_engine.mito.max_background_compactions` | Integer | Auto | Max number of running background compaction jobs (default: 1/4 of cpu cores). |
| `region_engine.mito.max_background_purges` | Integer | Auto | Max number of running background purge jobs (default: number of cpu cores). |
diff --git a/config/datanode.example.toml b/config/datanode.example.toml
index d378d55cc2..f1be288bdd 100644
--- a/config/datanode.example.toml
+++ b/config/datanode.example.toml
@@ -474,6 +474,17 @@ experimental_manifest_keep_removed_file_ttl = "1h"
## Whether to compress manifest and checkpoint file by gzip (default false).
compress_manifest = false
+## Under development; do not enable. Root directory for local series indexes.
+## Empty disables the feature. Relative paths resolve under `data_home`.
+#+ experimental_series_index_root = ""
+
+## Interval between series-index maintenance runs. Zero uses the default of 5 min.
+#+ experimental_series_index_maintenance_interval = "5m"
+
+## Requested minimum series-index bucket width (default: 5 days), rounded up to
+## an exact multiple of each region's compaction time window.
+#+ experimental_series_index_bucket_width = "5days"
+
## Max number of running background flush jobs (default: 1/2 of cpu cores).
## @toml2docs:none-default="Auto"
#+ max_background_flushes = 4
diff --git a/config/standalone.example.toml b/config/standalone.example.toml
index 0bef079324..8ec36a6f59 100644
--- a/config/standalone.example.toml
+++ b/config/standalone.example.toml
@@ -635,6 +635,17 @@ manifest_checkpoint_distance = 10
## Whether to compress manifest and checkpoint file by gzip (default false).
compress_manifest = false
+## Under development; do not enable. Root directory for local series indexes.
+## Empty disables the feature. Relative paths resolve under `data_home`.
+#+ experimental_series_index_root = ""
+
+## Interval between series-index maintenance runs. Zero uses the default of 5 min.
+#+ experimental_series_index_maintenance_interval = "5m"
+
+## Requested minimum series-index bucket width (default: 5 days), rounded up to
+## an exact multiple of each region's compaction time window.
+#+ experimental_series_index_bucket_width = "5days"
+
## Max number of running background flush jobs (default: 1/2 of cpu cores).
## @toml2docs:none-default="Auto"
#+ max_background_flushes = 4
diff --git a/src/mito2/AGENTS.md b/src/mito2/AGENTS.md
index faa581331b..583ec31bd4 100644
--- a/src/mito2/AGENTS.md
+++ b/src/mito2/AGENTS.md
@@ -26,7 +26,7 @@ snapshot isolation). It implements the `RegionEngine` trait from `store-api`.
| `compaction` | `src/mito2/src/compaction/` | Compaction scheduler (`scheduler.rs` + `scheduler/`), TWCS picker, strict-window manual picker, compactor, memory control |
| `access_layer` | `src/mito2/src/access_layer.rs` | SST read/write over the object store |
| `sst` | `src/mito2/src/sst/` | Parquet format, file metadata, index layout |
-| `series_index` | `src/mito2/src/series_index/` | Series index writer/searcher, bucket planning and SST builders, catalogs, immutable snapshots, file lifecycle, and background maintenance |
+| `series_index` | `src/mito2/src/series_index/` | Series index writer/searcher, bucket planning, catalogs, worker-owned maintenance and purge tasks, and index lifecycle |
| `read` | `src/mito2/src/read/` | `ScanRegion`, merge, dedup, projection, streaming |
| `manifest` | `src/mito2/src/manifest/` | `RegionManifestManager`, manifest actions/edits |
| `cache` | `src/mito2/src/cache.rs` | Write/file/page caches |
diff --git a/src/mito2/src/config.rs b/src/mito2/src/config.rs
index 6ec790c52b..18d1843371 100644
--- a/src/mito2/src/config.rs
+++ b/src/mito2/src/config.rs
@@ -89,14 +89,16 @@ pub struct MitoConfig {
// Background job configs:
/// Max number of running background index build jobs (default: 1/8 of cpu cores).
pub max_background_index_builds: usize,
- // TODO: Document both series-index settings in the example configs and regenerate configuration
- // docs before exposing the feature.
/// Under development; do not enable. Root directory for loading series indexes, currently stored
/// on the local filesystem. Empty disables the feature. Relative paths resolve under `data_home`.
pub experimental_series_index_root: String,
/// Interval between series-index maintenance runs (default 5 min). Zero uses the default.
#[serde(with = "humantime_serde")]
pub experimental_series_index_maintenance_interval: Duration,
+ /// Under development; do not enable. Requested minimum bucket width for series indexes.
+ /// It is rounded up to an exact multiple of each region's compaction time window.
+ #[serde(with = "humantime_serde")]
+ pub experimental_series_index_bucket_width: Duration,
/// Max number of running background flush jobs (default: 1/2 of cpu cores).
pub max_background_flushes: usize,
/// Max number of running background compaction jobs (default: 1/4 of cpu cores).
@@ -217,6 +219,7 @@ impl Default for MitoConfig {
experimental_series_index_root: String::new(),
experimental_series_index_maintenance_interval:
DEFAULT_SERIES_INDEX_MAINTENANCE_INTERVAL,
+ experimental_series_index_bucket_width: Duration::from_secs(5 * 24 * 60 * 60),
max_background_flushes: divide_num_cpus(2),
max_background_compactions: divide_num_cpus(4),
max_background_purges: get_total_cpu_cores(),
@@ -439,9 +442,14 @@ mod tests {
.experimental_series_index_root
.is_empty()
);
+ assert_eq!(
+ MitoConfig::default().experimental_series_index_bucket_width,
+ Duration::from_secs(5 * 24 * 60 * 60)
+ );
let mut config: MitoConfig = toml::from_str(
"experimental_series_index_root = 'indexes'
- experimental_series_index_maintenance_interval = '30s'",
+ experimental_series_index_maintenance_interval = '30s'
+ experimental_series_index_bucket_width = '2days'",
)
.unwrap();
config.sanitize("/data").unwrap();
@@ -450,6 +458,10 @@ mod tests {
config.experimental_series_index_maintenance_interval,
Duration::from_secs(30)
);
+ assert_eq!(
+ config.experimental_series_index_bucket_width,
+ Duration::from_secs(2 * 24 * 60 * 60)
+ );
let restored: MitoConfig = toml::from_str(&toml::to_string(&config).unwrap()).unwrap();
assert_eq!(config, restored);
diff --git a/src/mito2/src/engine/compaction_test.rs b/src/mito2/src/engine/compaction_test.rs
index c67afbf8c6..094e9188c3 100644
--- a/src/mito2/src/engine/compaction_test.rs
+++ b/src/mito2/src/engine/compaction_test.rs
@@ -38,11 +38,11 @@ use tokio::sync::{Notify, Semaphore};
use crate::config::MitoConfig;
use crate::engine::MitoEngine;
-use crate::engine::flush_test::MockTimeProvider;
use crate::engine::listener::{CompactionListener, EventListener};
use crate::test_util::{
CreateRequestBuilder, TestEnv, build_rows_for_key, column_metadata_to_column_schema, put_rows,
};
+use crate::time_provider::mock::MockTimeProvider;
pub(crate) async fn put_and_flush(
engine: &MitoEngine,
diff --git a/src/mito2/src/engine/edit_region_test.rs b/src/mito2/src/engine/edit_region_test.rs
index de9611ee2d..e25c7793ea 100644
--- a/src/mito2/src/engine/edit_region_test.rs
+++ b/src/mito2/src/engine/edit_region_test.rs
@@ -35,12 +35,12 @@ use tokio::sync::{Barrier, mpsc, oneshot};
use crate::config::MitoConfig;
use crate::engine::MitoEngine;
-use crate::engine::flush_test::MockTimeProvider;
use crate::engine::listener::EventListener;
use crate::manifest::action::RegionEdit;
use crate::region::{MitoRegionRef, RegionLeaderState, RegionRoleState};
use crate::sst::file::FileMeta;
use crate::test_util::{CreateRequestBuilder, TestEnv, build_rows, rows_schema};
+use crate::time_provider::mock::MockTimeProvider;
#[tokio::test]
async fn test_edit_region_schedule_compaction() {
diff --git a/src/mito2/src/engine/flush_test.rs b/src/mito2/src/engine/flush_test.rs
index d7155e1052..47d5486e6d 100644
--- a/src/mito2/src/engine/flush_test.rs
+++ b/src/mito2/src/engine/flush_test.rs
@@ -17,7 +17,7 @@
use std::assert_matches;
use std::collections::HashMap;
use std::sync::Arc;
-use std::sync::atomic::{AtomicBool, AtomicI64, AtomicUsize, Ordering};
+use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use std::time::Duration;
use api::v1::Rows;
@@ -54,7 +54,7 @@ use crate::test_util::{
prepare_test_for_kafka_log_store, put_rows, raft_engine_log_store_factory, reopen_region,
rows_schema, single_kafka_log_store_factory,
};
-use crate::time_provider::TimeProvider;
+use crate::time_provider::mock::MockTimeProvider;
use crate::worker::MAX_INITIAL_CHECK_DELAY_SECS;
async fn set_write_buffer_size_to_current_usage(engine: &MitoEngine, region_id: RegionId) {
@@ -1334,43 +1334,6 @@ fn kafka_wal_options(topic: &Option) -> HashMap {
.unwrap_or_default()
}
-#[derive(Debug)]
-pub(crate) struct MockTimeProvider {
- now: AtomicI64,
- elapsed: AtomicI64,
-}
-
-impl TimeProvider for MockTimeProvider {
- fn current_time_millis(&self) -> i64 {
- self.now.load(Ordering::Relaxed)
- }
-
- fn elapsed_since(&self, _current_millis: i64) -> i64 {
- self.elapsed.load(Ordering::Relaxed)
- }
-
- fn wait_duration(&self, _duration: Duration) -> Duration {
- Duration::from_millis(20)
- }
-}
-
-impl MockTimeProvider {
- pub(crate) fn new(now: i64) -> Self {
- Self {
- now: AtomicI64::new(now),
- elapsed: AtomicI64::new(0),
- }
- }
-
- pub(crate) fn set_now(&self, now: i64) {
- self.now.store(now, Ordering::Relaxed);
- }
-
- fn set_elapsed(&self, elapsed: i64) {
- self.elapsed.store(elapsed, Ordering::Relaxed);
- }
-}
-
#[tokio::test]
async fn test_auto_flush_engine() {
test_auto_flush_engine_with_format(false).await;
diff --git a/src/mito2/src/metrics.rs b/src/mito2/src/metrics.rs
index 2d2decbde7..49e91d93ef 100644
--- a/src/mito2/src/metrics.rs
+++ b/src/mito2/src/metrics.rs
@@ -319,6 +319,19 @@ lazy_static! {
// Index metrics.
lazy_static! {
// Index metrics.
+ /// Outcomes of series-index reconciliation passes.
+ pub static ref SERIES_INDEX_RECONCILE_TOTAL: IntCounterVec = register_int_counter_vec!(
+ "greptime_mito_series_index_reconcile_total",
+ "series-index reconciliation passes",
+ &["result"],
+ ).unwrap();
+ /// Elapsed time of series-index reconciliation phases.
+ pub static ref SERIES_INDEX_RECONCILE_ELAPSED: HistogramVec = register_histogram_vec!(
+ "greptime_mito_series_index_reconcile_elapsed",
+ "series-index reconciliation elapsed time",
+ &["phase"],
+ exponential_buckets(0.01, 10.0, 7).unwrap(),
+ ).unwrap();
/// Series and range index file operations.
pub static ref SERIES_INDEX_FILE_OPERATION_TOTAL: IntCounterVec = register_int_counter_vec!(
"greptime_mito_series_index_file_operation_total",
diff --git a/src/mito2/src/region.rs b/src/mito2/src/region.rs
index 10195f5680..bdf18550cd 100644
--- a/src/mito2/src/region.rs
+++ b/src/mito2/src/region.rs
@@ -242,7 +242,6 @@ impl StagingPartitionInfo {
impl MitoRegion {
/// Returns the current immutable series-index snapshot.
- #[allow(dead_code)] // Used by the upcoming query integration.
pub(crate) fn series_index_version(&self) -> Arc {
self.series_index_version_control.current()
}
diff --git a/src/mito2/src/series_index.rs b/src/mito2/src/series_index.rs
index 8173050f2a..6982e4398f 100644
--- a/src/mito2/src/series_index.rs
+++ b/src/mito2/src/series_index.rs
@@ -12,42 +12,37 @@
// See the License for the specific language governing permissions and
// limitations under the License.
-//! Series index writer and searcher.
+//! Series index construction, search, and maintenance.
+//!
+//! Under development. Index files are currently stored on the local filesystem.
-// These components are consumed by the upcoming query and maintenance integration.
-#[allow(dead_code)]
-mod catalog;
-// Consumed by the follow-up background-maintenance integration.
-#[allow(dead_code)]
mod bucket;
-// Consumed by the follow-up background-maintenance integration.
-#[allow(dead_code)]
mod builder;
-#[allow(dead_code)]
+mod catalog;
+mod maintenance;
mod purger;
mod searcher;
mod task;
#[cfg(test)]
mod tests;
-#[allow(dead_code)]
mod version;
mod writer;
use futures::stream::BoxStream;
-pub use searcher::SeriesIndexSearcher;
use store_api::metric_engine_consts::{
DATA_SCHEMA_TABLE_ID_COLUMN_NAME as TABLE_ID_COLUMN,
DATA_SCHEMA_TSID_COLUMN_NAME as TSID_COLUMN,
};
-pub use writer::{
- SeriesIndexWriter, SeriesIndexWriterMetrics, SeriesIndexWriterOptions, series_index_schema,
-};
use crate::error::Result;
pub(crate) use crate::series_index::catalog::{delete_catalogs, load_version_control};
pub(crate) use crate::series_index::purger::{IndexFilePurger, series_index_channel};
+pub use crate::series_index::searcher::SeriesIndexSearcher;
pub(crate) use crate::series_index::task::{SeriesIndexTaskState, spawn_series_index_tasks};
pub(crate) use crate::series_index::version::{SeriesIndexVersion, SeriesIndexVersionControl};
+pub use crate::series_index::writer::{
+ SeriesIndexWriter, SeriesIndexWriterMetrics, SeriesIndexWriterOptions, series_index_schema,
+};
pub(crate) const MIN_TS_COLUMN: &str = "__series_min_ts";
pub(crate) const MAX_TS_COLUMN: &str = "__series_max_ts";
diff --git a/src/mito2/src/series_index/maintenance.rs b/src/mito2/src/series_index/maintenance.rs
new file mode 100644
index 0000000000..91b323e2fc
--- /dev/null
+++ b/src/mito2/src/series_index/maintenance.rs
@@ -0,0 +1,305 @@
+// Copyright 2023 Greptime Team
+//
+// Licensed under the Apache License, Version 2.0 (the "License");
+// you may not use this file except in compliance with the License.
+// You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing, software
+// distributed under the License is distributed on an "AS IS" BASIS,
+// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+// See the License for the specific language governing permissions and
+// limitations under the License.
+
+//! Worker-owned series-index reconciliation and publication.
+
+use std::collections::HashSet;
+use std::sync::Arc;
+use std::time::{Duration, Instant};
+
+use common_telemetry::{debug, info};
+use object_store::ObjectStore;
+use store_api::storage::RegionId;
+
+use crate::error::Result;
+use crate::metrics::{SERIES_INDEX_RECONCILE_ELAPSED, SERIES_INDEX_RECONCILE_TOTAL};
+use crate::read::series_candidate::is_sparse_metric_metadata;
+use crate::region::MitoRegionRef;
+use crate::region::version::VersionRef;
+use crate::series_index::bucket::{
+ group_files_into_series_buckets, plan_series_indexes, rounded_bucket_width,
+};
+use crate::series_index::builder::{build_range_index, build_series_index};
+use crate::series_index::catalog::{
+ RangeIndexCatalog, SeriesIndexCatalog, range_catalog_path, series_catalog_path, store_catalog,
+};
+use crate::series_index::purger::IndexFilePurger;
+use crate::series_index::version::{SeriesIndexFileHandle, SeriesIndexVersion};
+
+/// Retires newly completed series files unless their snapshot is published.
+#[derive(Default)]
+struct UnpublishedSeriesFiles(Vec);
+
+impl UnpublishedSeriesFiles {
+ fn disarm(&mut self) {
+ self.0.clear();
+ }
+}
+
+impl Drop for UnpublishedSeriesFiles {
+ fn drop(&mut self) {
+ for handle in &self.0 {
+ handle.mark_deleted();
+ }
+ }
+}
+
+#[derive(Debug, Default)]
+pub(crate) struct ReconcileStats {
+ pub(crate) source_files: usize,
+ pub(crate) built_range: usize,
+ pub(crate) built_series: usize,
+ pub(crate) removed_range: usize,
+ pub(crate) removed_series: usize,
+ pub(crate) computed_buckets: usize,
+ pub(crate) skipped_buckets: usize,
+}
+
+impl ReconcileStats {
+ fn changed(&self) -> bool {
+ self.built_range + self.built_series + self.removed_range + self.removed_series > 0
+ }
+}
+
+/// Reconciles indexes for one region snapshot, persists catalogs, then atomically publishes it.
+pub(crate) async fn reconcile_series_indexes(
+ worker_id: u32,
+ store: ObjectStore,
+ region: MitoRegionRef,
+ requested_bucket_width: Duration,
+ now_ms: i64,
+ purger: IndexFilePurger,
+) -> Result {
+ let total_start = Instant::now();
+ // Use this snapshot throughout reconciliation, even if the region version advances.
+ let version = region.version_control.current().version;
+ if !is_sparse_metric_metadata(&version.metadata) {
+ SERIES_INDEX_RECONCILE_TOTAL
+ .with_label_values(&["noop"])
+ .inc();
+ return Ok(ReconcileStats::default());
+ }
+ let build_start = Instant::now();
+ let mut unpublished = UnpublishedSeriesFiles::default();
+ let (next, stats) = build_index_version(
+ worker_id,
+ &store,
+ ®ion,
+ &version,
+ requested_bucket_width,
+ now_ms,
+ &purger,
+ &mut unpublished,
+ )
+ .await?;
+ SERIES_INDEX_RECONCILE_ELAPSED
+ .with_label_values(&["build"])
+ .observe(build_start.elapsed().as_secs_f64());
+ // Persist both catalogs before making the new snapshot visible to readers.
+ if let Some(next) = next {
+ persist_index_catalogs(&store, region.region_id, &next).await?;
+ publish_index_version(®ion, Arc::new(next));
+ }
+ unpublished.disarm();
+ let result = if stats.changed() { "changed" } else { "noop" };
+ SERIES_INDEX_RECONCILE_TOTAL
+ .with_label_values(&[result])
+ .inc();
+ SERIES_INDEX_RECONCILE_ELAPSED
+ .with_label_values(&["total"])
+ .observe(total_start.elapsed().as_secs_f64());
+ if stats.changed() {
+ info!(
+ "Reconciled series-index snapshot, worker: {worker_id}, region: {}, elapsed: {:?}, stats: {:?}",
+ region.region_id,
+ total_start.elapsed(),
+ stats
+ );
+ } else {
+ debug!(
+ "Series-index reconciliation made no changes, worker: {worker_id}, region: {}",
+ region.region_id
+ );
+ }
+ Ok(stats)
+}
+
+/// Builds a changed snapshot, retaining reusable indexes and removing obsolete coverage.
+/// Returns `None` when no range or series indexes were added or removed.
+#[allow(clippy::too_many_arguments)]
+async fn build_index_version(
+ worker_id: u32,
+ store: &ObjectStore,
+ region: &MitoRegionRef,
+ version: &VersionRef,
+ requested_bucket_width: Duration,
+ now_ms: i64,
+ purger: &IndexFilePurger,
+ unpublished: &mut UnpublishedSeriesFiles,
+) -> Result<(Option, ReconcileStats)> {
+ let mut stats = ReconcileStats::default();
+ let files = version
+ .ssts
+ .levels()
+ .iter()
+ .flat_map(|level| level.files())
+ .cloned()
+ .collect::>();
+ stats.source_files = files.len();
+ let visible = files
+ .iter()
+ .map(|file| file.file_id().file_id())
+ .collect::>();
+ let current = region.series_index_version();
+ // The SST purger deletes companion range files after final handle release. Prune
+ // metadata here using the captured SST snapshot, independently of physical deletion;
+ // a later region-version change is picked up by the next reconciliation.
+ stats.removed_range = current.range_indexes.difference(&visible).count();
+ let buckets = match version.compaction_time_window {
+ Some(window) => rounded_bucket_width(requested_bucket_width, window)
+ // Successful rounding guarantees the window fits in i64; subsecond windows
+ // use the same one-second minimum as rounded_bucket_width.
+ .map(|width| {
+ group_files_into_series_buckets(&files, width, (window.as_secs() as i64).max(1))
+ })
+ .unwrap_or_default(),
+ None => {
+ debug!(
+ "Deferring series indexes without compaction window, worker: {worker_id}, region: {}",
+ region.region_id
+ );
+ Vec::new()
+ }
+ };
+ let plan = plan_series_indexes(
+ buckets,
+ current.index_buckets.clone(),
+ version.options.ttl,
+ now_ms,
+ );
+ stats.computed_buckets = plan.computed_buckets;
+ stats.skipped_buckets = plan.skipped_buckets;
+ if plan.builds.is_empty()
+ && plan.expired_index_ids.is_empty()
+ && stats.removed_range == 0
+ && current.range_indexes.len() == visible.len()
+ {
+ return Ok((None, stats));
+ }
+ let mut range_indexes = current.range_indexes.clone();
+ range_indexes.retain(|file_id| visible.contains(file_id));
+ let mut series_indexes = current.series_indexes.clone();
+ for id in plan
+ .expired_index_ids
+ .iter()
+ .chain(&plan.superseded_index_ids)
+ {
+ series_indexes.remove(id);
+ }
+ for (bucket, expected) in plan.builds {
+ // Complete companion indexes independently so a failed series build preserves them.
+ for file in &bucket.files {
+ if !range_indexes.contains(&file.file_id().file_id())
+ && let Some(file_id) =
+ build_range_index(store, region, version, file.clone()).await?
+ {
+ stats.built_range += 1;
+ range_indexes.insert(file_id);
+ }
+ }
+ let series_handle =
+ build_series_index(store, region, version, &bucket, &expected, purger).await?;
+ unpublished.0.push(series_handle.clone());
+ stats.built_series += 1;
+ series_indexes.insert(expected.index_uuid, series_handle);
+ }
+ // Cover SSTs outside planned aggregate builds, including skipped buckets.
+ for file in files {
+ let file_id = file.file_id().file_id();
+ if range_indexes.contains(&file_id) {
+ continue;
+ }
+ if let Some(file_id) = build_range_index(store, region, version, file).await? {
+ stats.built_range += 1;
+ range_indexes.insert(file_id);
+ }
+ }
+ stats.removed_series = current
+ .series_indexes
+ .keys()
+ .filter(|id| !series_indexes.contains_key(id))
+ .count();
+ // Bucket reconciliation is speculative until indexes change. Recompute its coverage
+ // next time rather than replacing the published snapshot on a no-op pass.
+ if !stats.changed() {
+ return Ok((None, stats));
+ }
+ let next = SeriesIndexVersion {
+ range_indexes,
+ series_indexes,
+ index_buckets: plan.index_buckets,
+ };
+ Ok((Some(next), stats))
+}
+
+/// Writes catalogs in a stable order; the two writes are not atomic together.
+async fn persist_index_catalogs(
+ store: &ObjectStore,
+ region_id: RegionId,
+ next: &SeriesIndexVersion,
+) -> Result<()> {
+ let mut range_entries = next.range_indexes.iter().copied().collect::>();
+ range_entries.sort_unstable_by(|left, right| left.as_bytes().cmp(right.as_bytes()));
+ let mut series_entries = next
+ .series_indexes
+ .values()
+ .map(|handle| handle.entry().clone())
+ .collect::>();
+ series_entries.sort_unstable_by_key(|entry| {
+ (
+ entry.bucket_start,
+ entry.bucket_end,
+ entry.min_file_sequence,
+ entry.max_file_sequence,
+ )
+ });
+ store_catalog(
+ store,
+ &range_catalog_path(region_id),
+ &RangeIndexCatalog {
+ indexes: range_entries,
+ },
+ )
+ .await?;
+ store_catalog(
+ store,
+ &series_catalog_path(region_id),
+ &SeriesIndexCatalog {
+ indexes: series_entries,
+ },
+ )
+ .await?;
+ Ok(())
+}
+
+/// Publishes the snapshot and retires series files absent from the new version.
+fn publish_index_version(region: &MitoRegionRef, next: Arc) {
+ let previous = region.series_index_version_control.publish(next.clone());
+ for (id, handle) in &previous.series_indexes {
+ if !next.series_indexes.contains_key(id) {
+ // Purge only after readers release their retained handles.
+ handle.mark_deleted();
+ }
+ }
+}
diff --git a/src/mito2/src/series_index/purger.rs b/src/mito2/src/series_index/purger.rs
index 039ca66ccc..5298e9e0d1 100644
--- a/src/mito2/src/series_index/purger.rs
+++ b/src/mito2/src/series_index/purger.rs
@@ -108,3 +108,56 @@ pub(crate) fn series_index_channel(
let (sender, receiver) = unbounded_channel();
(IndexFilePurger { store, sender }, receiver)
}
+
+#[cfg(test)]
+mod tests {
+ use std::time::Duration;
+
+ use object_store::ObjectStore;
+ use object_store::services::Memory;
+ use store_api::storage::{FileId, RegionId};
+
+ use super::*;
+
+ #[tokio::test]
+ async fn test_purge_drains_queue_after_last_sender_drops() {
+ let store = ObjectStore::new(Memory::default()).unwrap();
+ let (purger, receiver) = series_index_channel(store.clone());
+ let mut paths = Vec::new();
+ for _ in 0..3 {
+ let file_id = RegionFileId::new(RegionId::new(1, 1), FileId::random());
+ let path = series_index_path(file_id.region_id(), file_id.file_id());
+ store.write(&path, "index").await.unwrap();
+ paths.push(path);
+ purger.purge(PurgeRequest { file_id });
+ }
+ drop(purger);
+ tokio::time::timeout(
+ Duration::from_secs(10),
+ run_index_purge_task(0, store.clone(), receiver),
+ )
+ .await
+ .unwrap();
+ for path in paths {
+ assert!(!store.exists(&path).await.unwrap());
+ }
+ }
+
+ #[tokio::test]
+ async fn test_purge_falls_back_after_receiver_drops() {
+ let store = ObjectStore::new(Memory::default()).unwrap();
+ let (purger, receiver) = series_index_channel(store.clone());
+ drop(receiver);
+ let file_id = RegionFileId::new(RegionId::new(1, 1), FileId::random());
+ let path = series_index_path(file_id.region_id(), file_id.file_id());
+ store.write(&path, "index").await.unwrap();
+ purger.purge(PurgeRequest { file_id });
+ tokio::time::timeout(Duration::from_secs(10), async {
+ while store.exists(&path).await.unwrap() {
+ tokio::task::yield_now().await;
+ }
+ })
+ .await
+ .unwrap();
+ }
+}
diff --git a/src/mito2/src/series_index/task.rs b/src/mito2/src/series_index/task.rs
index 76b42dd62f..bef477a039 100644
--- a/src/mito2/src/series_index/task.rs
+++ b/src/mito2/src/series_index/task.rs
@@ -18,14 +18,18 @@ use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::time::Duration;
-use common_telemetry::info;
+use common_telemetry::{info, warn};
use object_store::ObjectStore;
use tokio::sync::Notify;
use tokio::sync::mpsc::UnboundedReceiver;
use tokio::task::JoinHandle;
use tokio::time::{Instant, MissedTickBehavior};
-use crate::series_index::purger::{PurgeRequest, run_index_purge_task};
+use crate::metrics::SERIES_INDEX_RECONCILE_TOTAL;
+use crate::region::{RegionLeaderState, RegionMapRef, RegionRoleState};
+use crate::series_index::maintenance::reconcile_series_indexes;
+use crate::series_index::purger::{IndexFilePurger, PurgeRequest, run_index_purge_task};
+use crate::time_provider::TimeProviderRef;
/// Shared lifecycle state for a worker's series-index task.
#[derive(Debug)]
@@ -46,6 +50,10 @@ impl SeriesIndexTaskState {
self.running.load(Ordering::Acquire)
}
+ pub(crate) fn wake(&self) {
+ self.notify.notify_one();
+ }
+
pub(crate) fn stop(&self) {
self.running.store(false, Ordering::Release);
// Retain a permit if maintenance has not started waiting yet.
@@ -58,20 +66,34 @@ impl SeriesIndexTaskState {
}
/// Starts both tasks, detaching purge and returning the maintenance handle.
+#[allow(clippy::too_many_arguments)]
pub(crate) fn spawn_series_index_tasks(
worker_id: u32,
store: ObjectStore,
+ regions: RegionMapRef,
state: Arc,
+ bucket_width: Duration,
+ purger: IndexFilePurger,
purge_receiver: UnboundedReceiver,
interval: Duration,
+ time_provider: TimeProviderRef,
) -> JoinHandle<()> {
// Snapshots may retain senders after the worker stops; purge until all senders drop.
- common_runtime::spawn_global(run_index_purge_task(worker_id, store, purge_receiver));
+ common_runtime::spawn_global(run_index_purge_task(
+ worker_id,
+ store.clone(),
+ purge_receiver,
+ ));
common_runtime::spawn_global(async move {
SeriesIndexTask {
worker_id,
+ store,
+ regions,
+ bucket_width,
+ purger,
state,
interval,
+ time_provider,
}
.run()
.await;
@@ -80,9 +102,14 @@ pub(crate) fn spawn_series_index_tasks(
/// Periodic series-index maintenance for one region worker.
struct SeriesIndexTask {
+ store: ObjectStore,
+ regions: RegionMapRef,
+ bucket_width: Duration,
+ purger: IndexFilePurger,
worker_id: u32,
state: Arc,
interval: Duration,
+ time_provider: TimeProviderRef,
}
impl SeriesIndexTask {
@@ -90,17 +117,17 @@ impl SeriesIndexTask {
async fn run(mut self) {
let worker_id = self.worker_id;
info!("Start series-index background task, worker: {worker_id}");
- let interval = self.interval;
+ let interval = self.time_provider.wait_duration(self.interval);
let mut timer = tokio::time::interval_at(Instant::now() + interval, interval);
- timer.set_missed_tick_behavior(MissedTickBehavior::Skip);
+ // Schedule future ticks from a late tick rather than the original cadence.
+ timer.set_missed_tick_behavior(MissedTickBehavior::Delay);
while self.state.is_running() {
tokio::select! {
_ = self.state.notified() => {}
- _ = timer.tick() => {
- if self.state.is_running() {
- self.maintain().await;
- }
- }
+ _ = timer.tick() => {}
+ }
+ if self.state.is_running() {
+ self.maintain().await;
}
}
info!("Stop series-index background task, worker: {worker_id}");
@@ -108,6 +135,94 @@ impl SeriesIndexTask {
/// Runs periodic maintenance independently of incoming deletion requests.
async fn maintain(&mut self) {
- // TODO: Reconcile indexes and perform other periodic maintenance here.
+ for region in self.regions.list_regions() {
+ if !self.state.is_running() {
+ break;
+ }
+ // Best effort: the region can still change state during reconciliation.
+ // Local indexes can be built on followers as well as writable leaders.
+ if !matches!(
+ region.state(),
+ RegionRoleState::Follower | RegionRoleState::Leader(RegionLeaderState::Writable)
+ ) {
+ continue;
+ }
+ if let Err(error) = reconcile_series_indexes(
+ self.worker_id,
+ self.store.clone(),
+ region.clone(),
+ self.bucket_width,
+ self.time_provider.current_time_millis(),
+ self.purger.clone(),
+ )
+ .await
+ {
+ SERIES_INDEX_RECONCILE_TOTAL
+ .with_label_values(&["failure"])
+ .inc();
+ warn!(error; "Failed to reconcile series indexes, worker: {}, region: {}", self.worker_id, region.region_id);
+ }
+ }
+ }
+}
+
+#[cfg(test)]
+mod tests {
+ use object_store::services::Memory;
+ use store_api::region_engine::RegionEngine;
+
+ use super::*;
+ use crate::region::RegionMap;
+ use crate::series_index::catalog::{range_catalog_path, series_catalog_path};
+ use crate::series_index::purger::series_index_channel;
+ use crate::series_index::tests::prepare_region;
+ use crate::test_util::TestEnv;
+
+ #[rstest::rstest]
+ #[case::follower(RegionRoleState::Follower, true)]
+ #[case::writable(RegionRoleState::Leader(RegionLeaderState::Writable), true)]
+ #[case::staging(RegionRoleState::Leader(RegionLeaderState::Staging), false)]
+ #[case::entering_staging(RegionRoleState::Leader(RegionLeaderState::EnteringStaging), false)]
+ #[case::altering(RegionRoleState::Leader(RegionLeaderState::Altering), false)]
+ #[case::dropping(RegionRoleState::Leader(RegionLeaderState::Dropping), false)]
+ #[case::truncating(RegionRoleState::Leader(RegionLeaderState::Truncating), false)]
+ #[case::editing(RegionRoleState::Leader(RegionLeaderState::Editing), false)]
+ #[case::downgrading(RegionRoleState::Leader(RegionLeaderState::Downgrading), false)]
+ #[tokio::test]
+ async fn test_maintenance_region_states(#[case] role: RegionRoleState, #[case] builds: bool) {
+ let mut env = TestEnv::with_prefix("series-maintenance-state").await;
+ let (engine, region) = prepare_region(&mut env).await;
+ // Install the desired state without triggering the corresponding DDL.
+ region.switch_state_to_staging(RegionLeaderState::Writable);
+ region
+ .manifest_ctx
+ .exit_staging(region.region_id, role)
+ .unwrap();
+ let store = ObjectStore::new(Memory::default()).unwrap();
+ let (purger, _receiver) = series_index_channel(store.clone());
+ let regions = Arc::new(RegionMap::default());
+ regions.insert_region(region.clone());
+ let mut task = SeriesIndexTask {
+ store: store.clone(),
+ regions,
+ bucket_width: Duration::from_secs(100),
+ purger,
+ worker_id: 0,
+ state: Arc::new(SeriesIndexTaskState::new()),
+ interval: Duration::from_secs(3600),
+ time_provider: Arc::new(crate::time_provider::StdTimeProvider),
+ };
+ task.maintain().await;
+ assert_eq!(
+ builds,
+ !region.series_index_version().series_indexes.is_empty()
+ );
+ for path in [
+ range_catalog_path(region.region_id),
+ series_catalog_path(region.region_id),
+ ] {
+ assert_eq!(builds, store.exists(&path).await.unwrap());
+ }
+ engine.stop().await.unwrap();
}
}
diff --git a/src/mito2/src/series_index/tests.rs b/src/mito2/src/series_index/tests.rs
index fa65dadcdd..50109f9d4d 100644
--- a/src/mito2/src/series_index/tests.rs
+++ b/src/mito2/src/series_index/tests.rs
@@ -12,25 +12,36 @@
// See the License for the specific language governing permissions and
// limitations under the License.
-//! Shared fixtures for series-index construction tests.
+//! Behavioral coverage for series-index reconciliation.
use std::sync::Arc;
+use std::time::Duration;
use api::v1::helper::row;
use api::v1::value::ValueData;
use api::v1::{ColumnDataType, Rows, SemanticType, WriteHint};
+use object_store::ObjectStore;
+use object_store::layers::mock::{self, MockLayerBuilder, oio};
+use object_store::services::Memory;
use store_api::codec::PrimaryKeyEncoding;
use store_api::metric_engine_consts::PRIMARY_KEY_ENCODING;
use store_api::region_engine::RegionEngine;
use store_api::region_request::{RegionPutRequest, RegionRequest};
-use store_api::storage::RegionId;
use store_api::storage::consts::PRIMARY_KEY_COLUMN_NAME;
+use store_api::storage::{FileId, RegionId};
+use super::catalog::{
+ delete_catalogs, load_version_control, range_catalog_path, range_index_path,
+ series_catalog_path, series_index_path,
+};
+use super::maintenance::reconcile_series_indexes;
+use super::purger::series_index_channel;
use crate::config::MitoConfig;
use crate::engine::MitoEngine;
-use crate::region::MitoRegionRef;
+use crate::region::{MitoRegionRef, RegionMap};
use crate::test_util::sst_util::{new_sparse_primary_key, sst_region_metadata_with_encoding};
use crate::test_util::{CreateRequestBuilder, TestEnv, flush_region, rows_schema};
+use crate::time_provider::StdTimeProvider;
/// Builds real sparse SSTs; background maintenance is disabled so tests control publication.
pub(super) async fn prepare_region(env: &mut TestEnv) -> (MitoEngine, MitoRegionRef) {
@@ -66,10 +77,12 @@ async fn prepare_region_with_timestamps(
"100s".to_string(),
);
// Keep the source SSTs stable while tests reconcile multiple buckets.
- request.options.insert(
- "compaction.twcs.trigger_file_num".to_string(),
- "100".to_string(),
- );
+ for window in ["active_window", "inactive_window"] {
+ request.options.insert(
+ format!("compaction.twcs.{window}.trigger_file_num"),
+ "100".to_string(),
+ );
+ }
let full_schema = rows_schema(&request);
let mut pk_column = full_schema[0].clone();
pk_column.column_name = PRIMARY_KEY_COLUMN_NAME.to_string();
@@ -112,3 +125,631 @@ async fn prepare_region_with_timestamps(
let region = engine.get_region(region_id).unwrap();
(engine, region)
}
+
+struct IndexTest {
+ region: MitoRegionRef,
+ store: ObjectStore,
+ purger: super::purger::IndexFilePurger,
+}
+
+impl IndexTest {
+ fn new(
+ region: MitoRegionRef,
+ ) -> (
+ Self,
+ tokio::sync::mpsc::UnboundedReceiver,
+ ) {
+ let store = ObjectStore::new(Memory::default()).unwrap();
+ let (purger, receiver) = series_index_channel(store.clone());
+ (
+ Self {
+ region,
+ store,
+ purger,
+ },
+ receiver,
+ )
+ }
+
+ async fn reconcile(
+ &self,
+ store: ObjectStore,
+ ) -> crate::error::Result {
+ reconcile_series_indexes(
+ 0,
+ store,
+ self.region.clone(),
+ Duration::from_secs(100),
+ 0,
+ self.purger.clone(),
+ )
+ .await
+ }
+}
+
+#[tokio::test]
+async fn test_reconcile_restores_and_reuses_indexes() {
+ let mut env = TestEnv::with_prefix("series-reconcile").await;
+ let (engine, region) = prepare_region(&mut env).await;
+ let (test, _receiver) = IndexTest::new(region.clone());
+ let store = &test.store;
+ let purger = &test.purger;
+ let stats = test.reconcile(store.clone()).await.unwrap();
+ assert_eq!((4, 1), (stats.built_range, stats.built_series));
+ let first = region.series_index_version();
+ let first_id = *first.series_indexes.keys().next().unwrap();
+
+ let missing_paths = [
+ range_index_path(
+ region.region_id,
+ *first.range_indexes.iter().next().unwrap(),
+ ),
+ series_index_path(region.region_id, first_id),
+ ];
+ for path in &missing_paths {
+ store.delete(path).await.unwrap();
+ }
+
+ // Restore catalog entries even when their index files are missing.
+ let restored = load_version_control(store, region.region_id, purger)
+ .await
+ .current();
+ assert_eq!(first.range_indexes, restored.range_indexes);
+ assert_eq!(first.index_buckets, restored.index_buckets);
+ assert_eq!(
+ first.series_indexes[&first_id].entry(),
+ restored.series_indexes[&first_id].entry()
+ );
+
+ // Catalog changes are not reloaded, and missing index files are not repaired.
+ for path in [
+ range_catalog_path(region.region_id),
+ series_catalog_path(region.region_id),
+ ] {
+ store.write(&path, "{}").await.unwrap();
+ }
+ let stats = test.reconcile(store.clone()).await.unwrap();
+ assert_eq!((0, 0), (stats.built_range, stats.built_series));
+ assert!(Arc::ptr_eq(&first, ®ion.series_index_version()));
+ for path in [
+ range_catalog_path(region.region_id),
+ series_catalog_path(region.region_id),
+ ] {
+ assert_eq!(
+ b"{}".as_slice(),
+ store.read(&path).await.unwrap().to_bytes().as_ref()
+ );
+ }
+ for path in &missing_paths {
+ assert!(!store.exists(path).await.unwrap());
+ }
+ engine.stop().await.unwrap();
+}
+
+#[tokio::test]
+async fn test_reconcile_range_only_build_and_removal() {
+ let mut env = TestEnv::with_prefix("series-range-lifecycle").await;
+ let (engine, region) = prepare_region_with_timestamps(&mut env, &[1000]).await;
+ let (test, _receiver) = IndexTest::new(region.clone());
+ let initial = region.series_index_version();
+ let stats = test.reconcile(test.store.clone()).await.unwrap();
+ assert_eq!((1, 0), (stats.built_range, stats.built_series));
+ let built = region.series_index_version();
+ assert!(!Arc::ptr_eq(&initial, &built));
+ let files = region
+ .version()
+ .ssts
+ .levels()
+ .iter()
+ .flat_map(|level| level.files())
+ .map(|file| file.meta_ref().clone())
+ .collect();
+ region.version_control.apply_edit(
+ Some(crate::manifest::action::RegionEdit {
+ files_to_remove: files,
+ files_to_add: Vec::new(),
+ timestamp_ms: None,
+ compaction_time_window: None,
+ flushed_entry_id: None,
+ flushed_sequence: None,
+ committed_sequence: None,
+ }),
+ &[],
+ crate::test_util::new_noop_file_purger(),
+ );
+ let stats = test.reconcile(test.store.clone()).await.unwrap();
+ assert_eq!(
+ (0, 0, 1),
+ (stats.built_range, stats.built_series, stats.removed_range)
+ );
+ let removed = region.series_index_version();
+ assert!(removed.range_indexes.is_empty());
+ assert!(!Arc::ptr_eq(&built, &removed));
+ let restored = load_version_control(&test.store, region.region_id, &test.purger).await;
+ assert!(restored.current().range_indexes.is_empty());
+ test.reconcile(test.store.clone()).await.unwrap();
+ assert!(Arc::ptr_eq(&removed, ®ion.series_index_version()));
+ engine.stop().await.unwrap();
+}
+
+struct FailingSeriesWriter {
+ inner: oio::Writer,
+ fail: bool,
+}
+
+#[rstest::rstest]
+#[case::series("series")]
+#[case::range_catalog("range-index.json")]
+#[case::series_catalog("series-index.json")]
+#[tokio::test]
+async fn test_reconcile_replaces_whole_bucket_after_failure(#[case] failure: &str) {
+ let mut env = TestEnv::with_prefix("series-replacement").await;
+ let (engine, region) =
+ prepare_region_with_timestamps(&mut env, &[1000, 2000, 3000, 4000, 5000]).await;
+ // Retain all source handles while hiding the newest SST for the initial build.
+ let sources = region.version();
+ let newest = sources
+ .ssts
+ .levels()
+ .iter()
+ .flat_map(|level| level.files())
+ .max_by_key(|file| file.meta_ref().sequence)
+ .unwrap()
+ .meta_ref()
+ .clone();
+ let edit = |add| crate::manifest::action::RegionEdit {
+ files_to_add: if add { vec![newest.clone()] } else { vec![] },
+ files_to_remove: if add { vec![] } else { vec![newest.clone()] },
+ timestamp_ms: None,
+ compaction_time_window: None,
+ flushed_entry_id: None,
+ flushed_sequence: None,
+ committed_sequence: None,
+ };
+ region.version_control.apply_edit(
+ Some(edit(false)),
+ &[],
+ crate::test_util::new_noop_file_purger(),
+ );
+ let (test, mut receiver) = IndexTest::new(region.clone());
+ test.reconcile(test.store.clone()).await.unwrap();
+ let previous = region.series_index_version();
+ let old_id = *previous.series_indexes.keys().next().unwrap();
+ assert_eq!(
+ 4,
+ previous.series_indexes[&old_id]
+ .entry()
+ .source_file_ids
+ .len()
+ );
+ region.version_control.apply_edit(
+ Some(edit(true)),
+ &[],
+ crate::test_util::new_noop_file_purger(),
+ );
+
+ let failure = failure.to_string();
+ let layer = MockLayerBuilder::default()
+ .writer_factory(Arc::new(move |path, _, inner| {
+ Box::new(FailingSeriesWriter {
+ inner,
+ fail: if failure == "series" {
+ path.contains("/series/")
+ } else {
+ path.ends_with(&failure)
+ },
+ })
+ }))
+ .build()
+ .unwrap();
+ assert!(
+ test.reconcile(test.store.clone().layer(layer))
+ .await
+ .is_err()
+ );
+ assert!(Arc::ptr_eq(&previous, ®ion.series_index_version()));
+ // Failed catalog publication may retire a completed replacement, never its predecessor.
+ while let Ok(request) = receiver.try_recv() {
+ assert_ne!(old_id, request.file_id.file_id());
+ }
+
+ let stats = test.reconcile(test.store.clone()).await.unwrap();
+ assert_eq!((1, 1), (stats.built_series, stats.removed_series));
+ let replacement = region.series_index_version();
+ assert_eq!(1, replacement.series_indexes.len());
+ assert!(!replacement.series_indexes.contains_key(&old_id));
+ let entry = replacement.series_indexes.values().next().unwrap().entry();
+ assert_eq!(5, entry.source_file_ids.len());
+ assert!(entry.source_file_ids.contains(&newest.file_id));
+ let restored = load_version_control(&test.store, region.region_id, &test.purger).await;
+ assert_eq!(replacement.index_buckets, restored.current().index_buckets);
+ assert_eq!(1, restored.current().series_indexes.len());
+ test.reconcile(test.store.clone()).await.unwrap();
+ assert!(Arc::ptr_eq(&replacement, ®ion.series_index_version()));
+
+ // Publication retires the old file only after the last reader releases its snapshot.
+ let old_path = series_index_path(region.region_id, old_id);
+ assert!(test.store.exists(&old_path).await.unwrap());
+ assert!(receiver.try_recv().is_err());
+ drop(previous);
+ assert_eq!(old_id, receiver.try_recv().unwrap().file_id.file_id());
+ assert!(receiver.try_recv().is_err());
+ engine.stop().await.unwrap();
+}
+
+impl mock::Write for FailingSeriesWriter {
+ async fn write(&mut self, buffer: mock::Buffer) -> mock::Result<()> {
+ self.inner.write(buffer).await
+ }
+
+ async fn close(&mut self) -> mock::Result {
+ if self.fail {
+ return Err(mock::Error::new(
+ mock::ErrorKind::Unexpected,
+ "injected series write failure",
+ ));
+ }
+ self.inner.close().await
+ }
+
+ async fn abort(&mut self) -> mock::Result<()> {
+ self.inner.abort().await
+ }
+}
+
+#[tokio::test]
+async fn test_failed_series_build_keeps_completed_sst_range_indexes() {
+ let mut env = TestEnv::with_prefix("series-build-failure").await;
+ let (engine, region) = prepare_region(&mut env).await;
+ let (test, mut receiver) = IndexTest::new(region.clone());
+ let store = &test.store;
+ let layer = MockLayerBuilder::default()
+ .writer_factory(Arc::new(|path, _, inner| {
+ Box::new(FailingSeriesWriter {
+ inner,
+ fail: path.contains("/series/"),
+ })
+ }))
+ .build()
+ .unwrap();
+ let failing_store = store.clone().layer(layer);
+ assert!(test.reconcile(failing_store).await.is_err());
+ assert!(region.series_index_version().range_indexes.is_empty());
+ assert!(receiver.try_recv().is_err());
+ for file in region
+ .version()
+ .ssts
+ .levels()
+ .iter()
+ .flat_map(|level| level.files())
+ {
+ let path = range_index_path(region.region_id, file.file_id().file_id());
+ assert!(store.exists(&path).await.unwrap());
+ }
+ // A retry can rebuild and publish the complete snapshot.
+ let stats = test.reconcile(store.clone()).await.unwrap();
+ assert_eq!((4, 1), (stats.built_range, stats.built_series));
+ engine.stop().await.unwrap();
+}
+
+#[tokio::test]
+async fn test_reconcile_publishes_after_region_version_changes() {
+ for during_catalog_write in [false, true] {
+ let mut env = TestEnv::with_prefix("series-version-change").await;
+ let (engine, region) = prepare_region(&mut env).await;
+ let initial_version = region.version();
+ let (test, mut receiver) = IndexTest::new(region.clone());
+ let store = &test.store;
+ let purger = &test.purger;
+ let target_region = region.clone();
+ let catalog_path = range_catalog_path(region.region_id);
+ let layer = MockLayerBuilder::default()
+ .writer_factory(Arc::new(move |path, _, inner| {
+ let should_advance = if during_catalog_write {
+ path == catalog_path
+ } else {
+ path.contains("/series/")
+ };
+ if should_advance {
+ // Replace the version while reconciliation is writing its captured snapshot.
+ target_region
+ .version_control
+ .alter_options(target_region.version().options.clone());
+ }
+ inner
+ }))
+ .build()
+ .unwrap();
+ let stats = test.reconcile(store.clone().layer(layer)).await.unwrap();
+
+ assert!(!Arc::ptr_eq(&initial_version, ®ion.version()));
+ assert_eq!((4, 1), (stats.built_range, stats.built_series));
+ let published = region.series_index_version();
+ assert_eq!(4, published.range_indexes.len());
+ assert_eq!(1, published.series_indexes.len());
+ let restored = load_version_control(store, region.region_id, purger).await;
+ assert_eq!(published.range_indexes, restored.current().range_indexes);
+ let handle = published.series_indexes.values().next().unwrap();
+ assert_eq!(
+ handle.entry(),
+ restored.current().series_indexes[&handle.entry().index_uuid].entry()
+ );
+ assert!(receiver.try_recv().is_err());
+ engine.stop().await.unwrap();
+ }
+}
+
+#[tokio::test]
+async fn test_failed_reconcile_retires_only_unpublished_series() {
+ use std::collections::{HashMap, HashSet};
+ use std::sync::Mutex;
+
+ use super::bucket::{group_files_into_series_buckets, plan_series_indexes};
+ use super::builder::{build_range_index, build_series_index};
+ use super::purger::run_index_purge_task;
+ use super::version::SeriesIndexVersion;
+
+ for failure in ["series", "range", "range-index.json", "series-index.json"] {
+ let mut env = TestEnv::with_prefix("series-unpublished-cleanup").await;
+ let mut timestamps = vec![
+ -99000, -98000, -97000, -96000, 1000, 2000, 3000, 4000, 101000,
+ ];
+ if failure != "range" {
+ timestamps.extend([102000, 103000, 104000]);
+ }
+ let (engine, region) = prepare_region_with_timestamps(&mut env, ×tamps).await;
+ let (test, mut receiver) = IndexTest::new(region.clone());
+ let store = &test.store;
+ let purger = &test.purger;
+
+ // Publish the first bucket so failed attempts must preserve a reused handle.
+ let version = region.version();
+ let files = version
+ .ssts
+ .levels()
+ .iter()
+ .flat_map(|level| level.files())
+ .cloned()
+ .collect::>();
+ let plan = plan_series_indexes(
+ group_files_into_series_buckets(
+ &files,
+ 100,
+ version.compaction_time_window.unwrap().as_secs() as i64,
+ ),
+ Default::default(),
+ None,
+ 0,
+ );
+ assert_eq!(if failure == "range" { 2 } else { 3 }, plan.builds.len());
+ let (bucket, entry) = &plan.builds[0];
+ let mut ranges = HashSet::new();
+ for file in &bucket.files {
+ let id = build_range_index(store, ®ion, &version, file.clone())
+ .await
+ .unwrap()
+ .unwrap();
+ ranges.insert(id);
+ }
+ let handle = build_series_index(store, ®ion, &version, bucket, entry, purger)
+ .await
+ .unwrap();
+ let reused_path = series_index_path(region.region_id, entry.index_uuid);
+ region
+ .series_index_version_control
+ .publish(Arc::new(SeriesIndexVersion::new(
+ ranges,
+ HashMap::from([(entry.index_uuid, handle)]),
+ )));
+ let previous = region.series_index_version();
+ let standalone_path = files
+ .iter()
+ .find(|file| file.time_range().0 == common_time::Timestamp::new_millisecond(101000))
+ .map(|file| range_index_path(region.region_id, file.file_id().file_id()))
+ .unwrap();
+
+ // Every completed unpublished output must be retired.
+ let attempted = Arc::new(Mutex::new(Vec::::new()));
+ let recorded = attempted.clone();
+ let layer = MockLayerBuilder::default()
+ .writer_factory(Arc::new(move |path, _, inner| {
+ let mut paths = recorded.lock().unwrap();
+ if path.contains("/series/") {
+ paths.push(path.to_string());
+ }
+ let fail = match failure {
+ "series" => path.contains("/series/") && paths.len() == 2,
+ "range" => path == standalone_path,
+ catalog => path.ends_with(catalog),
+ };
+ Box::new(FailingSeriesWriter { inner, fail })
+ }))
+ .build()
+ .unwrap();
+ let result = test.reconcile(store.clone().layer(layer)).await;
+ assert!(
+ result.is_err(),
+ "expected {failure} failure, got {result:?}; attempted series paths: {:?}",
+ attempted.lock().unwrap()
+ );
+ assert!(Arc::ptr_eq(&previous, ®ion.series_index_version()));
+
+ let completed = attempted
+ .lock()
+ .unwrap()
+ .iter()
+ .take(if failure.ends_with(".json") { 2 } else { 1 })
+ .cloned()
+ .collect::>();
+ assert_eq!(
+ if failure.ends_with(".json") { 2 } else { 1 },
+ completed.len()
+ );
+ let mut retired = HashSet::new();
+ // Drain through the real purge task using a finite channel.
+ let (drain_purger, drain_receiver) = series_index_channel(store.clone());
+ while let Ok(request) = receiver.try_recv() {
+ let path = series_index_path(request.file_id.region_id(), request.file_id.file_id());
+ assert!(store.exists(&path).await.unwrap());
+ assert!(retired.insert(path));
+ drain_purger.purge(request);
+ }
+ assert_eq!(completed, retired);
+ drop(drain_purger);
+ run_index_purge_task(0, store.clone(), drain_receiver).await;
+ let attempted_paths = attempted.lock().unwrap().clone();
+ for path in attempted_paths {
+ assert!(!store.exists(&path).await.unwrap());
+ }
+ assert!(store.exists(&reused_path).await.unwrap());
+
+ let stats = test.reconcile(store.clone()).await.unwrap();
+ assert_eq!(if failure == "range" { 1 } else { 2 }, stats.built_series);
+ let published = region.series_index_version();
+ for handle in published.series_indexes.values() {
+ assert!(
+ store
+ .exists(&series_index_path(
+ region.region_id,
+ handle.entry().index_uuid
+ ))
+ .await
+ .unwrap()
+ );
+ }
+ // Release all snapshots without retiring them: neither reused nor published files
+ // should have been marked deleted by the cleanup guard.
+ region.series_index_version_control.publish(Arc::default());
+ drop(published);
+ drop(previous);
+ assert!(receiver.try_recv().is_err());
+ engine.stop().await.unwrap();
+ }
+}
+
+/// Waits for a background state change without relying on exact scheduling delays.
+async fn wait_for(mut ready: impl FnMut() -> bool) {
+ tokio::time::timeout(Duration::from_secs(10), async {
+ while !ready() {
+ tokio::time::sleep(Duration::from_millis(10)).await;
+ }
+ })
+ .await
+ .unwrap();
+}
+
+#[tokio::test]
+async fn test_maintenance_wakeup_and_timer() {
+ let mut env = TestEnv::with_prefix("series-task").await;
+ let (engine, region) = prepare_region(&mut env).await;
+ let mut options = region.version().options.clone();
+ options.ttl = Some(common_time::TimeToLive::Duration(Duration::from_secs(100)));
+ region.version_control.alter_options(options);
+ let store = ObjectStore::new(Memory::default()).unwrap();
+ let (purger, receiver) = series_index_channel(store.clone());
+ let state = Arc::new(super::task::SeriesIndexTaskState::new());
+ let clock = Arc::new(crate::time_provider::mock::MockTimeProvider::new(0));
+ // A notification issued before the task starts must also trigger maintenance.
+ state.wake();
+ let regions = Arc::new(RegionMap::default());
+ regions.insert_region(region.clone());
+ let task = super::task::spawn_series_index_tasks(
+ 0,
+ store,
+ regions,
+ state.clone(),
+ Duration::from_secs(100),
+ purger,
+ receiver,
+ Duration::from_secs(3600),
+ clock.clone(),
+ );
+ wait_for(|| !region.series_index_version().series_indexes.is_empty()).await;
+ // Advance event time; periodic maintenance must expire coverage without a wakeup.
+ clock.set_now(201_000);
+ wait_for(|| region.series_index_version().series_indexes.is_empty()).await;
+ state.stop();
+ task.await.unwrap();
+ engine.stop().await.unwrap();
+}
+
+#[tokio::test]
+async fn test_drop_catalogs_and_retained_snapshot_after_task_stop() {
+ use std::collections::BTreeMap;
+
+ use common_time::Timestamp;
+
+ use super::catalog::{SeriesIndexCatalog, SeriesIndexEntry, WindowSequence, store_catalog};
+
+ let store = ObjectStore::new(Memory::default()).unwrap();
+ let region_id = RegionId::new(1, 1);
+ let entry = SeriesIndexEntry {
+ index_uuid: FileId::random(),
+ bucket_start: Timestamp::new_second(0),
+ bucket_end: Timestamp::new_second(100),
+ source_file_ids: vec![FileId::random()],
+ min_file_sequence: 1,
+ max_file_sequence: 2,
+ compaction_window_secs: 10,
+ window_sequences: BTreeMap::from([(
+ 0,
+ WindowSequence {
+ start: 0,
+ end: 100,
+ max_sequence: 2,
+ },
+ )]),
+ };
+ let path = series_index_path(region_id, entry.index_uuid);
+ store.write(&path, "index").await.unwrap();
+ store_catalog(
+ &store,
+ &series_catalog_path(region_id),
+ &SeriesIndexCatalog {
+ indexes: vec![entry.clone()],
+ },
+ )
+ .await
+ .unwrap();
+ store_catalog(
+ &store,
+ &range_catalog_path(region_id),
+ &super::catalog::RangeIndexCatalog::default(),
+ )
+ .await
+ .unwrap();
+ let (purger, receiver) = series_index_channel(store.clone());
+ let control = load_version_control(&store, region_id, &purger).await;
+ let snapshot = control.current();
+ assert_eq!(&entry, snapshot.series_indexes[&entry.index_uuid].entry());
+ let state = Arc::new(super::task::SeriesIndexTaskState::new());
+ state.stop();
+ super::task::spawn_series_index_tasks(
+ 0,
+ store.clone(),
+ Arc::new(RegionMap::default()),
+ state,
+ Duration::from_secs(100),
+ purger,
+ receiver,
+ Duration::from_secs(3600),
+ Arc::new(StdTimeProvider),
+ )
+ .await
+ .unwrap();
+ delete_catalogs(&store, region_id).await;
+ // Removing already absent catalogs is harmless.
+ delete_catalogs(&store, region_id).await;
+ assert!(!store.exists(&range_catalog_path(region_id)).await.unwrap());
+ assert!(!store.exists(&series_catalog_path(region_id)).await.unwrap());
+ control.mark_dropped();
+ assert!(store.exists(&path).await.unwrap());
+ drop(snapshot);
+ tokio::time::timeout(Duration::from_secs(10), async {
+ while store.exists(&path).await.unwrap() {
+ tokio::task::yield_now().await;
+ }
+ })
+ .await
+ .unwrap();
+}
diff --git a/src/mito2/src/sst/file_purger.rs b/src/mito2/src/sst/file_purger.rs
index 935c027258..5ff88c0874 100644
--- a/src/mito2/src/sst/file_purger.rs
+++ b/src/mito2/src/sst/file_purger.rs
@@ -288,8 +288,14 @@ mod tests {
use crate::sst::index::puffin_manager::PuffinManagerFactory;
use crate::sst::location;
+ #[rstest::rstest]
+ #[case(false)]
+ #[case(true)]
#[tokio::test]
- async fn test_file_purge() {
+ async fn test_range_index_purge_on_handle_release(
+ #[case] gc_enabled: bool,
+ #[values(false, true)] is_delete: bool,
+ ) {
common_telemetry::init_default_ut_logging();
let dir = create_temp_dir("file-purge");
@@ -320,7 +326,19 @@ mod tests {
let scheduler = Arc::new(LocalScheduler::new(3));
- let file_purger = Arc::new(LocalFilePurger::new(scheduler.clone(), layer, None));
+ let index_store = ObjectStore::new(object_store::services::Memory::default()).unwrap();
+ let owner = RegionId::new(9, 1);
+ let index_path = crate::sst::range_index::range_index_path(owner, sst_file_id.file_id());
+ index_store.write(&index_path, "range index").await.unwrap();
+ let file_purger = create_file_purger(
+ gc_enabled,
+ PathType::Bare,
+ scheduler.clone(),
+ layer,
+ None,
+ Arc::new(crate::sst::file_ref::FileReferenceManager::new(None)),
+ Some(RangeIndexDeleter::new(index_store.clone(), owner)),
+ );
{
let handle = FileHandle::new(
@@ -344,13 +362,22 @@ mod tests {
},
file_purger,
);
- // mark file as deleted and drop the handle, we expect the file is deleted.
- handle.mark_deleted();
+ if is_delete {
+ handle.mark_deleted();
+ }
+ let reader = handle.clone();
+ drop(handle);
+ assert!(index_store.exists(&index_path).await.unwrap());
+ drop(reader);
}
scheduler.stop(true).await.unwrap();
- assert!(!object_store.exists(&path).await.unwrap());
+ assert_eq!(
+ object_store.exists(&path).await.unwrap(),
+ gc_enabled || !is_delete
+ );
+ assert_eq!(index_store.exists(&index_path).await.unwrap(), !is_delete);
}
#[tokio::test]
diff --git a/src/mito2/src/time_provider.rs b/src/mito2/src/time_provider.rs
index 37cfd59e41..6acda19792 100644
--- a/src/mito2/src/time_provider.rs
+++ b/src/mito2/src/time_provider.rs
@@ -50,3 +50,49 @@ impl TimeProvider for StdTimeProvider {
current_time_millis() - current_millis
}
}
+
+/// Controllable event time with short waits for background-task tests.
+#[cfg(test)]
+pub(crate) mod mock {
+ use std::sync::atomic::{AtomicI64, Ordering};
+ use std::time::Duration;
+
+ use crate::time_provider::TimeProvider;
+
+ #[derive(Debug)]
+ pub(crate) struct MockTimeProvider {
+ now: AtomicI64,
+ elapsed: AtomicI64,
+ }
+
+ impl TimeProvider for MockTimeProvider {
+ fn current_time_millis(&self) -> i64 {
+ self.now.load(Ordering::Relaxed)
+ }
+
+ fn elapsed_since(&self, _current_millis: i64) -> i64 {
+ self.elapsed.load(Ordering::Relaxed)
+ }
+
+ fn wait_duration(&self, _duration: Duration) -> Duration {
+ Duration::from_millis(20)
+ }
+ }
+
+ impl MockTimeProvider {
+ pub(crate) fn new(now: i64) -> Self {
+ Self {
+ now: AtomicI64::new(now),
+ elapsed: AtomicI64::new(0),
+ }
+ }
+
+ pub(crate) fn set_now(&self, now: i64) {
+ self.now.store(now, Ordering::Relaxed);
+ }
+
+ pub(crate) fn set_elapsed(&self, elapsed: i64) {
+ self.elapsed.store(elapsed, Ordering::Relaxed);
+ }
+ }
+}
diff --git a/src/mito2/src/worker.rs b/src/mito2/src/worker.rs
index 8efcf47ab6..1aeaa59c5c 100644
--- a/src/mito2/src/worker.rs
+++ b/src/mito2/src/worker.rs
@@ -603,13 +603,17 @@ impl WorkerStarter {
.zip(series_index_task_state.clone())
.map(|(store, state)| {
let (purger, purge_receiver) = series_index_channel(store.clone());
- series_index_purger = Some(purger);
+ series_index_purger = Some(purger.clone());
spawn_series_index_tasks(
self.id,
store,
+ regions.clone(),
state,
+ self.config.experimental_series_index_bucket_width,
+ purger,
purge_receiver,
self.config.experimental_series_index_maintenance_interval,
+ self.time_provider.clone(),
)
});
let now = self.time_provider.current_time_millis();
@@ -636,6 +640,7 @@ impl WorkerStarter {
self.index_build_job_pool,
self.config.max_background_index_builds,
),
+ series_index_task_state: series_index_task_state.clone(),
series_index_store: self.series_index_store,
series_index_purger,
flush_scheduler: FlushScheduler::new(self.flush_job_pool),
@@ -700,9 +705,9 @@ pub(crate) struct RegionWorker {
sender: Sender,
/// Handle to the worker thread.
handle: Mutex