fix(ecstore): harden durable ILM cursor receipts

This commit is contained in:
overtrue
2026-08-22 05:08:46 +08:00
parent cf0a78ce11
commit 7428f3138c
6 changed files with 631 additions and 109 deletions
+2 -2
View File
@@ -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
@@ -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#"<?xml version="1.0" encoding="UTF-8"?>
@@ -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
@@ -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<manual_transition_job::ManualTransitionWorkerFailureReason, u64>,
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<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
cursor_marker: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
cursor_revision: Option<u64>,
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<Box<ManualTransitionJobProgressCheckpoint>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
progress_proof: Option<Box<ManualTransitionJobProgressProof>>,
},
ManualTransitionScope {
content_sha256: String,
@@ -163,13 +203,38 @@ impl DurableIlmRecordCheckpoint {
}
}
pub(crate) fn compacted(&self) -> Result<Self> {
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<u64>,
) -> Result<Self> {
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<ManualTransitionJobProgressProof> {
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<u64>, next: Option<u64>) -> 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<Val
job.max_duration,
job.created_at_unix_nanos,
))?;
let progress_proof = ManualTransitionJobProgressProof::new(&job.report, &job.queue_snapshot, job.cursor_revision)?;
(
"job_id",
job_id.to_string(),
@@ -573,10 +827,8 @@ pub(crate) fn validate_durable_ilm_record(path: &str, data: &[u8]) -> Result<Val
state: job.state,
scan_completed: job.scan_completed,
cancel_requested: job.cancel_requested,
progress: Some(Box::new(ManualTransitionJobProgressCheckpoint {
report: job.report,
queue_snapshot: job.queue_snapshot,
})),
progress: None,
progress_proof: Some(Box::new(progress_proof)),
},
)
}
@@ -652,13 +904,15 @@ pub(crate) fn validate_durable_ilm_record(path: &str, data: &[u8]) -> Result<Val
mod tests {
use super::*;
fn manual_job_checkpoint(job: &manual_transition_job::ManualTransitionJobRecord) -> DurableIlmRecordCheckpoint {
fn try_manual_job_checkpoint(job: &manual_transition_job::ManualTransitionJobRecord) -> Result<DurableIlmRecordCheckpoint> {
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());
}
}
@@ -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<i128>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub cursor_revision: Option<u64>,
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();
+78 -10
View File
@@ -943,12 +943,19 @@ impl DecommissionDurableIlmReceipt {
}
fn encode(&self) -> Result<Vec<u8>> {
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<String> {
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(
+105 -1
View File
@@ -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::<Vec<_>>()
.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)]