diff --git a/crates/heal-contracts/src/heal_channel.rs b/crates/heal-contracts/src/heal_channel.rs index cef8fa0ab..22618fa7c 100644 --- a/crates/heal-contracts/src/heal_channel.rs +++ b/crates/heal-contracts/src/heal_channel.rs @@ -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 { diff --git a/crates/heal/src/heal/channel.rs b/crates/heal/src/heal/channel.rs index 708352b14..724929c91 100644 --- a/crates/heal/src/heal/channel.rs +++ b/crates/heal/src/heal/channel.rs @@ -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>, ) -> 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(); diff --git a/crates/heal/src/heal/manager.rs b/crates/heal/src/heal/manager.rs index 22f22b644..5eb3c934f 100644 --- a/crates/heal/src/heal/manager.rs +++ b/crates/heal/src/heal/manager.rs @@ -1793,6 +1793,9 @@ impl HealManager { mrf_notice_target: Option, ) -> Result { 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. diff --git a/crates/scanner/src/scanner_folder/item_actions.rs b/crates/scanner/src/scanner_folder/item_actions.rs index 67482603d..614a49550 100644 --- a/crates/scanner/src/scanner_folder/item_actions.rs +++ b/crates/scanner/src/scanner_folder/item_actions.rs @@ -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() } } diff --git a/crates/scanner/src/scanner_folder/tests.rs b/crates/scanner/src/scanner_folder/tests.rs index 05ba24c72..1c934b3c1 100644 --- a/crates/scanner/src/scanner_folder/tests.rs +++ b/crates/scanner/src/scanner_folder/tests.rs @@ -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( diff --git a/rustfs/src/app/multipart_usecase/tests/body_read_tests.rs b/rustfs/src/app/multipart_usecase/tests/body_read_tests.rs index 4444d6acc..e0cdb37d8 100644 --- a/rustfs/src/app/multipart_usecase/tests/body_read_tests.rs +++ b/rustfs/src/app/multipart_usecase/tests/body_read_tests.rs @@ -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::().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!( diff --git a/rustfs/src/app/object/request_body/tests.rs b/rustfs/src/app/object/request_body/tests.rs index cd4adcd3e..8cb2cefab 100644 --- a/rustfs/src/app/object/request_body/tests.rs +++ b/rustfs/src/app/object/request_body/tests.rs @@ -27,6 +27,19 @@ mod protocol; type FrameSender = mpsc::UnboundedSender, 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();