Files
rustfs/crates/heal/src/heal/channel.rs
T
houseme 360bceafce feat(heal): add progress and trace observability (#6179)
* feat(heal): track erasure set progress baseline

Record erasure-set heal byte progress from per-object results and seed progress totals from complete usage-cache snapshots when available.

Keep usage-cache failures observational so heal execution continues without a baseline.

Co-Authored-By: heihutu <heihutu@gmail.com>

* feat(heal): skip filtered erasure set versions

Skip erasure-set versions written after the durable heal start time, and queue lifecycle-expired versions for expiry before skipping them.

Track new-version and ILM-expired skips separately so progress can explain completed baseline work without treating these skips as retry-blocking failures.

Co-Authored-By: heihutu <heihutu@gmail.com>

* feat(heal): wire abandoned data-dir cleanup check

Connect check_abandoned_parts through ECStore, pool, and set layers so heal can invoke the existing orphan data-dir reclaim path instead of returning NotImplemented.

Add dry-run support to the reclaim scan and cover dry-run plus scoped set behavior with regression tests.

Co-Authored-By: heihutu <heihutu@gmail.com>

* feat(obs): add heal scanner trace bus

Introduce an in-process broadcast trace bus with typed heal and scanner events, lazy event construction, and bounded lagged-subscriber behavior.

Cover zero-subscriber publishing, subscription delivery, drop accounting, and lagged receivers with focused common-crate tests.

Co-Authored-By: heihutu <heihutu@gmail.com>

* feat(obs): stream heal trace events from admin API

Wire the admin trace endpoint to the common trace bus for heal/scanner events, including kind, regex, and threshold filtering.

Co-Authored-By: heihutu <heihutu@gmail.com>

* feat(obs): emit heal trace events

Publish heal task lifecycle and abandoned-parts cleanup events through the common trace bus so the admin trace stream has live heal diagnostics.

Co-Authored-By: heihutu <heihutu@gmail.com>

* feat(obs): emit scanner trace events

Publish scanner folder, lifecycle action, and heal-candidate events through the common trace bus for live admin scanner diagnostics.

Co-Authored-By: heihutu <heihutu@gmail.com>

* fix(heal): route data usage loader through storage api

Keep ECStore data-usage facade access behind the heal storage_api boundary so architecture migration guards can validate the heal progress path.

Co-Authored-By: heihutu <heihutu@gmail.com>

* perf(heal): avoid lifecycle snapshots on ordinary heal pages

Only request lifecycle object snapshots when the heal pass has lifecycle expiry context. This keeps ordinary listing and disk-walk pages from cloning FileInfo/ObjectInfo payloads while preserving the skip path that queues expired versions.

Co-Authored-By: heihutu <heihutu@gmail.com>

* test(heal): update bug-fix mocks for lifecycle snapshots

Carry the lifecycle snapshot opt-in argument through the remaining heal bug-fix test mocks so all-targets clippy covers the updated storage trait.

Co-Authored-By: heihutu <heihutu@gmail.com>

* test(rustfs): sync heal storage mock signature

Update the rustfs storage RPC test mock for the lifecycle snapshot opt-in argument and cover it with rustfs all-targets clippy.

Co-Authored-By: heihutu <heihutu@gmail.com>

* test(e2e): allocate smoke ports across nextest processes

Serialize E2E port selection with a small /tmp allocator so nextest workers do not reuse the same just-released ephemeral port before RustFS binds it.

Co-Authored-By: heihutu <heihutu@gmail.com>

---------

Co-authored-by: heihutu <heihutu@gmail.com>
2026-08-18 08:29:29 +08:00

1873 lines
71 KiB
Rust

// Copyright 2024 RustFS 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.
use crate::heal::{
manager::{HealManager, HealTaskReport},
progress::HealProgress,
task::{HealOptions, HealPriority, HealRequest, HealTaskStatus, HealType},
utils,
};
use crate::{Error, Result};
use rustfs_common::heal_channel::{
HealAdmissionReceipt, HealAdmissionResult, HealChannelCommand, HealChannelPriority, HealChannelReceiver, HealChannelRequest,
HealChannelResponse, HealReceiptCommand, HealReceiptReceiver, HealRequestSource, HealScanMode, publish_heal_response,
};
use rustfs_madmin::heal_commands::HealResultItem;
use serde::Serialize;
use std::sync::Arc;
use tokio::sync::{mpsc, oneshot};
use tracing::{debug, error, info};
const LOG_COMPONENT_HEAL: &str = "heal";
const LOG_SUBSYSTEM_CHANNEL: &str = "channel";
const EVENT_HEAL_CHANNEL_STATE: &str = "heal_channel_state";
const EVENT_HEAL_CHANNEL_REQUEST: &str = "heal_channel_request";
const EVENT_HEAL_CHANNEL_RESPONSE: &str = "heal_channel_response";
const MAX_HEAL_STATUS_PAYLOAD_SIZE: usize = 8 * 1024 * 1024;
fn admission_response(request_id: String, admission: HealAdmissionResult) -> HealChannelResponse {
let (success, error) = match admission {
HealAdmissionResult::Accepted | HealAdmissionResult::Merged => (true, None),
HealAdmissionResult::Full => (false, Some("Heal request queue is full".to_string())),
HealAdmissionResult::Dropped(reason) => (false, Some(format!("Heal request dropped: {}", reason.as_str()))),
};
HealChannelResponse {
request_id,
success,
data: Some(format!("admission={},reason={}", admission.result_label(), admission.reason_label()).into_bytes()),
error,
}
}
/// Heal channel processor
pub struct HealChannelProcessor {
/// Heal manager
heal_manager: Arc<HealManager>,
/// Response sender
response_sender: mpsc::UnboundedSender<HealChannelResponse>,
/// Response receiver
response_receiver: mpsc::UnboundedReceiver<HealChannelResponse>,
}
#[derive(Serialize)]
struct HealTaskStatusPayload<'a> {
summary: &'a str,
items: &'a [HealResultItem],
truncated: bool,
#[serde(skip_serializing_if = "Option::is_none")]
progress: Option<&'a HealProgress>,
}
fn encode_heal_task_status_payload(
summary: &str,
mut items: Vec<HealResultItem>,
progress: Option<&HealProgress>,
mut truncated: bool,
) -> Result<(Vec<u8>, bool)> {
loop {
let data = serde_json::to_vec(&HealTaskStatusPayload {
summary,
items: &items,
truncated,
progress,
})
.map_err(|e| Error::Serialization(format!("failed to serialize heal task status: {e}")))?;
if data.len() <= MAX_HEAL_STATUS_PAYLOAD_SIZE {
return Ok((data, truncated));
}
if items.is_empty() {
return Err(Error::Serialization("heal task status metadata exceeds size limit".to_string()));
}
truncated = true;
items.truncate(items.len() / 2);
}
}
fn heal_status_detail(detail: Option<String>, truncated: bool) -> Option<String> {
if !truncated {
return detail;
}
let truncation = "heal result items were truncated";
Some(detail.map_or_else(|| truncation.to_string(), |detail| format!("{detail}; {truncation}")))
}
fn encode_heal_status_response(
summary: &str,
items: Vec<HealResultItem>,
progress: Option<&HealProgress>,
detail: Option<String>,
truncated: bool,
) -> Result<(Vec<u8>, Option<String>)> {
let (data, truncated) = encode_heal_task_status_payload(summary, items, progress, truncated)?;
Ok((data, heal_status_detail(detail, truncated)))
}
impl HealChannelProcessor {
/// Create new HealChannelProcessor
pub fn new(heal_manager: Arc<HealManager>) -> Self {
let (response_tx, response_rx) = mpsc::unbounded_channel();
Self {
heal_manager,
response_sender: response_tx,
response_receiver: response_rx,
}
}
/// Execute a start directly against the manager without entering the
/// process-global unbounded command queue.
pub async fn execute_start_request(&self, request: HealChannelRequest) -> Result<HealAdmissionReceipt> {
let (response_tx, response_rx) = oneshot::channel();
self.process_start_request(request, false, true, response_tx).await?;
response_rx
.await
.map_err(|err| Error::other(format!("heal receipt channel closed: {err}")))?
.map_err(Error::other)
}
/// Execute a token query directly against the manager.
pub async fn execute_query_request(&self, heal_path: String, client_token: String) -> Result<HealChannelResponse> {
let (response_tx, response_rx) = oneshot::channel();
self.process_query_request(heal_path, client_token, response_tx).await?;
response_rx
.await
.map_err(|err| Error::other(format!("heal query channel closed: {err}")))?
.map_err(Error::other)
}
/// Execute cancellation directly against the manager.
pub async fn execute_cancel_request(&self, heal_path: String, client_token: String) -> Result<HealChannelResponse> {
let (response_tx, response_rx) = oneshot::channel();
self.process_cancel_request(heal_path, client_token, response_tx).await?;
response_rx
.await
.map_err(|err| Error::other(format!("heal cancel channel closed: {err}")))?
.map_err(Error::other)
}
/// Start processing legacy heal channel requests.
pub async fn start(&mut self, receiver: HealChannelReceiver) -> Result<()> {
let (receipt_sender, receipt_receiver) = mpsc::unbounded_channel();
drop(receipt_sender);
self.start_with_receipts(receiver, receipt_receiver).await
}
/// Start processing legacy and canonical-receipt heal channel requests.
pub async fn start_with_receipts(
&mut self,
mut receiver: HealChannelReceiver,
mut receipt_receiver: HealReceiptReceiver,
) -> Result<()> {
info!(
target: "rustfs::heal::channel",
event = EVENT_HEAL_CHANNEL_STATE,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_CHANNEL,
state = "started",
"Heal channel started"
);
let mut receipt_channel_open = true;
loop {
tokio::select! {
command = receiver.recv() => {
match command {
Some(command) => {
if let Err(e) = self.process_command(command).await {
error!(
target: "rustfs::heal::channel",
event = EVENT_HEAL_CHANNEL_REQUEST,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_CHANNEL,
state = "process_failed",
error = %e,
"Heal channel processing failed"
);
}
}
None => {
debug!(
target: "rustfs::heal::channel",
event = EVENT_HEAL_CHANNEL_STATE,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_CHANNEL,
state = "receiver_closed",
"Heal channel receiver closed"
);
break;
}
}
}
command = receipt_receiver.recv(), if receipt_channel_open => {
let Some(HealReceiptCommand { request, response_tx }) = command else {
receipt_channel_open = false;
continue;
};
if let Err(e) = self.process_start_request(request, false, true, response_tx).await {
error!(
target: "rustfs::heal::channel",
event = EVENT_HEAL_CHANNEL_REQUEST,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_CHANNEL,
state = "receipt_process_failed",
error = %e,
"Heal receipt request processing failed"
);
}
}
response = self.response_receiver.recv() => {
if let Some(response) = response {
// Handle response if needed
debug!(
target: "rustfs::heal::channel",
event = EVENT_HEAL_CHANNEL_RESPONSE,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_CHANNEL,
request_id = %response.request_id,
success = response.success,
state = "received_local",
"Heal response observed"
);
}
}
}
}
info!(
target: "rustfs::heal::channel",
event = EVENT_HEAL_CHANNEL_STATE,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_CHANNEL,
state = "stopped",
"Heal channel stopped"
);
Ok(())
}
/// Process heal command
async fn process_command(&self, command: HealChannelCommand) -> Result<()> {
match command {
HealChannelCommand::Start { request, response_tx } => self.process_legacy_start_request(request, response_tx).await,
HealChannelCommand::Query {
heal_path,
client_token,
response_tx,
} => self.process_query_request(heal_path, client_token, response_tx).await,
HealChannelCommand::Cancel {
heal_path,
client_token,
response_tx,
} => self.process_cancel_request(heal_path, client_token, response_tx).await,
}
}
async fn process_legacy_start_request(
&self,
request: HealChannelRequest,
response_tx: oneshot::Sender<std::result::Result<HealAdmissionResult, String>>,
) -> Result<()> {
let (receipt_tx, receipt_rx) = oneshot::channel();
self.process_start_request(request, true, false, receipt_tx).await?;
let result = receipt_rx
.await
.map_err(|err| Error::other(format!("heal receipt channel closed: {err}")))?
.map(|receipt| receipt.result);
let _ = response_tx.send(result);
Ok(())
}
/// Process start request
async fn process_start_request(
&self,
request: HealChannelRequest,
preserve_alias: bool,
publish_canonical_id: bool,
response_tx: oneshot::Sender<std::result::Result<HealAdmissionReceipt, String>>,
) -> Result<()> {
debug!(
target: "rustfs::heal::channel",
event = EVENT_HEAL_CHANNEL_REQUEST,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_CHANNEL,
request_id = %request.id,
bucket = %request.bucket,
object_prefix = %request.object_prefix.as_deref().unwrap_or(""),
state = "start_received",
"Heal start received"
);
// Convert channel request to heal request
let heal_request = match self.convert_to_heal_request(request.clone()) {
Ok(heal_request) => heal_request,
Err(err) => {
let error_text = err.to_string();
let _ = response_tx.send(Err(error_text.clone()));
self.publish_response(HealChannelResponse {
request_id: request.id,
success: false,
data: None,
error: Some(error_text),
});
return Ok(());
}
};
// Submit to heal manager
match self
.heal_manager
.submit_heal_request_with_receipt_and_alias(heal_request, preserve_alias)
.await
{
Ok(receipt) => {
let admission = receipt.result;
debug!(
target: "rustfs::heal::channel",
event = EVENT_HEAL_CHANNEL_REQUEST,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_CHANNEL,
request_id = %request.id,
admission = admission.result_label(),
state = "admission_decided",
"Heal admission decided"
);
let response_id = if publish_canonical_id {
receipt.task_id.clone()
} else {
request.id.clone()
};
self.publish_response(admission_response(response_id, receipt.result));
let _ = response_tx.send(Ok(receipt));
}
Err(e) => {
let error_text = e.to_string();
error!(
target: "rustfs::heal::channel",
event = EVENT_HEAL_CHANNEL_REQUEST,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_CHANNEL,
request_id = %request.id,
state = "submit_failed",
error = %error_text,
"Heal start submission failed"
);
let _ = response_tx.send(Err(error_text.clone()));
// Send error response
let response = HealChannelResponse {
request_id: request.id,
success: false,
data: None,
error: Some(error_text),
};
self.publish_response(response);
}
}
Ok(())
}
/// Process query request
async fn process_query_request(
&self,
heal_path: String,
client_token: String,
response_tx: oneshot::Sender<std::result::Result<HealChannelResponse, String>>,
) -> Result<()> {
debug!(
target: "rustfs::heal::channel",
event = EVENT_HEAL_CHANNEL_REQUEST,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_CHANNEL,
request_id = %client_token,
heal_path = %heal_path,
state = "query_received",
"Heal query received"
);
let report = if heal_path.trim_matches('/').is_empty() {
self.heal_manager.get_task_report(&client_token).await
} else {
self.heal_manager.get_task_report_for_path(&heal_path, &client_token).await
};
let (summary, detail, items, truncated, progress) = match report {
Ok(HealTaskReport {
status: HealTaskStatus::Pending | HealTaskStatus::Running,
result_items,
result_items_truncated,
progress,
}) => ("running".to_string(), None, result_items, result_items_truncated, progress),
Ok(HealTaskReport {
status: HealTaskStatus::Retrying { error, retry_attempt },
result_items,
result_items_truncated,
progress,
}) => (
"running".to_string(),
Some(format!("heal task retrying after recoverable failure, attempt {retry_attempt}: {error}")),
result_items,
result_items_truncated,
progress,
),
Ok(HealTaskReport {
status: HealTaskStatus::Completed,
result_items,
result_items_truncated,
progress,
}) => ("finished".to_string(), None, result_items, result_items_truncated, progress),
Ok(HealTaskReport {
status: HealTaskStatus::Cancelled,
result_items,
result_items_truncated,
progress,
}) => (
"stopped".to_string(),
Some("heal task cancelled".to_string()),
result_items,
result_items_truncated,
progress,
),
Ok(HealTaskReport {
status: HealTaskStatus::Timeout,
result_items,
result_items_truncated,
progress,
}) => (
"stopped".to_string(),
Some("heal task timed out".to_string()),
result_items,
result_items_truncated,
progress,
),
Ok(HealTaskReport {
status: HealTaskStatus::Failed { error },
result_items,
result_items_truncated,
progress,
}) => ("stopped".to_string(), Some(error), result_items, result_items_truncated, progress),
Err(crate::Error::TaskNotFound { .. }) => (
"notFound".to_string(),
Some("heal task not found or expired".to_string()),
Vec::new(),
false,
None,
),
Err(crate::Error::InvalidClientToken) => {
let response = HealChannelResponse {
request_id: client_token,
success: false,
data: None,
error: Some("invalid heal client token".to_string()),
};
let _ = response_tx.send(Ok(response.clone()));
self.publish_response(response);
return Ok(());
}
Err(err) => {
let error_text = err.to_string();
let response = HealChannelResponse {
request_id: client_token,
success: false,
data: None,
error: Some(error_text),
};
let _ = response_tx.send(Ok(response.clone()));
self.publish_response(response);
return Ok(());
}
};
let (data, detail) = encode_heal_status_response(&summary, items, progress.as_ref(), detail, truncated)?;
let response = HealChannelResponse {
request_id: client_token,
success: true,
data: Some(data),
error: detail,
};
let _ = response_tx.send(Ok(response.clone()));
self.publish_response(response);
Ok(())
}
/// Process cancel request
async fn process_cancel_request(
&self,
heal_path: String,
client_token: String,
response_tx: oneshot::Sender<std::result::Result<HealChannelResponse, String>>,
) -> Result<()> {
debug!(
target: "rustfs::heal::channel",
event = EVENT_HEAL_CHANNEL_REQUEST,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_CHANNEL,
request_id = %client_token,
heal_path = %heal_path,
state = "cancel_received",
"Heal cancel received"
);
let request_id = if client_token.is_empty() {
heal_path.clone()
} else {
client_token.clone()
};
let cancel_result = if client_token.is_empty() {
self.heal_manager.cancel_tasks_for_path(&heal_path).await.map(|_| ())
} else {
self.heal_manager.cancel_task(&client_token).await
};
let response = match cancel_result {
Ok(()) => HealChannelResponse {
request_id,
success: true,
data: Some("stopped".as_bytes().to_vec()),
error: None,
},
Err(Error::TaskNotFound { .. }) if client_token.is_empty() => HealChannelResponse {
request_id,
success: true,
data: Some("stopped".as_bytes().to_vec()),
error: None,
},
Err(err) => HealChannelResponse {
request_id,
success: false,
data: None,
error: Some(err.to_string()),
},
};
let _ = response_tx.send(Ok(response.clone()));
self.publish_response(response);
Ok(())
}
/// Convert channel request to heal request
fn convert_to_heal_request(&self, request: HealChannelRequest) -> Result<HealRequest> {
let recursive = request.recursive.unwrap_or(false);
let heal_type = if let Some(disk_id) = &request.disk {
let set_disk_id = utils::normalize_set_disk_id(disk_id).ok_or_else(|| Error::InvalidHealType {
heal_type: format!("erasure-set({disk_id})"),
})?;
HealType::ErasureSet {
buckets: vec![],
set_disk_id,
}
} else if request.bucket.is_empty() {
HealType::Cluster
} else if let Some(prefix) = &request.object_prefix {
if !prefix.is_empty() {
if recursive {
HealType::Prefix {
bucket: request.bucket.clone(),
prefix: prefix.clone(),
}
} else {
HealType::Object {
bucket: request.bucket.clone(),
object: prefix.clone(),
version_id: request.object_version_id.clone(),
}
}
} else {
HealType::Bucket {
bucket: request.bucket.clone(),
}
}
} else {
HealType::Bucket {
bucket: request.bucket.clone(),
}
};
let priority = match request.priority {
HealChannelPriority::Low => HealPriority::Low,
HealChannelPriority::Normal => HealPriority::Normal,
HealChannelPriority::High => HealPriority::High,
HealChannelPriority::Critical => HealPriority::Urgent,
};
let recreate_missing = request.recreate_missing.unwrap_or(match request.source {
HealRequestSource::Scanner => false,
HealRequestSource::Admin
| HealRequestSource::AutoHeal
| HealRequestSource::Internal
| HealRequestSource::ReadRepair => true,
});
// Build HealOptions with all available fields
let options = HealOptions {
scan_mode: request.scan_mode.unwrap_or(HealScanMode::Normal),
remove_corrupted: request.remove_corrupted.unwrap_or(false),
recreate_missing,
update_parity: request.update_parity.unwrap_or(true),
recursive,
dry_run: request.dry_run.unwrap_or(false),
no_lock: request.no_lock.unwrap_or(false),
timeout: request.timeout_seconds.map(std::time::Duration::from_secs),
pool_index: request.pool_index,
set_index: request.set_index,
};
let mut heal_request = HealRequest::new(heal_type, options, priority);
heal_request.id = request.id;
heal_request.source = request.source;
// force_start controls admission/queue semantics only. Do not reinterpret it as
// destructive heal options: admin clients commonly pass forceStart=true together
// with remove=false, and turning that into remove_corrupted=true can delete the
// remaining healthy bucket volumes before object shards are rebuilt.
heal_request.force_start = request.force_start;
Ok(heal_request)
}
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.clone()) {
error!(
target: "rustfs::heal::channel",
event = EVENT_HEAL_CHANNEL_RESPONSE,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_CHANNEL,
request_id = %response.request_id,
state = "enqueue_local_failed",
error = %e,
"Heal response local enqueue failed"
);
}
// 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!(
target: "rustfs::heal::channel",
event = EVENT_HEAL_CHANNEL_RESPONSE,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_CHANNEL,
state = "broadcast_failed",
error = %e,
"Heal response broadcast failed"
);
}
}
/// Get response sender for external use
pub fn get_response_sender(&self) -> mpsc::UnboundedSender<HealChannelResponse> {
self.response_sender.clone()
}
}
#[cfg(test)]
mod tests {
use super::super::{DiskStore, Endpoint};
use super::*;
use crate::heal::manager::HealConfig;
use crate::heal::storage::{HealObjectInfo, HealStorageAPI};
use rustfs_common::heal_channel::{
HealAdmissionDropReason, HealAdmissionResult, HealChannelPriority, HealChannelRequest, HealRequestSource, HealScanMode,
};
use std::sync::Arc;
use std::time::Duration;
// Mock storage for testing
struct MockStorage;
#[async_trait::async_trait]
impl HealStorageAPI for MockStorage {
async fn get_object_meta(&self, _bucket: &str, _object: &str) -> crate::Result<Option<HealObjectInfo>> {
Ok(None)
}
async fn get_object_data(&self, _bucket: &str, _object: &str) -> crate::Result<Option<Vec<u8>>> {
Ok(None)
}
async fn put_object_data(&self, _bucket: &str, _object: &str, _data: &[u8]) -> crate::Result<()> {
Ok(())
}
async fn delete_object(&self, _bucket: &str, _object: &str) -> crate::Result<()> {
Ok(())
}
async fn verify_object_integrity(&self, _bucket: &str, _object: &str) -> crate::Result<bool> {
Ok(true)
}
async fn ec_decode_rebuild(&self, _bucket: &str, _object: &str) -> crate::Result<Vec<u8>> {
Ok(vec![])
}
async fn get_disk_status(&self, _endpoint: &Endpoint) -> crate::Result<crate::heal::storage::DiskStatus> {
Ok(crate::heal::storage::DiskStatus::Ok)
}
async fn format_disk(&self, _endpoint: &Endpoint) -> crate::Result<()> {
Ok(())
}
async fn get_bucket_info(&self, _bucket: &str) -> crate::Result<Option<crate::heal::storage_api::status::BucketInfo>> {
Ok(None)
}
async fn heal_bucket_metadata(&self, _bucket: &str) -> crate::Result<()> {
Ok(())
}
async fn list_buckets(&self) -> crate::Result<Vec<crate::heal::storage_api::status::BucketInfo>> {
Ok(vec![])
}
async fn object_exists(&self, _bucket: &str, _object: &str) -> crate::Result<bool> {
Ok(false)
}
async fn get_object_size(&self, _bucket: &str, _object: &str) -> crate::Result<Option<u64>> {
Ok(None)
}
async fn get_object_checksum(&self, _bucket: &str, _object: &str) -> crate::Result<Option<String>> {
Ok(None)
}
async fn heal_object(
&self,
_bucket: &str,
_object: &str,
_version_id: Option<&str>,
_opts: &rustfs_common::heal_channel::HealOpts,
) -> crate::Result<(rustfs_madmin::heal_commands::HealResultItem, Option<crate::Error>)> {
Ok((rustfs_madmin::heal_commands::HealResultItem::default(), None))
}
async fn heal_bucket(
&self,
_bucket: &str,
_opts: &rustfs_common::heal_channel::HealOpts,
) -> crate::Result<rustfs_madmin::heal_commands::HealResultItem> {
Ok(rustfs_madmin::heal_commands::HealResultItem::default())
}
async fn heal_format(
&self,
_dry_run: bool,
) -> crate::Result<(rustfs_madmin::heal_commands::HealResultItem, Option<crate::Error>)> {
Ok((rustfs_madmin::heal_commands::HealResultItem::default(), None))
}
async fn list_objects_for_heal(
&self,
_bucket: &str,
_prefix: &str,
) -> crate::Result<Vec<crate::heal::storage::HealListItem>> {
Ok(vec![])
}
async fn list_objects_for_heal_page(
&self,
_bucket: &str,
_prefix: &str,
_continuation_token: Option<&str>,
_include_lifecycle_object_info: bool,
) -> crate::Result<(Vec<crate::heal::storage::HealListItem>, Option<String>, bool)> {
Ok((vec![], None, false))
}
async fn get_disk_for_resume(&self, _set_disk_id: &str) -> crate::Result<DiskStore> {
Err(crate::Error::other("Not implemented in mock"))
}
}
fn create_test_heal_manager() -> Arc<HealManager> {
let storage: Arc<dyn HealStorageAPI> = Arc::new(MockStorage);
Arc::new(HealManager::new(storage, None))
}
#[test]
fn test_heal_channel_processor_new() {
let heal_manager = create_test_heal_manager();
let processor = HealChannelProcessor::new(heal_manager);
let sender = processor.get_response_sender();
sender
.send(HealChannelResponse {
request_id: "request-id".to_string(),
success: true,
data: None,
error: None,
})
.expect("a freshly constructed processor must accept responses on its channel");
}
#[test]
fn oversized_status_items_are_truncated_before_transport() {
let items = vec![HealResultItem {
detail: "x".repeat(MAX_HEAL_STATUS_PAYLOAD_SIZE + 1),
..Default::default()
}];
let (data, detail) = encode_heal_status_response("running", items, None, None, false).unwrap();
assert!(data.len() <= MAX_HEAL_STATUS_PAYLOAD_SIZE);
let payload: serde_json::Value = serde_json::from_slice(&data).unwrap();
assert_eq!(payload["truncated"], true);
assert!(payload["items"].as_array().unwrap().is_empty());
assert_eq!(detail.as_deref(), Some("heal result items were truncated"));
}
#[test]
fn admission_response_preserves_all_admission_outcomes() {
let cases = [
(HealAdmissionResult::Accepted, true, None, "admission=accepted,reason=none"),
(HealAdmissionResult::Merged, true, None, "admission=merged,reason=none"),
(
HealAdmissionResult::Full,
false,
Some("Heal request queue is full"),
"admission=full,reason=none",
),
(
HealAdmissionResult::Dropped(HealAdmissionDropReason::PolicyDropped),
false,
Some("Heal request dropped: policy_dropped"),
"admission=dropped,reason=policy_dropped",
),
];
for (admission, success, error, data) in cases {
let response = admission_response("request-id".to_string(), admission);
assert_eq!(response.request_id, "request-id");
assert_eq!(response.success, success);
assert_eq!(response.error.as_deref(), error);
assert_eq!(response.data.as_deref(), Some(data.as_bytes()));
}
}
#[tokio::test]
async fn test_convert_to_heal_request_bucket() {
let heal_manager = create_test_heal_manager();
let processor = HealChannelProcessor::new(heal_manager);
let channel_request = HealChannelRequest {
id: "test-id".to_string(),
bucket: "test-bucket".to_string(),
object_prefix: None,
object_version_id: None,
disk: None,
priority: HealChannelPriority::Normal,
scan_mode: None,
remove_corrupted: None,
recreate_missing: None,
update_parity: None,
recursive: None,
dry_run: None,
no_lock: None,
timeout_seconds: None,
pool_index: None,
set_index: None,
force_start: false,
source: HealRequestSource::Internal,
};
let heal_request = processor.convert_to_heal_request(channel_request).unwrap();
assert_eq!(heal_request.id, "test-id");
assert!(matches!(heal_request.heal_type, HealType::Bucket { .. }));
assert_eq!(heal_request.priority, HealPriority::Normal);
}
#[tokio::test]
async fn test_convert_to_heal_request_cluster() {
let heal_manager = create_test_heal_manager();
let processor = HealChannelProcessor::new(heal_manager);
let channel_request = HealChannelRequest {
id: "test-id".to_string(),
bucket: String::new(),
object_prefix: None,
object_version_id: None,
disk: None,
priority: HealChannelPriority::High,
scan_mode: Some(HealScanMode::Normal),
remove_corrupted: Some(false),
recreate_missing: Some(true),
update_parity: Some(true),
recursive: Some(true),
dry_run: Some(false),
no_lock: None,
timeout_seconds: None,
pool_index: None,
set_index: None,
force_start: false,
source: HealRequestSource::Admin,
};
let heal_request = processor.convert_to_heal_request(channel_request).unwrap();
assert!(matches!(heal_request.heal_type, HealType::Cluster));
assert!(heal_request.options.recursive);
}
#[tokio::test]
async fn test_convert_to_heal_request_object() {
let heal_manager = create_test_heal_manager();
let processor = HealChannelProcessor::new(heal_manager);
let channel_request = HealChannelRequest {
id: "test-id".to_string(),
bucket: "test-bucket".to_string(),
object_prefix: Some("test-object".to_string()),
object_version_id: None,
disk: None,
priority: HealChannelPriority::High,
scan_mode: Some(HealScanMode::Deep),
remove_corrupted: Some(true),
recreate_missing: Some(true),
update_parity: Some(true),
recursive: Some(false),
dry_run: Some(false),
no_lock: Some(true),
timeout_seconds: Some(300),
pool_index: Some(0),
set_index: Some(1),
force_start: false,
source: HealRequestSource::Scanner,
};
let heal_request = processor.convert_to_heal_request(channel_request).unwrap();
assert!(matches!(heal_request.heal_type, HealType::Object { .. }));
assert_eq!(heal_request.priority, HealPriority::High);
assert_eq!(heal_request.source, HealRequestSource::Scanner);
assert_eq!(heal_request.options.scan_mode, HealScanMode::Deep);
assert!(heal_request.options.remove_corrupted);
assert!(heal_request.options.recreate_missing);
assert!(heal_request.options.no_lock);
}
#[tokio::test]
async fn test_convert_to_heal_request_scanner_defaults_recreate_missing_false() {
let heal_manager = create_test_heal_manager();
let processor = HealChannelProcessor::new(heal_manager);
let channel_request = HealChannelRequest {
id: "test-id".to_string(),
bucket: "test-bucket".to_string(),
object_prefix: Some("test-object".to_string()),
object_version_id: None,
disk: None,
priority: HealChannelPriority::Low,
scan_mode: None,
remove_corrupted: None,
recreate_missing: None,
update_parity: None,
recursive: Some(false),
dry_run: None,
no_lock: None,
timeout_seconds: None,
pool_index: None,
set_index: None,
force_start: false,
source: HealRequestSource::Scanner,
};
let heal_request = processor.convert_to_heal_request(channel_request).unwrap();
assert!(!heal_request.options.recreate_missing);
}
#[tokio::test]
async fn test_convert_to_heal_request_non_scanner_defaults_recreate_missing_true() {
let heal_manager = create_test_heal_manager();
let processor = HealChannelProcessor::new(heal_manager);
for source in [
HealRequestSource::Admin,
HealRequestSource::AutoHeal,
HealRequestSource::Internal,
HealRequestSource::ReadRepair,
] {
let channel_request = HealChannelRequest {
id: format!("test-id-{source:?}"),
bucket: "test-bucket".to_string(),
object_prefix: Some("test-object".to_string()),
object_version_id: None,
disk: None,
priority: HealChannelPriority::Normal,
scan_mode: None,
remove_corrupted: None,
recreate_missing: None,
update_parity: None,
recursive: Some(false),
dry_run: None,
no_lock: None,
timeout_seconds: None,
pool_index: None,
set_index: None,
force_start: false,
source,
};
let heal_request = processor.convert_to_heal_request(channel_request).unwrap();
assert!(heal_request.options.recreate_missing);
}
}
#[tokio::test]
async fn test_convert_to_heal_request_keeps_explicit_recreate_missing() {
let heal_manager = create_test_heal_manager();
let processor = HealChannelProcessor::new(heal_manager);
for (source, recreate_missing) in [
(HealRequestSource::Scanner, true),
(HealRequestSource::Admin, false),
(HealRequestSource::AutoHeal, false),
(HealRequestSource::Internal, false),
(HealRequestSource::ReadRepair, false),
] {
let channel_request = HealChannelRequest {
id: format!("test-id-{source:?}-{recreate_missing}"),
bucket: "test-bucket".to_string(),
object_prefix: Some("test-object".to_string()),
object_version_id: None,
disk: None,
priority: HealChannelPriority::Normal,
scan_mode: None,
remove_corrupted: None,
recreate_missing: Some(recreate_missing),
update_parity: None,
recursive: Some(false),
dry_run: None,
no_lock: None,
timeout_seconds: None,
pool_index: None,
set_index: None,
force_start: false,
source,
};
let heal_request = processor.convert_to_heal_request(channel_request).unwrap();
assert_eq!(heal_request.options.recreate_missing, recreate_missing);
}
}
#[tokio::test]
async fn test_convert_to_heal_request_prefix_when_recursive() {
let heal_manager = create_test_heal_manager();
let processor = HealChannelProcessor::new(heal_manager);
let channel_request = HealChannelRequest {
id: "test-id".to_string(),
bucket: "test-bucket".to_string(),
object_prefix: Some("logs/".to_string()),
object_version_id: None,
disk: None,
priority: HealChannelPriority::High,
scan_mode: Some(HealScanMode::Normal),
remove_corrupted: Some(false),
recreate_missing: Some(true),
update_parity: Some(true),
recursive: Some(true),
dry_run: Some(false),
no_lock: None,
timeout_seconds: None,
pool_index: None,
set_index: None,
force_start: false,
source: HealRequestSource::Admin,
};
let heal_request = processor.convert_to_heal_request(channel_request).unwrap();
assert!(matches!(
heal_request.heal_type,
HealType::Prefix { ref bucket, ref prefix } if bucket == "test-bucket" && prefix == "logs/"
));
assert!(heal_request.options.recursive);
}
#[tokio::test]
async fn test_convert_to_heal_request_erasure_set() {
let heal_manager = create_test_heal_manager();
let processor = HealChannelProcessor::new(heal_manager);
let channel_request = HealChannelRequest {
id: "test-id".to_string(),
bucket: "test-bucket".to_string(),
object_prefix: None,
object_version_id: None,
disk: Some("pool_0_set_1".to_string()),
priority: HealChannelPriority::Critical,
scan_mode: None,
remove_corrupted: None,
recreate_missing: None,
update_parity: None,
recursive: None,
dry_run: None,
no_lock: None,
timeout_seconds: None,
pool_index: None,
set_index: None,
force_start: false,
source: HealRequestSource::Internal,
};
let heal_request = processor.convert_to_heal_request(channel_request).unwrap();
assert!(matches!(heal_request.heal_type, HealType::ErasureSet { .. }));
assert_eq!(heal_request.priority, HealPriority::Urgent);
}
#[tokio::test]
async fn test_convert_to_heal_request_invalid_disk_id() {
let heal_manager = create_test_heal_manager();
let processor = HealChannelProcessor::new(heal_manager);
let channel_request = HealChannelRequest {
id: "test-id".to_string(),
bucket: "test-bucket".to_string(),
object_prefix: None,
object_version_id: None,
disk: Some("invalid-disk-id".to_string()),
priority: HealChannelPriority::Normal,
scan_mode: None,
remove_corrupted: None,
recreate_missing: None,
update_parity: None,
recursive: None,
dry_run: None,
no_lock: None,
timeout_seconds: None,
pool_index: None,
set_index: None,
force_start: false,
source: HealRequestSource::Internal,
};
let result = processor.convert_to_heal_request(channel_request);
assert!(result.is_err());
}
#[tokio::test]
async fn test_convert_to_heal_request_priority_mapping() {
let heal_manager = create_test_heal_manager();
let processor = HealChannelProcessor::new(heal_manager);
let priorities = vec![
(HealChannelPriority::Low, HealPriority::Low),
(HealChannelPriority::Normal, HealPriority::Normal),
(HealChannelPriority::High, HealPriority::High),
(HealChannelPriority::Critical, HealPriority::Urgent),
];
for (channel_priority, expected_heal_priority) in priorities {
let channel_request = HealChannelRequest {
id: "test-id".to_string(),
bucket: "test-bucket".to_string(),
object_prefix: None,
object_version_id: None,
disk: None,
priority: channel_priority,
scan_mode: None,
remove_corrupted: None,
recreate_missing: None,
update_parity: None,
recursive: None,
dry_run: None,
no_lock: None,
timeout_seconds: None,
pool_index: None,
set_index: None,
force_start: false,
source: HealRequestSource::Internal,
};
let heal_request = processor.convert_to_heal_request(channel_request).unwrap();
assert_eq!(heal_request.priority, expected_heal_priority);
}
}
#[tokio::test]
async fn test_convert_to_heal_request_force_start() {
let heal_manager = create_test_heal_manager();
let processor = HealChannelProcessor::new(heal_manager);
let channel_request = HealChannelRequest {
id: "test-id".to_string(),
bucket: "test-bucket".to_string(),
object_prefix: None,
object_version_id: None,
disk: None,
priority: HealChannelPriority::Normal,
scan_mode: None,
remove_corrupted: Some(false),
recreate_missing: Some(false),
update_parity: Some(false),
recursive: None,
dry_run: None,
no_lock: None,
timeout_seconds: None,
pool_index: None,
set_index: None,
force_start: true, // Admission force only; must not override explicit heal options.
source: HealRequestSource::Internal,
};
let heal_request = processor.convert_to_heal_request(channel_request).unwrap();
assert!(heal_request.force_start);
assert!(!heal_request.options.remove_corrupted);
assert!(!heal_request.options.recreate_missing);
assert!(!heal_request.options.update_parity);
}
#[tokio::test]
async fn test_convert_to_heal_request_empty_object_prefix() {
let heal_manager = create_test_heal_manager();
let processor = HealChannelProcessor::new(heal_manager);
let channel_request = HealChannelRequest {
id: "test-id".to_string(),
bucket: "test-bucket".to_string(),
object_prefix: Some("".to_string()), // Empty prefix should be treated as bucket heal
object_version_id: None,
disk: None,
priority: HealChannelPriority::Normal,
scan_mode: None,
remove_corrupted: None,
recreate_missing: None,
update_parity: None,
recursive: None,
dry_run: None,
no_lock: None,
timeout_seconds: None,
pool_index: None,
set_index: None,
force_start: false,
source: HealRequestSource::Internal,
};
let heal_request = processor.convert_to_heal_request(channel_request).unwrap();
assert!(matches!(heal_request.heal_type, HealType::Bucket { .. }));
}
#[tokio::test]
async fn test_process_start_request_returns_admission_result() {
let storage: Arc<dyn HealStorageAPI> = Arc::new(MockStorage);
let manager = Arc::new(HealManager::new(
storage,
Some(HealConfig {
queue_size: 1,
..HealConfig::default()
}),
));
let processor = HealChannelProcessor::new(manager);
let request = HealChannelRequest {
id: "admission-id".to_string(),
bucket: "bucket".to_string(),
object_prefix: Some("object".to_string()),
object_version_id: None,
disk: None,
priority: HealChannelPriority::Low,
scan_mode: Some(HealScanMode::Normal),
remove_corrupted: None,
recreate_missing: None,
update_parity: None,
recursive: None,
dry_run: None,
no_lock: None,
timeout_seconds: None,
pool_index: None,
set_index: None,
force_start: false,
source: HealRequestSource::Internal,
};
let (tx, rx) = oneshot::channel();
processor
.process_start_request(request.clone(), false, true, tx)
.await
.expect("first admission should succeed");
let first = rx
.await
.expect("oneshot should resolve")
.expect("admission should be returned");
assert_eq!(first.result, HealAdmissionResult::Accepted);
assert_eq!(first.task_id, "admission-id");
let mut duplicate = request;
duplicate.id = "duplicate-id".to_string();
let (tx, rx) = oneshot::channel();
processor
.process_start_request(duplicate, false, true, tx)
.await
.expect("duplicate admission should succeed");
let merged = rx
.await
.expect("oneshot should resolve")
.expect("admission should be returned");
assert_eq!(merged.result, HealAdmissionResult::Merged);
assert_eq!(merged.task_id, "admission-id");
}
#[tokio::test]
async fn test_legacy_start_preserves_duplicate_alias_and_submitted_response_id() {
let manager = create_test_heal_manager();
let processor = HealChannelProcessor::new(manager.clone());
let request = HealChannelRequest {
id: "legacy-original-id".to_string(),
bucket: "bucket".to_string(),
object_prefix: Some("object".to_string()),
priority: HealChannelPriority::Low,
source: HealRequestSource::Admin,
..Default::default()
};
let mut responses = rustfs_common::heal_channel::subscribe_heal_responses();
let (tx, rx) = oneshot::channel();
processor
.process_command(HealChannelCommand::Start {
request: request.clone(),
response_tx: tx,
})
.await
.expect("legacy start should be processed");
assert_eq!(
rx.await
.expect("legacy response should arrive")
.expect("legacy start should succeed"),
HealAdmissionResult::Accepted
);
let mut duplicate = request;
duplicate.id = "legacy-duplicate-id".to_string();
let (tx, rx) = oneshot::channel();
processor
.process_command(HealChannelCommand::Start {
request: duplicate,
response_tx: tx,
})
.await
.expect("legacy duplicate should be processed");
assert_eq!(
rx.await
.expect("legacy duplicate response should arrive")
.expect("legacy duplicate should succeed"),
HealAdmissionResult::Merged
);
tokio::time::timeout(Duration::from_secs(1), async {
loop {
let response = responses.recv().await.expect("legacy broadcast should stay open");
if response.request_id == "legacy-duplicate-id" {
break;
}
}
})
.await
.expect("legacy broadcast must retain the submitted request id");
assert_eq!(
manager
.get_task_status_for_path("bucket/object", "legacy-duplicate-id")
.await
.expect("legacy alias should resolve"),
HealTaskStatus::Pending
);
manager
.cancel_task("legacy-duplicate-id")
.await
.expect("legacy alias should cancel the canonical task");
}
#[tokio::test]
async fn test_public_legacy_start_processes_commands_with_receipt_channel_closed() {
let manager = create_test_heal_manager();
let mut processor = HealChannelProcessor::new(manager);
let (command_tx, command_rx) = mpsc::unbounded_channel();
let processor_task = tokio::spawn(async move { processor.start(command_rx).await });
let (response_tx, response_rx) = oneshot::channel();
command_tx
.send(HealChannelCommand::Start {
request: HealChannelRequest {
id: "legacy-start-loop-id".to_string(),
bucket: "bucket".to_string(),
object_prefix: Some("object".to_string()),
source: HealRequestSource::Admin,
..Default::default()
},
response_tx,
})
.expect("legacy command should send");
assert_eq!(
tokio::time::timeout(Duration::from_secs(1), response_rx)
.await
.expect("legacy processor should not starve")
.expect("legacy response should arrive")
.expect("legacy start should succeed"),
HealAdmissionResult::Accepted
);
drop(command_tx);
tokio::time::timeout(Duration::from_secs(1), processor_task)
.await
.expect("legacy processor should stop when command channel closes")
.expect("legacy processor task should join")
.expect("legacy processor should stop cleanly");
}
#[tokio::test]
async fn test_receipt_channel_returns_canonical_id_end_to_end() {
let manager = create_test_heal_manager();
let mut processor = HealChannelProcessor::new(manager);
let (command_tx, command_rx) = mpsc::unbounded_channel();
let (receipt_tx, receipt_rx) = mpsc::unbounded_channel();
let processor_task = tokio::spawn(async move { processor.start_with_receipts(command_rx, receipt_rx).await });
let request = HealChannelRequest {
id: "receipt-original-id".to_string(),
bucket: "bucket".to_string(),
object_prefix: Some("object".to_string()),
source: HealRequestSource::Admin,
..Default::default()
};
let (response_tx, response_rx) = oneshot::channel();
receipt_tx
.send(HealReceiptCommand {
request: request.clone(),
response_tx,
})
.expect("receipt command should send");
let accepted = tokio::time::timeout(Duration::from_secs(1), response_rx)
.await
.expect("receipt processor should not starve")
.expect("receipt response should arrive")
.expect("receipt start should succeed");
assert_eq!(accepted.result, HealAdmissionResult::Accepted);
assert_eq!(accepted.task_id, "receipt-original-id");
let mut duplicate = request;
duplicate.id = "receipt-duplicate-id".to_string();
let (response_tx, response_rx) = oneshot::channel();
receipt_tx
.send(HealReceiptCommand {
request: duplicate,
response_tx,
})
.expect("duplicate receipt command should send");
let merged = tokio::time::timeout(Duration::from_secs(1), response_rx)
.await
.expect("duplicate receipt should not starve")
.expect("duplicate receipt response should arrive")
.expect("duplicate receipt should succeed");
assert_eq!(merged.result, HealAdmissionResult::Merged);
assert_eq!(merged.task_id, "receipt-original-id");
drop(receipt_tx);
drop(command_tx);
tokio::time::timeout(Duration::from_secs(1), processor_task)
.await
.expect("receipt processor should stop when channels close")
.expect("receipt processor task should join")
.expect("receipt processor should stop cleanly");
}
#[tokio::test]
async fn direct_control_execution_preserves_target_dedup_and_token_ownership() {
let manager = create_test_heal_manager();
let processor = HealChannelProcessor::new(manager);
let original_id = uuid::Uuid::new_v4().to_string();
let request = HealChannelRequest {
id: original_id.clone(),
bucket: "bucket-a".to_string(),
object_prefix: Some("object".to_string()),
priority: HealChannelPriority::High,
source: HealRequestSource::Admin,
..Default::default()
};
let accepted = processor
.execute_start_request(request.clone())
.await
.expect("first direct start should be admitted");
assert_eq!(accepted.result, HealAdmissionResult::Accepted);
assert_eq!(accepted.task_id, original_id);
let mut duplicate = request;
duplicate.id = uuid::Uuid::new_v4().to_string();
let merged = processor
.execute_start_request(duplicate)
.await
.expect("duplicate direct start should merge");
assert_eq!(merged.result, HealAdmissionResult::Merged);
assert_eq!(merged.task_id, accepted.task_id);
let other = processor
.execute_start_request(HealChannelRequest {
id: uuid::Uuid::new_v4().to_string(),
bucket: "bucket-b".to_string(),
object_prefix: Some("object".to_string()),
priority: HealChannelPriority::High,
source: HealRequestSource::Admin,
..Default::default()
})
.await
.expect("different target should be admitted independently");
assert_eq!(other.result, HealAdmissionResult::Accepted);
assert_ne!(other.task_id, accepted.task_id);
let status = processor
.execute_query_request(String::new(), accepted.task_id.clone())
.await
.expect("canonical token should be queryable");
assert!(status.success);
let cancelled = processor
.execute_cancel_request(String::new(), accepted.task_id.clone())
.await
.expect("canonical token should be cancellable");
assert!(cancelled.success);
let stopped = processor
.execute_query_request(String::new(), accepted.task_id)
.await
.expect("cancelled task status should remain queryable");
assert!(stopped.success);
assert_eq!(stopped.error.as_deref(), Some("heal task not found or expired"));
let status: serde_json::Value = serde_json::from_slice(stopped.data.as_deref().unwrap()).unwrap();
assert_eq!(status["summary"], "notFound");
}
#[tokio::test]
async fn test_process_start_request_returns_error_on_invalid_request() {
let heal_manager = create_test_heal_manager();
let processor = HealChannelProcessor::new(heal_manager);
let request = HealChannelRequest {
id: "invalid-id".to_string(),
bucket: "bucket".to_string(),
object_prefix: None,
object_version_id: None,
disk: Some("invalid".to_string()),
priority: HealChannelPriority::Normal,
scan_mode: None,
remove_corrupted: None,
recreate_missing: None,
update_parity: None,
recursive: None,
dry_run: None,
no_lock: None,
timeout_seconds: None,
pool_index: None,
set_index: None,
force_start: false,
source: HealRequestSource::Internal,
};
let (tx, rx) = oneshot::channel();
processor
.process_start_request(request, false, true, tx)
.await
.expect("processor should surface invalid request through response channel");
assert!(rx.await.expect("oneshot should resolve").is_err());
}
#[tokio::test]
async fn test_process_query_request_reports_not_found_when_task_is_unknown() {
let heal_manager = create_test_heal_manager();
let processor = HealChannelProcessor::new(heal_manager);
let (tx, rx) = oneshot::channel();
processor
.process_query_request("bucket".to_string(), "completed-token".to_string(), tx)
.await
.expect("query should process");
let response = rx
.await
.expect("oneshot should resolve")
.expect("query response should be returned");
assert!(response.success);
assert_eq!(response.request_id, "completed-token");
assert_eq!(response.error.as_deref(), Some("heal task not found or expired"));
let payload: serde_json::Value =
serde_json::from_slice(response.data.as_deref().expect("status payload should be present"))
.expect("status payload should be json");
assert_eq!(payload["summary"], "notFound");
assert_eq!(payload["items"].as_array().expect("items should be an array").len(), 0);
}
#[tokio::test]
async fn test_process_query_request_reports_running_for_queued_task() {
let heal_manager = create_test_heal_manager();
let request = HealRequest::bucket("bucket".to_string());
let task_id = request.id.clone();
assert_eq!(
heal_manager
.submit_heal_request(request)
.await
.expect("request should be accepted"),
HealAdmissionResult::Accepted
);
let processor = HealChannelProcessor::new(heal_manager);
let (tx, rx) = oneshot::channel();
processor
.process_query_request("bucket".to_string(), task_id.clone(), tx)
.await
.expect("query should process");
let response = rx
.await
.expect("oneshot should resolve")
.expect("query response should be returned");
assert!(response.success);
assert_eq!(response.request_id, task_id);
let payload: serde_json::Value =
serde_json::from_slice(response.data.as_deref().expect("status payload should be present"))
.expect("status payload should be json");
assert_eq!(payload["summary"], "running");
assert_eq!(payload["items"].as_array().expect("items should be an array").len(), 0);
}
#[tokio::test]
async fn test_process_query_request_rejects_wrong_token_for_active_path() {
let heal_manager = create_test_heal_manager();
let request = HealRequest::bucket("bucket".to_string());
assert_eq!(
heal_manager
.submit_heal_request(request)
.await
.expect("request should be accepted"),
HealAdmissionResult::Accepted
);
let processor = HealChannelProcessor::new(heal_manager);
let (tx, rx) = oneshot::channel();
processor
.process_query_request("bucket".to_string(), "wrong-token".to_string(), tx)
.await
.expect("query should process");
let response = rx
.await
.expect("oneshot should resolve")
.expect("query response should be returned");
assert!(!response.success);
assert_eq!(response.request_id, "wrong-token");
assert_eq!(response.error.as_deref(), Some("invalid heal client token"));
}
#[tokio::test]
async fn test_process_query_request_empty_path_ignores_unrelated_tasks() {
let heal_manager = create_test_heal_manager();
heal_manager
.submit_heal_request(HealRequest::bucket("bucket".to_string()))
.await
.expect("request should be accepted");
let processor = HealChannelProcessor::new(heal_manager);
let (tx, rx) = oneshot::channel();
processor
.process_query_request(String::new(), "wrong-token".to_string(), tx)
.await
.expect("query should process");
let response = rx
.await
.expect("oneshot should resolve")
.expect("query response should be returned");
assert!(response.success);
let payload: serde_json::Value =
serde_json::from_slice(response.data.as_deref().expect("status payload should be present"))
.expect("status payload should be json");
assert_eq!(payload["summary"], "notFound");
assert_eq!(payload["items"].as_array().expect("items should be an array").len(), 0);
}
#[tokio::test]
async fn test_process_query_request_empty_path_uses_client_token_directly() {
let heal_manager = create_test_heal_manager();
let request = HealRequest::new(
HealType::ErasureSet {
buckets: vec![],
set_disk_id: "pool_0_set_1".to_string(),
},
HealOptions::default(),
HealPriority::High,
);
let task_id = request.id.clone();
heal_manager
.submit_heal_request(request)
.await
.expect("request should be accepted");
let processor = HealChannelProcessor::new(heal_manager);
let (tx, rx) = oneshot::channel();
processor
.process_query_request(String::new(), task_id.clone(), tx)
.await
.expect("query should process");
let response = rx
.await
.expect("oneshot should resolve")
.expect("query response should be returned");
assert!(response.success);
assert_eq!(response.request_id, task_id);
assert!(response.error.is_none());
let payload: serde_json::Value =
serde_json::from_slice(response.data.as_deref().expect("status payload should be present"))
.expect("status payload should be json");
assert_eq!(payload["summary"], "running");
assert_eq!(payload["items"].as_array().expect("items should be an array").len(), 0);
}
#[tokio::test]
async fn test_process_cancel_request_cancels_queued_task_by_token() {
let heal_manager = create_test_heal_manager();
let request = HealRequest::bucket("bucket".to_string());
let task_id = request.id.clone();
heal_manager
.submit_heal_request(request)
.await
.expect("request should be accepted");
let processor = HealChannelProcessor::new(heal_manager.clone());
let (tx, rx) = oneshot::channel();
processor
.process_cancel_request("bucket".to_string(), task_id.clone(), tx)
.await
.expect("cancel should process");
let response = rx
.await
.expect("oneshot should resolve")
.expect("cancel response should be returned");
assert!(response.success);
assert_eq!(response.request_id, task_id);
assert_eq!(response.data.as_deref(), Some("stopped".as_bytes()));
assert!(matches!(
heal_manager.get_task_status(&response.request_id).await,
Err(crate::Error::TaskNotFound { .. })
));
}
#[tokio::test]
async fn test_process_cancel_request_cancels_queued_task_by_path() {
let heal_manager = create_test_heal_manager();
let request = HealRequest::bucket("bucket".to_string());
let task_id = request.id.clone();
heal_manager
.submit_heal_request(request)
.await
.expect("request should be accepted");
let processor = HealChannelProcessor::new(heal_manager.clone());
let (tx, rx) = oneshot::channel();
processor
.process_cancel_request("bucket".to_string(), String::new(), tx)
.await
.expect("cancel should process");
let response = rx
.await
.expect("oneshot should resolve")
.expect("cancel response should be returned");
assert!(response.success);
assert_eq!(response.request_id, "bucket");
assert_eq!(response.data.as_deref(), Some("stopped".as_bytes()));
assert!(matches!(
heal_manager.get_task_status(&task_id).await,
Err(crate::Error::TaskNotFound { .. })
));
}
#[tokio::test]
async fn test_process_cancel_request_cancels_cluster_task_for_legacy_root_path() {
let heal_manager = create_test_heal_manager();
let cluster_request = HealRequest::new(HealType::Cluster, HealOptions::default(), HealPriority::High);
let cluster_task_id = cluster_request.id.clone();
let bucket_request = HealRequest::bucket("bucket".to_string());
let bucket_task_id = bucket_request.id.clone();
heal_manager
.submit_heal_request(cluster_request)
.await
.expect("cluster request should be accepted");
heal_manager
.submit_heal_request(bucket_request)
.await
.expect("bucket request should be accepted");
let processor = HealChannelProcessor::new(heal_manager.clone());
let (tx, rx) = oneshot::channel();
processor
.process_cancel_request(".".to_string(), String::new(), tx)
.await
.expect("cancel should process");
let response = rx
.await
.expect("oneshot should resolve")
.expect("cancel response should be returned");
assert!(response.success);
assert_eq!(response.request_id, ".");
assert_eq!(response.data.as_deref(), Some("stopped".as_bytes()));
assert!(response.error.is_none());
assert!(matches!(
heal_manager.get_task_status(&cluster_task_id).await,
Err(crate::Error::TaskNotFound { .. })
));
assert_eq!(
heal_manager
.get_task_status(&bucket_task_id)
.await
.expect("bucket request should not match the root path"),
HealTaskStatus::Pending
);
}
#[tokio::test]
async fn test_process_cancel_request_treats_unknown_path_as_stopped() {
let heal_manager = create_test_heal_manager();
let processor = HealChannelProcessor::new(heal_manager);
let (tx, rx) = oneshot::channel();
processor
.process_cancel_request("missing".to_string(), String::new(), tx)
.await
.expect("cancel should process");
let response = rx
.await
.expect("oneshot should resolve")
.expect("cancel response should be returned");
assert!(response.success);
assert_eq!(response.request_id, "missing");
assert_eq!(response.data.as_deref(), Some("stopped".as_bytes()));
assert!(response.error.is_none());
}
#[tokio::test]
async fn test_process_cancel_request_reports_unknown_task() {
let heal_manager = create_test_heal_manager();
let processor = HealChannelProcessor::new(heal_manager);
let (tx, rx) = oneshot::channel();
processor
.process_cancel_request("missing".to_string(), "missing-token".to_string(), tx)
.await
.expect("cancel should process");
let response = rx
.await
.expect("oneshot should resolve")
.expect("cancel response should be returned");
assert!(!response.success);
assert_eq!(response.request_id, "missing-token");
assert!(response.error.unwrap_or_default().contains("Heal task not found"));
}
}