mirror of
https://github.com/rustfs/rustfs.git
synced 2026-10-03 20:20:29 +00:00
fix(heal): assign stable identities to scanner repair requests (#7823)
* fix(heal): assign stable identities to scanner repair requests * test(s3): isolate UploadPart inactivity timing
This commit is contained in:
@@ -383,6 +383,16 @@ pub struct HealChannelRequest {
|
||||
pub source: HealRequestSource,
|
||||
}
|
||||
|
||||
impl HealChannelRequest {
|
||||
/// Create a request with a stable identity before publishing it.
|
||||
pub fn new() -> Self {
|
||||
Self {
|
||||
id: Uuid::new_v4().to_string(),
|
||||
..Default::default()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Heal response from ahm to admin
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct HealChannelResponse {
|
||||
|
||||
@@ -344,11 +344,14 @@ impl HealChannelProcessor {
|
||||
/// Process start request
|
||||
async fn process_start_request(
|
||||
&self,
|
||||
request: HealChannelRequest,
|
||||
mut request: HealChannelRequest,
|
||||
preserve_alias: bool,
|
||||
publish_canonical_id: bool,
|
||||
response_tx: oneshot::Sender<std::result::Result<HealAdmissionReceipt, String>>,
|
||||
) -> Result<()> {
|
||||
if request.id.is_empty() {
|
||||
request.id = uuid::Uuid::new_v4().to_string();
|
||||
}
|
||||
debug!(
|
||||
target: "rustfs::heal::channel",
|
||||
event = EVENT_HEAL_CHANNEL_REQUEST,
|
||||
@@ -1790,6 +1793,145 @@ mod tests {
|
||||
.expect("receipt processor should stop cleanly");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn heal_empty_request_ids_keep_scanner_objects_independent() {
|
||||
let manager = create_test_heal_manager();
|
||||
let mut processor = HealChannelProcessor::new(manager.clone());
|
||||
let request = HealChannelRequest {
|
||||
bucket: "scanner-bucket".to_string(),
|
||||
object_prefix: Some("first".to_string()),
|
||||
source: HealRequestSource::Scanner,
|
||||
priority: HealChannelPriority::Normal,
|
||||
..Default::default()
|
||||
};
|
||||
let mut other = request.clone();
|
||||
other.object_prefix = Some("second".to_string());
|
||||
let (first, second) =
|
||||
tokio::join!(processor.execute_start_request(request.clone()), processor.execute_start_request(other),);
|
||||
let first = first.expect("first scanner object should be admitted");
|
||||
let second = second.expect("second scanner object should be admitted");
|
||||
assert_eq!(first.result, HealAdmissionResult::Accepted);
|
||||
assert_eq!(second.result, HealAdmissionResult::Accepted);
|
||||
assert!(!first.task_id.is_empty());
|
||||
assert!(!second.task_id.is_empty());
|
||||
assert_ne!(first.task_id, second.task_id);
|
||||
let mut published_ids = std::collections::HashSet::new();
|
||||
for _ in 0..2 {
|
||||
let response = processor
|
||||
.response_receiver
|
||||
.try_recv()
|
||||
.expect("admission should publish its response");
|
||||
assert!(response.success);
|
||||
published_ids.insert(response.request_id);
|
||||
}
|
||||
assert_eq!(
|
||||
published_ids,
|
||||
std::collections::HashSet::from([first.task_id.clone(), second.task_id.clone()])
|
||||
);
|
||||
|
||||
let merged = processor
|
||||
.execute_start_request(request.clone())
|
||||
.await
|
||||
.expect("a repeated object should retain its canonical owner");
|
||||
assert_eq!(merged.result, HealAdmissionResult::Merged);
|
||||
assert_eq!(merged.task_id, first.task_id);
|
||||
let mut conflict = request;
|
||||
conflict.id = first.task_id.clone();
|
||||
conflict.object_prefix = Some("conflicting-object".to_string());
|
||||
let conflict = processor
|
||||
.execute_start_request(conflict)
|
||||
.await
|
||||
.expect("conflicting ID should get a receipt");
|
||||
assert_eq!(conflict.result, HealAdmissionResult::Dropped(HealAdmissionDropReason::AlreadyRunning));
|
||||
assert_eq!(conflict.task_id, first.task_id);
|
||||
for id in [&first.task_id, &second.task_id] {
|
||||
let status = processor
|
||||
.execute_query_request(String::new(), id.clone())
|
||||
.await
|
||||
.expect("generated token should be queryable");
|
||||
assert!(status.success);
|
||||
assert_eq!(status.request_id, *id);
|
||||
}
|
||||
let cancelled = processor
|
||||
.execute_cancel_request(String::new(), first.task_id.clone())
|
||||
.await
|
||||
.expect("generated token should cancel its own request");
|
||||
assert!(cancelled.success);
|
||||
assert!(matches!(manager.get_task_status(&first.task_id).await, Err(Error::TaskNotFound { .. })));
|
||||
assert_eq!(
|
||||
manager
|
||||
.get_task_status(&second.task_id)
|
||||
.await
|
||||
.expect("other object must retain its queued owner"),
|
||||
HealTaskStatus::Pending
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn heal_empty_request_ids_legacy_response_matches_admitted_task() {
|
||||
let manager = create_test_heal_manager();
|
||||
let mut processor = HealChannelProcessor::new(manager.clone());
|
||||
let (tx, rx) = oneshot::channel();
|
||||
processor
|
||||
.process_command(HealChannelCommand::Start {
|
||||
request: HealChannelRequest {
|
||||
bucket: "bucket".to_string(),
|
||||
object_prefix: Some("object".to_string()),
|
||||
source: HealRequestSource::Scanner,
|
||||
..Default::default()
|
||||
},
|
||||
response_tx: tx,
|
||||
})
|
||||
.await
|
||||
.expect("legacy start should be processed");
|
||||
assert_eq!(
|
||||
rx.await.expect("start response should arrive").expect("start should succeed"),
|
||||
HealAdmissionResult::Accepted
|
||||
);
|
||||
let response = processor
|
||||
.response_receiver
|
||||
.try_recv()
|
||||
.expect("legacy response should be published");
|
||||
assert!(response.success);
|
||||
assert!(!response.request_id.is_empty());
|
||||
assert_eq!(
|
||||
manager
|
||||
.get_task_status_for_path("bucket/object", &response.request_id)
|
||||
.await
|
||||
.expect("published ID must resolve to the admitted object"),
|
||||
HealTaskStatus::Pending
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn heal_empty_request_ids_are_normalized_at_manager_admission() {
|
||||
let manager = create_test_heal_manager();
|
||||
let mut first = HealRequest::object("bucket".to_string(), "first".to_string(), None);
|
||||
first.id.clear();
|
||||
let mut second = HealRequest::object("bucket".to_string(), "second".to_string(), None);
|
||||
second.id.clear();
|
||||
let (first, second) = tokio::join!(
|
||||
manager.submit_heal_request_with_receipt(first),
|
||||
manager.submit_heal_request_with_receipt(second),
|
||||
);
|
||||
let first = first.expect("first direct request should be admitted");
|
||||
let second = second.expect("second direct request should be admitted");
|
||||
assert_eq!(first.result, HealAdmissionResult::Accepted);
|
||||
assert_eq!(second.result, HealAdmissionResult::Accepted);
|
||||
assert!(!first.task_id.is_empty());
|
||||
assert!(!second.task_id.is_empty());
|
||||
assert_ne!(first.task_id, second.task_id);
|
||||
for id in [&first.task_id, &second.task_id] {
|
||||
assert_eq!(
|
||||
manager
|
||||
.get_task_status(id)
|
||||
.await
|
||||
.expect("canonical owner should be queryable"),
|
||||
HealTaskStatus::Pending
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn direct_control_execution_preserves_target_dedup_and_token_ownership() {
|
||||
let manager = create_test_heal_manager();
|
||||
|
||||
@@ -1793,6 +1793,9 @@ impl HealManager {
|
||||
mrf_notice_target: Option<MrfRepairNoticeTarget>,
|
||||
) -> Result<HealAdmissionReceipt> {
|
||||
let admission_start = Instant::now();
|
||||
if request.id.is_empty() {
|
||||
request.id = uuid::Uuid::new_v4().to_string();
|
||||
}
|
||||
let source = request.source;
|
||||
let force_start = request.force_start;
|
||||
// Keep ordinary STARTs outside forceStart's cancel-then-admit window.
|
||||
|
||||
@@ -325,7 +325,7 @@ pub(super) fn build_bucket_heal_request(bucket: String, priority: HealChannelPri
|
||||
priority,
|
||||
recreate_missing: Some(false),
|
||||
source: HealRequestSource::Scanner,
|
||||
..Default::default()
|
||||
..HealChannelRequest::new()
|
||||
}
|
||||
}
|
||||
|
||||
@@ -345,7 +345,7 @@ pub(super) fn build_object_heal_request(
|
||||
remove_corrupted: Some(HEAL_DELETE_DANGLING),
|
||||
recreate_missing: Some(false),
|
||||
source: HealRequestSource::Scanner,
|
||||
..Default::default()
|
||||
..HealChannelRequest::new()
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -1072,6 +1072,41 @@ async fn test_prune_failed_objects_max_zero_keeps_fresh() {
|
||||
assert!(!scanner.new_cache.info.failed_objects.contains_key("expired"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn scanner_heal_request_builders_assign_distinct_identities() {
|
||||
let requests = [
|
||||
build_bucket_heal_request("bucket".to_string(), HealChannelPriority::Low),
|
||||
build_object_heal_request(
|
||||
"bucket".to_string(),
|
||||
"first".to_string(),
|
||||
None,
|
||||
HealScanMode::Deep,
|
||||
HealChannelPriority::Low,
|
||||
),
|
||||
build_object_heal_request(
|
||||
"bucket".to_string(),
|
||||
"second".to_string(),
|
||||
Some(uuid::Uuid::new_v4().to_string()),
|
||||
HealScanMode::Deep,
|
||||
HealChannelPriority::Low,
|
||||
),
|
||||
build_non_destructive_object_heal_request(
|
||||
"bucket".to_string(),
|
||||
"third".to_string(),
|
||||
HealScanMode::Deep,
|
||||
HealChannelPriority::Low,
|
||||
),
|
||||
];
|
||||
let mut ids = std::collections::HashSet::new();
|
||||
for request in requests {
|
||||
let id = uuid::Uuid::parse_str(&request.id).expect("scanner must assign an ID before publishing its request");
|
||||
assert!(!id.is_nil());
|
||||
assert!(ids.insert(id), "independent scanner requests must not share an identity");
|
||||
assert_eq!(request.clone().id, request.id, "replaying a captured request must preserve its identity");
|
||||
assert_eq!(request.source, HealRequestSource::Scanner);
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_build_object_heal_request_omits_nil_version_id() {
|
||||
let request = build_object_heal_request(
|
||||
|
||||
@@ -172,6 +172,7 @@ async fn assert_body_timeout_storage_lifecycle(part_size: usize, partial_size: u
|
||||
let baseline = temporary_entries(&disks).await;
|
||||
let _ = path_count();
|
||||
let (request, sender, waiting, _) = observed_request(&bucket, &upload.upload_id, part_size);
|
||||
let control = request.extensions.get::<BodyReadControl>().expect("raw body control").clone();
|
||||
sender
|
||||
.send(Ok(Frame::data(Bytes::from(vec![7; partial_size]))))
|
||||
.expect("partial body");
|
||||
@@ -184,16 +185,20 @@ async fn assert_body_timeout_storage_lifecycle(part_size: usize, partial_size: u
|
||||
})
|
||||
.await
|
||||
.expect("storage must request raw input");
|
||||
// Advance only after actual storage demand; filesystem setup and
|
||||
// cleanup run on a real clock and cannot race auto-advance.
|
||||
tokio::time::pause();
|
||||
tokio::time::advance(Duration::from_secs(300)).await;
|
||||
tokio::time::resume();
|
||||
// The producer can request more input while encoded blocks are
|
||||
// still being written. Elapse only the raw body's inactivity so
|
||||
// those concurrent disk writes retain their real deadlines.
|
||||
control.advance_wait_for_test(Duration::from_secs(300));
|
||||
sender.send(Ok(Frame::data(Bytes::new()))).expect("wake waiting raw body");
|
||||
let error = tokio::time::timeout(Duration::from_secs(10), upload_future)
|
||||
.await
|
||||
.expect("inline cleanup must complete")
|
||||
.expect_err("stalled part");
|
||||
assert_eq!(error.code(), &S3ErrorCode::RequestTimeout);
|
||||
assert_eq!(
|
||||
error.code(),
|
||||
&S3ErrorCode::RequestTimeout,
|
||||
"part_size={part_size}, capped={capped}, replacement={replacement}, retry_after_failure={retry_after_failure}: {error:?}"
|
||||
);
|
||||
assert_eq!(path_count(), 1, "failure must exercise {path}");
|
||||
assert!(sender.is_closed(), "producer must release the failed raw body");
|
||||
assert_eq!(
|
||||
|
||||
@@ -27,6 +27,19 @@ mod protocol;
|
||||
|
||||
type FrameSender = mpsc::UnboundedSender<Result<Frame<Bytes>, io::Error>>;
|
||||
|
||||
impl BodyReadControl {
|
||||
/// Elapse this body's active wait without advancing storage timers.
|
||||
/// The fixture must wake the raw body afterward so it polls the deadline.
|
||||
pub(crate) fn advance_wait_for_test(&self, elapsed: Duration) {
|
||||
let mut state = self.0.budget.lock();
|
||||
assert!(
|
||||
state.demand && state.waiting_since.is_some(),
|
||||
"body must already be waiting for raw input"
|
||||
);
|
||||
state.waited = state.waited.saturating_add(elapsed);
|
||||
}
|
||||
}
|
||||
|
||||
fn raw_reader(timeout: Duration) -> (FrameSender, DynReader, BodyReadControl) {
|
||||
let (sender, receiver) = mpsc::unbounded_channel();
|
||||
let control = BodyReadControl::default();
|
||||
|
||||
Reference in New Issue
Block a user