diff --git a/crates/heal/src/heal/channel.rs b/crates/heal/src/heal/channel.rs index 4ad655d60..322185cbd 100644 --- a/crates/heal/src/heal/channel.rs +++ b/crates/heal/src/heal/channel.rs @@ -281,6 +281,12 @@ impl HealChannelProcessor { 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, @@ -965,6 +971,27 @@ mod tests { )); } + #[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(".".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()); + } + #[tokio::test] async fn test_process_cancel_request_reports_unknown_task() { let heal_manager = create_test_heal_manager(); diff --git a/crates/scanner/src/scanner.rs b/crates/scanner/src/scanner.rs index 40cb02918..f15fe48f3 100644 --- a/crates/scanner/src/scanner.rs +++ b/crates/scanner/src/scanner.rs @@ -162,6 +162,59 @@ fn get_cycle_scan_mode(current_cycle: u64, bitrot_start_cycle: u64, bitrot_start } } +fn background_heal_info_for_scan_start( + mut info: BackgroundHealInfo, + current_cycle: u64, + scan_mode: HealScanMode, + now: DateTime, +) -> Option { + let reset_bitrot_start = scan_mode == HealScanMode::Deep && should_reset_bitrot_start(&info, current_cycle, now); + if info.current_scan_mode == scan_mode && !reset_bitrot_start { + return None; + } + + info.current_scan_mode = scan_mode; + if reset_bitrot_start { + info.bitrot_start_cycle = current_cycle; + info.bitrot_start_time = Some(now); + } + + Some(info) +} + +fn should_reset_bitrot_start(info: &BackgroundHealInfo, current_cycle: u64, now: DateTime) -> bool { + let Some(bitrot_start_time) = info.bitrot_start_time else { + return true; + }; + + let Some(bitrot_cycle) = bitrot_scan_cycle() else { + return false; + }; + + if bitrot_cycle.is_zero() { + return true; + } + + if current_cycle.saturating_sub(info.bitrot_start_cycle) < heal_object_select_prob() as u64 { + return false; + } + + let elapsed = now + .signed_duration_since(bitrot_start_time) + .to_std() + .unwrap_or(Duration::ZERO); + elapsed >= bitrot_cycle +} + +fn background_heal_info_for_scan_complete(mut info: BackgroundHealInfo, scan_mode: HealScanMode) -> Option { + if scan_mode != HealScanMode::Deep || info.current_scan_mode != HealScanMode::Deep { + return None; + } + + info.current_scan_mode = HealScanMode::Normal; + Some(info) +} + fn retain_recent_cycle_completions(cycle_completed: &mut Vec>) { let keep = data_usage_update_dir_cycles() as usize; if cycle_completed.len() > keep { @@ -246,22 +299,17 @@ async fn run_data_scanner_cycle(ctx: &CancellationToken, storeapi: &Arc global_metrics().set_cycle(Some(cycle_info.clone())).await; - let background_heal_info = read_background_heal_info(storeapi.clone()).await; + let mut background_heal_info = read_background_heal_info(storeapi.clone()).await; let scan_mode = get_cycle_scan_mode( cycle_info.current, background_heal_info.bitrot_start_cycle, background_heal_info.bitrot_start_time, ); - if background_heal_info.current_scan_mode != scan_mode { - let mut new_heal_info = background_heal_info.clone(); - new_heal_info.current_scan_mode = scan_mode; - - if scan_mode == HealScanMode::Deep { - new_heal_info.bitrot_start_cycle = cycle_info.current; - new_heal_info.bitrot_start_time = Some(Utc::now()); - } - + if let Some(new_heal_info) = + background_heal_info_for_scan_start(background_heal_info.clone(), cycle_info.current, scan_mode, Utc::now()) + { + background_heal_info = new_heal_info.clone(); save_background_heal_info(storeapi.clone(), new_heal_info).await; } @@ -281,10 +329,16 @@ async fn run_data_scanner_cycle(ctx: &CancellationToken, storeapi: &Arc { error!(duration = ?now.elapsed(), "Fail run data scanner cycle: {e}"); emit_scan_cycle_complete(false, cycle_start.elapsed()); + if let Some(new_heal_info) = background_heal_info_for_scan_complete(background_heal_info.clone(), scan_mode) { + save_background_heal_info(storeapi.clone(), new_heal_info).await; + } return; } done_cycle(); emit_scan_cycle_complete(true, cycle_start.elapsed()); + if let Some(new_heal_info) = background_heal_info_for_scan_complete(background_heal_info.clone(), scan_mode) { + save_background_heal_info(storeapi.clone(), new_heal_info).await; + } cycle_info.next += 1; cycle_info.current = 0; @@ -520,6 +574,68 @@ mod tests { }); } + #[test] + #[serial] + fn test_background_heal_info_for_scan_start_marks_deep_active() { + let now = Utc::now(); + let info = background_heal_info_for_scan_start(BackgroundHealInfo::default(), 7, HealScanMode::Deep, now) + .expect("deep scan should update background heal info"); + + assert_eq!(info.current_scan_mode, HealScanMode::Deep); + assert_eq!(info.bitrot_start_cycle, 7); + assert_eq!(info.bitrot_start_time, Some(now)); + } + + #[test] + #[serial] + fn test_background_heal_info_for_scan_start_keeps_deep_window_start() { + with_var_unset(ENV_SCANNER_BITROT_CYCLE_SECS, || { + let started_at = Utc::now(); + let info = BackgroundHealInfo { + bitrot_start_time: Some(started_at), + bitrot_start_cycle: 7, + current_scan_mode: HealScanMode::Normal, + }; + + let info = background_heal_info_for_scan_start(info, 8, HealScanMode::Deep, Utc::now()) + .expect("deep scan should mark active status"); + + assert_eq!(info.current_scan_mode, HealScanMode::Deep); + assert_eq!(info.bitrot_start_cycle, 7); + assert_eq!(info.bitrot_start_time, Some(started_at)); + }); + } + + #[test] + #[serial] + fn test_background_heal_info_for_scan_complete_marks_deep_idle() { + let started_at = Utc::now(); + let info = BackgroundHealInfo { + bitrot_start_time: Some(started_at), + bitrot_start_cycle: 7, + current_scan_mode: HealScanMode::Deep, + }; + + let info = background_heal_info_for_scan_complete(info, HealScanMode::Deep) + .expect("completed deep scan should update background heal info"); + + assert_eq!(info.current_scan_mode, HealScanMode::Normal); + assert_eq!(info.bitrot_start_cycle, 7); + assert_eq!(info.bitrot_start_time, Some(started_at)); + } + + #[test] + #[serial] + fn test_background_heal_info_for_scan_complete_leaves_normal_scan_unchanged() { + let info = BackgroundHealInfo { + bitrot_start_time: Some(Utc::now()), + bitrot_start_cycle: 7, + current_scan_mode: HealScanMode::Normal, + }; + + assert!(background_heal_info_for_scan_complete(info, HealScanMode::Normal).is_none()); + } + #[test] fn test_retain_recent_cycle_completions_keeps_last_entries() { let base = Utc::now();