From 7428f3138c40067a8e389d6768619c5f6a808756 Mon Sep 17 00:00:00 2001 From: overtrue Date: Sat, 22 Aug 2026 05:08:46 +0800 Subject: [PATCH] fix(ecstore): harden durable ILM cursor receipts --- .config/nextest.toml | 4 +- .../bucket/lifecycle/bucket_lifecycle_ops.rs | 26 + .../src/bucket/lifecycle/durable_namespace.rs | 483 ++++++++++++++---- .../bucket/lifecycle/manual_transition_job.rs | 33 +- crates/ecstore/src/core/pools.rs | 88 +++- crates/ecstore/src/store/init.rs | 106 +++- 6 files changed, 631 insertions(+), 109 deletions(-) diff --git a/.config/nextest.toml b/.config/nextest.toml index a65ebe9fe..8cc4a9017 100644 --- a/.config/nextest.toml +++ b/.config/nextest.toml @@ -81,7 +81,7 @@ test-group = 'ecstore-serial-flaky' # The durable ILM decommission regressions build isolated multi-pool stores and # deliberately take source or target disks offline while checking fencing. [[profile.default.overrides]] -filter = 'package(rustfs-ecstore) & (test(decommission_migrates_and_verifies_registered_durable_ilm_records) | test(decommission_durable_ilm_target_read_error_is_not_masked_by_peer_success) | test(decommission_durable_ilm_terminal_receipt_recovers_failed_source_cleanup))' +filter = 'package(rustfs-ecstore) & (test(decommission_migrates_and_verifies_registered_durable_ilm_records) | test(decommission_durable_ilm_target_read_error_is_not_masked_by_peer_success) | test(decommission_durable_ilm_terminal_receipt_recovers_failed_source_cleanup) | test(decommission_durable_ilm_receipt_pagination_fails_closed_on_second_page))' test-group = 'ecstore-serial-flaky' # Serialize the bucket-incarnation / lifecycle-fence tests. They drive @@ -197,7 +197,7 @@ filter = 'package(rustfs-ecstore) & test(manual_transition_page_checkpoint_persi test-group = 'ecstore-serial-flaky' [[profile.ci.overrides]] -filter = 'package(rustfs-ecstore) & (test(decommission_migrates_and_verifies_registered_durable_ilm_records) | test(decommission_durable_ilm_target_read_error_is_not_masked_by_peer_success) | test(decommission_durable_ilm_terminal_receipt_recovers_failed_source_cleanup))' +filter = 'package(rustfs-ecstore) & (test(decommission_migrates_and_verifies_registered_durable_ilm_records) | test(decommission_durable_ilm_target_read_error_is_not_masked_by_peer_success) | test(decommission_durable_ilm_terminal_receipt_recovers_failed_source_cleanup) | test(decommission_durable_ilm_receipt_pagination_fails_closed_on_second_page))' test-group = 'ecstore-serial-flaky' # Serialize the bucket-incarnation / lifecycle-fence tests under the ci profile diff --git a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs index 1162787f7..69bfe6416 100644 --- a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs +++ b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs @@ -9136,6 +9136,7 @@ mod tests { assert_eq!(loaded.report.scanned, 37); assert_eq!(loaded.report.eligible, 11); assert_eq!(loaded.report.enqueued, 5); + assert_eq!(loaded.cursor_revision, Some(1)); assert!(loaded.lease_expires_at_unix_nanos > 0); let token = loaded .report @@ -9154,6 +9155,30 @@ mod tests { assert_eq!(admission.lease_id, loaded.lease_id); assert_eq!(admission.lease_expires_at_unix_nanos, loaded.lease_expires_at_unix_nanos); + let mut same_marker_report = report.clone(); + same_marker_report.scanned += 1; + persist_manual_transition_page_checkpoint( + &checkpoint_options, + &same_marker_report, + Some("logs/page-end".to_string()), + Some("opaque-next-version".to_string()), + ) + .await + .expect("same-marker version checkpoint should persist through the durable progress sink"); + let same_marker_checkpointed = load_manual_transition_job_record(ecstore.clone(), job_id) + .await + .expect("same-marker version checkpoint should reload"); + assert_eq!(same_marker_checkpointed.cursor_revision, Some(2)); + let (_, version_marker) = decode_manual_transition_continuation_token( + same_marker_checkpointed + .report + .continuation_token + .as_deref() + .expect("same-marker version checkpoint should persist a cursor"), + ) + .expect("same-marker version cursor should decode"); + assert_eq!(version_marker.as_deref(), Some("opaque-next-version")); + create_test_bucket(&ecstore, &bucket).await; let lifecycle_xml = format!( r#" @@ -9210,6 +9235,7 @@ mod tests { assert_eq!(checkpointed.report.scanned, 1000); assert_eq!(checkpointed.report.eligible, 1000); assert_eq!(checkpointed.report.dry_run_eligible, 1000); + assert_eq!(checkpointed.cursor_revision, Some(3)); let token = checkpointed .report .continuation_token diff --git a/crates/ecstore/src/bucket/lifecycle/durable_namespace.rs b/crates/ecstore/src/bucket/lifecycle/durable_namespace.rs index 6d3ebff70..81a8f38df 100644 --- a/crates/ecstore/src/bucket/lifecycle/durable_namespace.rs +++ b/crates/ecstore/src/bucket/lifecycle/durable_namespace.rs @@ -12,7 +12,9 @@ // See the License for the specific language governing permissions and // limitations under the License. -use rustfs_utils::crypto::hex_sha256; +use std::collections::BTreeMap; + +use rustfs_utils::crypto::{hex_sha256, is_sha256_checksum}; use serde::{Deserialize, Serialize}; use uuid::Uuid; @@ -26,6 +28,7 @@ use crate::error::{Error, Result}; pub(crate) const ILM_META_PREFIX: &str = "ilm"; const ILM_META_OBJECT_PREFIX: &str = "ilm/"; +const MANUAL_TRANSITION_CURSOR_MARKER_PROOF_MAX_SIZE: usize = 1024; #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub(crate) enum DurableIlmRecordKind { @@ -106,6 +109,41 @@ pub(crate) struct ManualTransitionJobProgressCheckpoint { queue_snapshot: ManualTransitionQueueSnapshot, } +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub(crate) struct ManualTransitionJobProgressProof { + scope_sha256: String, + scanned: u64, + eligible: u64, + enqueued: u64, + dry_run_eligible: u64, + skipped_not_transition: u64, + skipped_tier: u64, + skipped_delete_marker: u64, + skipped_directory: u64, + skipped_replication: u64, + skipped_already_transitioned: u64, + skipped_already_in_flight: u64, + skipped_queue_full: u64, + skipped_queue_closed: u64, + skipped_queue_timeout: u64, + transition_completed: u64, + transition_failed: u64, + tier_failure: u64, + tier_failure_by_reason: BTreeMap, + lifecycle_config_found: bool, + truncated_by_limit: bool, + truncated_by_duration: bool, + cancelled: bool, + #[serde(default, skip_serializing_if = "Option::is_none")] + continuation_token_sha256: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + cursor_marker: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + cursor_revision: Option, + queue_snapshot: ManualTransitionQueueSnapshot, +} + impl ValidatedDurableIlmRecord { pub(crate) fn context(&self) -> String { format!("namespace `{}` {} `{}`", self.namespace, self.id_kind, self.id) @@ -137,6 +175,8 @@ pub(crate) enum DurableIlmRecordCheckpoint { cancel_requested: bool, #[serde(default, skip_serializing_if = "Option::is_none")] progress: Option>, + #[serde(default, skip_serializing_if = "Option::is_none")] + progress_proof: Option>, }, ManualTransitionScope { content_sha256: String, @@ -163,13 +203,38 @@ impl DurableIlmRecordCheckpoint { } } + pub(crate) fn compacted(&self) -> Result { + let mut checkpoint = self.clone(); + if let Self::ManualTransitionJob { + progress, + progress_proof, + .. + } = &mut checkpoint + { + match (progress.take(), progress_proof.take()) { + (Some(progress), None) => { + *progress_proof = Some(Box::new(ManualTransitionJobProgressProof::new( + &progress.report, + &progress.queue_snapshot, + None, + )?)); + } + (None, Some(proof)) if proof.is_valid() => *progress_proof = Some(proof), + (None, None) => {} + _ => return Err(Error::other("durable ILM manual transition checkpoint is invalid")), + } + } + Ok(checkpoint) + } + pub(crate) fn validate_successor(&self, next: &Self) -> Result<()> { if self == next { if let Self::ManualTransitionJob { - progress: Some(progress), + progress, + progress_proof, .. } = self - && !manual_job_progress_is_valid(progress) + && !manual_job_progress_checkpoint_is_valid(progress.as_deref(), progress_proof.as_deref()) { return Err(Error::other("durable ILM manual transition checkpoint is invalid")); } @@ -217,30 +282,53 @@ impl DurableIlmRecordCheckpoint { } ( Self::ManualTransitionJob { + content_sha256: previous_content, identity_sha256: previous_identity, updated_at_unix_nanos: previous_updated_at, state: previous_state, scan_completed: previous_scan_completed, cancel_requested: previous_cancel_requested, progress: previous_progress, + progress_proof: previous_progress_proof, .. }, Self::ManualTransitionJob { + content_sha256: next_content, identity_sha256: next_identity, updated_at_unix_nanos: next_updated_at, state: next_state, scan_completed: next_scan_completed, cancel_requested: next_cancel_requested, progress: next_progress, + progress_proof: next_progress_proof, .. }, ) => { - previous_identity == next_identity - && next_updated_at > previous_updated_at - && manual_job_state_reaches(*previous_state, *next_state) - && (!previous_scan_completed || *next_scan_completed) - && (!previous_cancel_requested || *next_cancel_requested) - && manual_job_progress_reaches(previous_progress.as_deref(), next_progress.as_deref(), *next_scan_completed) + let same_generation = previous_content == next_content + && previous_identity == next_identity + && previous_updated_at == next_updated_at + && previous_state == next_state + && previous_scan_completed == next_scan_completed + && previous_cancel_requested == next_cancel_requested + && manual_job_progress_equivalent( + previous_progress.as_deref(), + previous_progress_proof.as_deref(), + next_progress.as_deref(), + next_progress_proof.as_deref(), + ); + same_generation + || (previous_identity == next_identity + && next_updated_at > previous_updated_at + && manual_job_state_reaches(*previous_state, *next_state) + && (!previous_scan_completed || *next_scan_completed) + && (!previous_cancel_requested || *next_cancel_requested) + && manual_job_progress_reaches( + previous_progress.as_deref(), + previous_progress_proof.as_deref(), + next_progress.as_deref(), + next_progress_proof.as_deref(), + *next_scan_completed, + )) } ( Self::ManualTransitionScope { @@ -294,22 +382,148 @@ fn manual_job_state_reaches( from == to || from == manual_transition_job::ManualTransitionJobState::Running } +impl ManualTransitionJobProgressProof { + fn new( + report: &ManualTransitionRunReport, + queue_snapshot: &ManualTransitionQueueSnapshot, + cursor_revision: Option, + ) -> Result { + let cursor_marker = match report.continuation_token.as_deref() { + Some(token) => { + let marker = decode_manual_transition_continuation_token(token)? + .0 + .ok_or_else(|| Error::other("durable ILM manual transition cursor marker is missing"))?; + (marker.len() <= MANUAL_TRANSITION_CURSOR_MARKER_PROOF_MAX_SIZE).then_some(marker) + } + None => None, + }; + if !manual_job_worker_results_are_valid(report) + || !manual_job_queue_snapshot_is_valid(queue_snapshot) + || (report.continuation_token.is_some() && cursor_revision == Some(0)) + { + return Err(Error::other("durable ILM manual transition progress is invalid")); + } + Ok(Self { + scope_sha256: checkpoint_hash(&(report.bucket.as_str(), report.prefix.as_str(), &report.tier, report.dry_run))?, + scanned: report.scanned, + eligible: report.eligible, + enqueued: report.enqueued, + dry_run_eligible: report.dry_run_eligible, + skipped_not_transition: report.skipped_not_transition, + skipped_tier: report.skipped_tier, + skipped_delete_marker: report.skipped_delete_marker, + skipped_directory: report.skipped_directory, + skipped_replication: report.skipped_replication, + skipped_already_transitioned: report.skipped_already_transitioned, + skipped_already_in_flight: report.skipped_already_in_flight, + skipped_queue_full: report.skipped_queue_full, + skipped_queue_closed: report.skipped_queue_closed, + skipped_queue_timeout: report.skipped_queue_timeout, + transition_completed: report.transition_completed, + transition_failed: report.transition_failed, + tier_failure: report.tier_failure, + tier_failure_by_reason: report.tier_failure_by_reason.clone(), + lifecycle_config_found: report.lifecycle_config_found, + truncated_by_limit: report.truncated_by_limit, + truncated_by_duration: report.truncated_by_duration, + cancelled: report.cancelled, + continuation_token_sha256: report + .continuation_token + .as_deref() + .map(|token| hex_sha256(token.as_bytes(), ToOwned::to_owned)), + cursor_marker, + cursor_revision, + queue_snapshot: *queue_snapshot, + }) + } + + fn is_valid(&self) -> bool { + let reason_total = self + .tier_failure_by_reason + .values() + .try_fold(0u64, |total, count| total.checked_add(*count)); + is_sha256_checksum(&self.scope_sha256) + && self.continuation_token_sha256.as_deref().is_none_or(is_sha256_checksum) + && match (&self.continuation_token_sha256, &self.cursor_marker) { + (None, None) | (Some(_), None) => true, + (Some(_), Some(marker)) => !marker.is_empty() && marker.len() <= MANUAL_TRANSITION_CURSOR_MARKER_PROOF_MAX_SIZE, + (None, Some(_)) => false, + } + && !(self.continuation_token_sha256.is_some() && self.cursor_revision == Some(0)) + && self + .transition_completed + .checked_add(self.transition_failed) + .is_some_and(|total| total <= self.enqueued) + && self.transition_failed <= self.tier_failure + && reason_total.is_some_and(|total| total <= self.tier_failure) + && manual_job_queue_snapshot_is_valid(&self.queue_snapshot) + } +} + +fn manual_job_progress_checkpoint_is_valid( + progress: Option<&ManualTransitionJobProgressCheckpoint>, + proof: Option<&ManualTransitionJobProgressProof>, +) -> bool { + match (progress, proof) { + (Some(progress), None) => manual_job_progress_is_valid(progress), + (None, Some(proof)) => proof.is_valid(), + (None, None) => true, + (Some(_), Some(_)) => false, + } +} + +fn manual_job_progress_proof( + progress: Option<&ManualTransitionJobProgressCheckpoint>, + proof: Option<&ManualTransitionJobProgressProof>, +) -> Option { + match (progress, proof) { + (Some(progress), None) => ManualTransitionJobProgressProof::new(&progress.report, &progress.queue_snapshot, None).ok(), + (None, Some(proof)) if proof.is_valid() => Some(proof.clone()), + _ => None, + } +} + +fn manual_job_progress_equivalent( + previous: Option<&ManualTransitionJobProgressCheckpoint>, + previous_proof: Option<&ManualTransitionJobProgressProof>, + next: Option<&ManualTransitionJobProgressCheckpoint>, + next_proof: Option<&ManualTransitionJobProgressProof>, +) -> bool { + if !manual_job_progress_checkpoint_is_valid(previous, previous_proof) + || !manual_job_progress_checkpoint_is_valid(next, next_proof) + { + return false; + } + match ( + manual_job_progress_proof(previous, previous_proof), + manual_job_progress_proof(next, next_proof), + ) { + (Some(previous), Some(next)) => previous == next, + (None, None) => true, + _ => false, + } +} + fn manual_job_progress_reaches( previous: Option<&ManualTransitionJobProgressCheckpoint>, + previous_proof: Option<&ManualTransitionJobProgressProof>, next: Option<&ManualTransitionJobProgressCheckpoint>, + next_proof: Option<&ManualTransitionJobProgressProof>, next_scan_completed: bool, ) -> bool { - let (previous, next) = match (previous, next) { - (None, Some(next)) => return manual_job_progress_is_valid(next), - (Some(previous), Some(next)) => (previous, next), - _ => return false, + if previous.is_none() && previous_proof.is_none() { + return (next.is_some() || next_proof.is_some()) && manual_job_progress_checkpoint_is_valid(next, next_proof); + } + let (Some(previous_compact), Some(next_compact)) = ( + manual_job_progress_proof(previous, previous_proof), + manual_job_progress_proof(next, next_proof), + ) else { + return false; }; - let previous_report = &previous.report; - let next_report = &next.report; macro_rules! counters_do_not_regress { ($($field:ident),+ $(,)?) => { - $(previous_report.$field <= next_report.$field)&&+ + $(previous_compact.$field <= next_compact.$field)&&+ }; } @@ -332,28 +546,34 @@ fn manual_job_progress_reaches( transition_failed, tier_failure, ); - let failure_reasons_monotonic = previous_report.tier_failure_by_reason.iter().all(|(reason, previous_count)| { - next_report - .tier_failure_by_reason - .get(reason) - .is_some_and(|next_count| next_count >= previous_count) - }); - let flags_monotonic = (!previous_report.lifecycle_config_found || next_report.lifecycle_config_found) - && (!previous_report.truncated_by_limit || next_report.truncated_by_limit) - && (!previous_report.truncated_by_duration || next_report.truncated_by_duration) - && (!previous_report.cancelled || next_report.cancelled); - let cursor_monotonic = manual_job_cursor_reaches(previous_report, next_report, next_scan_completed); - let progress_valid = manual_job_progress_is_valid(previous) && manual_job_progress_is_valid(next); + let failure_reasons_monotonic = previous_compact + .tier_failure_by_reason + .iter() + .all(|(reason, previous_count)| { + next_compact + .tier_failure_by_reason + .get(reason) + .is_some_and(|next_count| next_count >= previous_count) + }); + let flags_monotonic = (!previous_compact.lifecycle_config_found || next_compact.lifecycle_config_found) + && (!previous_compact.truncated_by_limit || next_compact.truncated_by_limit) + && (!previous_compact.truncated_by_duration || next_compact.truncated_by_duration) + && (!previous_compact.cancelled || next_compact.cancelled); + let cursor_monotonic = manual_job_cursor_reaches( + &previous_compact, + &next_compact, + previous.map(|progress| &progress.report), + next.map(|progress| &progress.report), + next_scan_completed, + ); - previous_report.bucket == next_report.bucket - && previous_report.prefix == next_report.prefix - && previous_report.tier == next_report.tier - && previous_report.dry_run == next_report.dry_run + previous_compact.scope_sha256 == next_compact.scope_sha256 && counters_monotonic && failure_reasons_monotonic && flags_monotonic && cursor_monotonic - && progress_valid + && previous_compact.is_valid() + && next_compact.is_valid() } fn manual_job_progress_is_valid(progress: &ManualTransitionJobProgressCheckpoint) -> bool { @@ -376,29 +596,62 @@ fn manual_job_worker_results_are_valid(report: &ManualTransitionRunReport) -> bo } fn manual_job_cursor_reaches( - previous: &ManualTransitionRunReport, - next: &ManualTransitionRunReport, + previous: &ManualTransitionJobProgressProof, + next: &ManualTransitionJobProgressProof, + previous_legacy: Option<&ManualTransitionRunReport>, + next_legacy: Option<&ManualTransitionRunReport>, next_scan_completed: bool, ) -> bool { - if previous.continuation_token == next.continuation_token { - return manual_job_cursor_is_valid(previous.continuation_token.as_deref()); + if previous.continuation_token_sha256 == next.continuation_token_sha256 { + return previous.cursor_marker == next.cursor_marker && previous.cursor_revision == next.cursor_revision; } - match (&previous.continuation_token, &next.continuation_token) { - (None, Some(next_token)) => next.scanned > previous.scanned && manual_job_cursor_is_valid(Some(next_token)), + match (&previous.continuation_token_sha256, &next.continuation_token_sha256) { + (None, Some(_)) => { + next.scanned > previous.scanned + && (manual_job_cursor_revision_advances(previous.cursor_revision, next.cursor_revision) + || (previous.cursor_revision.is_none() && next.cursor_revision.is_none())) + } (Some(_), None) => next_scan_completed, - (Some(previous_token), Some(next_token)) if next.scanned > previous.scanned => { - let (Ok((Some(previous_marker), _)), Ok((Some(next_marker), _))) = ( - decode_manual_transition_continuation_token(previous_token), - decode_manual_transition_continuation_token(next_token), - ) else { - return false; - }; - next_marker > previous_marker + (Some(_), Some(_)) if next.scanned > previous.scanned => { + manual_job_cursor_revision_advances(previous.cursor_revision, next.cursor_revision) + || manual_job_legacy_cursor_reaches(previous, next, previous_legacy, next_legacy) } _ => false, } } +fn manual_job_cursor_revision_advances(previous: Option, next: Option) -> bool { + match (previous, next) { + (Some(previous), Some(next)) => next > previous, + (None, Some(next)) => next > 0, + _ => false, + } +} + +fn manual_job_legacy_cursor_reaches( + previous_proof: &ManualTransitionJobProgressProof, + next_proof: &ManualTransitionJobProgressProof, + previous_legacy: Option<&ManualTransitionRunReport>, + next_legacy: Option<&ManualTransitionRunReport>, +) -> bool { + if let (Some(previous_marker), Some(next_marker)) = (&previous_proof.cursor_marker, &next_proof.cursor_marker) { + return next_marker > previous_marker; + } + let (Some(previous_token), Some(next_token)) = ( + previous_legacy.and_then(|report| report.continuation_token.as_deref()), + next_legacy.and_then(|report| report.continuation_token.as_deref()), + ) else { + return false; + }; + let (Ok((Some(previous_marker), _)), Ok((Some(next_marker), _))) = ( + decode_manual_transition_continuation_token(previous_token), + decode_manual_transition_continuation_token(next_token), + ) else { + return false; + }; + next_marker > previous_marker +} + fn manual_job_cursor_is_valid(token: Option<&str>) -> bool { let Some(token) = token else { return true; @@ -563,6 +816,7 @@ pub(crate) fn validate_durable_ilm_record(path: &str, data: &[u8]) -> Result Result Result DurableIlmRecordCheckpoint { + fn try_manual_job_checkpoint(job: &manual_transition_job::ManualTransitionJobRecord) -> Result { let path = manual_transition_job::manual_transition_job_record_object_name(job.job_id).expect("manual job path should build"); let encoded = job.encode().expect("manual job should encode"); - validate_durable_ilm_record(&path, &encoded) - .expect("manual job checkpoint should validate") - .checkpoint + Ok(validate_durable_ilm_record(&path, &encoded)?.checkpoint) + } + + fn manual_job_checkpoint(job: &manual_transition_job::ManualTransitionJobRecord) -> DurableIlmRecordCheckpoint { + try_manual_job_checkpoint(job).expect("manual job checkpoint should validate") } fn continuation_token_with_version(marker: &str, version_marker: Option<&str>) -> String { @@ -692,37 +946,77 @@ mod tests { } } + #[test] + fn manual_transition_job_checkpoint_compacts_legacy_progress_compatibly() { + let options = super::super::bucket_lifecycle_ops::ManualTransitionRunOptions::default(); + let mut job = + manual_transition_job::ManualTransitionJobRecord::new(Uuid::new_v4(), "legacy-checkpoint-bucket", &options, "owner"); + job.cursor_revision = None; + job.updated_at_unix_nanos += 1; + job.report.scanned = 1; + job.report.continuation_token = Some(continuation_token("logs/a")); + let compact = manual_job_checkpoint(&job); + let mut legacy = compact.clone(); + let DurableIlmRecordCheckpoint::ManualTransitionJob { + progress, + progress_proof, + .. + } = &mut legacy + else { + panic!("manual job should produce a manual checkpoint"); + }; + *progress = Some(Box::new(ManualTransitionJobProgressCheckpoint { + report: job.report.clone(), + queue_snapshot: job.queue_snapshot, + })); + *progress_proof = None; + + compact + .validate_successor(&legacy) + .expect("bounded checkpoints should accept the same legacy generation"); + legacy + .validate_successor(&compact) + .expect("legacy checkpoints should accept the same bounded generation"); + assert_eq!(legacy.compacted().expect("legacy checkpoint should compact"), compact); + } + #[test] fn manual_transition_job_checkpoint_rejects_progress_poison() { let options = super::super::bucket_lifecycle_ops::ManualTransitionRunOptions::default(); let mut initial = manual_transition_job::ManualTransitionJobRecord::new(Uuid::new_v4(), "manual-checkpoint-bucket", &options, "owner"); let initial_checkpoint = manual_job_checkpoint(&initial); - initial.updated_at_unix_nanos += 1; - initial.report.scanned = 1; - initial.report.continuation_token = Some(continuation_token("logs/a")); + let mut first_page = initial.report.clone(); + first_page.scanned = 1; + first_page.continuation_token = Some(continuation_token("logs/a")); + initial.update_running_progress(first_page, ManualTransitionQueueSnapshot::default()); let first_page_checkpoint = manual_job_checkpoint(&initial); initial_checkpoint .validate_successor(&first_page_checkpoint) .expect("the first durable cursor should advance from no cursor"); let mut legacy_checkpoint = initial_checkpoint; - let DurableIlmRecordCheckpoint::ManualTransitionJob { progress, .. } = &mut legacy_checkpoint else { + let DurableIlmRecordCheckpoint::ManualTransitionJob { + progress, + progress_proof, + .. + } = &mut legacy_checkpoint + else { panic!("manual job should produce a manual checkpoint"); }; *progress = None; + *progress_proof = None; legacy_checkpoint .validate_successor(&first_page_checkpoint) .expect("legacy checkpoints should upgrade to validated progress"); let mut previous = initial; - previous.updated_at_unix_nanos += 1; - previous.report.scanned = 10; - previous.report.eligible = 8; - previous.report.enqueued = 2; - previous.report.transition_completed = 1; - previous.report.continuation_token = Some(continuation_token("logs/b")); - previous.queue_snapshot = ManualTransitionQueueSnapshot { + let mut previous_report = previous.report.clone(); + previous_report.scanned = 10; + previous_report.eligible = 8; + previous_report.enqueued = 2; + previous_report.continuation_token = Some(continuation_token("logs/b")); + let previous_queue = ManualTransitionQueueSnapshot { queue_capacity: 10, queued: 1, active: 1, @@ -731,17 +1025,21 @@ mod tests { queue_send_timeout: 1, ..Default::default() }; + previous.update_running_progress(previous_report, previous_queue); + previous.report.transition_completed = 1; let previous_checkpoint = manual_job_checkpoint(&previous); let mut next = previous.clone(); - next.updated_at_unix_nanos += 1; - next.report.scanned = 11; - next.report.eligible = 9; + let mut next_report = next.report.clone(); + next_report.scanned = 11; + next_report.eligible = 9; + next_report.continuation_token = Some(continuation_token("logs/c")); + let mut next_queue = next.queue_snapshot; + next_queue.queued = 0; + next_queue.active = 0; + next_queue.queue_full = 3; + next.update_running_progress(next_report, next_queue); next.report.transition_completed = 2; - next.report.continuation_token = Some(continuation_token("logs/c")); - next.queue_snapshot.queued = 0; - next.queue_snapshot.active = 0; - next.queue_snapshot.queue_full = 3; let next_checkpoint = manual_job_checkpoint(&next); previous_checkpoint .validate_successor(&next_checkpoint) @@ -756,9 +1054,8 @@ mod tests { .is_err() ); - let mut cursor_rollback = next.clone(); + let mut cursor_rollback = previous.clone(); cursor_rollback.updated_at_unix_nanos += 1; - cursor_rollback.report.scanned = previous.report.scanned; cursor_rollback.report.scanned += 1; cursor_rollback.report.continuation_token = Some(continuation_token("logs/a")); assert!( @@ -768,19 +1065,29 @@ mod tests { ); let mut same_marker_version_previous = previous.clone(); - same_marker_version_previous.report.continuation_token = - Some(continuation_token_with_version("logs/b", Some("opaque-newer-version"))); + let mut same_marker_report = same_marker_version_previous.report.clone(); + same_marker_report.continuation_token = Some(continuation_token_with_version("logs/b", Some("opaque-z-version"))); + same_marker_version_previous.update_running_progress(same_marker_report, same_marker_version_previous.queue_snapshot); let same_marker_version_previous_checkpoint = manual_job_checkpoint(&same_marker_version_previous); + let mut same_marker_version_next = same_marker_version_previous.clone(); + let mut same_marker_next_report = same_marker_version_next.report.clone(); + same_marker_next_report.scanned += 1; + same_marker_next_report.continuation_token = Some(continuation_token_with_version("logs/b", Some("opaque-a-version"))); + same_marker_version_next.update_running_progress(same_marker_next_report, same_marker_version_next.queue_snapshot); + same_marker_version_previous_checkpoint + .validate_successor(&manual_job_checkpoint(&same_marker_version_next)) + .expect("producer cursor revision should prove same-marker version progress"); + let mut same_marker_version_rollback = same_marker_version_previous.clone(); same_marker_version_rollback.updated_at_unix_nanos += 1; same_marker_version_rollback.report.scanned += 1; same_marker_version_rollback.report.continuation_token = - Some(continuation_token_with_version("logs/b", Some("opaque-stale-version"))); + Some(continuation_token_with_version("logs/b", Some("opaque-arbitrary-version"))); assert!( same_marker_version_previous_checkpoint .validate_successor(&manual_job_checkpoint(&same_marker_version_rollback)) .is_err(), - "opaque version markers must not be treated as ordered progress" + "a different opaque version marker without producer evidence must fail closed" ); let mut worker_result_rollback = next.clone(); @@ -798,28 +1105,16 @@ mod tests { worker_result_overflow.report.transition_completed = u64::MAX; worker_result_overflow.report.transition_failed = 1; worker_result_overflow.report.tier_failure = 1; - assert!( - previous_checkpoint - .validate_successor(&manual_job_checkpoint(&worker_result_overflow)) - .is_err() - ); + assert!(try_manual_job_checkpoint(&worker_result_overflow).is_err()); let mut invalid_cursor = next.clone(); invalid_cursor.updated_at_unix_nanos += 1; invalid_cursor.report.continuation_token = Some("not-base64".to_string()); - assert!( - previous_checkpoint - .validate_successor(&manual_job_checkpoint(&invalid_cursor)) - .is_err() - ); + assert!(try_manual_job_checkpoint(&invalid_cursor).is_err()); let mut queue_state_poison = next; queue_state_poison.updated_at_unix_nanos += 1; queue_state_poison.queue_snapshot.queued = queue_state_poison.queue_snapshot.queue_capacity + 1; - assert!( - previous_checkpoint - .validate_successor(&manual_job_checkpoint(&queue_state_poison)) - .is_err() - ); + assert!(try_manual_job_checkpoint(&queue_state_poison).is_err()); } } diff --git a/crates/ecstore/src/bucket/lifecycle/manual_transition_job.rs b/crates/ecstore/src/bucket/lifecycle/manual_transition_job.rs index 9461c150a..2706665f7 100644 --- a/crates/ecstore/src/bucket/lifecycle/manual_transition_job.rs +++ b/crates/ecstore/src/bucket/lifecycle/manual_transition_job.rs @@ -199,6 +199,8 @@ pub struct ManualTransitionJobRecord { pub updated_at_unix_nanos: i128, #[serde(default, skip_serializing_if = "Option::is_none")] pub completed_at_unix_nanos: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub cursor_revision: Option, pub report: ManualTransitionRunReport, pub queue_snapshot: ManualTransitionQueueSnapshot, #[serde(default, skip_serializing_if = "Option::is_none")] @@ -228,6 +230,7 @@ impl ManualTransitionJobRecord { created_at_unix_nanos: now, updated_at_unix_nanos: now, completed_at_unix_nanos: None, + cursor_revision: Some(0), report: ManualTransitionRunReport { bucket: bucket.to_string(), prefix: options.prefix.clone(), @@ -242,7 +245,7 @@ impl ManualTransitionJobRecord { pub fn complete(&mut self, report: ManualTransitionRunReport, queue_snapshot: ManualTransitionQueueSnapshot) { self.scan_completed = true; - self.report.merge_scan_report_preserving_worker(&report); + self.merge_scan_report(&report); self.queue_snapshot = queue_snapshot; self.error = None; self.mark_terminal_if_worker_drained(); @@ -442,11 +445,18 @@ impl ManualTransitionJobRecord { pub fn update_running_progress(&mut self, report: ManualTransitionRunReport, queue_snapshot: ManualTransitionQueueSnapshot) { if self.state == ManualTransitionJobState::Running { - self.report.merge_scan_report_preserving_worker(&report); + self.merge_scan_report(&report); self.renew_lease(queue_snapshot); } } + fn merge_scan_report(&mut self, report: &ManualTransitionRunReport) { + if self.report.continuation_token != report.continuation_token { + self.cursor_revision = Some(self.cursor_revision.unwrap_or(0).saturating_add(1)); + } + self.report.merge_scan_report_preserving_worker(report); + } + pub fn mark_unknown_if_unowned(&mut self) { if self.state == ManualTransitionJobState::Running { self.state = ManualTransitionJobState::Unknown; @@ -2591,6 +2601,25 @@ mod tests { assert!(decoded.report.tier_failure_by_reason.is_empty()); } + #[test] + fn manual_transition_job_record_decodes_legacy_cursor_without_revision() { + let options = ManualTransitionRunOptions::default(); + let record = ManualTransitionJobRecord::new(Uuid::new_v4(), "bucket", &options, TEST_OWNER); + let encoded = record.encode().expect("job record should encode"); + let mut value: serde_json::Value = serde_json::from_slice(&encoded).expect("encoded job should be json"); + value["job"] + .as_object_mut() + .expect("job should be object") + .remove("cursor_revision"); + let record_bytes = serde_json::to_vec(&value["job"]).expect("legacy job should encode"); + value["content_sha256"] = serde_json::Value::String(hex_sha256(&record_bytes, ToOwned::to_owned)); + let legacy = serde_json::to_vec(&value).expect("legacy envelope should encode"); + + let decoded = ManualTransitionJobRecord::decode(record.job_id, &legacy).expect("legacy job should decode"); + + assert_eq!(decoded.cursor_revision, None); + } + #[test] fn manual_transition_job_record_rejects_unknown_report_fields() { let options = ManualTransitionRunOptions::default(); diff --git a/crates/ecstore/src/core/pools.rs b/crates/ecstore/src/core/pools.rs index 8c905620b..cc6cda6e6 100644 --- a/crates/ecstore/src/core/pools.rs +++ b/crates/ecstore/src/core/pools.rs @@ -943,12 +943,19 @@ impl DecommissionDurableIlmReceipt { } fn encode(&self) -> Result> { - self.validate()?; - let receipt_bytes = serde_json::to_vec(self)?; + let mut receipt = self.clone(); + receipt.checkpoint = receipt.checkpoint.compacted()?; + receipt.terminal_checkpoint = receipt + .terminal_checkpoint + .as_ref() + .map(DurableIlmRecordCheckpoint::compacted) + .transpose()?; + receipt.validate()?; + let receipt_bytes = serde_json::to_vec(&receipt)?; let persisted = PersistedDecommissionDurableIlmReceipt { schema: DECOMMISSION_DURABLE_ILM_RECEIPT_SCHEMA.to_string(), content_sha256: hex_sha256(&receipt_bytes, ToOwned::to_owned), - receipt: self.clone(), + receipt, }; let encoded = serde_json::to_vec(&persisted)?; if encoded.len() > DECOMMISSION_DURABLE_ILM_RECEIPT_MAX_SIZE { @@ -5642,6 +5649,30 @@ impl ECStore { self.list_decommission_durable_ilm_receipts(source_pool_idx).await } + #[cfg(test)] + pub(crate) async fn persist_decommission_durable_ilm_receipt_for_test( + &self, + source_pool_idx: usize, + target_pool_idx: usize, + source_path: &str, + record: &ValidatedDurableIlmRecord, + terminal: bool, + ) -> Result { + let mut receipt = DecommissionDurableIlmReceipt::new(source_path, record); + if terminal { + receipt.terminal_checkpoint = Some(record.checkpoint.clone()); + } + self.persist_decommission_durable_ilm_receipt(source_pool_idx, target_pool_idx, &receipt) + .await?; + let run_token = self.durable_ilm_receipt_run_token(source_pool_idx).await?; + Ok(decommission_durable_ilm_receipt_path(&run_token, source_path, record.id_kind, &record.id)) + } + + #[cfg(test)] + pub(crate) async fn persist_decommission_durable_ilm_manifest_for_test(&self, source_pool_idx: usize) -> Result<()> { + self.persist_decommission_durable_ilm_manifest(source_pool_idx).await + } + #[cfg(test)] pub(crate) async fn cleanup_decommission_durable_ilm_receipts_for_test(&self, source_pool_idx: usize) -> Result<()> { self.cleanup_decommission_durable_ilm_receipts(source_pool_idx).await @@ -6828,12 +6859,12 @@ pub(crate) fn fallback_free_capacity_dedup(disks: &[rustfs_madmin::Disk]) -> usi #[cfg(test)] mod pools_tests { use super::{ - DECOMMISSION_META_PREFIXES, DECOMMISSION_PROGRESS_SAVE_INTERVAL, DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD, - DecomBucketInfo, DecommissionDurableIlmReceipt, DecommissionStartPoolState, DecommissionTerminalState, ListCallback, - PoolDecommissionInfo, PoolMeta, PoolSpaceInfo, PoolStatus, apply_decommission_status_space_info, - bind_decommission_cancelers, bind_missing_decommission_cancelers, cancel_decommission_canceler, - classify_decommission_terminal_state, count_decommission_item, decommission_cancel_signal_result, - decommission_durable_ilm_receipt_path, decommission_durable_ilm_receipt_run_prefix, + DECOMMISSION_DURABLE_ILM_RECEIPT_MAX_SIZE, DECOMMISSION_META_PREFIXES, DECOMMISSION_PROGRESS_SAVE_INTERVAL, + DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD, DecomBucketInfo, DecommissionDurableIlmReceipt, DecommissionStartPoolState, + DecommissionTerminalState, ListCallback, PoolDecommissionInfo, PoolMeta, PoolSpaceInfo, PoolStatus, + apply_decommission_status_space_info, bind_decommission_cancelers, bind_missing_decommission_cancelers, + cancel_decommission_canceler, classify_decommission_terminal_state, count_decommission_item, + decommission_cancel_signal_result, decommission_durable_ilm_receipt_path, decommission_durable_ilm_receipt_run_prefix, decommission_durable_ilm_receipt_run_token, decommission_item_size, decommission_meta_bucket_options, decommission_start_pool_state, dedup_indices, default_decommission_bucket_concurrency, ensure_decommission_cancel_allowed, ensure_decommission_clear_allowed, ensure_decommission_listing_disks_available, @@ -6862,7 +6893,12 @@ mod pools_tests { track_decommission_current_object, track_decommission_current_object_stage, validate_start_decommission_request, wait_decommission_listing_retry, wait_decommission_worker_drain, with_decommission_entry_context, }; - use crate::bucket::lifecycle::DurableIlmRecordCheckpoint; + use crate::bucket::lifecycle::{ + DurableIlmRecordCheckpoint, + bucket_lifecycle_ops::{ManualTransitionQueueSnapshot, ManualTransitionRunOptions}, + manual_transition_job::{ManualTransitionJobRecord, manual_transition_job_record_object_name}, + validate_durable_ilm_record, + }; use crate::data_movement; use crate::disk::endpoint::Endpoint; use crate::error::{Error, StorageError}; @@ -6964,6 +7000,38 @@ mod pools_tests { assert_eq!(merged.terminal_checkpoint, Some(terminal_checkpoint)); } + #[test] + fn decommission_manual_job_receipt_compacts_large_progress() { + let prefix = "p".repeat(12 * 1024); + let options = ManualTransitionRunOptions { + prefix: prefix.clone(), + ..Default::default() + }; + let mut job = ManualTransitionJobRecord::new(uuid::Uuid::new_v4(), "bounded-receipt-bucket", &options, "owner"); + let token_bytes = serde_json::to_vec(&serde_json::json!({ + "marker": "m".repeat(12 * 1024), + "version_marker": "opaque-version" + })) + .expect("large continuation token should encode"); + let mut report = job.report.clone(); + report.scanned = 1; + report.continuation_token = Some(base64_simd::URL_SAFE_NO_PAD.encode_to_string(&token_bytes)); + job.update_running_progress(report, ManualTransitionQueueSnapshot::default()); + let path = manual_transition_job_record_object_name(job.job_id).expect("manual job path should build"); + let job_bytes = job.encode().expect("large manual job should remain within its record limit"); + assert!(job_bytes.len() > DECOMMISSION_DURABLE_ILM_RECEIPT_MAX_SIZE); + let record = validate_durable_ilm_record(&path, &job_bytes).expect("large manual job should validate"); + let mut receipt = DecommissionDurableIlmReceipt::new(&path, &record); + receipt.terminal_checkpoint = Some(record.checkpoint); + + let encoded = receipt.encode().expect("bounded progress proof should fit the receipt limit"); + let decoded = DecommissionDurableIlmReceipt::decode(&encoded).expect("bounded receipt should round trip"); + + assert!(encoded.len() <= DECOMMISSION_DURABLE_ILM_RECEIPT_MAX_SIZE); + assert_eq!(decoded.source_path, path); + assert!(decoded.terminal_checkpoint.is_some()); + } + #[test] fn test_apply_decommission_status_space_info_adds_idle_pool_usage() { let status = apply_decommission_status_space_info( diff --git a/crates/ecstore/src/store/init.rs b/crates/ecstore/src/store/init.rs index 1f25757a6..de7a6afe4 100644 --- a/crates/ecstore/src/store/init.rs +++ b/crates/ecstore/src/store/init.rs @@ -555,7 +555,7 @@ mod tests { #[cfg(feature = "test-util")] use crate::{ bucket::lifecycle::{ - ILM_META_PREFIX, + DurableIlmRecordCheckpoint, ILM_META_PREFIX, ValidatedDurableIlmRecord, bucket_lifecycle_ops::{ManualTransitionRunOptions, recover_manual_transition_jobs_once}, lifecycle::{TRANSITION_PENDING, TransitionOptions}, manual_transition_job::{ @@ -619,6 +619,8 @@ mod tests { range::HTTPRangeSpec, }, }; + #[cfg(feature = "test-util")] + use futures::{StreamExt as _, TryStreamExt as _}; use http::HeaderMap; use rustfs_config::server_config::KVS; use rustfs_filemeta::ObjectPartInfo; @@ -3222,6 +3224,108 @@ mod tests { assert!(backend.remove_versions().await.contains(&(entry.obj_name, entry.version_id))); } + #[cfg(feature = "test-util")] + #[tokio::test] + #[serial_test::serial(storage_class_env)] + async fn decommission_durable_ilm_receipt_pagination_fails_closed_on_second_page() { + const RECEIPT_COUNT: usize = 1001; + + let temp_dir = tempfile::tempdir().expect("create paginated receipt store dir"); + let (_ctx, store, _shutdown) = + without_storage_class_env(build_isolated_test_store(temp_dir.path(), "durable-ilm-receipt-pages", &[4, 4])).await; + store.pool_meta.write().await.pools[0].decommission = Some(PoolDecommissionInfo { + start_time: Some(OffsetDateTime::now_utc()), + ..Default::default() + }); + + futures::stream::iter(0..RECEIPT_COUNT) + .map(|index| { + let store = store.clone(); + async move { + let id = format!("{index:064x}"); + let source_path = format!("ilm/tier-delete-journal/{id}.json"); + let record = ValidatedDurableIlmRecord { + namespace: "tier-delete-journal", + id_kind: "operation_id", + id, + checkpoint: DurableIlmRecordCheckpoint::TierDeleteJournal { + content_sha256: format!("{:064x}", index + RECEIPT_COUNT), + identity_sha256: "f".repeat(64), + committed: false, + }, + }; + store + .persist_decommission_durable_ilm_receipt_for_test(0, 0, &source_path, &record, true) + .await?; + store + .persist_decommission_durable_ilm_receipt_for_test(0, 1, &source_path, &record, true) + .await?; + Ok::<(), Error>(()) + } + }) + .buffer_unordered(32) + .try_collect::>() + .await + .expect("more than one receipt page should persist"); + store + .persist_decommission_durable_ilm_manifest_for_test(0) + .await + .expect("paginated source receipts should produce a manifest"); + + let target_receipts = store + .decommission_durable_ilm_receipt_paths_for_test(0) + .await + .expect("paginated target receipts should list"); + assert_eq!(target_receipts.len(), RECEIPT_COUNT); + let (target_pool_idx, second_page_path) = target_receipts + .get(1000) + .cloned() + .expect("the real 1000-item page boundary should expose a second page receipt"); + let receipt_bytes = com::read_config(store.pools[target_pool_idx].clone(), &second_page_path) + .await + .expect("second page receipt should be readable"); + + com::delete_config(store.pools[target_pool_idx].clone(), &second_page_path) + .await + .expect("second page receipt should delete"); + let missing = store + .complete_decommission(0) + .await + .expect_err("a missing second page receipt must block completion") + .to_string(); + assert!(missing.contains(&second_page_path)); + assert!( + !store.pool_meta.read().await.pools[0] + .decommission + .as_ref() + .expect("source pool should remain in decommission") + .complete + ); + assert!(com::read_config(store.pools[0].clone(), &second_page_path).await.is_ok()); + + com::save_config(store.pools[target_pool_idx].clone(), &second_page_path, receipt_bytes.clone()) + .await + .expect("second page receipt should restore"); + com::save_config(store.pools[target_pool_idx].clone(), &second_page_path, b"{corrupt".to_vec()) + .await + .expect("second page receipt should corrupt deterministically"); + let corrupt = store + .complete_decommission(0) + .await + .expect_err("a corrupt second page receipt must block completion") + .to_string(); + assert!(corrupt.contains(&second_page_path)); + assert!(corrupt.contains("invalid")); + assert!( + !store.pool_meta.read().await.pools[0] + .decommission + .as_ref() + .expect("source pool should remain in decommission") + .complete + ); + assert!(com::read_config(store.pools[0].clone(), &second_page_path).await.is_ok()); + } + #[cfg(feature = "test-util")] #[tokio::test] #[serial_test::serial(storage_class_env)]