From 1d5b7bc0336da7c69f0422b58e82d96ba2d62ed8 Mon Sep 17 00:00:00 2001 From: evenyag Date: Wed, 9 Sep 2026 21:49:12 +0800 Subject: [PATCH] feat(mito2): reconcile series indexes in background Signed-off-by: evenyag --- config/config.md | 6 + config/datanode.example.toml | 11 + config/standalone.example.toml | 11 + src/mito2/AGENTS.md | 2 +- src/mito2/src/config.rs | 18 +- src/mito2/src/engine/compaction_test.rs | 2 +- src/mito2/src/engine/edit_region_test.rs | 2 +- src/mito2/src/engine/flush_test.rs | 41 +- src/mito2/src/metrics.rs | 13 + src/mito2/src/region.rs | 1 - src/mito2/src/series_index.rs | 23 +- src/mito2/src/series_index/maintenance.rs | 305 ++++++++++ src/mito2/src/series_index/purger.rs | 53 ++ src/mito2/src/series_index/task.rs | 137 ++++- src/mito2/src/series_index/tests.rs | 655 +++++++++++++++++++++- src/mito2/src/sst/file_purger.rs | 37 +- src/mito2/src/time_provider.rs | 46 ++ src/mito2/src/worker.rs | 14 +- src/mito2/src/worker/handle_catchup.rs | 3 + src/mito2/src/worker/handle_create.rs | 3 + src/mito2/src/worker/handle_open.rs | 4 + tests-integration/tests/http.rs | 1 + 22 files changed, 1302 insertions(+), 86 deletions(-) create mode 100644 src/mito2/src/series_index/maintenance.rs 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>>, - /// Handle to the series-index maintenance task. + /// Handle to the sequential series-index task. series_index_handle: Mutex>>, - /// Controls the worker-owned series-index maintenance task. + /// Controls the worker-owned series-index task. series_index_task_state: Option>, /// Whether to run the worker thread. running: Arc, @@ -935,6 +940,8 @@ struct RegionWorkerLoop { write_buffer_manager: WriteBufferManagerRef, /// Scheduler for index build task. index_build_scheduler: IndexBuildScheduler, + /// Controls the worker-owned series-index task. + series_index_task_state: Option>, /// Store for companion range indexes deleted by the region SST purger. series_index_store: Option, series_index_purger: Option, @@ -1625,6 +1632,7 @@ mod tests { .create_worker_group(MitoConfig { num_workers: 4, experimental_series_index_root: "series-index".to_string(), + experimental_series_index_maintenance_interval: Duration::from_secs(3600), ..Default::default() }) .await; diff --git a/src/mito2/src/worker/handle_catchup.rs b/src/mito2/src/worker/handle_catchup.rs index 5ef45f82f7..576ffed7ba 100644 --- a/src/mito2/src/worker/handle_catchup.rs +++ b/src/mito2/src/worker/handle_catchup.rs @@ -134,6 +134,9 @@ impl RegionWorkerLoop { .await?; debug_assert!(!reopened_region.is_writable()); self.regions.insert_region(reopened_region.clone()); + if let Some(state) = &self.series_index_task_state { + state.wake(); + } Ok(reopened_region) } diff --git a/src/mito2/src/worker/handle_create.rs b/src/mito2/src/worker/handle_create.rs index 3acb5da76d..2263cfa2ff 100644 --- a/src/mito2/src/worker/handle_create.rs +++ b/src/mito2/src/worker/handle_create.rs @@ -96,6 +96,9 @@ impl RegionWorkerLoop { // Insert the MitoRegion into the RegionMap. self.regions.insert_region(region); + if let Some(state) = &self.series_index_task_state { + state.wake(); + } Ok(0) } diff --git a/src/mito2/src/worker/handle_open.rs b/src/mito2/src/worker/handle_open.rs index 825a677b76..52b902ffe9 100644 --- a/src/mito2/src/worker/handle_open.rs +++ b/src/mito2/src/worker/handle_open.rs @@ -227,6 +227,7 @@ impl RegionWorkerLoop { let opening_regions = self.opening_regions.clone(); let region_count = self.region_count.clone(); let worker_id = self.id; + let series_index_task_state = self.series_index_task_state.clone(); opening_regions.insert_sender(region_id, sender); common_runtime::spawn_global(async move { match opener.open(&config, &wal).await { @@ -249,6 +250,9 @@ impl RegionWorkerLoop { // Insert the Region into the RegionMap. regions.insert_region(region); + if let Some(state) = &series_index_task_state { + state.wake(); + } let senders = opening_regions.remove_sender(region_id); for sender in senders { diff --git a/tests-integration/tests/http.rs b/tests-integration/tests/http.rs index 84d1f7d957..23b5129051 100644 --- a/tests-integration/tests/http.rs +++ b/tests-integration/tests/http.rs @@ -2383,6 +2383,7 @@ experimental_manifest_keep_removed_file_ttl = "1h" compress_manifest = false experimental_series_index_root = "" experimental_series_index_maintenance_interval = "5m" +experimental_series_index_bucket_width = "5days" experimental_compaction_memory_limit = "unlimited" experimental_compaction_on_exhausted = "wait" auto_flush_interval = "10m"