mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-22 20:36:38 +00:00
fix(ecstore): harden durable ILM cursor receipts
This commit is contained in:
@@ -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();
|
||||
|
||||
@@ -888,12 +888,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 {
|
||||
@@ -5492,6 +5499,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
|
||||
@@ -6634,12 +6665,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,
|
||||
@@ -6669,7 +6700,12 @@ mod pools_tests {
|
||||
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};
|
||||
@@ -6769,6 +6805,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(
|
||||
|
||||
@@ -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)]
|
||||
|
||||
Reference in New Issue
Block a user