From e6472fd12a744efd911e290f04e4f49900bf3b6f Mon Sep 17 00:00:00 2001 From: discord9 Date: Thu, 9 Jul 2026 18:09:37 +0800 Subject: [PATCH] fix: pause GC during maintenance mode (#8450) Skip scheduled meta GC while cluster maintenance mode is enabled and reject manual GC requests explicitly instead of returning an empty success report. Also increase mito GC's default lingering time to 1h and update generated config docs and config API expectations. Signed-off-by: discord9 --- config/config.md | 2 +- config/datanode.example.toml | 2 +- src/meta-srv/src/error.rs | 7 + src/meta-srv/src/gc/mock.rs | 1 + src/meta-srv/src/gc/mock/basic.rs | 1 + src/meta-srv/src/gc/mock/candidate_select.rs | 4 + src/meta-srv/src/gc/mock/concurrent.rs | 5 + src/meta-srv/src/gc/mock/config.rs | 3 + src/meta-srv/src/gc/mock/err_handle.rs | 3 + src/meta-srv/src/gc/mock/full_list.rs | 3 + src/meta-srv/src/gc/mock/integration.rs | 2 + src/meta-srv/src/gc/mock/misc.rs | 2 + src/meta-srv/src/gc/scheduler.rs | 162 ++++++++++++++++++- src/meta-srv/src/metasrv/builder.rs | 1 + src/mito2/src/gc.rs | 10 +- tests-integration/tests/http.rs | 2 +- 16 files changed, 204 insertions(+), 6 deletions(-) diff --git a/config/config.md b/config/config.md index 3a0c1729af5..b4fb40c4778 100644 --- a/config/config.md +++ b/config/config.md @@ -608,7 +608,7 @@ | `region_engine.mito.bloom_filter_index.mem_threshold_on_create` | String | `auto` | Memory threshold for the index creation.
- `auto`: automatically determine the threshold based on the system memory size (default)
- `unlimited`: no memory limit
- `[size]` e.g. `64MB`: fixed memory threshold | | `region_engine.mito.gc` | -- | -- | -- | | `region_engine.mito.gc.enable` | Bool | `false` | Whether GC is enabled. Need to be the same with metasrv's `gc.enable` or unexpected behavior will occur | -| `region_engine.mito.gc.lingering_time` | String | `1m` | Lingering time before deleting files.
Should be long enough to allow long running queries to finish.
If set to None, then unused files will be deleted immediately. | +| `region_engine.mito.gc.lingering_time` | String | `1h` | Lingering time before deleting files.
Should be long enough to allow long running queries to finish.
If set to None, then unused files will be deleted immediately. | | `region_engine.mito.gc.unknown_file_lingering_time` | String | `1d` | Lingering time before deleting unknown files (files with undetermined expel time).
Only applies during full file listing GC.
This uses the object's last-modified timestamp as a heuristic (strict less-than comparison);
do not configure this value too small in production to avoid deleting pre-manifest files
from in-progress compaction or flush.
For active/open regions, an unknown file is deleted only if its object last-modified time exceeds this TTL.
If the object store does not provide a last-modified timestamp, the file is conservatively kept.
For dropped regions, unknown files are deleted immediately. | | `region_engine.file` | -- | -- | Enable the file engine. | | `region_engine.metric` | -- | -- | Metric engine options. | diff --git a/config/datanode.example.toml b/config/datanode.example.toml index 0c0fdbee8f7..678af443768 100644 --- a/config/datanode.example.toml +++ b/config/datanode.example.toml @@ -678,7 +678,7 @@ enable = false ## Lingering time before deleting files. ## Should be long enough to allow long running queries to finish. ## If set to None, then unused files will be deleted immediately. -lingering_time = "1m" +lingering_time = "1h" ## Lingering time before deleting unknown files (files with undetermined expel time). ## Only applies during full file listing GC. diff --git a/src/meta-srv/src/error.rs b/src/meta-srv/src/error.rs index 003bee94d3e..aa689080cd3 100644 --- a/src/meta-srv/src/error.rs +++ b/src/meta-srv/src/error.rs @@ -426,6 +426,12 @@ pub enum Error { location: Location, }, + #[snafu(display("Manual GC is rejected because maintenance mode is enabled"))] + ManualGcRejectedByMaintenanceMode { + #[snafu(implicit)] + location: Location, + }, + #[cfg(feature = "mysql_kvbackend")] #[snafu(display("Failed to parse mysql url: {}", mysql_url))] ParseMySqlUrl { @@ -1194,6 +1200,7 @@ impl ErrorExt for Error { | Error::ParseAddr { .. } | Error::UnsupportedSelectorType { .. } | Error::InvalidArguments { .. } + | Error::ManualGcRejectedByMaintenanceMode { .. } | Error::ProcedureNotFound { .. } | Error::TooManyPartitions { .. } | Error::TomlFormat { .. } diff --git a/src/meta-srv/src/gc/mock.rs b/src/meta-srv/src/gc/mock.rs index ac257576a41..0686d9a4653 100644 --- a/src/meta-srv/src/gc/mock.rs +++ b/src/meta-srv/src/gc/mock.rs @@ -344,6 +344,7 @@ impl TestEnv { let scheduler = GcScheduler { ctx: ctx.clone(), + runtime_switch_manager: crate::gc::scheduler::new_test_runtime_switch_manager(), receiver: rx, config, region_gc_tracker: Arc::new(tokio::sync::Mutex::new(HashMap::new())), diff --git a/src/meta-srv/src/gc/mock/basic.rs b/src/meta-srv/src/gc/mock/basic.rs index fb395d899e8..87278abcbf9 100644 --- a/src/meta-srv/src/gc/mock/basic.rs +++ b/src/meta-srv/src/gc/mock/basic.rs @@ -151,6 +151,7 @@ async fn test_handle_tick() { let scheduler = GcScheduler { ctx: ctx.clone(), + runtime_switch_manager: crate::gc::scheduler::new_test_runtime_switch_manager(), receiver: GcScheduler::channel().1, config: GcSchedulerOptions::default(), region_gc_tracker: Arc::new(tokio::sync::Mutex::new(HashMap::new())), diff --git a/src/meta-srv/src/gc/mock/candidate_select.rs b/src/meta-srv/src/gc/mock/candidate_select.rs index 73da83802a9..9aca46bd516 100644 --- a/src/meta-srv/src/gc/mock/candidate_select.rs +++ b/src/meta-srv/src/gc/mock/candidate_select.rs @@ -70,6 +70,7 @@ async fn test_gc_candidate_filtering_by_role() { let scheduler = GcScheduler { ctx: ctx.clone(), + runtime_switch_manager: crate::gc::scheduler::new_test_runtime_switch_manager(), receiver: GcScheduler::channel().1, config: GcSchedulerOptions::default(), region_gc_tracker: Arc::new(tokio::sync::Mutex::new(HashMap::new())), @@ -149,6 +150,7 @@ async fn test_gc_candidate_size_threshold() { let scheduler = GcScheduler { ctx: ctx.clone(), + runtime_switch_manager: crate::gc::scheduler::new_test_runtime_switch_manager(), receiver: GcScheduler::channel().1, config, region_gc_tracker: Arc::new(tokio::sync::Mutex::new(HashMap::new())), @@ -241,6 +243,7 @@ async fn test_gc_candidate_scoring() { let scheduler = GcScheduler { ctx: ctx.clone(), + runtime_switch_manager: crate::gc::scheduler::new_test_runtime_switch_manager(), receiver: GcScheduler::channel().1, config, region_gc_tracker: Arc::new(tokio::sync::Mutex::new(HashMap::new())), @@ -327,6 +330,7 @@ async fn test_gc_candidate_regions_per_table_threshold() { let scheduler = GcScheduler { ctx: ctx.clone(), + runtime_switch_manager: crate::gc::scheduler::new_test_runtime_switch_manager(), receiver: GcScheduler::channel().1, config, region_gc_tracker: Arc::new(tokio::sync::Mutex::new(HashMap::new())), diff --git a/src/meta-srv/src/gc/mock/concurrent.rs b/src/meta-srv/src/gc/mock/concurrent.rs index 2554c7046be..ddaf8128431 100644 --- a/src/meta-srv/src/gc/mock/concurrent.rs +++ b/src/meta-srv/src/gc/mock/concurrent.rs @@ -77,6 +77,7 @@ async fn test_concurrent_table_processing_limits() { let scheduler = GcScheduler { ctx: ctx.clone(), + runtime_switch_manager: crate::gc::scheduler::new_test_runtime_switch_manager(), receiver: GcScheduler::channel().1, config, region_gc_tracker: Arc::new(tokio::sync::Mutex::new(HashMap::new())), @@ -159,6 +160,7 @@ async fn test_datanode_processes_tables_with_partial_gc_failures() { let scheduler = GcScheduler { ctx: ctx.clone(), + runtime_switch_manager: crate::gc::scheduler::new_test_runtime_switch_manager(), receiver: GcScheduler::channel().1, config: GcSchedulerOptions::default(), region_gc_tracker: Arc::new(tokio::sync::Mutex::new(HashMap::new())), @@ -269,6 +271,7 @@ async fn test_region_gc_concurrency_limit() { let scheduler = GcScheduler { ctx: ctx.clone(), + runtime_switch_manager: crate::gc::scheduler::new_test_runtime_switch_manager(), receiver: GcScheduler::channel().1, config, region_gc_tracker: Arc::new(tokio::sync::Mutex::new(HashMap::new())), @@ -381,6 +384,7 @@ async fn test_region_gc_concurrency_with_partial_failures() { let scheduler = GcScheduler { ctx: ctx.clone(), + runtime_switch_manager: crate::gc::scheduler::new_test_runtime_switch_manager(), receiver: GcScheduler::channel().1, config, region_gc_tracker: Arc::new(tokio::sync::Mutex::new(HashMap::new())), @@ -520,6 +524,7 @@ async fn test_region_gc_concurrency_with_retryable_errors() { let scheduler = GcScheduler { ctx: ctx.clone(), + runtime_switch_manager: crate::gc::scheduler::new_test_runtime_switch_manager(), receiver: GcScheduler::channel().1, config, region_gc_tracker: Arc::new(tokio::sync::Mutex::new(HashMap::new())), diff --git a/src/meta-srv/src/gc/mock/config.rs b/src/meta-srv/src/gc/mock/config.rs index f4ec9be9480..f3a83d5ab1a 100644 --- a/src/meta-srv/src/gc/mock/config.rs +++ b/src/meta-srv/src/gc/mock/config.rs @@ -59,6 +59,7 @@ async fn test_different_gc_weights() { let scheduler1 = GcScheduler { ctx: ctx.clone(), + runtime_switch_manager: crate::gc::scheduler::new_test_runtime_switch_manager(), receiver: GcScheduler::channel().1, config: config1, region_gc_tracker: Arc::new(tokio::sync::Mutex::new(HashMap::new())), @@ -83,6 +84,7 @@ async fn test_different_gc_weights() { let scheduler2 = GcScheduler { ctx: ctx.clone(), + runtime_switch_manager: crate::gc::scheduler::new_test_runtime_switch_manager(), receiver: GcScheduler::channel().1, config: config2, region_gc_tracker: Arc::new(tokio::sync::Mutex::new(HashMap::new())), @@ -161,6 +163,7 @@ async fn test_regions_per_table_threshold() { let scheduler = GcScheduler { ctx: ctx.clone(), + runtime_switch_manager: crate::gc::scheduler::new_test_runtime_switch_manager(), receiver: GcScheduler::channel().1, config, region_gc_tracker: Arc::new(tokio::sync::Mutex::new(HashMap::new())), diff --git a/src/meta-srv/src/gc/mock/err_handle.rs b/src/meta-srv/src/gc/mock/err_handle.rs index 0dd9b3d115a..aa07e21f036 100644 --- a/src/meta-srv/src/gc/mock/err_handle.rs +++ b/src/meta-srv/src/gc/mock/err_handle.rs @@ -83,6 +83,7 @@ async fn test_gc_regions_failure_handling() { let scheduler = GcScheduler { ctx: ctx.clone(), + runtime_switch_manager: crate::gc::scheduler::new_test_runtime_switch_manager(), receiver: GcScheduler::channel().1, config: GcSchedulerOptions::default(), region_gc_tracker: Arc::new(tokio::sync::Mutex::new(HashMap::new())), @@ -178,6 +179,7 @@ async fn test_get_file_references_failure() { let scheduler = GcScheduler { ctx: ctx.clone(), + runtime_switch_manager: crate::gc::scheduler::new_test_runtime_switch_manager(), receiver: GcScheduler::channel().1, config: GcSchedulerOptions { retry_backoff_duration: Duration::from_millis(10), // shorten for test @@ -256,6 +258,7 @@ async fn test_get_table_route_failure() { let scheduler = GcScheduler { ctx: ctx.clone(), + runtime_switch_manager: crate::gc::scheduler::new_test_runtime_switch_manager(), receiver: GcScheduler::channel().1, config: GcSchedulerOptions::default(), region_gc_tracker: Arc::new(tokio::sync::Mutex::new(HashMap::new())), diff --git a/src/meta-srv/src/gc/mock/full_list.rs b/src/meta-srv/src/gc/mock/full_list.rs index 6b188c0869c..847c0e73e40 100644 --- a/src/meta-srv/src/gc/mock/full_list.rs +++ b/src/meta-srv/src/gc/mock/full_list.rs @@ -68,6 +68,7 @@ async fn test_full_file_listing_first_time_gc() { let scheduler = GcScheduler { ctx: ctx.clone(), + runtime_switch_manager: crate::gc::scheduler::new_test_runtime_switch_manager(), receiver: GcScheduler::channel().1, config, region_gc_tracker: Arc::new(tokio::sync::Mutex::new(HashMap::new())), @@ -140,6 +141,7 @@ async fn test_full_file_listing_interval_enforcement() { let scheduler = GcScheduler { ctx: ctx.clone(), + runtime_switch_manager: crate::gc::scheduler::new_test_runtime_switch_manager(), receiver: GcScheduler::channel().1, config, region_gc_tracker: Arc::new(tokio::sync::Mutex::new(HashMap::new())), @@ -240,6 +242,7 @@ async fn test_full_file_listing_no_interval_passed() { let scheduler = GcScheduler { ctx: ctx.clone(), + runtime_switch_manager: crate::gc::scheduler::new_test_runtime_switch_manager(), receiver: GcScheduler::channel().1, config, region_gc_tracker: Arc::new(tokio::sync::Mutex::new(HashMap::new())), diff --git a/src/meta-srv/src/gc/mock/integration.rs b/src/meta-srv/src/gc/mock/integration.rs index 8aa8f977ad5..76fdbac35f8 100644 --- a/src/meta-srv/src/gc/mock/integration.rs +++ b/src/meta-srv/src/gc/mock/integration.rs @@ -76,6 +76,7 @@ async fn test_full_gc_workflow() { let scheduler = GcScheduler { ctx: ctx.clone(), + runtime_switch_manager: crate::gc::scheduler::new_test_runtime_switch_manager(), receiver: GcScheduler::channel().1, config: GcSchedulerOptions::default(), region_gc_tracker: Arc::new(tokio::sync::Mutex::new(HashMap::new())), @@ -205,6 +206,7 @@ async fn test_tracker_cleanup() { let scheduler = GcScheduler { ctx: ctx.clone(), + runtime_switch_manager: crate::gc::scheduler::new_test_runtime_switch_manager(), receiver: GcScheduler::channel().1, config, region_gc_tracker: Arc::new(tokio::sync::Mutex::new(old_region_gc_tracker)), diff --git a/src/meta-srv/src/gc/mock/misc.rs b/src/meta-srv/src/gc/mock/misc.rs index 76d14136e44..c8b2b43155e 100644 --- a/src/meta-srv/src/gc/mock/misc.rs +++ b/src/meta-srv/src/gc/mock/misc.rs @@ -51,6 +51,7 @@ async fn test_empty_file_refs_manifest() { let scheduler = GcScheduler { ctx: ctx.clone(), + runtime_switch_manager: crate::gc::scheduler::new_test_runtime_switch_manager(), receiver: GcScheduler::channel().1, config: GcSchedulerOptions::default(), region_gc_tracker: Arc::new(tokio::sync::Mutex::new(HashMap::new())), @@ -139,6 +140,7 @@ async fn test_multiple_regions_per_table() { let scheduler = GcScheduler { ctx: ctx.clone(), + runtime_switch_manager: crate::gc::scheduler::new_test_runtime_switch_manager(), receiver: GcScheduler::channel().1, config: GcSchedulerOptions::default(), region_gc_tracker: Arc::new(tokio::sync::Mutex::new(HashMap::new())), diff --git a/src/meta-srv/src/gc/scheduler.rs b/src/meta-srv/src/gc/scheduler.rs index 772c7e7cb69..4f732f52bc0 100644 --- a/src/meta-srv/src/gc/scheduler.rs +++ b/src/meta-srv/src/gc/scheduler.rs @@ -18,16 +18,18 @@ use std::time::{Duration, Instant}; use common_meta::DatanodeId; use common_meta::key::TableMetadataManagerRef; +use common_meta::key::runtime_switch::RuntimeSwitchManagerRef; use common_procedure::ProcedureManagerRef; use common_telemetry::tracing::Instrument as _; use common_telemetry::{error, info}; +use snafu::ResultExt; use store_api::storage::{GcReport, RegionId}; use tokio::sync::mpsc::{Receiver, Sender}; use tokio::sync::{Mutex, oneshot}; use crate::cluster::MetaPeerClientRef; use crate::define_ticker; -use crate::error::{Error, Result}; +use crate::error::{self, Error, Result}; use crate::gc::Region2Peers; use crate::gc::ctx::{DefaultGcSchedulerCtx, SchedulerCtx}; use crate::gc::dropped::DroppedRegionCollector; @@ -115,6 +117,8 @@ define_ticker!( /// [`GcScheduler`] is used to periodically trigger garbage collection on datanodes. pub struct GcScheduler { pub(crate) ctx: Arc, + /// Runtime switch manager to check maintenance mode. + pub(crate) runtime_switch_manager: RuntimeSwitchManagerRef, /// The receiver of events. pub(crate) receiver: Receiver, /// GC configuration. @@ -133,6 +137,7 @@ impl GcScheduler { meta_peer_client: MetaPeerClientRef, mailbox: MailboxRef, server_addr: String, + runtime_switch_manager: RuntimeSwitchManagerRef, config: GcSchedulerOptions, ) -> Result<(Self, GcTicker)> { // Validate configuration before creating the scheduler @@ -148,6 +153,7 @@ impl GcScheduler { mailbox, server_addr, )?), + runtime_switch_manager, receiver: rx, config, region_gc_tracker: Arc::new(Mutex::new(HashMap::new())), @@ -192,7 +198,11 @@ impl GcScheduler { .instrument(span) .await; if let Err(e) = &result { - error!(e; "Failed to handle manual gc"); + if matches!(e, Error::ManualGcRejectedByMaintenanceMode { .. }) { + info!("Rejected manual gc request: {}", e); + } else { + error!(e; "Failed to handle manual gc"); + } } let _ = sender.send(result); } @@ -204,6 +214,10 @@ impl GcScheduler { METRIC_META_GC_SCHEDULER_CYCLES_TOTAL.inc(); let _timer = METRIC_META_GC_SCHEDULER_DURATION_SECONDS.start_timer(); info!("Start to trigger gc"); + if self.is_maintenance_mode_enabled().await? { + info!("Skip gc trigger because maintenance mode is enabled"); + return Ok(GcJobReport::default()); + } let span = common_telemetry::tracing::info_span!("meta_gc_handle_tick"); let report = self.trigger_gc().instrument(span).await?; @@ -227,6 +241,11 @@ impl GcScheduler { ) -> Result { info!("Start to handle manual gc request"); + if self.is_maintenance_mode_enabled().await? { + info!("Skip manual gc request because maintenance mode is enabled"); + return error::ManualGcRejectedByMaintenanceModeSnafu {}.fail(); + } + // No specific regions, use default tick behavior let Some(regions) = region_ids else { let report = self.trigger_gc().await?; @@ -294,11 +313,26 @@ impl GcScheduler { info!("Finished manual gc request"); Ok(report) } + + pub(crate) async fn is_maintenance_mode_enabled(&self) -> Result { + self.runtime_switch_manager + .maintenance_mode() + .await + .context(error::RuntimeSwitchManagerSnafu) + } +} + +#[cfg(test)] +pub(crate) fn new_test_runtime_switch_manager() -> RuntimeSwitchManagerRef { + Arc::new(common_meta::key::runtime_switch::RuntimeSwitchManager::new( + Arc::new(common_meta::kv_backend::memory::MemoryKvBackend::new()), + )) } #[cfg(test)] mod tests { use std::collections::HashMap; + use std::sync::atomic::{AtomicUsize, Ordering}; use std::time::Duration; use common_meta::datanode::RegionStat; @@ -309,6 +343,72 @@ mod tests { use super::*; + #[derive(Default)] + struct CountingSchedulerCtx { + get_table_to_region_stats_calls: AtomicUsize, + get_table_reparts_calls: AtomicUsize, + gc_regions_calls: AtomicUsize, + } + + impl CountingSchedulerCtx { + fn assert_no_scheduler_work(&self) { + assert_eq!( + 0, + self.get_table_to_region_stats_calls.load(Ordering::Relaxed), + "get_table_to_region_stats should not be called" + ); + assert_eq!( + 0, + self.get_table_reparts_calls.load(Ordering::Relaxed), + "get_table_reparts should not be called" + ); + assert_eq!( + 0, + self.gc_regions_calls.load(Ordering::Relaxed), + "gc_regions should not be called" + ); + } + } + + #[async_trait::async_trait] + impl SchedulerCtx for CountingSchedulerCtx { + async fn get_table_to_region_stats(&self) -> Result>> { + self.get_table_to_region_stats_calls + .fetch_add(1, Ordering::Relaxed); + panic!("get_table_to_region_stats should not be called in maintenance mode") + } + + async fn get_table_reparts(&self) -> Result> { + self.get_table_reparts_calls.fetch_add(1, Ordering::Relaxed); + panic!("get_table_reparts should not be called in maintenance mode") + } + + async fn get_table_route( + &self, + _table_id: TableId, + ) -> Result<(TableId, PhysicalTableRouteValue)> { + unreachable!("get_table_route should not be called in this test") + } + + async fn batch_get_table_route( + &self, + _table_ids: &[TableId], + ) -> Result> { + unreachable!("batch_get_table_route should not be called in this test") + } + + async fn gc_regions( + &self, + _region_ids: &[RegionId], + _full_file_listing: bool, + _timeout: Duration, + _region_routes_override: Region2Peers, + ) -> Result { + self.gc_regions_calls.fetch_add(1, Ordering::Relaxed); + panic!("gc_regions should not be called in maintenance mode") + } + } + struct ErrorMockSchedulerCtx; #[async_trait::async_trait] @@ -356,6 +456,7 @@ mod tests { let scheduler = GcScheduler { ctx: Arc::new(ErrorMockSchedulerCtx), + runtime_switch_manager: new_test_runtime_switch_manager(), receiver: rx, config: GcSchedulerOptions::default(), region_gc_tracker: Arc::new(Mutex::new(HashMap::new())), @@ -372,4 +473,61 @@ mod tests { assert!(result.is_err()); } + + #[tokio::test] + async fn test_maintenance_mode_skips_manual_gc() { + let (tx, rx) = GcScheduler::channel(); + drop(tx); + let runtime_switch_manager = new_test_runtime_switch_manager(); + runtime_switch_manager.set_maintenance_mode().await.unwrap(); + + let ctx = Arc::new(CountingSchedulerCtx::default()); + let scheduler = GcScheduler { + ctx: ctx.clone(), + runtime_switch_manager, + receiver: rx, + config: GcSchedulerOptions::default(), + region_gc_tracker: Arc::new(Mutex::new(HashMap::new())), + last_tracker_cleanup: Arc::new(Mutex::new(Instant::now())), + }; + + let result = scheduler + .handle_manual_gc( + Some(vec![RegionId::new(1, 0)]), + Some(false), + Some(Duration::from_secs(1)), + ) + .await; + + let err = result.unwrap_err(); + assert!(matches!( + err, + error::Error::ManualGcRejectedByMaintenanceMode { .. } + )); + assert!(err.to_string().contains("maintenance mode is enabled")); + ctx.assert_no_scheduler_work(); + } + + #[tokio::test] + async fn test_maintenance_mode_skips_tick_gc() { + let (tx, rx) = GcScheduler::channel(); + drop(tx); + let runtime_switch_manager = new_test_runtime_switch_manager(); + runtime_switch_manager.set_maintenance_mode().await.unwrap(); + + let ctx = Arc::new(CountingSchedulerCtx::default()); + let scheduler = GcScheduler { + ctx: ctx.clone(), + runtime_switch_manager, + receiver: rx, + config: GcSchedulerOptions::default(), + region_gc_tracker: Arc::new(Mutex::new(HashMap::new())), + last_tracker_cleanup: Arc::new(Mutex::new(Instant::now())), + }; + + let result = scheduler.handle_tick().await; + + assert!(result.is_ok()); + ctx.assert_no_scheduler_work(); + } } diff --git a/src/meta-srv/src/metasrv/builder.rs b/src/meta-srv/src/metasrv/builder.rs index 69bfd3f34bb..90ce54142f1 100644 --- a/src/meta-srv/src/metasrv/builder.rs +++ b/src/meta-srv/src/metasrv/builder.rs @@ -516,6 +516,7 @@ impl MetasrvBuilder { meta_peer_client.clone(), mailbox.clone(), options.grpc.server_addr.clone(), + runtime_switch_manager.clone(), options.gc.clone(), )?; gc_scheduler.try_start()?; diff --git a/src/mito2/src/gc.rs b/src/mito2/src/gc.rs index 8f82c4159e1..61ba362a7bd 100644 --- a/src/mito2/src/gc.rs +++ b/src/mito2/src/gc.rs @@ -173,7 +173,7 @@ impl Default for GcConfig { Self { enable: false, // expect long running queries to be finished(or at least be able to notify it's using a deleted file) within a reasonable time - lingering_time: Some(Duration::from_secs(60)), + lingering_time: Some(Duration::from_secs(60 * 60)), // 1 day, for unknown expel time, which is when this file get removed from manifest. // Only applies to full-listing GC for active/open regions. A long default avoids // accidentally deleting pre-manifest files (e.g. compaction/flush still in progress). @@ -1135,4 +1135,12 @@ mod tests { ] ); } + + #[test] + fn test_gc_config_default_lingering_time() { + assert_eq!( + GcConfig::default().lingering_time, + Some(Duration::from_secs(60 * 60)) + ); + } } diff --git a/tests-integration/tests/http.rs b/tests-integration/tests/http.rs index 8b6dcd27aff..c6366f48efe 100644 --- a/tests-integration/tests/http.rs +++ b/tests-integration/tests/http.rs @@ -1982,7 +1982,7 @@ apply_on_query = "auto" mem_threshold_on_create = "auto" {vector_index_config}[region_engine.mito.gc] enable = false -lingering_time = "1m" +lingering_time = "1h" unknown_file_lingering_time = "1day" max_concurrent_lister_per_gc_job = 32 max_concurrent_gc_job = 4