diff --git a/crates/heal/src/heal/channel.rs b/crates/heal/src/heal/channel.rs index 4681d4a44..708352b14 100644 --- a/crates/heal/src/heal/channel.rs +++ b/crates/heal/src/heal/channel.rs @@ -79,6 +79,8 @@ struct HealTaskStatusPayload<'a> { progress: Option<&'a HealProgress>, #[serde(skip_serializing_if = "Option::is_none")] outcome: Option<&'a super::outcome::HealTaskOutcome>, + #[serde(rename = "outcomeStatus", skip_serializing_if = "Option::is_none")] + outcome_status: Option<&'static str>, #[serde(skip_serializing_if = "Option::is_none")] settings: Option<&'a HealOpts>, } @@ -105,6 +107,7 @@ fn encode_heal_task_status_payload( min_seq: sequence.1, progress, outcome, + outcome_status: (outcome.is_none() && matches!(summary, "finished" | "stopped")).then_some("unavailable"), settings, }) .map_err(|e| Error::Serialization(format!("failed to serialize heal task status: {e}")))?; @@ -907,6 +910,35 @@ mod tests { .expect("a freshly constructed processor must accept responses on its channel"); } + #[test] + fn terminal_without_outcome_is_explicitly_unavailable() { + for summary in ["finished", "stopped", "running"] { + let (bytes, _) = encode_heal_status_response(summary, Vec::new(), HealStatusResponseContext::default()) + .expect("status without canonical outcome"); + let payload: serde_json::Value = serde_json::from_slice(&bytes).expect("status JSON"); + assert!(payload.get("outcome").is_none(), "missing counters cannot become zero counters"); + if summary == "running" { + assert!(payload.get("outcomeStatus").is_none()); + } else { + assert_eq!(payload["outcomeStatus"], "unavailable"); + } + } + let mut outcome = super::super::outcome::HealTaskOutcome::default(); + outcome.finish(None); + let (bytes, _) = encode_heal_status_response( + "finished", + Vec::new(), + HealStatusResponseContext { + outcome: Some(&outcome), + ..Default::default() + }, + ) + .expect("known terminal outcome"); + let payload: serde_json::Value = serde_json::from_slice(&bytes).expect("known status JSON"); + assert!(payload.get("outcomeStatus").is_none()); + assert_eq!(payload["outcome"]["execution"]["state"], "completed"); + } + #[test] fn oversized_status_items_are_truncated_before_transport() { let items = vec![HealResultItem { diff --git a/crates/heal/src/heal/manager.rs b/crates/heal/src/heal/manager.rs index 3c5db47ac..3780ea445 100644 --- a/crates/heal/src/heal/manager.rs +++ b/crates/heal/src/heal/manager.rs @@ -2294,7 +2294,8 @@ impl HealManager { source: HealRequestSource, options: &HealOptions, ) -> Result { - let completed = CompletedHealStatus { + let previous = self.completed_heals.lock().await.get(task_id).cloned(); + let mut completed = CompletedHealStatus { outcome: None, progress: None, retained_bytes: std::sync::OnceLock::new(), @@ -2307,6 +2308,17 @@ impl HealManager { next_seq: 0, min_seq: 0, }; + if let Some(previous) = previous.filter(|previous| previous.heal_type == *heal_type) { + completed.progress = previous.progress.clone(); + completed.outcome = previous.outcome.as_deref().cloned().map(|mut outcome| { + outcome.finish(Some(crate::heal::outcome::HealAbortReason::Cancelled)); + Arc::new(outcome) + }); + completed.seqed_items = previous.seqed_items.clone(); + completed.next_seq = previous.next_seq; + completed.min_seq = previous.min_seq; + completed.result_items_truncated = previous.result_items_truncated; + } self.publish_admin_terminal(task_id, heal_type, source, &completed).await } diff --git a/crates/heal/src/heal/manager/queue.rs b/crates/heal/src/heal/manager/queue.rs index c5b01c55d..2607f4256 100644 --- a/crates/heal/src/heal/manager/queue.rs +++ b/crates/heal/src/heal/manager/queue.rs @@ -207,14 +207,22 @@ impl CompletedHealStatus { } pub(super) async fn snapshot(task: &HealTask, status: HealTaskStatus) -> Self { - let seqed_items = task.get_seqed_result_items().await; - let (next_seq, min_seq) = task.result_seq_cursors(); + let (seqed_items, (next_seq, min_seq)) = { + // Freeze the window and its cursors together while a cancelled + // worker may still append its last result. + let items = task.result_items.read().await; + (items.iter().cloned().collect(), task.result_seq_cursors()) + }; + let mut outcome = task.get_outcome().await; + if status == HealTaskStatus::Cancelled { + outcome.finish(Some(crate::heal::outcome::HealAbortReason::Cancelled)); + } let mut snapshot = Self { heal_type: task.heal_type.clone(), options: task.options.clone(), status, progress: Some(task.get_progress().await), - outcome: Some(Arc::new(task.get_outcome().await)), + outcome: Some(Arc::new(outcome)), retained_bytes: std::sync::OnceLock::new(), result_items_truncated: task.result_items_truncated(), completed_at: SystemTime::now(), diff --git a/crates/heal/src/heal/manager/root_recovery.rs b/crates/heal/src/heal/manager/root_recovery.rs index 49c4daaa1..3715f945f 100644 --- a/crates/heal/src/heal/manager/root_recovery.rs +++ b/crates/heal/src/heal/manager/root_recovery.rs @@ -24,6 +24,9 @@ use crate::heal::{DiskStore, RUSTFS_META_BUCKET}; use serde::{Deserialize, Serialize}; use uuid::Uuid; +mod report; +use report::{ROOT_REPORT_PREFIX, RootHealReport, read_report, remove_report, write_report}; + // The metadata bucket already exists and its parent is durable. Creating a // nested journal directory here would also require syncing every ancestor. const ROOT_RECOVERY_PREFIX: &str = "root-heal-"; @@ -205,7 +208,7 @@ struct RootHealIntent { created_at: SystemTime, } -#[derive(Debug, PartialEq, Serialize, Deserialize)] +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] #[serde(deny_unknown_fields)] struct RootHealTerminal { schema: u32, @@ -219,6 +222,21 @@ struct RootHealTerminal { } impl RootHealTerminal { + fn validate(&self, task_id: &str) -> Result<()> { + let _ = terminal_path(task_id)?; + if self.schema != ROOT_TERMINAL_SCHEMA || self.task_id != task_id { + return Err(Error::Other(format!("Unsupported or mismatched root heal terminal record {task_id}"))); + } + self.heal_type.validate()?; + if !matches!( + self.status, + HealTaskStatus::Completed | HealTaskStatus::Cancelled | HealTaskStatus::Failed { .. } + ) { + return Err(Error::Other(format!("Non-terminal root heal receipt {task_id}"))); + } + Ok(()) + } + fn from_completed(task_id: &str, completed: &CompletedHealStatus) -> Self { Self { schema: ROOT_TERMINAL_SCHEMA, @@ -298,6 +316,8 @@ pub(super) struct RootHealRecovery { disabled_for_tests: bool, #[cfg(test)] disks: Option>, + #[cfg(test)] + pub(super) fail_after_terminal_write: std::sync::atomic::AtomicBool, } pub(super) fn is_admin_heal_recovery(heal_type: &HealType, source: HealRequestSource) -> bool { @@ -384,16 +404,7 @@ fn decode_terminal(task_id: &str, bytes: &[u8]) -> Result { let _ = terminal_path(task_id)?; let terminal: RootHealTerminal = serde_json::from_slice(bytes) .map_err(|error| Error::Other(format!("Invalid root heal terminal record {task_id}: {error}")))?; - if terminal.schema != ROOT_TERMINAL_SCHEMA || terminal.task_id != task_id { - return Err(Error::Other(format!("Unsupported or mismatched root heal terminal record {task_id}"))); - } - terminal.heal_type.validate()?; - if !matches!( - terminal.status, - HealTaskStatus::Completed | HealTaskStatus::Cancelled | HealTaskStatus::Failed { .. } - ) { - return Err(Error::Other(format!("Non-terminal root heal receipt {task_id}"))); - } + terminal.validate(task_id)?; Ok(terminal) } @@ -403,6 +414,7 @@ pub(super) struct RootTerminalGcReport { pub(super) retained: usize, pub(super) pending_removed: usize, pub(super) terminals_removed: usize, + pub(super) reports_removed: usize, pub(super) budget_exhausted: bool, } @@ -413,6 +425,7 @@ impl RootHealRecovery { mutation: Mutex::new(()), disabled_for_tests: false, disks: Some(disks), + fail_after_terminal_write: Default::default(), } } @@ -423,6 +436,8 @@ impl RootHealRecovery { disabled_for_tests: true, #[cfg(test)] disks: None, + #[cfg(test)] + fail_after_terminal_write: Default::default(), } } @@ -504,14 +519,22 @@ impl RootHealRecovery { disks: &[DiskStore], task_id: &str, terminal: RootHealTerminal, + completed: &CompletedHealStatus, ) -> Result> { let path = terminal_path(task_id)?; - if let Some((_, bytes)) = Self::find_terminal(disks, task_id).await? { + if let Some((disk, bytes)) = Self::find_terminal(disks, task_id).await? { let current = decode_terminal(task_id, &bytes)?; - if current == terminal { - return Self::find(disks, task_id).await; + // Only cancellation has a second publication: the worker may finish + // recording its last object after cancellation retired active ownership. + let cancelled_refinement = current.status == HealTaskStatus::Cancelled + && terminal.status == HealTaskStatus::Cancelled + && current.heal_type == terminal.heal_type + && current.options == terminal.options; + if current != terminal && !cancelled_refinement { + return Err(Error::Other(format!("Root heal terminal record changed for {task_id}"))); } - return Err(Error::Other(format!("Root heal terminal record changed for {task_id}"))); + write_report(&disk, task_id, &RootHealReport::from_completed(current, completed)).await?; + return Self::find(disks, task_id).await; } let pending = Self::find(disks, task_id).await?; let disk = pending @@ -521,6 +544,9 @@ impl RootHealRecovery { .ok_or_else(|| Error::Other("No local disk available for root heal terminal receipt".to_string()))?; let bytes = serde_json::to_vec(&terminal) .map_err(|error| Error::Other(format!("Serialize root heal terminal receipt: {error}")))?; + // RUSTFS_COMPAT_TODO(backlog-2519): preserve rollback fences. Remove after supported readers accept standalone reports. + // A crash before the marker leaves the pending responsibility intact. + write_report(&disk, task_id, &RootHealReport::from_completed(terminal, completed)).await?; match EcstoreDiskAPI::compare_and_update_file(disk.as_ref(), RUSTFS_META_BUCKET, &path, None, Some(bytes.into())).await? { EcstoreConditionalFileUpdate::Updated => Ok(pending), _ => Err(Error::Other(format!("Root heal terminal record changed for {task_id}"))), @@ -697,7 +723,8 @@ impl RootHealRecovery { let pending = decode_intent(task_id, &bytes)?; let heal_type = HealType::from(pending.heal_type); let terminal = RootHealTerminal::cancelled(task_id, &heal_type, pending.options); - let _ = Self::persist_terminal_locked(&disks, task_id, terminal).await?; + let completed = terminal.clone().into_completed(); + let _ = Self::persist_terminal_locked(&disks, task_id, terminal, &completed).await?; match EcstoreDiskAPI::compare_and_update_file( disk.as_ref(), RUSTFS_META_BUCKET, @@ -729,7 +756,15 @@ impl RootHealRecovery { let _guard = self.mutation.lock().await; let disks = self.disks().await?; let pending = - Self::persist_terminal_locked(&disks, task_id, RootHealTerminal::from_completed(task_id, completed)).await?; + Self::persist_terminal_locked(&disks, task_id, RootHealTerminal::from_completed(task_id, completed), completed) + .await?; + #[cfg(test)] + if self + .fail_after_terminal_write + .swap(false, std::sync::atomic::Ordering::SeqCst) + { + return Err(Error::Io(std::io::Error::other("injected failure after terminal report publication"))); + } if let Some((disk, bytes)) = pending { match EcstoreDiskAPI::compare_and_update_file( disk.as_ref(), @@ -761,10 +796,20 @@ impl RootHealRecovery { } let _guard = self.mutation.lock().await; let disks = self.disks().await?; - let Some((_, bytes)) = Self::find_retained_terminal(&disks, task_id, SystemTime::now()).await? else { + let Some((disk, bytes)) = Self::find_retained_terminal(&disks, task_id, SystemTime::now()).await? else { return Ok(None); }; - Ok(Some(decode_terminal(task_id, &bytes)?.into_completed())) + let terminal = decode_terminal(task_id, &bytes)?; + if let Some((report, _)) = read_report(&disk, task_id).await? { + if *report.terminal() != terminal { + return Err(Error::Other(format!("Root heal report does not match its terminal receipt {task_id}"))); + } + return Ok(Some(report.into_completed())); + } + // Legacy receipts cannot reconstruct counters or the retained result window. + let mut completed = terminal.into_completed(); + completed.result_items_truncated = true; + Ok(Some(completed)) } pub(super) async fn completed_matches_path(&self, heal_path: &str) -> Result { @@ -831,6 +876,7 @@ impl RootHealRecovery { report.scanned += 1; let Some(task_id) = entry .strip_prefix(ROOT_TERMINAL_PREFIX) + .or_else(|| entry.strip_prefix(ROOT_REPORT_PREFIX)) .and_then(|entry| entry.strip_suffix(".json")) else { continue; @@ -849,6 +895,23 @@ impl RootHealRecovery { break; } let Some((terminal_disk, terminal_bytes)) = Self::find_terminal(&disks, &task_id).await? else { + // A crash before terminal publication, or an old-version GC, + // can leave an uncommitted report. It never authorizes replay. + for disk in &disks { + if deletes >= ROOT_TERMINAL_GC_DELETE_BUDGET { + report.budget_exhausted = true; + break; + } + if let Some((snapshot, bytes)) = read_report(disk, &task_id).await? { + if snapshot.terminal().retained_at(now) { + report.retained += 1; + } else { + remove_report(disk, &task_id, bytes).await?; + deletes += 1; + report.reports_removed += 1; + } + } + } continue; }; let terminal = decode_terminal(&task_id, &terminal_bytes)?; @@ -879,6 +942,19 @@ impl RootHealRecovery { } } } + if let Some((snapshot, bytes)) = read_report(&terminal_disk, &task_id).await? { + if *snapshot.terminal() != terminal { + return Err(Error::Other(format!("Root heal report does not match expired receipt {task_id}"))); + } + remove_report(&terminal_disk, &task_id, bytes).await?; + deletes += 1; + report.reports_removed += 1; + if deletes >= ROOT_TERMINAL_GC_DELETE_BUDGET { + report.retained += 1; + report.budget_exhausted = true; + break; + } + } match EcstoreDiskAPI::compare_and_update_file( terminal_disk.as_ref(), RUSTFS_META_BUCKET, diff --git a/crates/heal/src/heal/manager/root_recovery/report.rs b/crates/heal/src/heal/manager/root_recovery/report.rs new file mode 100644 index 000000000..22d0e711a --- /dev/null +++ b/crates/heal/src/heal/manager/root_recovery/report.rs @@ -0,0 +1,221 @@ +// Copyright 2026 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 super::*; +use crate::heal::outcome::{HealAbortReason, HealExecutionOutcome}; +use std::io::Write; +use tokio::io::AsyncReadExt; + +// This prefix must not overlap either namespace understood by schema-1 readers. +pub(super) const ROOT_REPORT_PREFIX: &str = "heal-terminal-report-"; +const ROOT_REPORT_SCHEMA: u32 = 1; +pub(super) const MAX_ROOT_REPORT_BYTES: usize = 8 * 1024 * 1024; + +/// The unchanged legacy terminal is the commit marker and rollback fence. +/// Its original timestamp also anchors refinements of a cancelled worker's report. +#[derive(Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub(super) struct RootHealReport { + schema: u32, + terminal: RootHealTerminal, + outcome: Option, + progress: Option, + result_items_truncated: bool, + seqed_items: Vec<(u64, HealResultItem)>, + next_seq: u64, + min_seq: u64, +} + +impl RootHealReport { + pub(super) fn from_completed(terminal: RootHealTerminal, completed: &CompletedHealStatus) -> Self { + Self { + schema: ROOT_REPORT_SCHEMA, + terminal, + outcome: completed.outcome.as_deref().cloned(), + progress: completed.progress.clone(), + result_items_truncated: completed.result_items_truncated, + seqed_items: completed.seqed_items.clone(), + next_seq: completed.next_seq, + min_seq: completed.min_seq, + } + } + + pub(super) fn terminal(&self) -> &RootHealTerminal { + &self.terminal + } + + pub(super) fn into_completed(self) -> CompletedHealStatus { + CompletedHealStatus { + outcome: self.outcome.map(Arc::new), + progress: self.progress, + result_items_truncated: self.result_items_truncated, + seqed_items: self.seqed_items, + next_seq: self.next_seq, + min_seq: self.min_seq, + ..self.terminal.into_completed() + } + } + + pub(super) fn encode(&self) -> Result { + struct BoundedWriter(Vec); + impl Write for BoundedWriter { + fn write(&mut self, bytes: &[u8]) -> std::io::Result { + if bytes.len() > MAX_ROOT_REPORT_BYTES.saturating_sub(self.0.len()) { + return Err(std::io::Error::other("root heal report exceeds its size limit")); + } + self.0.extend_from_slice(bytes); + Ok(bytes.len()) + } + + fn flush(&mut self) -> std::io::Result<()> { + Ok(()) + } + } + let mut writer = BoundedWriter(Vec::new()); + serde_json::to_writer(&mut writer, self) + .map_err(|error| Error::Serialization(format!("Serialize root heal report: {error}")))?; + Ok(writer.0.into()) + } +} + +pub(super) fn report_path(task_id: &str) -> Result { + let _ = terminal_path(task_id)?; + Ok(format!("{ROOT_REPORT_PREFIX}{task_id}.json")) +} + +fn decode_report(task_id: &str, bytes: &[u8]) -> Result { + let mut report: RootHealReport = serde_json::from_slice(bytes) + .map_err(|error| Error::Serialization(format!("Invalid root heal report {task_id}: {error}")))?; + report.terminal.validate(task_id)?; + if report.schema != ROOT_REPORT_SCHEMA { + return Err(Error::Other(format!("Unsupported root heal report schema for {task_id}"))); + } + if let Some(outcome) = &report.outcome + && (matches!(outcome.execution, HealExecutionOutcome::Pending | HealExecutionOutcome::Running) + || (report.terminal.status == HealTaskStatus::Cancelled + && outcome.execution != HealExecutionOutcome::Aborted(HealAbortReason::Cancelled))) + { + return Err(Error::Other(format!("Non-terminal or mismatched heal outcome for {task_id}"))); + } + let mut previous = None; + for (seq, item) in &mut report.seqed_items { + if *seq < report.min_seq || *seq >= report.next_seq || previous.is_some_and(|previous| previous >= *seq) { + return Err(Error::Other(format!("Invalid root heal report cursor for {task_id}"))); + } + previous = Some(*seq); + // Serde's string growth must not inflate the original retention accounting. + for value in [ + &mut item.heal_item_type, + &mut item.bucket, + &mut item.object, + &mut item.version_id, + &mut item.detail, + ] { + value.shrink_to_fit(); + } + for infos in [&mut item.before, &mut item.after] { + infos.drives.shrink_to_fit(); + for drive in &mut infos.drives { + drive.uuid.shrink_to_fit(); + drive.endpoint.shrink_to_fit(); + drive.state.shrink_to_fit(); + } + } + } + if report.min_seq > report.next_seq { + return Err(Error::Other(format!("Invalid root heal report cursor range for {task_id}"))); + } + let mut window = CompletedHealStatus { + seqed_items: std::mem::take(&mut report.seqed_items), + ..report.terminal.clone().into_completed() + }; + let count = window.seqed_items.len(); + window.bound_result_window(); + if window.seqed_items.len() != count { + return Err(Error::Other(format!("Root heal report result window exceeds its limit for {task_id}"))); + } + report.seqed_items = window.seqed_items; + Ok(report) +} + +pub(super) async fn read_report(disk: &DiskStore, task_id: &str) -> Result> { + let path = report_path(task_id)?; + let reader = match EcstoreDiskAPI::read_file(disk.as_ref(), RUSTFS_META_BUCKET, &path).await { + Ok(reader) => reader, + Err(DiskError::FileNotFound) => return Ok(None), + Err(error) => return Err(error.into()), + }; + let mut bytes = Vec::new(); + reader + .take(u64::try_from(MAX_ROOT_REPORT_BYTES + 1).map_err(Error::other)?) + .read_to_end(&mut bytes) + .await?; + if bytes.len() > MAX_ROOT_REPORT_BYTES { + return Err(Error::Other(format!("Root heal report exceeds its size limit for {task_id}"))); + } + Ok(Some((decode_report(task_id, &bytes)?, bytes.into()))) +} + +pub(super) async fn write_report(disk: &DiskStore, task_id: &str, report: &RootHealReport) -> Result<()> { + let bytes = report.encode()?; + let previous = read_report(disk, task_id).await?; + if previous.as_ref().is_some_and(|(_, current)| *current == bytes) { + return Ok(()); + } + match EcstoreDiskAPI::compare_and_update_file( + disk.as_ref(), + RUSTFS_META_BUCKET, + &report_path(task_id)?, + previous.map(|(_, bytes)| bytes), + Some(bytes), + ) + .await? + { + EcstoreConditionalFileUpdate::Updated => Ok(()), + _ => Err(Error::Other(format!("Root heal report changed while publishing {task_id}"))), + } +} + +pub(super) async fn remove_report(disk: &DiskStore, task_id: &str, bytes: EcstoreDiskBytes) -> Result<()> { + match EcstoreDiskAPI::compare_and_update_file(disk.as_ref(), RUSTFS_META_BUCKET, &report_path(task_id)?, Some(bytes), None) + .await? + { + EcstoreConditionalFileUpdate::Updated => Ok(()), + _ => Err(Error::Other(format!("Root heal report changed while pruning {task_id}"))), + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn root_report_encoding_rejects_oversized_payload() { + let task_id = Uuid::new_v4().to_string(); + let terminal = RootHealTerminal::cancelled(&task_id, &HealType::Cluster, HealOptions::default()); + let mut completed = terminal.clone().into_completed(); + completed.progress = Some(HealProgress { + current_object: Some("x".repeat(MAX_ROOT_REPORT_BYTES)), + ..Default::default() + }); + let report = RootHealReport::from_completed(terminal, &completed); + assert!( + report + .encode() + .expect_err("bounded writer must reject oversized JSON") + .to_string() + .contains("size limit") + ); + } +} diff --git a/crates/heal/src/heal/manager/scheduler.rs b/crates/heal/src/heal/manager/scheduler.rs index 6a05e913a..bd0f4e80f 100644 --- a/crates/heal/src/heal/manager/scheduler.rs +++ b/crates/heal/src/heal/manager/scheduler.rs @@ -302,8 +302,8 @@ impl HealManager { tests::pause_completed_retention_before_publish(&task_id, &completed_status).await; let mut active_heals_guard = active_heals_clone.lock().await; let owns_completion = active_heals_guard.contains_key(&task_id); - let cancelled_completion = if owns_completion { - false + let cancelled_snapshot = if owns_completion { + None } else { // Cancellation can win while a finished worker waits // for active ownership. It must not resurrect a retry @@ -313,18 +313,21 @@ impl HealManager { .lock() .await .get(&task_id) - .is_some_and(|completed| completed.status == HealTaskStatus::Cancelled) + .filter(|completed| completed.status == HealTaskStatus::Cancelled) + .cloned() }; - if cancelled_completion { + let cancelled_completion = cancelled_snapshot.is_some(); + if let Some(cancelled) = &cancelled_snapshot { completed_status = HealTaskStatus::Cancelled; completed_status_entry.status = HealTaskStatus::Cancelled; completed_status_entry.outcome = Some(Arc::new(task.get_outcome().await)); + completed_status_entry.completed_at = cancelled.completed_at; } let terminal_completion = matches!( completed_status, HealTaskStatus::Completed | HealTaskStatus::Cancelled | HealTaskStatus::Failed { .. } ); - if owns_completion + if (owns_completion || cancelled_completion) && terminal_completion && root_recovery::is_admin_heal_recovery(&task.heal_type, task.source) && let Err(error) = root_recovery_clone @@ -344,6 +347,24 @@ impl HealManager { error = %error, "Failed to publish heal terminal receipt" ); + // A failed write can already have replaced the report. + // Resolve it from disk instead of assuming the old snapshot won. + if let Some(cancelled) = &cancelled_snapshot { + completed_status_entry = match root_recovery_clone.completed(&task_id).await { + Ok(Some(persisted)) => persisted, + Ok(None) | Err(_) => { + let mut unavailable = (**cancelled).clone(); + unavailable.outcome = None; + unavailable.progress = None; + unavailable.seqed_items = Vec::new(); + unavailable.next_seq = 0; + unavailable.min_seq = 0; + unavailable.result_items_truncated = true; + unavailable.retained_bytes.take(); + unavailable + } + }; + } } if owns_completion && !terminal_completion diff --git a/crates/heal/src/heal/manager/tests/root_recovery.rs b/crates/heal/src/heal/manager/tests/root_recovery.rs index 463ad5750..a1db5f447 100644 --- a/crates/heal/src/heal/manager/tests/root_recovery.rs +++ b/crates/heal/src/heal/manager/tests/root_recovery.rs @@ -134,6 +134,600 @@ fn completed_admin_status(heal_type: &HealType, completed_at: SystemTime) -> Com } } +#[tokio::test] +async fn root_recovery_terminal_outcome_survives_repeated_restart() { + use crate::heal::outcome::{HealExecutionOutcome, HealTraversalCoverage}; + + let (_temp, disk) = recovery_disk().await; + let manager = recovery_manager(vec![disk.clone()]); + let request = root_request(); + manager.root_recovery.persist(&request).await.expect("admit root heal"); + let mut completed = completed_admin_status(&request.heal_type, SystemTime::now()); + let mut outcome = HealTaskOutcome::default(); + outcome.execution = HealExecutionOutcome::Completed; + outcome.coverage = HealTraversalCoverage::Complete; + outcome.counters.processed = 257; + outcome.counters.healed = 255; + outcome.counters.unchanged = 2; + completed.outcome = Some(Arc::new(outcome)); + let expected = serde_json::to_value(completed.outcome.as_deref()).expect("expected outcome"); + manager + .publish_admin_terminal(&request.id, &request.heal_type, request.source, &completed) + .await + .expect("publish root outcome"); + drop(manager); + + for _ in 0..2 { + let restarted = recovery_manager(vec![disk.clone()]); + restarted.replay_root_heals().await.expect("replay terminal receipt"); + assert_eq!(restarted.get_queue_length().await, 0); + let report = restarted.get_task_report(&request.id).await.expect("retained report"); + assert_eq!(report.status, HealTaskStatus::Completed); + assert_eq!(serde_json::to_value(report.outcome.as_deref()).expect("restored outcome"), expected); + let retained = restarted + .root_recovery + .completed(&request.id) + .await + .expect("receipt") + .expect("retained"); + assert_eq!(retained.completed_at, completed.completed_at, "restart cannot renew retention"); + } +} + +fn terminal_with_outcome(heal_type: &HealType, status: HealTaskStatus) -> CompletedHealStatus { + use crate::heal::outcome::{ + HealAbortReason, HealFailureClass, HealObjectDisposition, HealObjectIdentity, HealObjectKind, HealObjectOutcome, + }; + let mut completed = completed_admin_status(heal_type, SystemTime::now()); + let mut outcome = HealTaskOutcome::default(); + outcome.start(); + for index in 0..257 { + let disposition = if index < 255 { + HealObjectDisposition::Repaired + } else { + HealObjectDisposition::VerifiedHealthy + }; + outcome.record(HealObjectOutcome { + identity: HealObjectIdentity { + kind: HealObjectKind::Object, + bucket: "bucket".to_string(), + object: format!("object-{index}"), + version_id: None, + bucket_incarnation_id: None, + pool_index: None, + set_index: None, + }, + disposition, + detail: Some("verified repair".to_string()), + }); + } + if matches!(status, HealTaskStatus::Failed { .. }) { + let mut failure = outcome.objects.back().expect("last object").clone(); + failure.disposition = HealObjectDisposition::Failed(HealFailureClass::Permanent); + outcome.record(failure); + outcome.attempt_failed(); + } + outcome.finish((status == HealTaskStatus::Cancelled).then_some(HealAbortReason::Cancelled)); + completed.outcome = Some(Arc::new(outcome)); + completed.status = status; + completed.seqed_items = vec![( + 9, + HealResultItem { + object: "retained-object".to_string(), + ..Default::default() + }, + )]; + completed.next_seq = 10; + completed.min_seq = 9; + completed.result_items_truncated = true; + completed +} + +#[tokio::test] +async fn root_recovery_terminal_report_preserves_states_windows_and_legacy_marker() { + for status in [ + HealTaskStatus::Completed, + HealTaskStatus::Cancelled, + HealTaskStatus::Failed { + error: "permanent failure".to_string(), + }, + ] { + let (_temp, disk) = recovery_disk().await; + let manager = recovery_manager(vec![disk.clone()]); + let request = root_request(); + let completed = terminal_with_outcome(&request.heal_type, status.clone()); + manager.root_recovery.persist(&request).await.expect("admission"); + manager + .publish_admin_terminal(&request.id, &request.heal_type, request.source, &completed) + .await + .expect("terminal report"); + let marker = disk + .read_all(RUSTFS_META_BUCKET, &format!("terminal-root-heal-{}.json", request.id)) + .await + .expect("legacy marker"); + let marker_json: serde_json::Value = serde_json::from_slice(&marker).expect("legacy JSON"); + let keys = marker_json + .as_object() + .expect("terminal object") + .keys() + .map(String::as_str) + .collect::>(); + assert_eq!( + keys, + HashSet::from([ + "schema", + "task_id", + "heal_type", + "status", + "options", + "progress", + "completed_at" + ]) + ); + assert_eq!(marker_json["schema"], 1, "rollback readers retain their original format"); + let expected = serde_json::to_value(completed.outcome.as_deref()).expect("outcome JSON"); + assert_eq!(expected["objectsTruncated"], true); + assert_eq!(expected["objects"].as_array().expect("window").len(), 128); + drop(manager); + for _ in 0..2 { + let restarted = recovery_manager(vec![disk.clone()]); + restarted.replay_root_heals().await.expect("restart"); + let report = restarted + .get_task_report_since(&request.id, Some(8)) + .await + .expect("incremental report"); + assert_eq!(report.status, status); + assert_eq!(serde_json::to_value(report.outcome.as_deref()).expect("restored outcome"), expected); + assert_eq!(report.result_items.len(), 1); + assert_eq!(report.result_items[0].object, "retained-object"); + assert_eq!((report.next_seq, report.min_seq, report.result_items_truncated), (10, 9, true)); + assert_eq!( + serde_json::to_value(report.progress).expect("progress"), + serde_json::to_value(&completed.progress).expect("expected progress") + ); + } + assert_eq!( + disk.read_all(RUSTFS_META_BUCKET, &format!("terminal-root-heal-{}.json", request.id)) + .await + .expect("unchanged marker"), + marker + ); + } +} + +#[tokio::test] +async fn root_recovery_legacy_terminal_does_not_invent_outcome() { + let (_temp, disk) = recovery_disk().await; + let task_id = "00000000-0000-0000-0000-000000002519"; + let mut marker: serde_json::Value = serde_json::from_str(r#"{"schema":1,"task_id":"00000000-0000-0000-0000-000000002519","heal_type":{"type":"cluster"},"status":"Cancelled","progress":null,"completed_at":{"secs_since_epoch":1,"nanos_since_epoch":0}}"#).expect("pinned legacy terminal"); + marker["completed_at"] = serde_json::to_value(SystemTime::now()).expect("current retention epoch"); + disk.write_all( + RUSTFS_META_BUCKET, + &format!("terminal-root-heal-{task_id}.json"), + serde_json::to_vec(&marker).expect("legacy marker").into(), + ) + .await + .expect("write legacy marker"); + let manager = recovery_manager(vec![disk]); + let report = manager.get_task_report(task_id).await.expect("legacy report"); + assert_eq!(report.status, HealTaskStatus::Cancelled); + assert!(report.outcome.is_none(), "historical counters are unavailable"); + assert!(report.result_items_truncated, "missing historical detail must be explicit"); +} + +#[cfg(unix)] +#[tokio::test] +async fn root_recovery_cancelled_worker_publishes_durable_refinement_or_keeps_previous_report() { + use crate::heal::outcome::{HealAbortReason, HealExecutionOutcome}; + for failure_mode in 0..4 { + let (temp, disk) = recovery_disk().await; + let manager = recovery_manager(vec![disk.clone()]); + let bucket = format!("terminal-report-cancel-{}", Uuid::new_v4()); + let request = admin_request(HealType::Object { + bucket: bucket.clone(), + object: "object".to_string(), + version_id: None, + }); + let task_id = request.id.clone(); + let hook = Arc::new(CompletedRetentionHook { + pause_before_publish: true, + ..Default::default() + }); + { + let mut hooks = COMPLETED_RETENTION_HOOKS.lock().await; + hooks.insert(bucket.clone(), hook.clone()); + hooks.insert(task_id.clone(), hook.clone()); + } + manager.submit_heal_request(request).await.expect("admit cancellable object"); + process_manager_queue_once(&manager).await; + tokio::time::timeout(Duration::from_secs(10), hook.started.notified()) + .await + .expect("worker entered storage"); + manager.cancel_task(&task_id).await.expect("durable cancellation"); + let initial = manager.get_task_report(&task_id).await.expect("initial cancelled report"); + assert_eq!( + initial.outcome.as_ref().expect("initial outcome").execution, + HealExecutionOutcome::Aborted(HealAbortReason::Cancelled) + ); + let completed_at = manager + .root_recovery + .completed(&task_id) + .await + .expect("receipt") + .expect("retained") + .completed_at; + hook.execute.notify_one(); + tokio::time::timeout(Duration::from_secs(10), hook.before_publish.notified()) + .await + .expect("worker finalized outcome"); + let read_only = + matches!(failure_mode, 1 | 3).then(|| RestoreDirectoryMode::read_only(temp.path().join(RUSTFS_META_BUCKET))); + if failure_mode == 3 { + use std::os::unix::fs::PermissionsExt as _; + std::fs::set_permissions(temp.path().join(RUSTFS_META_BUCKET), std::fs::Permissions::from_mode(0o0)) + .expect("make the report owner unreadable as well as unwritable"); + } + manager + .root_recovery + .fail_after_terminal_write + .store(failure_mode == 2, std::sync::atomic::Ordering::SeqCst); + hook.publish.notify_one(); + tokio::time::timeout(Duration::from_secs(10), hook.handoff.notified()) + .await + .expect("worker published completion"); + let final_report = manager.get_task_report(&task_id).await.expect("final cancelled report"); + let expected = if failure_mode == 3 { + assert!( + final_report.outcome.is_none(), + "an uncertain unreadable report must be marked unavailable" + ); + assert!(final_report.result_items_truncated); + assert!(final_report.result_items.is_empty()); + serde_json::to_value(initial.outcome.as_deref()).expect("last durable report after I/O recovers") + } else { + serde_json::to_value(final_report.outcome.as_deref()).expect("final outcome") + }; + if failure_mode == 1 { + assert_eq!( + expected, + serde_json::to_value(initial.outcome.as_deref()).expect("previous durable outcome") + ); + } else if failure_mode != 3 { + assert!( + final_report.outcome.as_ref().expect("final counters").counters.processed + > initial.outcome.as_ref().expect("initial counters").counters.processed + ); + } + drop(read_only); + for _ in 0..2 { + let restarted = recovery_manager(vec![disk.clone()]); + restarted.replay_root_heals().await.expect("cancelled restart"); + assert_eq!(restarted.get_queue_length().await, 0); + let restored = restarted.get_task_report(&task_id).await.expect("restored cancellation"); + assert_eq!(serde_json::to_value(restored.outcome.as_deref()).expect("restored outcome"), expected); + assert_eq!( + restarted + .root_recovery + .completed(&task_id) + .await + .expect("receipt") + .expect("retained") + .completed_at, + completed_at + ); + } + hook.finish.notify_one(); + COMPLETED_RETENTION_HOOKS + .lock() + .await + .retain(|key, _| key != &bucket && key != &task_id); + } +} + +#[tokio::test] +async fn root_recovery_retry_cancellation_preserves_previous_attempt_outcome() { + use crate::heal::outcome::{HealAbortReason, HealExecutionOutcome}; + let (_temp, disk) = recovery_disk().await; + let manager = recovery_manager(vec![disk.clone()]); + let request = root_request(); + manager.root_recovery.persist(&request).await.expect("durable retry intent"); + let previous = terminal_with_outcome( + &request.heal_type, + HealTaskStatus::Failed { + error: "retryable".to_string(), + }, + ); + let counters = previous.outcome.as_ref().expect("attempt outcome").counters.clone(); + manager + .completed_heals + .lock() + .await + .insert(request.id.clone(), Arc::new(previous)); + manager.retrying_heals.lock().await.insert( + request.id.clone(), + RetryingHeal { + request: request.clone(), + error: "retryable".to_string(), + cancel_token: CancellationToken::new(), + }, + ); + manager.cancel_task(&request.id).await.expect("cancel retry"); + drop(manager); + let restarted = recovery_manager(vec![disk]); + let report = restarted + .get_task_report(&request.id) + .await + .expect("retained retry cancellation"); + let outcome = report.outcome.expect("previous attempt retained"); + assert_eq!(outcome.execution, HealExecutionOutcome::Aborted(HealAbortReason::Cancelled)); + assert_eq!(outcome.counters, counters); + assert_eq!(report.result_items.len(), 1); +} + +#[tokio::test] +async fn root_recovery_corrupt_or_oversize_reports_do_not_resurrect_terminal_work() { + let (_temp, disk) = recovery_disk().await; + let manager = recovery_manager(vec![disk.clone()]); + let request = root_request(); + let completed = terminal_with_outcome(&request.heal_type, HealTaskStatus::Completed); + manager.root_recovery.persist(&request).await.expect("admission"); + manager + .publish_admin_terminal(&request.id, &request.heal_type, request.source, &completed) + .await + .expect("terminal"); + let path = format!("heal-terminal-report-{}.json", request.id); + let original = disk.read_all(RUSTFS_META_BUCKET, &path).await.expect("valid report"); + let original_json: serde_json::Value = serde_json::from_slice(&original).expect("report JSON"); + for (pointer, value) in [ + ("/schema", serde_json::json!(2)), + ("/terminal/task_id", serde_json::json!(Uuid::new_v4().to_string())), + ("/terminal/completed_at/secs_since_epoch", serde_json::json!(1)), + ("/outcome/execution", serde_json::json!({"state":"running"})), + ("/outcome/counters/processed", serde_json::json!(0)), + ("/outcome/objects/0/detail", serde_json::json!("x".repeat(1025))), + ("/min_seq", serde_json::json!(11)), + ] { + let mut corrupted = original_json.clone(); + *corrupted.pointer_mut(pointer).expect("existing field") = value; + disk.write_all(RUSTFS_META_BUCKET, &path, serde_json::to_vec(&corrupted).expect("corrupt JSON").into()) + .await + .expect("inject corruption"); + assert!(manager.get_task_report(&request.id).await.is_err(), "must reject {pointer}"); + assert!( + manager + .root_recovery + .pending() + .await + .expect("terminal still fences replay") + .is_empty() + ); + } + disk.write_all(RUSTFS_META_BUCKET, &path, b"{".to_vec().into()) + .await + .expect("inject malformed JSON"); + assert!(manager.get_task_report(&request.id).await.is_err()); + let mut boundary = original.to_vec(); + boundary.resize(8 * 1024 * 1024, b' '); + disk.write_all(RUSTFS_META_BUCKET, &path, boundary.clone().into()) + .await + .expect("exact-size valid JSON"); + assert!( + manager + .get_task_report(&request.id) + .await + .expect("exact byte limit is accepted") + .outcome + .is_some() + ); + boundary.push(b' '); + disk.write_all(RUSTFS_META_BUCKET, &path, boundary.into()) + .await + .expect("limit plus one"); + assert!( + manager + .get_task_report(&request.id) + .await + .expect_err("oversized otherwise-valid JSON") + .to_string() + .contains("size limit") + ); + disk.write_all(RUSTFS_META_BUCKET, &path, original) + .await + .expect("restore valid report"); + assert!( + manager + .get_task_report(&request.id) + .await + .expect("valid report restored") + .outcome + .is_some() + ); +} + +#[cfg(unix)] +#[tokio::test] +async fn root_recovery_report_write_failure_preserves_pending_owner() { + let (temp, disk) = recovery_disk().await; + let (_other_temp, other) = recovery_disk().await; + let manager = recovery_manager(vec![disk.clone(), other.clone()]); + let request = root_request(); + manager.root_recovery.persist(&request).await.expect("admission"); + let intent_path = format!("root-heal-{}.json", request.id); + let original = disk + .read_all(RUSTFS_META_BUCKET, &intent_path) + .await + .expect("original responsibility"); + let completed = terminal_with_outcome(&request.heal_type, HealTaskStatus::Completed); + let read_only = RestoreDirectoryMode::read_only(temp.path().join(RUSTFS_META_BUCKET)); + assert!( + manager + .publish_admin_terminal(&request.id, &request.heal_type, request.source, &completed) + .await + .is_err() + ); + assert_eq!( + disk.read_all(RUSTFS_META_BUCKET, &intent_path) + .await + .expect("responsibility retained"), + original + ); + for disk in [&disk, &other] { + assert!(matches!( + disk.read_all(RUSTFS_META_BUCKET, &format!("terminal-root-heal-{}.json", request.id)) + .await, + Err(DiskError::FileNotFound) + )); + assert!(matches!( + disk.read_all(RUSTFS_META_BUCKET, &format!("heal-terminal-report-{}.json", request.id)) + .await, + Err(DiskError::FileNotFound) + )); + } + drop(read_only); + manager + .publish_admin_terminal(&request.id, &request.heal_type, request.source, &completed) + .await + .expect("retry on original owner"); +} + +#[tokio::test] +async fn root_recovery_old_cancelled_token_cannot_stop_successor() { + let (_temp, disk) = recovery_disk().await; + let manager = recovery_manager(vec![disk.clone()]); + let old = root_request(); + let cancelled = terminal_with_outcome(&old.heal_type, HealTaskStatus::Cancelled); + manager + .publish_admin_terminal(&old.id, &old.heal_type, old.source, &cancelled) + .await + .expect("old cancellation"); + let mut successor = root_request(); + successor.force_start = true; + manager.submit_heal_request(successor.clone()).await.expect("admit successor"); + let path = format!("root-heal-{}.json", successor.id); + let original = disk.read_all(RUSTFS_META_BUCKET, &path).await.expect("successor owner"); + manager.cancel_task(&old.id).await.expect("repeat old STOP"); + assert_eq!( + manager.get_task_status(&successor.id).await.expect("successor still pending"), + HealTaskStatus::Pending + ); + assert_eq!(disk.read_all(RUSTFS_META_BUCKET, &path).await.expect("same successor intent"), original); + drop(manager); + let restarted = recovery_manager(vec![disk]); + restarted.replay_root_heals().await.expect("successor restart"); + assert_eq!( + restarted + .get_task_status(&successor.id) + .await + .expect("successor survives restart"), + HealTaskStatus::Pending + ); + assert_eq!( + serde_json::to_value( + restarted + .get_task_report(&old.id) + .await + .expect("old report") + .outcome + .as_deref() + ) + .expect("restored"), + serde_json::to_value(cancelled.outcome.as_deref()).expect("original") + ); +} + +#[tokio::test] +async fn root_recovery_uncertain_terminal_publication_preserves_commit_and_pending_fence() { + let (_temp, disk) = recovery_disk().await; + let manager = recovery_manager(vec![disk.clone()]); + let request = root_request(); + manager + .root_recovery + .persist(&request) + .await + .expect("original responsibility"); + let completed = terminal_with_outcome(&request.heal_type, HealTaskStatus::Cancelled); + manager + .root_recovery + .fail_after_terminal_write + .store(true, std::sync::atomic::Ordering::SeqCst); + assert!( + manager + .publish_admin_terminal(&request.id, &request.heal_type, request.source, &completed) + .await + .is_err() + ); + assert!( + disk.read_all(RUSTFS_META_BUCKET, &format!("root-heal-{}.json", request.id)) + .await + .is_ok(), + "uncertain publication must not retire pending ownership" + ); + let restarted = recovery_manager(vec![disk]); + restarted + .replay_root_heals() + .await + .expect("commit marker fences stale pending"); + assert_eq!(restarted.get_queue_length().await, 0); + let report = restarted + .get_task_report(&request.id) + .await + .expect("committed report survives uncertain response"); + assert_eq!( + serde_json::to_value(report.outcome.as_deref()).expect("restored"), + serde_json::to_value(completed.outcome.as_deref()).expect("committed") + ); +} + +#[tokio::test] +async fn root_recovery_orphan_report_never_commits_or_retires_pending_work() { + let (_temp, disk) = recovery_disk().await; + let manager = recovery_manager(vec![disk.clone()]); + let request = root_request(); + let mut completed = terminal_with_outcome(&request.heal_type, HealTaskStatus::Completed); + completed.completed_at = SystemTime::now() - KEEP_HEAL_TASK_STATUS_DURATION - Duration::from_secs(1); + manager + .publish_admin_terminal(&request.id, &request.heal_type, request.source, &completed) + .await + .expect("report fixture"); + let path = format!("terminal-root-heal-{}.json", request.id); + let marker = disk.read_all(RUSTFS_META_BUCKET, &path).await.expect("marker"); + assert_eq!( + crate::heal::storage_api::owner::EcstoreDiskAPI::compare_and_update_file( + disk.as_ref(), + RUSTFS_META_BUCKET, + &path, + Some(marker), + None + ) + .await + .expect("simulate missing commit marker"), + crate::heal::storage_api::owner::EcstoreConditionalFileUpdate::Updated + ); + manager + .root_recovery + .persist(&request) + .await + .expect("original pending responsibility"); + assert_eq!( + manager + .root_recovery + .pending() + .await + .expect("uncommitted report cannot mask pending") + .len(), + 1 + ); + assert!(matches!(manager.get_task_report(&request.id).await, Err(Error::TaskNotFound { .. }))); + let gc = manager + .root_recovery + .gc_terminal_receipts_once(SystemTime::now()) + .await + .expect("orphan GC"); + assert_eq!(gc.reports_removed, 1); + assert_eq!((gc.pending_removed, gc.terminals_removed), (0, 0)); + assert_eq!(manager.root_recovery.pending().await.expect("pending survives GC").len(), 1); +} + #[cfg(unix)] #[tokio::test] async fn root_recovery_new_intent_skips_prepublication_read_only_owner() { @@ -792,7 +1386,8 @@ async fn root_recovery_terminal_gc_is_delete_budget_bounded() { .gc_terminal_receipts_once(now) .await .expect("budgeted terminal GC"); - assert_eq!(report.terminals_removed, 64); + assert_eq!(report.terminals_removed, 32); + assert_eq!(report.reports_removed, 32, "report deletion shares the 64-operation budget"); assert!(report.budget_exhausted); let terminal_entries = disk .list_dir("", RUSTFS_META_BUCKET, "", -1) @@ -801,7 +1396,7 @@ async fn root_recovery_terminal_gc_is_delete_budget_bounded() { .into_iter() .filter(|entry| entry.starts_with("terminal-root-heal-")) .count(); - assert_eq!(terminal_entries, 1); + assert_eq!(terminal_entries, 33); } #[tokio::test] diff --git a/crates/heal/src/heal/outcome.rs b/crates/heal/src/heal/outcome.rs index 613a649b9..1ad85fa5a 100644 --- a/crates/heal/src/heal/outcome.rs +++ b/crates/heal/src/heal/outcome.rs @@ -23,7 +23,7 @@ const MAX_OUTCOME_ITEMS: usize = 128; const MAX_OUTCOME_BYTES: usize = 64 * 1024; const MAX_OUTCOME_DETAIL_BYTES: usize = 1024; -#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)] +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] #[serde(rename_all = "snake_case")] pub enum HealObjectKind { Object, @@ -31,8 +31,8 @@ pub enum HealObjectKind { Decode, } -#[derive(Debug, Clone, PartialEq, Eq, Serialize)] -#[serde(rename_all = "camelCase")] +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "camelCase", deny_unknown_fields)] pub struct HealObjectIdentity { pub kind: HealObjectKind, pub bucket: String, @@ -44,7 +44,7 @@ pub struct HealObjectIdentity { pub set_index: Option, } -#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)] +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] #[serde(rename_all = "snake_case")] pub enum HealDeferredReason { DanglingDeleteGrace, @@ -54,7 +54,7 @@ pub enum HealDeferredReason { Deadline, } -#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)] +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] #[serde(rename_all = "snake_case")] pub enum HealFailureClass { Recoverable, @@ -62,12 +62,13 @@ pub enum HealFailureClass { Permanent, } -#[derive(Debug, Clone, PartialEq, Eq, Serialize)] +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] #[serde( tag = "state", content = "details", rename_all = "snake_case", - rename_all_fields = "camelCase" + rename_all_fields = "camelCase", + deny_unknown_fields )] pub enum HealObjectDisposition { /// The legacy storage response does not prove the requested check or commit. @@ -84,8 +85,8 @@ pub enum HealObjectDisposition { DryRunObserved, } -#[derive(Debug, Clone, PartialEq, Eq, Serialize)] -#[serde(rename_all = "camelCase")] +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "camelCase", deny_unknown_fields)] pub struct HealObjectOutcome { pub identity: HealObjectIdentity, pub disposition: HealObjectDisposition, @@ -126,7 +127,7 @@ impl HealObjectOutcome { } } -#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize)] +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)] #[serde(rename_all = "snake_case")] pub enum HealTraversalCoverage { #[default] @@ -183,6 +184,61 @@ pub struct HealTaskOutcome { untraversable: bool, } +impl<'de> Deserialize<'de> for HealTaskOutcome { + fn deserialize>(deserializer: D) -> Result { + #[derive(Deserialize)] + #[serde(rename_all = "camelCase", deny_unknown_fields)] + struct Snapshot { + execution: HealExecutionOutcome, + coverage: HealTraversalCoverage, + counters: HealOutcomeCounters, + objects: VecDeque, + objects_truncated: bool, + } + + let mut snapshot = Snapshot::deserialize(deserializer)?; + if snapshot.objects.len() > MAX_OUTCOME_ITEMS { + return Err(serde::de::Error::custom("heal outcome object window exceeds its limit")); + } + let mut retained_object_bytes = 0usize; + for object in &mut snapshot.objects { + object.identity.bucket.shrink_to_fit(); + object.identity.object.shrink_to_fit(); + if let Some(version) = &mut object.identity.version_id { + version.shrink_to_fit(); + } + if let Some(detail) = &mut object.detail { + if detail.len() > MAX_OUTCOME_DETAIL_BYTES { + return Err(serde::de::Error::custom("heal outcome detail exceeds its limit")); + } + detail.shrink_to_fit(); + } + retained_object_bytes = retained_object_bytes.saturating_add(object.retained_bytes()); + } + if retained_object_bytes > MAX_OUTCOME_BYTES { + return Err(serde::de::Error::custom("heal outcome bytes exceed their limit")); + } + let counters = &snapshot.counters; + let total = counters + .healed + .checked_add(counters.unchanged) + .and_then(|total| total.checked_add(counters.skipped)) + .and_then(|total| total.checked_add(counters.failed)); + if !counters.overflowed && (total != Some(counters.processed) || counters.unknown > counters.skipped) { + return Err(serde::de::Error::custom("heal outcome counters are inconsistent")); + } + Ok(Self { + execution: snapshot.execution, + coverage: snapshot.coverage, + counters: snapshot.counters, + objects: snapshot.objects, + objects_truncated: snapshot.objects_truncated, + retained_object_bytes, + untraversable: snapshot.execution == HealExecutionOutcome::Aborted(HealAbortReason::Untraversable), + }) + } +} + #[derive(Debug, thiserror::Error)] pub enum HealOutcomeWireError { #[error("heal outcome is missing execution or counters")] @@ -435,6 +491,38 @@ mod canonical_outcome_tests { } } + #[test] + fn persisted_outcome_rebuilds_accounting_and_enforces_window_boundaries() { + let mut outcome = HealTaskOutcome::default(); + for _ in 0..MAX_OUTCOME_ITEMS { + outcome.record(item(HealObjectDisposition::Repaired)); + } + outcome.finish(None); + let value = serde_json::to_value(&outcome).expect("bounded outcome"); + let mut restored: HealTaskOutcome = serde_json::from_value(value.clone()).expect("restore bounded window"); + assert_eq!(restored.retained_object_bytes, outcome.retained_object_bytes); + restored.record(item(HealObjectDisposition::Repaired)); + assert_eq!(restored.objects.len(), MAX_OUTCOME_ITEMS); + assert!(restored.objects_truncated); + let mut oversized = value; + let extra = oversized["objects"][0].clone(); + oversized["objects"].as_array_mut().expect("objects").push(extra); + assert!(serde_json::from_value::(oversized).is_err()); + + let mut object = item(HealObjectDisposition::Repaired); + let fixed_bytes = object.retained_bytes() - object.identity.object.capacity(); + object.identity.object = "x".repeat(MAX_OUTCOME_BYTES - fixed_bytes); + let mut outcome = HealTaskOutcome::default(); + outcome.record(object); + outcome.finish(None); + assert_eq!(outcome.retained_object_bytes, MAX_OUTCOME_BYTES); + let mut value = serde_json::to_value(outcome).expect("exact byte boundary"); + let restored: HealTaskOutcome = serde_json::from_value(value.clone()).expect("exact bound remains readable"); + assert_eq!(restored.retained_object_bytes, MAX_OUTCOME_BYTES); + value["objects"][0]["identity"]["object"] = serde_json::json!("x".repeat(MAX_OUTCOME_BYTES - fixed_bytes + 1)); + assert!(serde_json::from_value::(value).is_err()); + } + #[test] fn canonical_outcome_categories_have_one_terminal_count() { let mut outcome = HealTaskOutcome::default(); diff --git a/docs/architecture/README.md b/docs/architecture/README.md index c151711cf..678bd65ba 100644 --- a/docs/architecture/README.md +++ b/docs/architecture/README.md @@ -35,6 +35,7 @@ Required headings and strings in these files are asserted by `scripts/check_arch | [erasure-coding.md](erasure-coding.md) | changing anything under `crates/ecstore/src/erasure/`, `crates/filemeta/`, `crates/ecstore/src/set_disk/`, storage-class or layout code, or any decode, quorum, or heal boundary (normative spec) | | [placement-repair-invariants.md](placement-repair-invariants.md) | changing anything that resolves an object to a pool, set, or disk, or that admits scanner or heal work | | [heal-concurrency-model.md](heal-concurrency-model.md) | changing heal, PUT/multipart commit, delete, lifecycle expiry, or data-movement code that shares the `(bucket, object)` commit surface, or asking whether RustFS needs a persistent healing marker | +| [heal-terminal-reports.md](heal-terminal-reports.md) | changing retained admin heal outcomes, terminal report publication, cancellation refinement, or report upgrade/rollback behavior | | [unified-object-generation.md](unified-object-generation.md) | adding or changing anything that fences a commit, scopes a read lease, gates old-directory cleanup, binds prepared pool reads, or settles quota against the current object version | | [ilm-tiering-persistence-contracts.md](ilm-tiering-persistence-contracts.md) | changing an ILM transition, tier configuration mutation, manual transition job, tier-delete recovery path, pool decommission, or any code that can create, transfer, or destroy ownership of a remote-tier object | | [decommission-compatibility.md](decommission-compatibility.md) | changing pool decommission or rebalance behavior, its admin API shape, the persisted `PoolMeta` fields, or how tier free versions move between pools | diff --git a/docs/architecture/compat-cleanup-register.md b/docs/architecture/compat-cleanup-register.md index 5fda12709..43317c1db 100644 --- a/docs/architecture/compat-cleanup-register.md +++ b/docs/architecture/compat-cleanup-register.md @@ -11,6 +11,7 @@ ## Open Items +- `backlog-2519` retained admin heal reports: keep the schema-1 terminal as the commit and replay fence, with a bounded versioned report in a separate namespace on the same disk. Older rollback readers can still query the terminal and suppress replay, but cannot expose its outcome; newer readers mark missing outcomes as unavailable. Remove the legacy marker and missing-report adapter only after all supported direct-upgrade and rollback readers understand the report format and retained schema-1-only receipts have expired. - `odm-list-bare-envelope` historical ODM continuation tokens: preserve complete bare v1/v2 envelopes. Framed issuance defaults on for the deployed framed-only generation; upgrades from older bare-only readers must explicitly disable it before starting new nodes and keep it off until reader convergence. Remove the legacy classifier and framing issuance override only after every supported reader accepts framing and outstanding bare listings have drained or clients explicitly restarted them; tokens have no automatic expiry. Exact full-envelope object keys remain intrinsically ambiguous during this compatibility period. - `backlog-2263` legacy heal MRF inspection: retained per-record journals remain readable while committed-snapshot ownership and writer activation are staged. Remove legacy import only after all supported direct-upgrade and rollback readers understand committed snapshots and migration tooling confirms that no retained or restorable legacy journal requires it. This does not enable a new writer or change the automatic legacy consumer. - `backlog-1337` legacy restore orphan recovery: releases that predate the restore worker-lock marker can leave a valid operation-id and `ongoing-request="true"` after cancellation or process failure, with no durable liveness proof. New servers allow an exact, non-nil legacy generation to be superseded only when its consistently parsed request date is at least 24 hours old. Remove the clock-based legacy fallback after the minimum supported direct-upgrade release writes the v1 worker-lock marker on every restore and operators have resolved every retained pre-v1 ongoing generation. diff --git a/docs/architecture/heal-terminal-reports.md b/docs/architecture/heal-terminal-reports.md new file mode 100644 index 000000000..1d7a09140 --- /dev/null +++ b/docs/architecture/heal-terminal-reports.md @@ -0,0 +1,29 @@ +# Retained Admin Heal Reports + +Admin heal tokens retain their terminal status for ten minutes from the original completion time. The persistence owner is `RootHealRecovery` in `crates/heal/src/heal/manager/root_recovery.rs`; the bounded report codec is in `crates/heal/src/heal/manager/root_recovery/report.rs`. + +## Publication and Recovery + +The original schema-1 `terminal-root-heal-.json` remains the commit marker. Its separate `heal-terminal-report-.json` report has its own schema version and embeds the exact terminal identity. Both files belong to the same coordinator disk as the pending intent. + +Publication writes the report through the storage owner's conditional file update, writes the terminal marker, and then conditionally removes the pending intent. Recovery never treats a report without its marker as committed. A committed terminal continues to suppress stale pending work even when its report is unreadable. An initial publication failure leaves the existing responsibility available for recovery. + +The report preserves the canonical execution, traversal coverage, cumulative counters, diagnostic object window, progress, legacy result window, truncation flags, and incremental cursors. Restoring it neither recounts objects nor derives counters from progress. Its embedded terminal must match the retained marker before the report can be returned. + +## Cancellation + +An active cancellation first publishes an `aborted/cancelled` partial snapshot before retiring active ownership. The worker can subsequently finish recording its last object. Its final report replaces the earlier report, while the terminal marker and original completion timestamp remain unchanged. Cancelling a retry preserves any retained preceding attempt's outcome; unavailable historical counters remain unavailable. + +If the final report update fails, the scheduler rereads the disk because an error can occur after publication. It returns that persisted report when readable. If the result cannot be resolved, the cached response exposes no canonical outcome or result window and marks the detail as truncated. It does not assume that the previous report won an uncertain write. + +## Bounds and Retention + +Report encoding and reading each enforce an 8 MiB JSON limit. The canonical diagnostic window retains at most 128 objects and 64 KiB of object storage accounting, with at most 1 KiB per detail. The legacy result window retains the existing 1 MiB memory budget. Decoding validates the report version, terminal identity, terminal execution, counters, object bounds, and cursor ordering, and reconstructs memory accounting. + +GC uses the original completion timestamp, including after repeated restarts or cancellation refinement. For an expired terminal, it removes stale pending work before the report and marker. Report deletions share the existing 64-deletion budget. Expired orphan reports left before marker publication or by rollback GC are removed without deleting pending responsibility. Corrupt records are retained and reported as errors. + +## Upgrade and Rollback + +Schema-1 terminals without reports remain queryable. A terminal response without a canonical outcome adds `outcomeStatus: "unavailable"`; it does not fabricate counters or complete coverage. Missing historical result windows are marked truncated. Responses with an outcome and running responses do not add this field. + +Older binaries retain their original terminal decoder and ignore the distinct report namespace. They continue to query status and prevent replay, but do not expose the new report's outcome. A later upgrade can read a retained report again; reports whose markers were collected by the old binary are handled as uncommitted orphans. Rollback therefore preserves the old status/replay contract, not the new outcome capability.