From 51ef87ed19b3ba98f600450c1b4b4ff90a7c67ed Mon Sep 17 00:00:00 2001 From: Henry Guo Date: Fri, 12 Jun 2026 07:05:52 +0800 Subject: [PATCH] fix(heal): recover renewed disk health checks (#3366) * fix(heal): recover renewed disk health checks * test(heal): cover replaced remote disk rebuild --------- Co-authored-by: Henry Guo Co-authored-by: houseme --- crates/e2e_test/src/common.rs | 53 +++++++++ .../src/heal_erasure_disk_rebuild_test.rs | 106 +++++++++++++++++- crates/ecstore/src/disk/disk_store.rs | 5 + crates/ecstore/src/disk/mod.rs | 8 ++ crates/ecstore/src/rpc/remote_disk.rs | 5 + crates/ecstore/src/set_disk/lock.rs | 49 +++++++- rustfs/src/server/layer.rs | 24 ++++ 7 files changed, 247 insertions(+), 3 deletions(-) diff --git a/crates/e2e_test/src/common.rs b/crates/e2e_test/src/common.rs index 1c3cef90e..e5a83c387 100644 --- a/crates/e2e_test/src/common.rs +++ b/crates/e2e_test/src/common.rs @@ -722,6 +722,13 @@ impl RustFSTestClusterEnvironment { self.extra_env.push((key.into(), value.into())); } + fn ensure_node_index(&self, node_idx: usize) -> Result<(), Box> { + if node_idx >= self.nodes.len() { + return Err(format!("node_idx {node_idx} is invalid").into()); + } + Ok(()) + } + /// Build the volumes argument string for RustFS binary (internal helper method). /// /// Concatenates the address and data directory of all cluster nodes into a single string @@ -781,6 +788,39 @@ impl RustFSTestClusterEnvironment { Ok(()) } + /// Start one node process using the cluster's existing volume layout. + pub async fn start_node(&mut self, node_idx: usize) -> Result<(), Box> { + self.ensure_node_index(node_idx)?; + if self.nodes[node_idx].process.is_some() { + return Err(format!("cluster node {node_idx} is already running").into()); + } + + let binary_path = rustfs_binary_path(); + let volumes_arg = self.build_volumes_arg(); + let node = &mut self.nodes[node_idx]; + info!("Starting cluster node {} on {}", node_idx, node.address); + + let mut command = Command::new(&binary_path); + command + .env("RUSTFS_VOLUMES", &volumes_arg) + .env("RUSTFS_ADDRESS", &node.address) + .env("RUSTFS_ACCESS_KEY", &self.access_key) + .env("RUSTFS_SECRET_KEY", &self.secret_key) + .env("RUSTFS_CONSOLE_ENABLE", "false") + .env("RUST_LOG", "rustfs=info,rustfs_notify=debug"); + + for (key, value) in &self.extra_env { + command.env(key, value); + } + + let process = command.current_dir(&node.data_dir).spawn()?; + node.process = Some(process); + + self.wait_for_node_ready(&self.nodes[node_idx].address, node_idx).await?; + self.wait_for_node_service_ready(node_idx).await?; + Ok(()) + } + /// Wait for a single cluster node's TCP port to become reachable (internal helper method). /// /// Attempts to establish a TCP connection to the node's address, retries up to 60 times @@ -911,6 +951,19 @@ impl RustFSTestClusterEnvironment { } } } + + /// Stop a single cluster node and wait for its process to exit. + pub fn stop_node(&mut self, node_idx: usize) -> Result<(), Box> { + self.ensure_node_index(node_idx)?; + let Some(mut process) = self.nodes[node_idx].process.take() else { + return Ok(()); + }; + + info!("Stopping cluster node {}", node_idx); + process.kill()?; + process.wait()?; + Ok(()) + } } impl Drop for RustFSTestClusterEnvironment { diff --git a/crates/e2e_test/src/heal_erasure_disk_rebuild_test.rs b/crates/e2e_test/src/heal_erasure_disk_rebuild_test.rs index 67d81c259..e5435f3d5 100644 --- a/crates/e2e_test/src/heal_erasure_disk_rebuild_test.rs +++ b/crates/e2e_test/src/heal_erasure_disk_rebuild_test.rs @@ -16,12 +16,13 @@ #[cfg(test)] mod tests { - use crate::common::{RustFSTestEnvironment, execute_awscurl, init_logging}; + use crate::common::{RustFSTestClusterEnvironment, RustFSTestEnvironment, execute_awscurl, init_logging}; use aws_sdk_s3::primitives::ByteStream; use serial_test::serial; use std::collections::HashSet; + use std::error::Error; use std::path::{Path, PathBuf}; - use tokio::time::{Duration, sleep}; + use tokio::time::{Duration, sleep, timeout}; use tracing::info; fn has_file_under(path: &Path) -> bool { @@ -328,4 +329,105 @@ mod tests { panic!("admin deep heal did not rebuild all files on the wiped disk within timeout"); } + + #[tokio::test(flavor = "multi_thread")] + #[serial] + async fn test_cluster_admin_heal_rebuilds_replaced_remote_disk() -> Result<(), Box> { + init_logging(); + info!("Admin deep heal should rebuild data on a remote node after its disk is replaced and the node rejoins"); + + let mut cluster = RustFSTestClusterEnvironment::new(4).await?; + cluster.set_env("RUSTFS_UNSAFE_BYPASS_DISK_CHECK", "true"); + cluster.set_env("RUSTFS_HEAL_ENABLED", "true"); + cluster.set_env("RUSTFS_SCANNER_ENABLED", "true"); + cluster.start().await?; + let clients = cluster.create_all_clients()?; + + let bucket = "heal-replaced-remote-disk"; + clients[0].create_bucket().bucket(bucket).send().await?; + + let online_key = "cluster/online-before-replacement.bin"; + let online_body = b"object written while all cluster nodes are online".to_vec(); + clients[0] + .put_object() + .bucket(bucket) + .key(online_key) + .body(ByteStream::from(online_body.clone())) + .send() + .await?; + + let replaced_disk = PathBuf::from(&cluster.nodes[1].data_dir); + assert!( + object_metadata_exists_on_disk(&replaced_disk, bucket, online_key), + "node 1 should contain metadata before disk replacement" + ); + + cluster.stop_node(1)?; + std::fs::remove_dir_all(&replaced_disk)?; + std::fs::create_dir_all(&replaced_disk)?; + assert!(!has_file_under(&replaced_disk), "replacement disk must start empty"); + + let outage_key = "cluster/written-while-node-down.bin"; + let outage_body = b"object written while one remote node is offline".to_vec(); + timeout(Duration::from_secs(30), async { + clients[0] + .put_object() + .bucket(bucket) + .key(outage_key) + .body(ByteStream::from(outage_body.clone())) + .send() + .await + }) + .await??; + + cluster.start_node(1).await?; + + let status_url = format!("{}/rustfs/admin/v3/background-heal/status", cluster.nodes[0].url); + let status_body = execute_awscurl(&status_url, "POST", None, &cluster.access_key, &cluster.secret_key).await?; + assert!( + !status_body.contains("MissingContentLength"), + "background heal status should not fail without an explicit Content-Length: {status_body}" + ); + + let heal_body = r#"{"recursive":true,"dryRun":false,"remove":false,"recreate":true,"scanMode":2,"updateParity":false,"nolock":false}"#; + let heal_url = format!("{}/rustfs/admin/v3/heal/{}?forceStart=true", cluster.nodes[0].url, bucket); + execute_awscurl(&heal_url, "POST", Some(heal_body), &cluster.access_key, &cluster.secret_key).await?; + + let expected_objects = [(online_key, online_body.as_slice()), (outage_key, outage_body.as_slice())]; + let mut remaining_rebuild_keys: HashSet<&str> = expected_objects.iter().map(|(key, _)| *key).collect(); + let heal_timeout_secs = std::env::var("RUSTFS_HEAL_REPLACED_DISK_TIMEOUT_SECS") + .ok() + .and_then(|value| value.parse::().ok()) + .unwrap_or(90); + + for _ in 0..heal_timeout_secs { + for (key, body) in &expected_objects { + let response = clients[0].get_object().bucket(bucket).key(*key).send().await?; + let actual = response.body.collect().await?.into_bytes(); + assert_eq!(actual.as_ref(), *body, "object body changed for {key}"); + } + + if !remaining_rebuild_keys.is_empty() { + let rebuilt = remaining_rebuild_keys + .iter() + .copied() + .filter(|key| object_metadata_exists_on_disk(&replaced_disk, bucket, key)) + .collect::>(); + for key in rebuilt { + let _ = remaining_rebuild_keys.remove(key); + } + } + + if remaining_rebuild_keys.is_empty() { + return Ok(()); + } + + sleep(Duration::from_secs(1)).await; + } + + Err(format!( + "admin deep heal did not rebuild replaced remote disk metadata for {remaining_rebuild_keys:?} within timeout" + ) + .into()) + } } diff --git a/crates/ecstore/src/disk/disk_store.rs b/crates/ecstore/src/disk/disk_store.rs index 226b67f6b..def933476 100644 --- a/crates/ecstore/src/disk/disk_store.rs +++ b/crates/ecstore/src/disk/disk_store.rs @@ -649,6 +649,11 @@ impl LocalDiskWrapper { self.health.reset_for_store_init_retry(&self.disk.endpoint()); } + #[cfg(test)] + pub fn health_check_enabled_for_test(&self) -> bool { + self.health_check + } + /// Enable health monitoring after disk creation. /// Used to defer health checks until after startup format loading completes. pub fn enable_health_check(&self) { diff --git a/crates/ecstore/src/disk/mod.rs b/crates/ecstore/src/disk/mod.rs index 111885315..fe6f4100f 100644 --- a/crates/ecstore/src/disk/mod.rs +++ b/crates/ecstore/src/disk/mod.rs @@ -435,6 +435,14 @@ impl Disk { } } + #[cfg(test)] + pub fn health_check_enabled_for_test(&self) -> bool { + match self { + Disk::Local(local_disk) => local_disk.health_check_enabled_for_test(), + Disk::Remote(remote_disk) => remote_disk.health_check_enabled_for_test(), + } + } + pub fn record_capacity_probe(&self, total: u64, used: u64, free: u64) { match self { Disk::Local(local_disk) => local_disk.record_capacity_probe(total, used, free), diff --git a/crates/ecstore/src/rpc/remote_disk.rs b/crates/ecstore/src/rpc/remote_disk.rs index f9b84787d..5f6e8c37b 100644 --- a/crates/ecstore/src/rpc/remote_disk.rs +++ b/crates/ecstore/src/rpc/remote_disk.rs @@ -221,6 +221,11 @@ impl RemoteDisk { self.health.reset_for_store_init_retry(&self.endpoint); } + #[cfg(test)] + pub fn health_check_enabled_for_test(&self) -> bool { + self.health_check + } + fn spawn_recovery_monitor_if_needed(&self) { if !self.health_check { return; diff --git a/crates/ecstore/src/set_disk/lock.rs b/crates/ecstore/src/set_disk/lock.rs index bd51d05f9..0c0e0deff 100644 --- a/crates/ecstore/src/set_disk/lock.rs +++ b/crates/ecstore/src/set_disk/lock.rs @@ -300,6 +300,7 @@ impl SetDisks { // Check that the endpoint matches let _ = new_disk.set_disk_id(Some(fm.erasure.this)).await; + new_disk.enable_health_check(); if new_disk.is_local() { let mut global_local_disk_map = GLOBAL_LOCAL_DISK_MAP.write().await; @@ -337,7 +338,14 @@ impl SetDisks { } pub(super) async fn connect_endpoint(ep: &Endpoint) -> disk::error::Result<(DiskStore, FormatV3)> { - let disk = new_disk(ep, &DiskOption::default()).await?; + let disk = new_disk( + ep, + &DiskOption { + cleanup: false, + health_check: true, + }, + ) + .await?; let fm = load_format_erasure(&disk, false).await?; @@ -604,4 +612,43 @@ mod tests { drop(temp_dirs); } + + #[tokio::test] + async fn renew_disk_enables_health_monitoring_for_recovered_disk() { + let disk_count = 4; + let format = FormatV3::new(1, disk_count); + + let mut temp_dirs = Vec::with_capacity(disk_count); + let mut endpoints = Vec::with_capacity(disk_count); + + for disk_idx in 0..disk_count { + let (temp_dir, endpoint, _) = make_formatted_local_disk(disk_idx, &format).await; + temp_dirs.push(temp_dir); + endpoints.push(endpoint); + } + + let set_disks = SetDisks::new( + "test-owner".to_string(), + Arc::new(RwLock::new(vec![None; disk_count])), + disk_count, + disk_count / 2, + 0, + 0, + endpoints.clone(), + format, + Vec::new(), + ) + .await; + + set_disks.renew_disk(&endpoints[0]).await; + + let disks = set_disks.get_disks_internal().await; + let renewed_disk = disks[0].as_ref().expect("renew_disk should attach the recovered disk"); + assert!( + renewed_disk.health_check_enabled_for_test(), + "renewed disks must keep health monitoring enabled so later faulty marks can recover" + ); + + drop(temp_dirs); + } } diff --git a/rustfs/src/server/layer.rs b/rustfs/src/server/layer.rs index cf68f5910..566abfe40 100644 --- a/rustfs/src/server/layer.rs +++ b/rustfs/src/server/layer.rs @@ -484,10 +484,12 @@ fn is_empty_body_admin_path(method: &Method, uri: &http::Uri) -> bool { path, "/minio/admin/v3/rebalance/start" | "/minio/admin/v3/rebalance/stop" + | "/minio/admin/v3/background-heal/status" | "/minio/admin/v3/pools/decommission" | "/minio/admin/v3/pools/cancel" | "/rustfs/admin/v3/rebalance/start" | "/rustfs/admin/v3/rebalance/stop" + | "/rustfs/admin/v3/background-heal/status" | "/rustfs/admin/v3/pools/decommission" | "/rustfs/admin/v3/pools/cancel" ) || is_heal_status_query(path, uri.query()) @@ -1850,6 +1852,28 @@ mod tests { } } + #[test] + fn admin_background_heal_status_without_content_length_is_normalized() { + let paths = [ + format!("{MINIO_ADMIN_V3_PREFIX}/background-heal/status"), + format!("{ADMIN_PREFIX}/v3/background-heal/status"), + ]; + + for path in paths { + let request = Request::builder() + .method(Method::POST) + .uri(path.clone()) + .header(http::header::TRANSFER_ENCODING, "chunked") + .body(()) + .expect("request"); + + assert!( + should_force_zero_content_length_for_empty_body_route(&request), + "{path} should force Content-Length: 0" + ); + } + } + #[test] fn admin_heal_start_without_status_token_is_not_normalized() { let request = Request::builder()