mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-27 15:37:02 +00:00
chore(heal): remove unused global response broadcast (#1991)
This commit is contained in:
@@ -18,7 +18,7 @@ use std::{
|
|||||||
fmt::{self, Display},
|
fmt::{self, Display},
|
||||||
sync::OnceLock,
|
sync::OnceLock,
|
||||||
};
|
};
|
||||||
use tokio::sync::{broadcast, mpsc};
|
use tokio::sync::mpsc;
|
||||||
use uuid::Uuid;
|
use uuid::Uuid;
|
||||||
|
|
||||||
pub const HEAL_DELETE_DANGLING: bool = true;
|
pub const HEAL_DELETE_DANGLING: bool = true;
|
||||||
@@ -290,11 +290,6 @@ pub type HealChannelReceiver = mpsc::UnboundedReceiver<HealChannelCommand>;
|
|||||||
/// Global heal channel sender
|
/// Global heal channel sender
|
||||||
static GLOBAL_HEAL_CHANNEL_SENDER: OnceLock<HealChannelSender> = OnceLock::new();
|
static GLOBAL_HEAL_CHANNEL_SENDER: OnceLock<HealChannelSender> = OnceLock::new();
|
||||||
|
|
||||||
type HealResponseSender = broadcast::Sender<HealChannelResponse>;
|
|
||||||
|
|
||||||
/// Global heal response broadcaster
|
|
||||||
static GLOBAL_HEAL_RESPONSE_SENDER: OnceLock<HealResponseSender> = OnceLock::new();
|
|
||||||
|
|
||||||
/// Initialize global heal channel
|
/// Initialize global heal channel
|
||||||
pub fn init_heal_channel() -> HealChannelReceiver {
|
pub fn init_heal_channel() -> HealChannelReceiver {
|
||||||
let (tx, rx) = mpsc::unbounded_channel();
|
let (tx, rx) = mpsc::unbounded_channel();
|
||||||
@@ -321,23 +316,6 @@ pub async fn send_heal_command(command: HealChannelCommand) -> Result<(), String
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
fn heal_response_sender() -> &'static HealResponseSender {
|
|
||||||
GLOBAL_HEAL_RESPONSE_SENDER.get_or_init(|| {
|
|
||||||
let (tx, _rx) = broadcast::channel(1024);
|
|
||||||
tx
|
|
||||||
})
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Publish a heal response to subscribers.
|
|
||||||
pub fn publish_heal_response(response: HealChannelResponse) -> Result<(), broadcast::error::SendError<HealChannelResponse>> {
|
|
||||||
heal_response_sender().send(response).map(|_| ())
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Subscribe to heal responses.
|
|
||||||
pub fn subscribe_heal_responses() -> broadcast::Receiver<HealChannelResponse> {
|
|
||||||
heal_response_sender().subscribe()
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Send heal start request
|
/// Send heal start request
|
||||||
pub async fn send_heal_request(request: HealChannelRequest) -> Result<(), String> {
|
pub async fn send_heal_request(request: HealChannelRequest) -> Result<(), String> {
|
||||||
send_heal_command(HealChannelCommand::Start(request)).await
|
send_heal_command(HealChannelCommand::Start(request)).await
|
||||||
@@ -536,20 +514,3 @@ pub async fn send_heal_disk(set_disk_id: String, priority: Option<HealChannelPri
|
|||||||
};
|
};
|
||||||
send_heal_request(req).await
|
send_heal_request(req).await
|
||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(test)]
|
|
||||||
mod tests {
|
|
||||||
use super::*;
|
|
||||||
|
|
||||||
#[tokio::test]
|
|
||||||
async fn heal_response_broadcast_reaches_subscriber() {
|
|
||||||
let mut receiver = subscribe_heal_responses();
|
|
||||||
let response = create_heal_response("req-1".to_string(), true, None, None);
|
|
||||||
|
|
||||||
publish_heal_response(response.clone()).expect("publish should succeed");
|
|
||||||
|
|
||||||
let received = receiver.recv().await.expect("should receive heal response");
|
|
||||||
assert_eq!(received.request_id, response.request_id);
|
|
||||||
assert!(received.success);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|||||||
@@ -20,7 +20,6 @@ use crate::heal::{
|
|||||||
use crate::{Error, Result};
|
use crate::{Error, Result};
|
||||||
use rustfs_common::heal_channel::{
|
use rustfs_common::heal_channel::{
|
||||||
HealChannelCommand, HealChannelPriority, HealChannelReceiver, HealChannelRequest, HealChannelResponse, HealScanMode,
|
HealChannelCommand, HealChannelPriority, HealChannelReceiver, HealChannelRequest, HealChannelResponse, HealScanMode,
|
||||||
publish_heal_response,
|
|
||||||
};
|
};
|
||||||
use std::sync::Arc;
|
use std::sync::Arc;
|
||||||
use tokio::sync::mpsc;
|
use tokio::sync::mpsc;
|
||||||
@@ -227,15 +226,9 @@ impl HealChannelProcessor {
|
|||||||
}
|
}
|
||||||
|
|
||||||
fn publish_response(&self, response: HealChannelResponse) {
|
fn publish_response(&self, response: HealChannelResponse) {
|
||||||
// Try to send to local channel first, but don't block broadcast on failure
|
if let Err(e) = self.response_sender.send(response) {
|
||||||
if let Err(e) = self.response_sender.send(response.clone()) {
|
|
||||||
error!("Failed to enqueue heal response locally: {}", e);
|
error!("Failed to enqueue heal response locally: {}", e);
|
||||||
}
|
}
|
||||||
// Always attempt to broadcast, even if local send failed
|
|
||||||
// Use the original response for broadcast; local send uses a clone
|
|
||||||
if let Err(e) = publish_heal_response(response) {
|
|
||||||
error!("Failed to broadcast heal response: {}", e);
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Get response sender for external use
|
/// Get response sender for external use
|
||||||
|
|||||||
@@ -27,7 +27,7 @@ use rustfs_heal::heal::{
|
|||||||
};
|
};
|
||||||
use serial_test::serial;
|
use serial_test::serial;
|
||||||
use std::{
|
use std::{
|
||||||
path::PathBuf,
|
path::{Path, PathBuf},
|
||||||
sync::{Arc, Once, OnceLock},
|
sync::{Arc, Once, OnceLock},
|
||||||
time::Duration,
|
time::Duration,
|
||||||
};
|
};
|
||||||
@@ -49,6 +49,17 @@ pub fn init_tracing() {
|
|||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
|
async fn wait_for_path_exists(path: &Path, timeout: Duration) -> bool {
|
||||||
|
let deadline = std::time::Instant::now() + timeout;
|
||||||
|
while std::time::Instant::now() < deadline {
|
||||||
|
if path.exists() {
|
||||||
|
return true;
|
||||||
|
}
|
||||||
|
tokio::time::sleep(Duration::from_millis(200)).await;
|
||||||
|
}
|
||||||
|
path.exists()
|
||||||
|
}
|
||||||
|
|
||||||
/// Test helper: Create test environment with ECStore
|
/// Test helper: Create test environment with ECStore
|
||||||
async fn setup_test_env() -> (Vec<PathBuf>, Arc<ECStore>, Arc<ECStoreHealStorage>) {
|
async fn setup_test_env() -> (Vec<PathBuf>, Arc<ECStore>, Arc<ECStoreHealStorage>) {
|
||||||
init_tracing();
|
init_tracing();
|
||||||
@@ -329,25 +340,11 @@ mod serial_tests {
|
|||||||
let (_result, error) = heal_storage.heal_format(false).await.expect("Failed to heal format");
|
let (_result, error) = heal_storage.heal_format(false).await.expect("Failed to heal format");
|
||||||
assert!(error.is_none(), "Heal format returned error: {error:?}");
|
assert!(error.is_none(), "Heal format returned error: {error:?}");
|
||||||
|
|
||||||
// ─── 2️⃣ wait for format.json to be restored with polling + timeout ───────
|
// Wait for task completion
|
||||||
// The minimal scanner interval is clamped to 10s in manager.rs, so we set timeout to 20s
|
let restored = wait_for_path_exists(&format_path, Duration::from_secs(20)).await;
|
||||||
let timeout_duration = Duration::from_secs(20);
|
|
||||||
let poll_interval = Duration::from_millis(200);
|
|
||||||
|
|
||||||
let result = tokio::time::timeout(timeout_duration, async {
|
// ─── 2️⃣ verify format.json is restored ───────
|
||||||
loop {
|
assert!(restored, "format.json does not exist on disk after heal");
|
||||||
if format_path.exists() {
|
|
||||||
break;
|
|
||||||
}
|
|
||||||
tokio::time::sleep(poll_interval).await;
|
|
||||||
}
|
|
||||||
})
|
|
||||||
.await;
|
|
||||||
|
|
||||||
assert!(result.is_ok(), "format.json was not restored within timeout period");
|
|
||||||
|
|
||||||
// ─── 3️⃣ verify format.json is restored ───────
|
|
||||||
assert!(format_path.exists(), "format.json does not exist on disk after heal");
|
|
||||||
|
|
||||||
info!("Heal format basic test passed");
|
info!("Heal format basic test passed");
|
||||||
}
|
}
|
||||||
@@ -365,51 +362,37 @@ mod serial_tests {
|
|||||||
create_test_bucket(&ecstore, bucket_name).await;
|
create_test_bucket(&ecstore, bucket_name).await;
|
||||||
upload_test_object(&ecstore, bucket_name, object_name, test_data).await;
|
upload_test_object(&ecstore, bucket_name, object_name, test_data).await;
|
||||||
|
|
||||||
let obj_dir = disk_paths[0].join(bucket_name).join(object_name);
|
|
||||||
let target_part = WalkDir::new(&obj_dir)
|
|
||||||
.min_depth(2)
|
|
||||||
.max_depth(2)
|
|
||||||
.into_iter()
|
|
||||||
.filter_map(Result::ok)
|
|
||||||
.find(|e| e.file_type().is_file() && e.file_name().to_str().map(|n| n.starts_with("part.")).unwrap_or(false))
|
|
||||||
.map(|e| e.into_path())
|
|
||||||
.expect("Failed to locate part file to delete");
|
|
||||||
|
|
||||||
// ─── 1️⃣ delete format.json on one disk ──────────────
|
// ─── 1️⃣ delete format.json on one disk ──────────────
|
||||||
let format_path = disk_paths[0].join(".rustfs.sys").join("format.json");
|
let format_path = disk_paths[0].join(".rustfs.sys").join("format.json");
|
||||||
std::fs::remove_dir_all(&disk_paths[0]).expect("failed to delete all contents under disk_paths[0]");
|
std::fs::remove_dir_all(&disk_paths[0]).expect("failed to delete all contents under disk_paths[0]");
|
||||||
std::fs::create_dir_all(&disk_paths[0]).expect("failed to recreate disk_paths[0] directory");
|
std::fs::create_dir_all(&disk_paths[0]).expect("failed to recreate disk_paths[0] directory");
|
||||||
println!("✅ Deleted format.json on disk: {:?}", disk_paths[0]);
|
println!("✅ Deleted format.json on disk: {:?}", disk_paths[0]);
|
||||||
|
|
||||||
// Create heal manager with faster interval
|
let (_result, error) = heal_storage.heal_format(false).await.expect("Failed to heal format");
|
||||||
let cfg = HealConfig {
|
assert!(error.is_none(), "Heal format returned error: {error:?}");
|
||||||
heal_interval: Duration::from_secs(1),
|
|
||||||
..Default::default()
|
// Wait for task completion
|
||||||
|
let restored = wait_for_path_exists(&format_path, Duration::from_secs(20)).await;
|
||||||
|
|
||||||
|
// ─── 2️⃣ verify format.json is restored ───────
|
||||||
|
assert!(restored, "format.json does not exist on disk after heal");
|
||||||
|
// ─── 3 verify each part file is restored ───────
|
||||||
|
let heal_opts = HealOpts {
|
||||||
|
recursive: false,
|
||||||
|
dry_run: false,
|
||||||
|
remove: false,
|
||||||
|
recreate: true,
|
||||||
|
scan_mode: HealScanMode::Normal,
|
||||||
|
update_parity: true,
|
||||||
|
no_lock: false,
|
||||||
|
pool: None,
|
||||||
|
set: None,
|
||||||
};
|
};
|
||||||
let heal_manager = HealManager::new(heal_storage.clone(), Some(cfg));
|
let (_result, error) = heal_storage
|
||||||
heal_manager.start().await.unwrap();
|
.heal_object(bucket_name, object_name, None, &heal_opts)
|
||||||
|
.await
|
||||||
// ─── 2️⃣ wait for format.json and part file to be restored with polling + timeout ───────
|
.expect("Failed to heal object");
|
||||||
// The minimal scanner interval is clamped to 10s in manager.rs, so we set timeout to 20s
|
assert!(error.is_none(), "Heal object returned error: {error:?}");
|
||||||
let timeout_duration = Duration::from_secs(20);
|
|
||||||
let poll_interval = Duration::from_millis(200);
|
|
||||||
|
|
||||||
let result = tokio::time::timeout(timeout_duration, async {
|
|
||||||
loop {
|
|
||||||
if format_path.exists() && target_part.exists() {
|
|
||||||
break;
|
|
||||||
}
|
|
||||||
tokio::time::sleep(poll_interval).await;
|
|
||||||
}
|
|
||||||
})
|
|
||||||
.await;
|
|
||||||
|
|
||||||
assert!(result.is_ok(), "format.json or part file was not restored within timeout period");
|
|
||||||
|
|
||||||
// ─── 3️⃣ verify format.json is restored ───────
|
|
||||||
assert!(format_path.exists(), "format.json does not exist on disk after heal");
|
|
||||||
// ─── 4️⃣ verify each part file is restored ───────
|
|
||||||
assert!(target_part.exists());
|
|
||||||
|
|
||||||
// Verify object metadata is accessible
|
// Verify object metadata is accessible
|
||||||
let obj_info = ecstore
|
let obj_info = ecstore
|
||||||
|
|||||||
Reference in New Issue
Block a user