From b513dbaf4a3b21add12c031e5abc965251f80423 Mon Sep 17 00:00:00 2001 From: discord9 Date: Tue, 14 Jul 2026 20:48:13 +0800 Subject: [PATCH] fix: reject datanode startup on GC config mismatch (#8509) * fix: reject datanode gc config mismatch Signed-off-by: discord9 * refactor: minimize datanode gc startup check Signed-off-by: discord9 * chore: update greptime-proto revision Signed-off-by: discord9 * chore: use merged greptime-proto revision Signed-off-by: discord9 --------- Signed-off-by: discord9 --- Cargo.lock | 2 +- Cargo.toml | 2 +- src/datanode/src/error.rs | 13 +++++++++++++ src/datanode/src/heartbeat.rs | 17 +++++++++++++++++ src/frontend/src/frontend.rs | 1 + src/meta-client/src/client/heartbeat.rs | 8 +++++--- src/meta-srv/src/handler.rs | 7 ++++--- src/meta-srv/src/handler/test_utils.rs | 1 + src/meta-srv/src/metasrv.rs | 3 +++ 9 files changed, 46 insertions(+), 8 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 7a9ec59760..cb6cb42c2e 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -5938,7 +5938,7 @@ dependencies = [ [[package]] name = "greptime-proto" version = "0.1.0" -source = "git+https://github.com/GreptimeTeam/greptime-proto.git?rev=6fc4d88851c5ad16846698e4dff6e74c1f7a48ba#6fc4d88851c5ad16846698e4dff6e74c1f7a48ba" +source = "git+https://github.com/GreptimeTeam/greptime-proto.git?rev=c6f92ee9e12ee7bd5ee3c0956b1c9caa39be3f76#c6f92ee9e12ee7bd5ee3c0956b1c9caa39be3f76" dependencies = [ "prost 0.14.1", "prost-types 0.14.1", diff --git a/Cargo.toml b/Cargo.toml index 5fecc43fa5..31093a556c 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -158,7 +158,7 @@ fs2 = "0.4" fst = "0.4.7" futures = "0.3" futures-util = "0.3" -greptime-proto = { git = "https://github.com/GreptimeTeam/greptime-proto.git", rev = "6fc4d88851c5ad16846698e4dff6e74c1f7a48ba" } +greptime-proto = { git = "https://github.com/GreptimeTeam/greptime-proto.git", rev = "c6f92ee9e12ee7bd5ee3c0956b1c9caa39be3f76" } hex = "0.4" http = "1" humantime = "2.1" diff --git a/src/datanode/src/error.rs b/src/datanode/src/error.rs index 0e2e1da959..6c176547ed 100644 --- a/src/datanode/src/error.rs +++ b/src/datanode/src/error.rs @@ -287,6 +287,18 @@ pub enum Error { location: Location, }, + #[snafu(display( + "GC configuration mismatch: metasrv.gc.enable={}, datanode.region_engine.mito.gc.enable={}", + metasrv_gc_enabled, + datanode_gc_enabled, + ))] + GcConfigMismatch { + metasrv_gc_enabled: bool, + datanode_gc_enabled: bool, + #[snafu(implicit)] + location: Location, + }, + #[snafu(display("Unsupported output type, expected: {}", expected))] UnsupportedOutput { expected: String, @@ -456,6 +468,7 @@ impl ErrorExt for Error { | ColumnNoneDefaultValue { .. } | MissingRequiredField { .. } | RegionEngineNotFound { .. } + | GcConfigMismatch { .. } | ParseAddr { .. } | TomlFormat { .. } | BuildDatanode { .. } => StatusCode::InvalidArguments, diff --git a/src/datanode/src/heartbeat.rs b/src/datanode/src/heartbeat.rs index bfdaf32a1c..0e0c51624e 100644 --- a/src/datanode/src/heartbeat.rs +++ b/src/datanode/src/heartbeat.rs @@ -121,6 +121,7 @@ impl HeartbeatTask { pub async fn create_streams( meta_client: &MetaClient, + local_gc_enabled: bool, running: Arc, handler_executor: HeartbeatResponseHandlerExecutorRef, mailbox: MailboxRef, @@ -129,6 +130,13 @@ impl HeartbeatTask { ) -> Result<(HeartbeatSender, HeartbeatConfig)> { let client_id = meta_client.id(); let (tx, mut rx, config) = meta_client.heartbeat().await.context(MetaClientInitSnafu)?; + if config.gc_enabled != local_gc_enabled { + return error::GcConfigMismatchSnafu { + metasrv_gc_enabled: config.gc_enabled, + datanode_gc_enabled: local_gc_enabled, + } + .fail(); + } let mut last_received_lease = Instant::now(); @@ -220,9 +228,17 @@ impl HeartbeatTask { let mailbox = Arc::new(HeartbeatMailbox::new(outgoing_tx)); let quit_signal = Arc::new(Notify::new()); + let local_gc_enabled = self + .region_server + .mito_engine() + .context(RegionEngineNotFoundSnafu { name: "mito" })? + .mito_config() + .gc + .enable; let (mut tx, config) = Self::create_streams( &meta_client, + local_gc_enabled, running.clone(), handler_executor.clone(), mailbox.clone(), @@ -368,6 +384,7 @@ impl HeartbeatTask { error!(e; "Failed to send heartbeat to metasrv"); match Self::create_streams( &meta_client, + local_gc_enabled, running.clone(), handler_executor.clone(), mailbox.clone(), diff --git a/src/frontend/src/frontend.rs b/src/frontend/src/frontend.rs index 918185cb8f..c7f7e71d34 100644 --- a/src/frontend/src/frontend.rs +++ b/src/frontend/src/frontend.rs @@ -237,6 +237,7 @@ mod tests { is_handshake.then_some(api::v1::meta::HeartbeatConfig { heartbeat_interval_ms, retry_interval_ms: heartbeat_interval_ms, + gc_enabled: false, }); is_handshake = false; let response = HeartbeatResponse { diff --git a/src/meta-client/src/client/heartbeat.rs b/src/meta-client/src/client/heartbeat.rs index c31369dac7..add7bd54de 100644 --- a/src/meta-client/src/client/heartbeat.rs +++ b/src/meta-client/src/client/heartbeat.rs @@ -39,6 +39,7 @@ use crate::error::{InvalidResponseHeaderSnafu, Result}; pub struct HeartbeatConfig { pub interval: Duration, pub retry_interval: Duration, + pub gc_enabled: bool, } impl Default for HeartbeatConfig { @@ -46,6 +47,7 @@ impl Default for HeartbeatConfig { Self { interval: BASE_HEARTBEAT_INTERVAL, retry_interval: BASE_HEARTBEAT_INTERVAL, + gc_enabled: false, } } } @@ -54,8 +56,8 @@ impl fmt::Display for HeartbeatConfig { fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { write!( f, - "interval={:?}, retry={:?}", - self.interval, self.retry_interval + "interval={:?}, retry={:?}, gc_enabled={}", + self.interval, self.retry_interval, self.gc_enabled ) } } @@ -68,6 +70,7 @@ impl HeartbeatConfig { Self { interval: Duration::from_millis(cfg.heartbeat_interval_ms), retry_interval: Duration::from_millis(cfg.retry_interval_ms), + gc_enabled: cfg.gc_enabled, } } else { let fallback = Self::default(); @@ -259,7 +262,6 @@ impl Inner { .map_err(error::Error::from)? .context(error::CreateHeartbeatStreamSnafu)?; - // Extract heartbeat configuration from handshake response let config = HeartbeatConfig::from_response(&res); info!( diff --git a/src/meta-srv/src/handler.rs b/src/meta-srv/src/handler.rs index 9cfd4e6079..72d5fcdab8 100644 --- a/src/meta-srv/src/handler.rs +++ b/src/meta-srv/src/handler.rs @@ -21,8 +21,8 @@ use std::time::{Duration, Instant}; use api::v1::meta::mailbox_message::Payload; use api::v1::meta::{ - HeartbeatRequest, HeartbeatResponse, MailboxMessage, PROTOCOL_VERSION, RegionLease, - ResponseHeader, Role, + HeartbeatConfig, HeartbeatRequest, HeartbeatResponse, MailboxMessage, PROTOCOL_VERSION, + RegionLease, ResponseHeader, Role, }; use check_leader_handler::CheckLeaderHandler; use collect_cluster_info_handler::{ @@ -387,7 +387,8 @@ impl HeartbeatHandlerGroup { // Populate heartbeat_config during handshake let heartbeat_config = if is_handshake { - let config = ctx.heartbeat_options_for(role).into(); + let mut config: HeartbeatConfig = ctx.heartbeat_options_for(role).into(); + config.gc_enabled = ctx.gc_enabled; info!( "Handshake with {:?} node, sending config: {:?}", diff --git a/src/meta-srv/src/handler/test_utils.rs b/src/meta-srv/src/handler/test_utils.rs index 7bfdfbee79..d8241c785e 100644 --- a/src/meta-srv/src/handler/test_utils.rs +++ b/src/meta-srv/src/handler/test_utils.rs @@ -92,6 +92,7 @@ impl TestEnv { leader_region_registry: self.leader_region_registry.clone(), topic_stats_registry: self.topic_stats_registry.clone(), heartbeat_interval: BASE_HEARTBEAT_INTERVAL, + gc_enabled: false, is_handshake: false, } } diff --git a/src/meta-srv/src/metasrv.rs b/src/meta-srv/src/metasrv.rs index f7f5bbf77d..7480b8cb50 100644 --- a/src/meta-srv/src/metasrv.rs +++ b/src/meta-srv/src/metasrv.rs @@ -175,6 +175,7 @@ impl From for HeartbeatConfig { Self { heartbeat_interval_ms: opts.interval.as_millis() as u64, retry_interval_ms: opts.retry_interval.as_millis() as u64, + gc_enabled: false, } } } @@ -438,6 +439,7 @@ pub struct Context { pub leader_region_registry: LeaderRegionRegistryRef, pub topic_stats_registry: TopicStatsRegistryRef, pub heartbeat_interval: Duration, + pub gc_enabled: bool, pub is_handshake: bool, } @@ -932,6 +934,7 @@ impl Metasrv { leader_region_registry, topic_stats_registry, heartbeat_interval: self.options().heartbeat_interval, + gc_enabled: self.options().gc.enable, is_handshake: false, } }