Compare commits

..

6 Commits

Author SHA1 Message Date
Zhengchao An 72e89dc9d7 Merge branch 'main' into fix/bounded-transition-recovery 2026-09-06 03:32:10 +08:00
Zhengchao An e56de7d783 Merge branch 'main' into fix/bounded-transition-recovery 2026-09-06 01:27:52 +08:00
overtrue 36f4118cd3 test(ilm): keep expiry sentinel within control range 2026-09-06 01:15:21 +08:00
Zhengchao An 42fb655863 Merge branch 'main' into fix/bounded-transition-recovery 2026-09-05 23:56:32 +08:00
Zhengchao An d6d0bea4a9 Merge branch 'main' into fix/bounded-transition-recovery 2026-09-05 22:21:17 +08:00
cxymds 847d3f2433 fix(ilm): persist bounded transition recovery controls 2026-09-05 21:57:57 +08:00
9 changed files with 2544 additions and 86 deletions
+7
View File
@@ -69,6 +69,13 @@ pub mod bucket {
};
}
pub mod recovery_control {
pub use crate::bucket::lifecycle::recovery_control::{
IlmRecoveryClassification, IlmRecoveryControlPage, IlmRecoveryControlView, IlmRecoveryProtocol,
inspect_recovery_control, list_recovery_controls,
};
}
pub mod transition_transaction {
pub use crate::bucket::lifecycle::transition_transaction::{
TransitionOperatorDeleteResult, TransitionOperatorError, TransitionOperatorProbe, TransitionOperatorStatus,
@@ -22,7 +22,7 @@ use super::{
bucket_lifecycle_ops::{
ManualTransitionQueueSnapshot, ManualTransitionRunReport, decode_manual_transition_continuation_token,
},
manual_transition_job, tier_delete_journal, transition_transaction,
manual_transition_job, recovery_control, tier_delete_journal, transition_transaction,
};
use crate::error::{Error, Result};
use crate::services::tier::tier_probe_intent;
@@ -41,6 +41,7 @@ pub(crate) enum DurableIlmRecordKind {
ManualTransitionScope,
ManualTransitionTask,
ManualTransitionWorkerResult,
RecoveryControl,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
@@ -105,8 +106,14 @@ pub(crate) const MANUAL_TRANSITION_WORKER_RESULT_NAMESPACE: DurableIlmNamespace
max_record_size: manual_transition_job::MAX_MANUAL_TRANSITION_WORKER_RESULT_RECORD_SIZE,
kind: DurableIlmRecordKind::ManualTransitionWorkerResult,
};
pub(crate) const RECOVERY_CONTROL_NAMESPACE: DurableIlmNamespace = DurableIlmNamespace {
name: "recovery-control",
prefix: recovery_control::ILM_RECOVERY_CONTROL_PREFIX,
max_record_size: recovery_control::MAX_ILM_RECOVERY_CONTROL_SIZE,
kind: DurableIlmRecordKind::RecoveryControl,
};
pub(crate) const DURABLE_ILM_NAMESPACES: [DurableIlmNamespace; 9] = [
pub(crate) const DURABLE_ILM_NAMESPACES: [DurableIlmNamespace; 10] = [
TIER_DELETE_JOURNAL_NAMESPACE,
TIER_DELETE_JOURNAL_V6_NAMESPACE,
TIER_DELETE_DISPATCH_MANIFEST_NAMESPACE,
@@ -116,6 +123,7 @@ pub(crate) const DURABLE_ILM_NAMESPACES: [DurableIlmNamespace; 9] = [
MANUAL_TRANSITION_SCOPE_NAMESPACE,
MANUAL_TRANSITION_TASK_NAMESPACE,
MANUAL_TRANSITION_WORKER_RESULT_NAMESPACE,
RECOVERY_CONTROL_NAMESPACE,
];
#[derive(Debug, Clone, PartialEq, Eq)]
@@ -241,6 +249,18 @@ pub(crate) enum DurableIlmRecordCheckpoint {
ManualTransitionWorkerResult {
content_sha256: String,
},
RecoveryControl {
content_sha256: String,
identity_sha256: String,
source_generation_sha256: String,
first_seen_at_unix_nanos: i64,
revision: u64,
classification: recovery_control::IlmRecoveryClassification,
attempt_count: u64,
consecutive_failure_count: u32,
#[serde(default, skip_serializing_if = "Option::is_none")]
owner_fence_sha256: Option<String>,
},
}
impl DurableIlmRecordCheckpoint {
@@ -254,7 +274,8 @@ impl DurableIlmRecordCheckpoint {
| Self::ManualTransitionJob { content_sha256, .. }
| Self::ManualTransitionScope { content_sha256, .. }
| Self::ManualTransitionTask { content_sha256 }
| Self::ManualTransitionWorkerResult { content_sha256 } => content_sha256,
| Self::ManualTransitionWorkerResult { content_sha256 }
| Self::RecoveryControl { content_sha256, .. } => content_sha256,
}
}
@@ -528,6 +549,51 @@ impl DurableIlmRecordCheckpoint {
..
},
) => previous_identity == next_identity && next_updated_at > previous_updated_at,
(
Self::RecoveryControl {
identity_sha256: previous_identity,
source_generation_sha256: previous_generation,
first_seen_at_unix_nanos: previous_first_seen,
revision: previous_revision,
classification: previous_classification,
attempt_count: previous_attempts,
consecutive_failure_count: previous_failures,
owner_fence_sha256: previous_owner,
..
},
Self::RecoveryControl {
identity_sha256: next_identity,
source_generation_sha256: next_generation,
first_seen_at_unix_nanos: next_first_seen,
revision: next_revision,
classification: next_classification,
attempt_count: next_attempts,
consecutive_failure_count: next_failures,
owner_fence_sha256: next_owner,
..
},
) => {
let adjacent = previous_identity == next_identity
&& previous_first_seen == next_first_seen
&& previous_revision.checked_add(1) == Some(*next_revision);
let claim = next_owner.is_some()
&& *previous_classification == recovery_control::IlmRecoveryClassification::Retrying
&& *next_classification == recovery_control::IlmRecoveryClassification::Retrying
&& previous_attempts.checked_add(1) == Some(*next_attempts)
&& previous_failures == next_failures;
let source_refresh = previous_owner.is_some()
&& previous_owner == next_owner
&& *previous_classification == recovery_control::IlmRecoveryClassification::Retrying
&& *next_classification == recovery_control::IlmRecoveryClassification::Retrying
&& previous_attempts == next_attempts
&& previous_failures == next_failures
&& previous_generation != next_generation;
let completion = previous_owner.is_some()
&& next_owner.is_none()
&& previous_generation == next_generation
&& previous_attempts == next_attempts;
adjacent && (claim || source_refresh || completion)
}
_ => false,
};
@@ -553,6 +619,14 @@ impl DurableIlmRecordCheckpoint {
{
return false;
}
if let Self::RecoveryControl { classification, .. } = terminal
&& !matches!(
classification,
recovery_control::IlmRecoveryClassification::Terminal | recovery_control::IlmRecoveryClassification::Abandoned
)
{
return false;
}
if self == terminal || self.validate_successor(terminal).is_ok() {
return true;
}
@@ -652,6 +726,32 @@ impl DurableIlmRecordCheckpoint {
.is_some_and(|distance| tier_probe_state_reaches(*previous_state, *terminal_state, distance))
&& (!previous_remote_version_known || previous_remote_version == terminal_remote_version)
}
(
Self::RecoveryControl {
identity_sha256: previous_identity,
source_generation_sha256: previous_generation,
first_seen_at_unix_nanos: previous_first_seen,
revision: previous_revision,
attempt_count: previous_attempts,
..
},
Self::RecoveryControl {
identity_sha256: terminal_identity,
source_generation_sha256: terminal_generation,
first_seen_at_unix_nanos: terminal_first_seen,
revision: terminal_revision,
attempt_count: terminal_attempts,
classification:
recovery_control::IlmRecoveryClassification::Terminal | recovery_control::IlmRecoveryClassification::Abandoned,
..
},
) => {
previous_identity == terminal_identity
&& (previous_generation == terminal_generation || terminal_attempts > previous_attempts)
&& previous_first_seen == terminal_first_seen
&& terminal_revision > previous_revision
&& terminal_attempts >= previous_attempts
}
_ => false,
}
}
@@ -1219,6 +1319,35 @@ pub(crate) fn validate_durable_ilm_record(path: &str, data: &[u8]) -> Result<Val
},
)
}
DurableIlmRecordKind::RecoveryControl => {
let (protocol, control_id) = recovery_control::recovery_control_id_from_record_object_name(path)
.map_err(|err| Error::other(err.to_string()))?;
let control =
recovery_control::IlmRecoveryControl::decode(&control_id, data).map_err(|err| Error::other(err.to_string()))?;
let canonical = recovery_control::recovery_control_record_object_name(protocol, &control_id)
.map_err(|err| Error::other(err.to_string()))?;
if canonical != path || control.identity.protocol != protocol {
return Err(Error::other("ILM recovery control path is not canonical"));
}
let identity_sha256 = checkpoint_hash(&control.identity)?;
let source_generation_sha256 = checkpoint_hash(&control.observed_source_generation)?;
let owner_fence_sha256 = control.owner.as_ref().map(checkpoint_hash).transpose()?;
(
"control_id",
control_id,
DurableIlmRecordCheckpoint::RecoveryControl {
content_sha256,
identity_sha256,
source_generation_sha256,
first_seen_at_unix_nanos: control.first_seen_at_unix_nanos,
revision: control.revision,
classification: control.classification,
attempt_count: control.attempt_count,
consecutive_failure_count: control.consecutive_failure_count,
owner_fence_sha256,
},
)
}
DurableIlmRecordKind::ManualTransitionJob => {
let job_id = manual_transition_job::manual_transition_job_id_from_record_object_name(path)
.map_err(|err| Error::other(err.to_string()))?;
@@ -1412,6 +1541,94 @@ mod tests {
.checkpoint
}
fn recovery_control_fixture() -> recovery_control::IlmRecoveryControl {
let source_path = "ilm/transition-transactions/records/12/34/1234567890abcdef1234567890abcdef.json";
let generation = recovery_control::IlmRecoverySourceGeneration::new(
transition_transaction::TRANSITION_TRANSACTION_SCHEMA,
"source-etag",
"a".repeat(64),
vec![recovery_control::IlmRecoverySourceCopy {
authority: "pool-0/set-0".to_string(),
canonical_path: source_path.to_string(),
etag: "source-etag".to_string(),
encoded_len: 128,
content_sha256: "a".repeat(64),
}],
)
.expect("source generation should build");
recovery_control::IlmRecoveryControl::new(
recovery_control::IlmRecoveryControlIdentity {
protocol: recovery_control::IlmRecoveryProtocol::TransitionTransaction,
canonical_source_path: source_path.to_string(),
stable_operation_identity: "12345678-90ab-cdef-1234-567890abcdef".to_string(),
record_class: "transition_transaction_v1".to_string(),
},
generation,
recovery_control::IlmRecoveryClassification::Retrying,
1_000_000_000,
recovery_control::IlmRecoveryErrorCode::None,
)
.expect("recovery control should build")
}
fn recovery_control_checkpoint(control: &recovery_control::IlmRecoveryControl) -> DurableIlmRecordCheckpoint {
let control_id = control.identity.source_operation_digest().expect("control id should derive");
let path = recovery_control::recovery_control_record_object_name(control.identity.protocol, &control_id)
.expect("control path should build");
let encoded = control.encode().expect("control should encode");
let namespace = classify_durable_ilm_record(&path)
.expect("recovery control namespace should classify")
.expect("recovery control should be durable");
assert_eq!(namespace, &RECOVERY_CONTROL_NAMESPACE);
validate_durable_ilm_record(&path, &encoded)
.expect("recovery control should validate")
.checkpoint
}
#[test]
fn recovery_control_checkpoint_tracks_claim_retry_and_terminal_generations() {
let initial_control = recovery_control_fixture();
let initial = recovery_control_checkpoint(&initial_control);
let mut claimed_control = initial_control;
let mut advanced_generation = claimed_control.observed_source_generation.clone();
advanced_generation.source_schema = "rustfs-transition-transaction-v2".to_string();
claimed_control
.claim_for_source_generation("node-a", Uuid::new_v4(), 2_000_000_000, 300_000_000_000, advanced_generation)
.expect("control should claim");
let claimed = recovery_control_checkpoint(&claimed_control);
initial.validate_successor(&claimed).expect("claim should advance receipt");
let mut retry_control = claimed_control;
retry_control
.record_retryable_failure(3_000_000_000, recovery_control::IlmRecoveryErrorCode::BackendTimeout)
.expect("retry should persist");
let retry = recovery_control_checkpoint(&retry_control);
claimed.validate_successor(&retry).expect("retry should advance receipt");
let ready_at = retry_control
.next_attempt_at_unix_nanos
.expect("retry deadline should persist");
let mut terminal_control = retry_control;
terminal_control
.claim("node-b", Uuid::new_v4(), ready_at, 300_000_000_000)
.expect("retry should claim");
let reclaimed = recovery_control_checkpoint(&terminal_control);
retry.validate_successor(&reclaimed).expect("reclaim should advance receipt");
terminal_control
.finish_attempt(
recovery_control::IlmRecoveryClassification::Terminal,
recovery_control::IlmRecoveryErrorCode::None,
)
.expect("control should terminate");
let terminal = recovery_control_checkpoint(&terminal_control);
reclaimed
.validate_successor(&terminal)
.expect("terminal state should advance receipt");
assert!(initial.is_predecessor_of_terminal(&terminal));
assert!(!initial.is_predecessor_of_terminal(&retry));
}
#[test]
fn tier_probe_intent_checkpoint_tracks_exact_monotonic_generations() {
let initial_intent = tier_probe_intent_fixture();
@@ -24,6 +24,7 @@ pub(crate) use metadata_boundary::{LifecycleExpiryConfigs, get_expiry_configs, g
mod object_handlers_common;
mod object_lock_boundary;
pub use self::core as lifecycle;
pub mod recovery_control;
mod replication_sink;
pub mod rule;
mod runtime_boundary;
File diff suppressed because it is too large Load Diff
@@ -23,6 +23,11 @@ use uuid::Uuid;
use crate::bucket::lifecycle::config_boundary;
use crate::bucket::lifecycle::durable_namespace::TRANSITION_TRANSACTION_NAMESPACE;
use crate::bucket::lifecycle::lifecycle::TRANSITION_COMPLETE;
use crate::bucket::lifecycle::recovery_control::{
IlmRecoveryClassification, IlmRecoveryControl, IlmRecoveryControlIdentity, IlmRecoveryErrorCode, IlmRecoveryProtocol,
ObservedIlmRecoveryControl, load_recovery_control, observe_recovery_source, recovery_control_record_object_name,
save_recovery_control_if_absent, save_recovery_control_if_current,
};
use crate::bucket::lifecycle::tier_sweeper::{
delete_confirmed_transition_candidate_exact_with_lease_idempotent,
delete_object_from_remote_tier_idempotent_with_manager_and_identity,
@@ -44,6 +49,7 @@ const EVENT_LIFECYCLE_TRANSITION_TRANSACTION_RECOVERY: &str = "lifecycle_transit
pub const DEFAULT_TRANSITION_TRANSACTION_RECOVERY_LIMIT: usize = 1_000;
const TRANSITION_TRANSACTION_RECOVERY_INTERVAL: Duration = Duration::from_secs(60);
const TRANSITION_TRANSACTION_RECOVERY_TIMEOUT: Duration = Duration::from_secs(300);
const TRANSITION_RECOVERY_CONTROL_LEASE_NANOS: i64 = 15 * 60 * 1_000_000_000;
pub const TRANSITION_TRANSACTION_SCHEMA: &str = "rustfs-transition-transaction-v1";
pub const TRANSITION_TRANSACTION_PREFIX: &str = "ilm/transition-transactions";
pub const TRANSITION_TRANSACTION_RECORD_PREFIX: &str = TRANSITION_TRANSACTION_NAMESPACE.prefix;
@@ -737,6 +743,8 @@ pub enum TransitionTransactionRecoveryOutcome {
RemoteCandidateDeleted,
RecordDeleted,
Retained,
RetainedAmbiguous(IlmRecoveryErrorCode),
OperatorRequired(IlmRecoveryErrorCode),
}
#[cfg(test)]
@@ -817,6 +825,80 @@ async fn pause_before_transition_recovery_claim(transaction_id: Uuid) {
}
}
#[cfg(test)]
#[derive(Default)]
struct TransitionRecoveryTerminalBarrierState {
transaction_id: Uuid,
arrived: tokio::sync::Notify,
release: tokio::sync::Notify,
}
#[cfg(test)]
pub(crate) struct TransitionRecoveryTerminalBarrier {
state: Arc<TransitionRecoveryTerminalBarrierState>,
}
#[cfg(test)]
static TRANSITION_RECOVERY_TERMINAL_BARRIER: std::sync::OnceLock<
std::sync::Mutex<Option<Arc<TransitionRecoveryTerminalBarrierState>>>,
> = std::sync::OnceLock::new();
#[cfg(test)]
impl TransitionRecoveryTerminalBarrier {
pub(crate) fn install(transaction_id: Uuid) -> Self {
let state = Arc::new(TransitionRecoveryTerminalBarrierState {
transaction_id,
..Default::default()
});
let mut slot = TRANSITION_RECOVERY_TERMINAL_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("transition recovery terminal barrier mutex should not poison");
assert!(
slot.is_none(),
"transition recovery terminal barrier must be installed by one test at a time"
);
*slot = Some(Arc::clone(&state));
drop(slot);
Self { state }
}
pub(crate) async fn wait_until_paused(&self) {
tokio::time::timeout(Duration::from_secs(30), self.state.arrived.notified())
.await
.expect("transition recovery should persist terminal control before source cleanup");
}
}
#[cfg(test)]
impl Drop for TransitionRecoveryTerminalBarrier {
fn drop(&mut self) {
self.state.release.notify_one();
let mut slot = TRANSITION_RECOVERY_TERMINAL_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("transition recovery terminal barrier mutex should not poison");
if slot.as_ref().is_some_and(|state| Arc::ptr_eq(state, &self.state)) {
*slot = None;
}
}
}
#[cfg(test)]
async fn pause_after_transition_recovery_terminal(transaction_id: Uuid) {
let barrier = TRANSITION_RECOVERY_TERMINAL_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("transition recovery terminal barrier mutex should not poison")
.as_ref()
.filter(|barrier| barrier.transaction_id == transaction_id)
.cloned();
if let Some(barrier) = barrier {
barrier.arrived.notify_one();
barrier.release.notified().await;
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
#[serde(rename_all = "snake_case")]
pub enum TransitionOperatorProbe {
@@ -1020,17 +1102,35 @@ fn transition_transaction_id_from_record_object_name(object: &str) -> Result<Uui
let suffix = object
.strip_prefix(&prefix)
.ok_or(TransitionTransactionError::Corrupt("transaction record path has wrong prefix"))?;
let file_name = suffix
.rsplit('/')
let mut parts = suffix.split('/');
let shard_a = parts
.next()
.ok_or(TransitionTransactionError::Corrupt("transaction record path is incomplete"))?;
let shard_b = parts
.next()
.ok_or(TransitionTransactionError::Corrupt("transaction record path is incomplete"))?;
let file_name = parts
.next()
.ok_or(TransitionTransactionError::Corrupt("transaction record path is incomplete"))?;
if parts.next().is_some() {
return Err(TransitionTransactionError::Corrupt("transaction record path is not canonical"));
}
let transaction_key = file_name
.strip_suffix(".json")
.ok_or(TransitionTransactionError::Corrupt("transaction record path has wrong suffix"))?;
if transaction_key.len() != 32 || !transaction_key.bytes().all(|byte| byte.is_ascii_hexdigit()) {
if transaction_key.len() != 32
|| !transaction_key
.bytes()
.all(|byte| byte.is_ascii_hexdigit() && !byte.is_ascii_uppercase())
|| shard_a != &transaction_key[..2]
|| shard_b != &transaction_key[2..4]
{
return Err(TransitionTransactionError::Corrupt("transaction record path has invalid transaction id"));
}
Uuid::parse_str(transaction_key).map_err(|_| TransitionTransactionError::Corrupt("transaction record path has invalid uuid"))
Uuid::parse_str(transaction_key)
.ok()
.filter(|transaction_id| !transaction_id.is_nil())
.ok_or(TransitionTransactionError::Corrupt("transaction record path has invalid uuid"))
}
pub async fn process_transition_transaction_record(
@@ -1055,6 +1155,27 @@ async fn process_transition_transaction_record_at(
) -> EcstoreResult<TransitionTransactionRecoveryOutcome> {
let record_name =
transition_transaction_record_object_name(observed.transaction_id).map_err(transition_transaction_store_error)?;
let now_unix_nanos =
i64::try_from(now_unix_nanos).map_err(|_| Error::other("transition transaction recovery timestamp does not fit i64"))?;
let recovery_control_identity = transition_recovery_control_identity(observed, &record_name);
let recovery_control_id = recovery_control_identity
.source_operation_digest()
.map_err(|err| Error::other(err.to_string()))?;
let control_record_name =
recovery_control_record_object_name(IlmRecoveryProtocol::TransitionTransaction, &recovery_control_id)
.map_err(|err| Error::other(err.to_string()))?;
let control_lock = if transition_state_needs_recovery_control(observed, now_unix_nanos) {
Some(
api.new_ns_lock(RUSTFS_META_BUCKET, &format!("{control_record_name}.recovery-lock"))
.await?,
)
} else {
None
};
let _control_guard = match &control_lock {
Some(lock) => Some(lock.get_write_lock(crate::set_disk::get_lock_acquire_timeout()).await?),
None => None,
};
// The synthetic key avoids nesting the recovery lock with the config
// object's own I/O lock. Holding it across the bounded source proof and
// remote DELETE elects one destructive recovery worker across nodes.
@@ -1073,55 +1194,400 @@ async fn process_transition_transaction_record_at(
return Ok(TransitionTransactionRecoveryOutcome::Retained);
}
match current.state {
let mut recovery_control = if transition_state_needs_recovery_control(&current, now_unix_nanos) {
if cleanup_terminal_transition_recovery_control(
api.clone(),
&current,
&record_name,
&recovery_control_identity,
&recovery_control_id,
)
.await?
{
return Ok(TransitionTransactionRecoveryOutcome::RecordDeleted);
}
match claim_transition_recovery_control(
api.clone(),
&current,
&record_name,
recovery_control_identity,
&recovery_control_id,
now_unix_nanos,
)
.await?
{
Some(control) => Some(control),
None => return Ok(TransitionTransactionRecoveryOutcome::Retained),
}
} else {
None
};
let recovery = match current.state {
TransitionTransactionState::Uploaded => {
if transition_transaction_ownership_is_active(&current, now_unix_nanos) {
return Ok(TransitionTransactionRecoveryOutcome::Retained);
}
let mut cleanup = current.clone();
cleanup
.mark_cleanup_pending(
current.fence(),
TransitionCleanupProof {
transaction_id: current.transaction_id,
write_id: current.write_id,
remote_object: current.remote_object.clone(),
remote_version: current.remote_version.clone(),
backend_fingerprint: current.backend_fingerprint,
decision: TransitionCleanupDecision::UploadAbortedBeforeLocalCommit,
},
)
.map_err(transition_transaction_store_error)?;
#[cfg(test)]
pause_before_transition_recovery_claim(current.transaction_id).await;
match save_transition_transaction_record_if_current(api.clone(), &current, &cleanup).await {
Ok(()) => recover_cleanup_pending(api, &cleanup).await,
Err(Error::PreconditionFailed) | Err(Error::ConfigNotFound) => Ok(TransitionTransactionRecoveryOutcome::Retained),
Err(err) => Err(err),
if transition_transaction_ownership_is_active(&current, i128::from(now_unix_nanos)) {
Ok(TransitionTransactionRecoveryOutcome::Retained)
} else {
let mut cleanup = current.clone();
cleanup
.mark_cleanup_pending(
current.fence(),
TransitionCleanupProof {
transaction_id: current.transaction_id,
write_id: current.write_id,
remote_object: current.remote_object.clone(),
remote_version: current.remote_version.clone(),
backend_fingerprint: current.backend_fingerprint,
decision: TransitionCleanupDecision::UploadAbortedBeforeLocalCommit,
},
)
.map_err(transition_transaction_store_error)?;
#[cfg(test)]
pause_before_transition_recovery_claim(current.transaction_id).await;
match save_transition_transaction_record_if_current(api.clone(), &current, &cleanup).await {
Ok(()) => recover_cleanup_pending(api.clone(), &cleanup).await,
Err(Error::PreconditionFailed) | Err(Error::ConfigNotFound) => {
Ok(TransitionTransactionRecoveryOutcome::Retained)
}
Err(err) => Err(err),
}
}
}
TransitionTransactionState::CleanupPending => recover_cleanup_pending(api, &current).await,
TransitionTransactionState::CleanupPending => recover_cleanup_pending(api.clone(), &current).await,
TransitionTransactionState::LocalCommitStarted => match local_commit_matches_transaction(api.clone(), &current).await {
Ok(true) => {
delete_transition_transaction_record(api, &current).await?;
Ok(TransitionTransactionRecoveryOutcome::RecordDeleted)
}
Ok(false) => Ok(TransitionTransactionRecoveryOutcome::Retained),
Err(err) if transition_source_is_missing(&err) => Ok(TransitionTransactionRecoveryOutcome::Retained),
Ok(true) => Ok(TransitionTransactionRecoveryOutcome::RecordDeleted),
Ok(false) => Ok(TransitionTransactionRecoveryOutcome::OperatorRequired(
IlmRecoveryErrorCode::LocalCommitAmbiguous,
)),
Err(err) if transition_source_is_missing(&err) => Ok(TransitionTransactionRecoveryOutcome::OperatorRequired(
IlmRecoveryErrorCode::LocalCommitAmbiguous,
)),
Err(err) => Err(err),
},
TransitionTransactionState::AbortedNoRemote | TransitionTransactionState::Committed => {
delete_transition_transaction_record(api, &current).await?;
Ok(TransitionTransactionRecoveryOutcome::RecordDeleted)
}
TransitionTransactionState::UploadOutcomeUnknown => {
if transition_transaction_ownership_is_active(&current, now_unix_nanos) {
if transition_transaction_ownership_is_active(&current, i128::from(now_unix_nanos)) {
Ok(TransitionTransactionRecoveryOutcome::Retained)
} else {
recover_unknown_upload_outcome(api, &current).await
recover_unknown_upload_outcome(api.clone(), &current).await
}
}
TransitionTransactionState::UploadStarted => Ok(TransitionTransactionRecoveryOutcome::Retained),
TransitionTransactionState::UploadStarted => {
if transition_transaction_ownership_is_active(&current, i128::from(now_unix_nanos)) {
Ok(TransitionTransactionRecoveryOutcome::Retained)
} else {
Ok(TransitionTransactionRecoveryOutcome::RetainedAmbiguous(
IlmRecoveryErrorCode::RemoteVersionUnknown,
))
}
}
};
if let Some(mut control) = recovery_control.take() {
let source_to_delete = if matches!(
recovery,
Ok(TransitionTransactionRecoveryOutcome::RemoteCandidateDeleted
| TransitionTransactionRecoveryOutcome::RecordDeleted)
) {
let refreshed =
refresh_transition_recovery_control_source(api.clone(), control, &record_name, current.transaction_id).await?;
control = refreshed.0;
refreshed.1
} else {
None
};
persist_transition_recovery_result(api.clone(), control, &recovery, now_unix_nanos).await?;
if let Some(source) = source_to_delete {
#[cfg(test)]
pause_after_transition_recovery_terminal(source.transaction_id).await;
delete_transition_transaction_record(api, &source).await?;
}
} else if matches!(
recovery,
Ok(TransitionTransactionRecoveryOutcome::RemoteCandidateDeleted | TransitionTransactionRecoveryOutcome::RecordDeleted)
) {
delete_transition_transaction_record(api, &current).await?;
}
recovery
}
fn transition_recovery_control_identity(transaction: &TransitionTransaction, record_name: &str) -> IlmRecoveryControlIdentity {
IlmRecoveryControlIdentity {
protocol: IlmRecoveryProtocol::TransitionTransaction,
canonical_source_path: record_name.to_string(),
stable_operation_identity: transaction.transaction_id.to_string(),
record_class: "transition_transaction_v1".to_string(),
}
}
#[cfg(test)]
pub(crate) fn transition_recovery_control_id(transaction: &TransitionTransaction) -> Result<String> {
let record_name = transition_transaction_record_object_name(transaction.transaction_id)?;
transition_recovery_control_identity(transaction, &record_name)
.source_operation_digest()
.map_err(|_| TransitionTransactionError::Corrupt("transition recovery control identity is invalid"))
}
fn transition_state_needs_recovery_control(transaction: &TransitionTransaction, now_unix_nanos: i64) -> bool {
now_unix_nanos >= transaction.not_after_unix_nanos
&& !matches!(
transaction.state,
TransitionTransactionState::AbortedNoRemote | TransitionTransactionState::Committed
)
}
async fn cleanup_terminal_transition_recovery_control(
api: Arc<ECStore>,
transaction: &TransitionTransaction,
record_name: &str,
identity: &IlmRecoveryControlIdentity,
control_id: &str,
) -> EcstoreResult<bool> {
let observed = match load_recovery_control(api.clone(), IlmRecoveryProtocol::TransitionTransaction, control_id).await {
Ok(observed) => observed,
Err(Error::ConfigNotFound) => return Ok(false),
Err(err) => return Err(err),
};
if observed.control.classification != IlmRecoveryClassification::Terminal {
return Ok(false);
}
let source = observe_recovery_source(api.clone(), record_name, TRANSITION_TRANSACTION_SCHEMA).await?;
let exact_source = source.is_consistent()
&& source.generation == observed.control.observed_source_generation
&& source.canonical_data.as_deref().is_some_and(|data| {
TransitionTransaction::decode(transaction.transaction_id, data).is_ok_and(|decoded| decoded == *transaction)
});
if observed.control.identity != *identity || !exact_source {
return Ok(false);
}
delete_transition_transaction_record(api, transaction).await?;
Ok(true)
}
async fn claim_transition_recovery_control(
api: Arc<ECStore>,
transaction: &TransitionTransaction,
record_name: &str,
identity: IlmRecoveryControlIdentity,
control_id: &str,
now_unix_nanos: i64,
) -> EcstoreResult<Option<ObservedIlmRecoveryControl>> {
let existing = match load_recovery_control(api.clone(), IlmRecoveryProtocol::TransitionTransaction, control_id).await {
Ok(control) => Some(control),
Err(Error::ConfigNotFound) => None,
Err(err) => return Err(err),
};
if let Some(observed) = existing.as_ref() {
if observed.control.identity != identity {
return Ok(None);
}
if observed
.control
.owner
.as_ref()
.is_some_and(|owner| owner.lease_expires_at_unix_nanos <= now_unix_nanos)
{
let mut expired = observed.control.clone();
expired
.record_expired_attempt(now_unix_nanos)
.map_err(|err| Error::other(err.to_string()))?;
save_recovery_control_if_current(api, observed, &expired).await?;
return Ok(None);
}
if !observed.control.should_attempt_at(now_unix_nanos) {
return Ok(None);
}
}
let source = match observe_recovery_source(api.clone(), record_name, TRANSITION_TRANSACTION_SCHEMA).await {
Ok(source) => source,
Err(err) => {
if let Some(observed) = existing {
persist_transition_recovery_source_failure(api, observed, now_unix_nanos).await?;
return Ok(None);
}
return Err(err);
}
};
let source_matches = source.is_consistent()
&& source.canonical_data.as_deref().is_some_and(|data| {
TransitionTransaction::decode(transaction.transaction_id, data).is_ok_and(|observed| observed == *transaction)
});
let source_error = if source_matches {
IlmRecoveryErrorCode::None
} else if source.canonical_data.is_some() {
IlmRecoveryErrorCode::SourceGenerationChanged
} else {
IlmRecoveryErrorCode::SourceDivergent
};
let mut observed = match existing {
Some(control) => control,
None => {
let candidate = IlmRecoveryControl::new(
identity.clone(),
source.generation.clone(),
if source_matches {
IlmRecoveryClassification::Retrying
} else {
IlmRecoveryClassification::Corrupt
},
now_unix_nanos,
source_error,
)
.map_err(|err| Error::other(err.to_string()))?;
match save_recovery_control_if_absent(api.clone(), &candidate).await {
Ok(()) | Err(Error::PreconditionFailed) => {}
Err(err) => return Err(err),
}
load_recovery_control(api.clone(), IlmRecoveryProtocol::TransitionTransaction, control_id).await?
}
};
if observed.control.identity != identity || !observed.control.should_attempt_at(now_unix_nanos) {
return Ok(None);
}
let mut claimed = observed.control.clone();
claimed
.claim_for_source_generation(
api.id.to_string(),
Uuid::new_v4(),
now_unix_nanos,
TRANSITION_RECOVERY_CONTROL_LEASE_NANOS,
source.generation,
)
.map_err(|err| Error::other(err.to_string()))?;
save_recovery_control_if_current(api.clone(), &observed, &claimed).await?;
observed = load_recovery_control(api.clone(), IlmRecoveryProtocol::TransitionTransaction, control_id).await?;
if observed.control != claimed {
return Err(Error::PreconditionFailed);
}
if !source_matches {
let mut corrupt = observed.control.clone();
corrupt
.finish_attempt(IlmRecoveryClassification::Corrupt, source_error)
.map_err(|err| Error::other(err.to_string()))?;
save_recovery_control_if_current(api, &observed, &corrupt).await?;
return Ok(None);
}
Ok(Some(observed))
}
async fn persist_transition_recovery_source_failure(
api: Arc<ECStore>,
observed: ObservedIlmRecoveryControl,
now_unix_nanos: i64,
) -> EcstoreResult<()> {
let mut claimed = observed.control.clone();
claimed
.claim(
api.id.to_string(),
Uuid::new_v4(),
now_unix_nanos,
TRANSITION_RECOVERY_CONTROL_LEASE_NANOS,
)
.map_err(|err| Error::other(err.to_string()))?;
save_recovery_control_if_current(api.clone(), &observed, &claimed).await?;
let claimed = load_recovery_control(
api.clone(),
IlmRecoveryProtocol::TransitionTransaction,
&claimed
.identity
.source_operation_digest()
.map_err(|err| Error::other(err.to_string()))?,
)
.await?;
let mut failed = claimed.control.clone();
failed
.record_retryable_failure(now_unix_nanos, IlmRecoveryErrorCode::SourceUnavailable)
.map_err(|err| Error::other(err.to_string()))?;
save_recovery_control_if_current(api, &claimed, &failed).await
}
async fn refresh_transition_recovery_control_source(
api: Arc<ECStore>,
mut observed: ObservedIlmRecoveryControl,
record_name: &str,
transaction_id: Uuid,
) -> EcstoreResult<(ObservedIlmRecoveryControl, Option<TransitionTransaction>)> {
let transaction = match load_transition_transaction_record(api.clone(), transaction_id).await {
Ok(transaction) => transaction,
Err(Error::ConfigNotFound) => return Ok((observed, None)),
Err(err) => return Err(err),
};
let source = observe_recovery_source(api.clone(), record_name, TRANSITION_TRANSACTION_SCHEMA).await?;
let exact_source = source.is_consistent()
&& source
.canonical_data
.as_deref()
.is_some_and(|data| TransitionTransaction::decode(transaction_id, data).is_ok_and(|decoded| decoded == transaction));
if !exact_source {
return Err(Error::PreconditionFailed);
}
if observed.control.observed_source_generation != source.generation {
let mut refreshed = observed.control.clone();
refreshed
.refresh_owned_source_generation(source.generation)
.map_err(|err| Error::other(err.to_string()))?;
save_recovery_control_if_current(api.clone(), &observed, &refreshed).await?;
observed = load_recovery_control(
api,
IlmRecoveryProtocol::TransitionTransaction,
&refreshed
.identity
.source_operation_digest()
.map_err(|err| Error::other(err.to_string()))?,
)
.await?;
if observed.control != refreshed {
return Err(Error::PreconditionFailed);
}
}
Ok((observed, Some(transaction)))
}
async fn persist_transition_recovery_result(
api: Arc<ECStore>,
observed: ObservedIlmRecoveryControl,
recovery: &EcstoreResult<TransitionTransactionRecoveryOutcome>,
now_unix_nanos: i64,
) -> EcstoreResult<()> {
let mut next = observed.control.clone();
match recovery {
Ok(
TransitionTransactionRecoveryOutcome::RemoteCandidateDeleted | TransitionTransactionRecoveryOutcome::RecordDeleted,
) => next
.finish_attempt(IlmRecoveryClassification::Terminal, IlmRecoveryErrorCode::None)
.map_err(|err| Error::other(err.to_string()))?,
Ok(TransitionTransactionRecoveryOutcome::Retained) => next
.record_retryable_failure(now_unix_nanos, IlmRecoveryErrorCode::SourceGenerationChanged)
.map_err(|err| Error::other(err.to_string()))?,
Ok(TransitionTransactionRecoveryOutcome::RetainedAmbiguous(code)) => next
.finish_attempt(IlmRecoveryClassification::RetainedAmbiguous, *code)
.map_err(|err| Error::other(err.to_string()))?,
Ok(TransitionTransactionRecoveryOutcome::OperatorRequired(code)) => next
.finish_attempt(IlmRecoveryClassification::OperatorRequired, *code)
.map_err(|err| Error::other(err.to_string()))?,
Err(err) => next
.record_retryable_failure(now_unix_nanos, transition_recovery_error_code(err))
.map_err(|err| Error::other(err.to_string()))?,
}
save_recovery_control_if_current(api, &observed, &next).await
}
fn transition_recovery_error_code(err: &Error) -> IlmRecoveryErrorCode {
match err {
Error::PreconditionFailed => IlmRecoveryErrorCode::CasConflict,
Error::ConfigNotFound
| Error::FileNotFound
| Error::FileVersionNotFound
| Error::ObjectNotFound(_, _)
| Error::VersionNotFound(_, _, _)
| Error::BucketNotFound(_) => IlmRecoveryErrorCode::SourceUnavailable,
Error::SlowDown => IlmRecoveryErrorCode::BackendThrottled,
_ => IlmRecoveryErrorCode::Unknown,
}
}
@@ -1134,10 +1600,7 @@ async fn recover_cleanup_pending(
transaction: &TransitionTransaction,
) -> EcstoreResult<TransitionTransactionRecoveryOutcome> {
match local_commit_matches_transaction(api.clone(), transaction).await {
Ok(true) => {
delete_transition_transaction_record(api, transaction).await?;
Ok(TransitionTransactionRecoveryOutcome::RecordDeleted)
}
Ok(true) => Ok(TransitionTransactionRecoveryOutcome::RecordDeleted),
Ok(false) => delete_unreferenced_transition_candidate(api, transaction).await,
Err(err) if transition_source_is_missing(&err) => delete_unreferenced_transition_candidate(api, transaction).await,
Err(err) => Err(err),
@@ -1157,7 +1620,6 @@ async fn delete_unreferenced_transition_candidate(
return Ok(TransitionTransactionRecoveryOutcome::Retained);
}
delete_transition_remote_candidate(api.clone(), &current).await?;
delete_transition_transaction_record(api, &current).await?;
Ok(TransitionTransactionRecoveryOutcome::RemoteCandidateDeleted)
}
@@ -1178,24 +1640,26 @@ async fn recover_unknown_upload_outcome(
.await
.map_err(Error::other)?
{
TransitionCandidateProbe::Missing => {
delete_transition_transaction_record(api, transaction).await?;
Ok(TransitionTransactionRecoveryOutcome::RecordDeleted)
}
TransitionCandidateProbe::Missing => Ok(TransitionTransactionRecoveryOutcome::RecordDeleted),
TransitionCandidateProbe::UnversionedPresent => {
cleanup_recovered_unknown_upload_candidate(api, transaction, TransitionRemoteVersion::unversioned()).await
}
TransitionCandidateProbe::VersionedPresent(version_id)
if Uuid::parse_str(&version_id).is_ok_and(|version_id| version_id.is_nil()) =>
{
Ok(TransitionTransactionRecoveryOutcome::Retained)
Ok(TransitionTransactionRecoveryOutcome::RetainedAmbiguous(
IlmRecoveryErrorCode::RemoteVersionUnknown,
))
}
TransitionCandidateProbe::VersionedPresent(version_id) => {
cleanup_recovered_unknown_upload_candidate(api, transaction, TransitionRemoteVersion::versioned(version_id)).await
}
TransitionCandidateProbe::Ambiguous | TransitionCandidateProbe::Unsupported => {
Ok(TransitionTransactionRecoveryOutcome::Retained)
}
TransitionCandidateProbe::Ambiguous => Ok(TransitionTransactionRecoveryOutcome::RetainedAmbiguous(
IlmRecoveryErrorCode::RemoteProbeAmbiguous,
)),
TransitionCandidateProbe::Unsupported => Ok(TransitionTransactionRecoveryOutcome::RetainedAmbiguous(
IlmRecoveryErrorCode::RemoteProbeUnsupported,
)),
}
}
@@ -1323,6 +1787,11 @@ async fn recover_transition_transaction_records_with_now(
false,
)
.await?;
if list.is_truncated && list.next_continuation_token.is_none() {
return Err(Error::other(
"transition transaction recovery returned a truncated page without a continuation marker",
));
}
let mut stats = TransitionTransactionRecoveryStats {
scanned: 0,
@@ -1381,7 +1850,11 @@ async fn recover_transition_transaction_records_with_now(
) => {
stats.recovered += 1;
}
Ok(TransitionTransactionRecoveryOutcome::Retained) => {
Ok(
TransitionTransactionRecoveryOutcome::Retained
| TransitionTransactionRecoveryOutcome::RetainedAmbiguous(_)
| TransitionTransactionRecoveryOutcome::OperatorRequired(_),
) => {
stats.retained += 1;
debug!(
event = EVENT_LIFECYCLE_TRANSITION_TRANSACTION_RECOVERY,
@@ -1509,11 +1982,74 @@ fn state_requires_known_remote_version(state: TransitionTransactionState) -> boo
#[cfg(test)]
mod tests {
use std::collections::HashMap;
use std::sync::atomic::{AtomicBool, Ordering};
use super::*;
const BACKEND_FINGERPRINT: [u8; 32] = [7; 32];
struct RecoveryAttemptDropGuard(Arc<AtomicBool>);
impl Drop for RecoveryAttemptDropGuard {
fn drop(&mut self) {
self.0.store(true, Ordering::SeqCst);
}
}
async fn pending_recovery_attempt(started: Arc<tokio::sync::Notify>, dropped: Arc<AtomicBool>) -> EcstoreResult<()> {
let _drop_guard = RecoveryAttemptDropGuard(dropped);
started.notify_one();
std::future::pending().await
}
#[tokio::test(start_paused = true)]
async fn transition_recovery_timeout_and_cancellation_drop_inflight_attempts() {
let timeout_started = Arc::new(tokio::sync::Notify::new());
let timeout_dropped = Arc::new(AtomicBool::new(false));
let timeout_task = tokio::spawn({
let started = Arc::clone(&timeout_started);
let dropped = Arc::clone(&timeout_dropped);
async move {
await_transition_transaction_recovery(
&CancellationToken::new(),
TRANSITION_TRANSACTION_RECOVERY_TIMEOUT,
pending_recovery_attempt(started, dropped),
)
.await
}
});
timeout_started.notified().await;
tokio::time::advance(TRANSITION_TRANSACTION_RECOVERY_TIMEOUT).await;
let timed_out = timeout_task.await.expect("timeout wrapper task should join");
assert!(matches!(timed_out, Some(Err(_))), "outer timeout should fail the recovery pass");
assert!(timeout_dropped.load(Ordering::SeqCst), "outer timeout must drop its in-flight attempt");
let cancel_token = CancellationToken::new();
let cancel_started = Arc::new(tokio::sync::Notify::new());
let cancel_dropped = Arc::new(AtomicBool::new(false));
let cancel_task = tokio::spawn({
let cancel_token = cancel_token.clone();
let started = Arc::clone(&cancel_started);
let dropped = Arc::clone(&cancel_dropped);
async move {
await_transition_transaction_recovery(
&cancel_token,
TRANSITION_TRANSACTION_RECOVERY_TIMEOUT,
pending_recovery_attempt(started, dropped),
)
.await
}
});
cancel_started.notified().await;
cancel_token.cancel();
let cancelled = cancel_task.await.expect("cancellation wrapper task should join");
assert!(cancelled.is_none(), "outer cancellation should stop the recovery loop");
assert!(
cancel_dropped.load(Ordering::SeqCst),
"outer cancellation must drop its in-flight attempt"
);
}
#[derive(Default)]
struct MemoryTransactionStore {
records: HashMap<Uuid, Vec<u8>>,
@@ -1968,5 +2504,19 @@ mod tests {
transition_transaction_record_object_name(Uuid::nil()),
Err(TransitionTransactionError::Corrupt("transaction_id is nil"))
));
assert_eq!(
transition_transaction_id_from_record_object_name(&object).expect("canonical record path should parse"),
transaction_id
);
for malformed in [
object.to_ascii_uppercase(),
object.replace("/aa/aa/", "/ff/aa/"),
object.replace("/aa/aa/", "/aa/aa/extra/"),
] {
assert!(matches!(
transition_transaction_id_from_record_object_name(&malformed),
Err(TransitionTransactionError::Corrupt(_))
));
}
}
}
+283 -22
View File
@@ -825,6 +825,11 @@ mod tests {
manual_transition_scope_record_object_name, manual_transition_task_object_name,
manual_transition_worker_result_object_name, manual_transition_worker_result_task_key,
},
recovery_control::{
IlmRecoveryClassification, IlmRecoveryControl, IlmRecoveryControlIdentity, IlmRecoveryErrorCode,
IlmRecoveryProtocol, MAX_RECOVERY_ATTEMPTS, load_recovery_control, observe_recovery_source,
save_recovery_control_if_absent,
},
tier_delete_journal::{
DecommissionCheckpointTargetFailureHook, TIER_DELETE_DISPATCH_MANIFEST_PREFIX, TIER_DELETE_JOURNAL_PREFIX,
TierDeleteChunkTestBarrier, TierDeleteChunkTestStage, TierDeleteDispatchBatchLimitGuard,
@@ -844,12 +849,13 @@ mod tests {
},
transition_transaction::{
TRANSITION_TRANSACTION_RECORD_PREFIX, TransitionCleanupDecision, TransitionCleanupProof, TransitionOperatorError,
TransitionOperatorProbe, TransitionRecoveryClaimBarrier, TransitionRemoteVersion, TransitionSourceIdentity,
TransitionSourceVersionMode, TransitionTransaction, TransitionTransactionInit, TransitionTransactionState,
delete_transition_candidate_for_operator, finalize_missing_transition_transaction_for_operator,
inspect_transition_transaction_for_operator, load_transition_transaction_record,
recover_transition_transaction_records, recover_transition_transaction_records_at,
save_transition_transaction_record, save_transition_transaction_record_if_current,
TransitionOperatorProbe, TransitionRecoveryClaimBarrier, TransitionRecoveryTerminalBarrier,
TransitionRemoteVersion, TransitionSourceIdentity, TransitionSourceVersionMode, TransitionTransaction,
TransitionTransactionInit, TransitionTransactionState, delete_transition_candidate_for_operator,
finalize_missing_transition_transaction_for_operator, inspect_transition_transaction_for_operator,
load_transition_transaction_record, recover_transition_transaction_records,
recover_transition_transaction_records_at, save_transition_transaction_record,
save_transition_transaction_record_if_current, transition_recovery_control_id,
transition_transaction_record_object_name,
},
validate_durable_ilm_record,
@@ -19239,6 +19245,105 @@ mod tests {
assert!(!Arc::ptr_eq(&ctx_a, &ctx_b), "the regression requires two distinct instance contexts");
}
#[cfg(feature = "test-util")]
#[tokio::test]
#[serial_test::serial(storage_class_env)]
async fn transition_transaction_recovery_expires_abandoned_attempt_at_budget_bound() {
let temp_dir = tempfile::tempdir().expect("create temp store dir");
let (ctx, store, _shutdown) = without_storage_class_env(build_isolated_test_store(
temp_dir.path(),
"transition-transaction-expired-attempt-budget",
&[4],
))
.await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
let transaction = TransitionTransaction::new(TransitionTransactionInit {
deployment_id: ctx.deployment_id().expect("test store should initialize deployment id"),
transaction_id: uuid::Uuid::new_v4(),
owner_epoch: uuid::Uuid::new_v4(),
write_id: uuid::Uuid::new_v4(),
source: TransitionSourceIdentity {
bucket: "source-bucket".to_string(),
object: "source-object".to_string(),
version_id: Some(uuid::Uuid::new_v4()),
data_dir: uuid::Uuid::new_v4(),
mod_time_unix_nanos: 1_770_000_000_000_000_000,
size: 42,
etag: "source-etag".to_string(),
version_mode: TransitionSourceVersionMode::Versioned,
},
tier_name: "UNUSEDABANDONEDTIER".to_string(),
backend_fingerprint: [7; 32],
not_after_unix_nanos: 1,
})
.expect("transaction should build");
save_transition_transaction_record(store.clone(), &transaction)
.await
.expect("transaction record should persist");
let record_name =
transition_transaction_record_object_name(transaction.transaction_id).expect("transaction record name should derive");
let source = observe_recovery_source(
store.clone(),
&record_name,
crate::bucket::lifecycle::transition_transaction::TRANSITION_TRANSACTION_SCHEMA,
)
.await
.expect("transaction source generation should be observable");
let mut control = IlmRecoveryControl::new(
IlmRecoveryControlIdentity {
protocol: IlmRecoveryProtocol::TransitionTransaction,
canonical_source_path: record_name,
stable_operation_identity: transaction.transaction_id.to_string(),
record_class: "transition_transaction_v1".to_string(),
},
source.generation,
IlmRecoveryClassification::Retrying,
2_000_000_000,
IlmRecoveryErrorCode::None,
)
.expect("recovery control should build");
let mut now = 3_000_000_000;
for _ in 1..MAX_RECOVERY_ATTEMPTS {
now = now.max(control.next_attempt_at_unix_nanos.unwrap_or(now));
control
.claim("cancelled-or-timed-out-owner", uuid::Uuid::new_v4(), now, 1)
.expect("abandoned attempt should claim");
control
.record_expired_attempt(now + 1)
.expect("expired attempt should consume retry budget");
now += 2;
}
now = now.max(control.next_attempt_at_unix_nanos.expect("last retry should have a backoff"));
control
.claim("cancelled-or-timed-out-owner", uuid::Uuid::new_v4(), now, 1)
.expect("final abandoned attempt should claim");
save_recovery_control_if_absent(store.clone(), &control)
.await
.expect("claimed recovery control should persist");
let stats = recover_transition_transaction_records_at(store.clone(), 100, None, i128::from(now + 1))
.await
.expect("recovery should account for the expired attempt");
assert_eq!((stats.scanned, stats.recovered, stats.retained, stats.failed), (1, 0, 1, 0));
let control_id = transition_recovery_control_id(&transaction).expect("control id should derive");
let persisted = load_recovery_control(store.clone(), IlmRecoveryProtocol::TransitionTransaction, &control_id)
.await
.expect("expired recovery control should remain inspectable");
assert_eq!(persisted.control.classification, IlmRecoveryClassification::OperatorRequired);
assert_eq!(persisted.control.attempt_count, u64::from(MAX_RECOVERY_ATTEMPTS));
assert_eq!(persisted.control.consecutive_failure_count, MAX_RECOVERY_ATTEMPTS);
assert_eq!(persisted.control.last_error_code, IlmRecoveryErrorCode::AttemptLeaseExpired);
assert!(persisted.control.owner.is_none());
assert_eq!(
transition_transaction_record_count(store).await,
1,
"budget exhaustion must retain the source record"
);
}
#[cfg(feature = "test-util")]
#[tokio::test]
#[serial_test::serial(storage_class_env)]
@@ -19351,6 +19456,7 @@ mod tests {
),
];
let mut expected_removes = Vec::new();
let mut recovery_control_ids = Vec::new();
for (case, put_version, remote_version, source_mode) in cases {
let mut transaction = TransitionTransaction::new(TransitionTransactionInit {
deployment_id: ctx.deployment_id().expect("test store should initialize deployment id"),
@@ -19388,6 +19494,8 @@ mod tests {
save_transition_transaction_record(store.clone(), &transaction)
.await
.expect("transaction record should persist");
recovery_control_ids
.push(transition_recovery_control_id(&transaction).expect("transition recovery control id should derive"));
expected_removes.push((transaction.remote_object, put_version));
}
@@ -19403,6 +19511,14 @@ mod tests {
assert_eq!(actual_removes, expected_removes, "recovery must preserve each remote version shape");
assert_eq!(backend.exact_remove_count(), 2);
assert_eq!(backend.object_count().await, 0);
for control_id in recovery_control_ids {
let control = load_recovery_control(store.clone(), IlmRecoveryProtocol::TransitionTransaction, &control_id)
.await
.expect("completed recovery control should remain inspectable");
assert_eq!(control.control.classification, IlmRecoveryClassification::Terminal);
assert_eq!(control.control.attempt_count, 1);
assert!(control.control.owner.is_none());
}
let replay = recover_transition_transaction_records(store, 100, None)
.await
@@ -19415,6 +19531,94 @@ mod tests {
);
}
#[cfg(feature = "test-util")]
#[tokio::test]
#[serial_test::serial(storage_class_env)]
async fn transition_transaction_recovery_resumes_source_cleanup_after_terminal_crash() {
let temp_dir = tempfile::tempdir().expect("create temp store dir");
let (ctx, store, _shutdown) =
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "transition-transaction-terminal-crash", &[4]))
.await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
let tier_name = "TXTERMINALCRASH";
let backend = register_mock_tier(&ctx.tier_config_mgr(), tier_name).await;
let backend_identity = TierConfigMgr::acquire_operation_lease(&ctx.tier_config_mgr(), tier_name)
.await
.expect("tier lease should resolve")
.backend_identity();
let remote_version = uuid::Uuid::new_v4().to_string();
let mut transaction = TransitionTransaction::new(TransitionTransactionInit {
deployment_id: ctx.deployment_id().expect("test store should initialize deployment id"),
transaction_id: uuid::Uuid::new_v4(),
owner_epoch: uuid::Uuid::new_v4(),
write_id: uuid::Uuid::new_v4(),
source: TransitionSourceIdentity {
bucket: "source-bucket".to_string(),
object: "source-object".to_string(),
version_id: Some(uuid::Uuid::new_v4()),
data_dir: uuid::Uuid::new_v4(),
mod_time_unix_nanos: 1_770_000_000_000_000_000,
size: 42,
etag: "source-etag".to_string(),
version_mode: TransitionSourceVersionMode::Versioned,
},
tier_name: tier_name.to_string(),
backend_fingerprint: backend_identity,
not_after_unix_nanos: 1,
})
.expect("transaction should build");
transaction
.advance(
transaction.fence(),
TransitionTransactionState::Uploaded,
Some(TransitionRemoteVersion::versioned(remote_version.clone())),
)
.expect("transaction should enter uploaded state");
backend.set_put_remote_version(Some(remote_version)).await;
let candidate = bytes::Bytes::from_static(b"terminal crash candidate");
backend
.put(
&transaction.remote_object,
ReaderImpl::Body(candidate.clone()),
i64::try_from(candidate.len()).expect("test candidate length should fit i64"),
)
.await
.expect("mock backend should accept candidate");
save_transition_transaction_record(store.clone(), &transaction)
.await
.expect("transaction record should persist");
let control_id = transition_recovery_control_id(&transaction).expect("control id should derive");
let barrier = TransitionRecoveryTerminalBarrier::install(transaction.transaction_id);
let recovery_store = store.clone();
let recovery = tokio::spawn(async move { recover_transition_transaction_records(recovery_store, 100, None).await });
barrier.wait_until_paused().await;
let terminal = load_recovery_control(store.clone(), IlmRecoveryProtocol::TransitionTransaction, &control_id)
.await
.expect("terminal control should persist before source cleanup");
assert_eq!(terminal.control.classification, IlmRecoveryClassification::Terminal);
assert_eq!(transition_transaction_record_count(store.clone()).await, 1);
assert_eq!(backend.object_count().await, 0);
assert_eq!(backend.exact_remove_count(), 1);
recovery.abort();
assert!(
recovery
.await
.expect_err("recovery should be cancelled at the crash boundary")
.is_cancelled()
);
drop(barrier);
let replay = recover_transition_transaction_records(store.clone(), 100, None)
.await
.expect("terminal control should resume source cleanup without another remote delete");
assert_eq!((replay.scanned, replay.recovered, replay.retained, replay.failed), (1, 1, 0, 0));
assert_eq!(transition_transaction_record_count(store).await, 0);
assert_eq!(backend.exact_remove_count(), 1, "terminal replay must not repeat the remote delete");
}
#[cfg(feature = "test-util")]
#[tokio::test]
#[serial_test::serial(storage_class_env)]
@@ -19472,6 +19676,8 @@ mod tests {
save_transition_transaction_record(store.clone(), &uploaded)
.await
.expect("transaction record should persist");
let recovery_control_id =
transition_recovery_control_id(&uploaded).expect("transition recovery control id should derive");
let barrier = TransitionRecoveryClaimBarrier::install(uploaded.transaction_id);
let recovery_store = store.clone();
@@ -19493,11 +19699,17 @@ mod tests {
.expect("recovery should treat the lost CAS as a retained transaction");
assert_eq!((stats.scanned, stats.recovered, stats.retained, stats.failed), (1, 0, 1, 0));
assert_eq!(
load_transition_transaction_record(store, uploaded.transaction_id)
load_transition_transaction_record(store.clone(), uploaded.transaction_id)
.await
.expect("newer transaction revision must remain"),
active
);
let control = load_recovery_control(store, IlmRecoveryProtocol::TransitionTransaction, &recovery_control_id)
.await
.expect("lost source CAS should retain a retryable recovery control");
assert_eq!(control.control.classification, IlmRecoveryClassification::Retrying);
assert_eq!(control.control.consecutive_failure_count, 1);
assert_eq!(control.control.last_error_code, IlmRecoveryErrorCode::SourceGenerationChanged);
assert_eq!(backend.object_count().await, 1, "a stale recovery must not delete the candidate");
assert_eq!(backend.remove_count().await, 0);
}
@@ -19799,27 +20011,15 @@ mod tests {
not_after_unix_nanos: 1_780_000_000_000_000_000,
})
.expect("transaction should build");
let uploaded_fence = transaction
transaction
.advance(
transaction.fence(),
TransitionTransactionState::Uploaded,
Some(TransitionRemoteVersion::versioned(remote_version)),
Some(TransitionRemoteVersion::versioned(remote_version.clone())),
)
.expect("transaction should enter uploaded state");
transaction
.mark_cleanup_pending(
uploaded_fence,
TransitionCleanupProof {
transaction_id: transaction.transaction_id,
write_id: transaction.write_id,
remote_object: transaction.remote_object.clone(),
remote_version: transaction.remote_version.clone(),
backend_fingerprint: transaction.backend_fingerprint,
decision: TransitionCleanupDecision::UploadAbortedBeforeLocalCommit,
},
)
.expect("transaction should enter cleanup pending state");
let candidate = bytes::Bytes::from_static(b"cleanup pending candidate retained after failure");
backend.set_put_remote_version(Some(remote_version)).await;
backend
.put(
&transaction.remote_object,
@@ -19831,6 +20031,8 @@ mod tests {
save_transition_transaction_record(store.clone(), &transaction)
.await
.expect("transaction record should persist");
let recovery_control_id =
transition_recovery_control_id(&transaction).expect("transition recovery control id should derive");
backend.set_remove_failure(true);
let stats = recover_transition_transaction_records(store.clone(), 100, None)
@@ -19846,6 +20048,42 @@ mod tests {
assert_eq!(backend.remove_versions().await, Vec::<(String, String)>::new());
assert_eq!(backend.exact_remove_count(), 1);
assert_eq!(backend.object_count().await, 1);
let control = load_recovery_control(store.clone(), IlmRecoveryProtocol::TransitionTransaction, &recovery_control_id)
.await
.expect("failed recovery control should persist");
assert_eq!(control.control.classification, IlmRecoveryClassification::Retrying);
assert_eq!(control.control.attempt_count, 1);
assert_eq!(control.control.consecutive_failure_count, 1);
assert!(
control
.control
.next_attempt_at_unix_nanos
.is_some_and(|next| next > OffsetDateTime::now_utc().unix_timestamp_nanos() as i64)
);
backend.set_remove_failure(false);
let replay = recover_transition_transaction_records(store.clone(), 100, None)
.await
.expect("recovery before the persisted deadline should be skipped");
assert_eq!((replay.scanned, replay.recovered, replay.retained, replay.failed), (1, 0, 1, 0));
assert_eq!(backend.exact_remove_count(), 1, "persisted backoff must prevent an immediate retry");
let retry_at = control
.control
.next_attempt_at_unix_nanos
.expect("retry deadline should persist");
let retried = recover_transition_transaction_records_at(store.clone(), 100, None, i128::from(retry_at) + 1)
.await
.expect("recovery at the persisted deadline should retry the advanced source generation");
assert_eq!((retried.scanned, retried.recovered, retried.retained, retried.failed), (1, 1, 0, 0));
assert_eq!(backend.exact_remove_count(), 2);
assert_eq!(backend.object_count().await, 0);
assert_eq!(transition_transaction_record_count(store.clone()).await, 0);
let terminal = load_recovery_control(store, IlmRecoveryProtocol::TransitionTransaction, &recovery_control_id)
.await
.expect("completed retry control should remain inspectable");
assert_eq!(terminal.control.classification, IlmRecoveryClassification::Terminal);
assert_eq!(terminal.control.attempt_count, 2);
}
#[cfg(feature = "test-util")]
@@ -20146,6 +20384,10 @@ mod tests {
local_commit_started
.advance(local_commit_started.fence(), TransitionTransactionState::LocalCommitStarted, None)
.expect("transaction should enter local commit state");
let upload_started_control_id =
transition_recovery_control_id(&upload_started).expect("upload-started control id should derive");
let local_commit_control_id =
transition_recovery_control_id(&local_commit_started).expect("local-commit control id should derive");
backend.set_put_remote_version(Some(remote_version)).await;
for transaction in [&upload_started, &local_commit_started] {
@@ -20176,6 +20418,19 @@ mod tests {
assert_eq!(backend.object_count().await, 2, "recovery must not delete an unproven remote candidate");
assert_eq!(backend.remove_count().await, 0);
assert_eq!(backend.exact_remove_count(), 0);
let upload_started_control =
load_recovery_control(store.clone(), IlmRecoveryProtocol::TransitionTransaction, &upload_started_control_id)
.await
.expect("upload-started control should persist");
assert_eq!(
upload_started_control.control.classification,
IlmRecoveryClassification::RetainedAmbiguous
);
let local_commit_control =
load_recovery_control(store, IlmRecoveryProtocol::TransitionTransaction, &local_commit_control_id)
.await
.expect("local-commit control should persist");
assert_eq!(local_commit_control.control.classification, IlmRecoveryClassification::OperatorRequired);
}
#[cfg(feature = "test-util")]
@@ -20526,6 +20781,12 @@ mod tests {
"an unsupported provider probe must retain the unknown upload"
);
assert_eq!(transition_transaction_record_count(store.clone()).await, 1);
let recovery_control_id =
transition_recovery_control_id(&transaction).expect("transition recovery control id should derive");
let control = load_recovery_control(store.clone(), IlmRecoveryProtocol::TransitionTransaction, &recovery_control_id)
.await
.expect("unsupported probe control should persist");
assert_eq!(control.control.classification, IlmRecoveryClassification::RetainedAmbiguous);
assert!(
backend.contains(&transaction.remote_object).await,
"unsupported recovery must not delete the candidate"
@@ -1020,6 +1020,7 @@ mod serial_tests {
}
let (_disk_paths, ecstore) = setup_isolated_test_env(false).await;
let expired_recovery_time = i128::from(i64::MAX / 2);
for case in [
CleanupCase::Persisted,
@@ -1117,7 +1118,7 @@ mod serial_tests {
.await
.expect("active unknown ownership must remain fenced after the transaction store was offline");
assert_eq!((retained.scanned, retained.recovered, retained.retained, retained.failed), (1, 0, 1, 0));
let recovered = recover_transition_transaction_records_at(ecstore.clone(), 100, None, i128::MAX)
let recovered = recover_transition_transaction_records_at(ecstore.clone(), 100, None, expired_recovery_time)
.await
.expect("expired unknown ownership may use the provider's missing proof");
assert_eq!(
@@ -1166,7 +1167,7 @@ mod serial_tests {
assert_eq!(retained.recovered, 0);
assert_eq!(retained.retained + retained.failed, 1);
backend.set_remove_failure(false);
let recovered = recover_transition_transaction_records_at(ecstore.clone(), 100, None, i128::MAX)
let recovered = recover_transition_transaction_records_at(ecstore.clone(), 100, None, expired_recovery_time)
.await
.expect("expired recovery should delete the candidate after the backend becomes available");
assert_eq!(
+144 -6
View File
@@ -18,18 +18,20 @@ use crate::admin::runtime_sources::object_store_from_extensions;
use crate::admin::storage_api::bucket::is_reserved_or_invalid_bucket;
use crate::admin::storage_api::error::StorageError;
use crate::admin::storage_api::lifecycle::{
ManualTransitionCancelCheck, ManualTransitionJobRecord, ManualTransitionJobState, ManualTransitionProgressSink,
ManualTransitionQueueSnapshot, ManualTransitionRunOptions, ManualTransitionRunReport, ManualTransitionScopeAdmission,
ManualTransitionScopeAdmissionClaim, TransitionOperatorDeleteResult, TransitionOperatorError,
claim_manual_transition_scope_admission, delete_manual_transition_scope_admission_if_current,
delete_transition_candidate_for_operator, enqueue_transition_for_existing_objects_scoped,
finalize_missing_transition_transaction_for_operator, inspect_transition_transaction_for_operator,
IlmRecoveryClassification, IlmRecoveryProtocol, ManualTransitionCancelCheck, ManualTransitionJobRecord,
ManualTransitionJobState, ManualTransitionProgressSink, ManualTransitionQueueSnapshot, ManualTransitionRunOptions,
ManualTransitionRunReport, ManualTransitionScopeAdmission, ManualTransitionScopeAdmissionClaim,
TransitionOperatorDeleteResult, TransitionOperatorError, claim_manual_transition_scope_admission,
delete_manual_transition_scope_admission_if_current, delete_transition_candidate_for_operator,
enqueue_transition_for_existing_objects_scoped, finalize_missing_transition_transaction_for_operator,
inspect_recovery_control, inspect_transition_transaction_for_operator, list_recovery_controls,
load_manual_transition_job_record, load_manual_transition_scope_admission, manual_transition_job_lease_expired,
manual_transition_queue_snapshot, manual_transition_scope_admission_lease_expired,
persist_manual_transition_job_progress_if_owned, renew_manual_transition_job_lease_if_owned,
request_manual_transition_job_cancel, save_manual_transition_job_record, update_manual_transition_job_record,
};
use crate::admin::storage_api::runtime::ECStore;
use crate::admin::storage_api::s3::{S3ErrorCode as AdminS3ErrorCode, error as admin_s3_error};
use crate::admin::utils::json_response;
use crate::server::{ADMIN_PREFIX, RemoteAddr};
use http::HeaderMap;
@@ -230,9 +232,48 @@ pub fn register_ilm_transition_route(r: &mut S3Router<AdminOperation>) -> std::i
format!("{ADMIN_PREFIX}/v3/ilm/transition/reconcile/{{transaction_id}}").as_str(),
AdminOperation(&TransitionReconcileApplyHandler {}),
)?;
r.insert(
Method::GET,
format!("{ADMIN_PREFIX}/v3/ilm/recovery/records").as_str(),
AdminOperation(&IlmRecoveryControlListHandler {}),
)?;
r.insert(
Method::GET,
format!("{ADMIN_PREFIX}/v3/ilm/recovery/records/{{control_id}}").as_str(),
AdminOperation(&IlmRecoveryControlInspectHandler {}),
)?;
Ok(())
}
#[derive(Debug, Deserialize)]
#[serde(deny_unknown_fields)]
struct IlmRecoveryControlListQuery {
protocol: IlmRecoveryProtocol,
#[serde(default)]
classification: Option<IlmRecoveryClassification>,
#[serde(default = "default_recovery_control_list_limit")]
limit: usize,
#[serde(default)]
marker: Option<String>,
}
const fn default_recovery_control_list_limit() -> usize {
100
}
fn parse_recovery_control_list_query(query: Option<&str>) -> S3Result<IlmRecoveryControlListQuery> {
let query = query.ok_or_else(|| admin_s3_error(AdminS3ErrorCode::InvalidRequest, "protocol is required"))?;
let parsed: IlmRecoveryControlListQuery = serde_urlencoded::from_bytes(query.as_bytes())
.map_err(|_| admin_s3_error(AdminS3ErrorCode::InvalidArgument, "invalid ILM recovery control query"))?;
if !(1..=1_000).contains(&parsed.limit) {
return Err(admin_s3_error(AdminS3ErrorCode::InvalidArgument, "limit must be between 1 and 1000"));
}
if parsed.marker.as_ref().is_some_and(|marker| marker.is_empty()) {
return Err(admin_s3_error(AdminS3ErrorCode::InvalidArgument, "marker must not be empty"));
}
Ok(parsed)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum ManualTransitionRunMode {
EnqueueOnly,
@@ -423,6 +464,26 @@ fn transition_transaction_id_from_params(params: &Params<'_, '_>) -> S3Result<Uu
.map_err(|_| s3_error!(InvalidArgument, "invalid transition transaction id"))
}
fn recovery_control_id_from_params(params: &Params<'_, '_>) -> S3Result<String> {
let control_id = params.get("control_id").unwrap_or("");
if control_id.len() != 64
|| !control_id
.bytes()
.all(|byte| byte.is_ascii_hexdigit() && !byte.is_ascii_uppercase())
{
return Err(admin_s3_error(AdminS3ErrorCode::InvalidArgument, "invalid ILM recovery control id"));
}
Ok(control_id.to_string())
}
fn map_recovery_control_error(err: StorageError) -> S3Error {
if err == StorageError::ConfigNotFound {
admin_s3_error(AdminS3ErrorCode::NoSuchKey, "ILM recovery control not found")
} else {
admin_s3_error(AdminS3ErrorCode::InternalError, "ILM recovery control request failed")
}
}
fn map_transition_operator_error(err: TransitionOperatorError) -> S3Error {
match err {
TransitionOperatorError::NotFound => s3_error!(NoSuchKey, "transition transaction not found"),
@@ -1031,6 +1092,40 @@ impl Operation for TransitionReconcileInspectHandler {
}
}
pub struct IlmRecoveryControlListHandler {}
#[async_trait::async_trait]
impl Operation for IlmRecoveryControlListHandler {
async fn call(&self, req: S3Request<Body>, _params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
authorize_transition_admin_request(&req, AdminAction::ListTierAction).await?;
let query = parse_recovery_control_list_query(req.uri.query())?;
let Some(store) = object_store_from_extensions(&req.extensions) else {
return Err(admin_s3_error(AdminS3ErrorCode::InternalError, "object store is not initialized"));
};
let page = list_recovery_controls(store, query.protocol, query.classification, query.limit, query.marker)
.await
.map_err(map_recovery_control_error)?;
json_response(StatusCode::OK, &page)
}
}
pub struct IlmRecoveryControlInspectHandler {}
#[async_trait::async_trait]
impl Operation for IlmRecoveryControlInspectHandler {
async fn call(&self, req: S3Request<Body>, params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
authorize_transition_admin_request(&req, AdminAction::ListTierAction).await?;
let control_id = recovery_control_id_from_params(&params)?;
let Some(store) = object_store_from_extensions(&req.extensions) else {
return Err(admin_s3_error(AdminS3ErrorCode::InternalError, "object store is not initialized"));
};
let control = inspect_recovery_control(store, &control_id)
.await
.map_err(map_recovery_control_error)?;
json_response(StatusCode::OK, &control)
}
}
pub struct TransitionReconcileApplyHandler {}
#[async_trait::async_trait]
@@ -1104,6 +1199,49 @@ mod tests {
f(&matched.params)
}
fn with_recovery_control_params<T>(path: &str, f: impl FnOnce(&Params<'_, '_>) -> T) -> T {
let mut router = Router::new();
router
.insert("/rustfs/admin/v3/ilm/recovery/records/{control_id}", ())
.expect("route should insert");
let matched = router.at(path).expect("route should match");
f(&matched.params)
}
#[test]
fn recovery_control_query_is_bounded_and_strict() {
let query = parse_recovery_control_list_query(Some("protocol=transition_transaction"))
.expect("minimal recovery query should parse");
assert_eq!(query.protocol, IlmRecoveryProtocol::TransitionTransaction);
assert_eq!(query.classification, None);
assert_eq!(query.limit, 100);
let filtered = parse_recovery_control_list_query(Some(
"protocol=tier_delete_journal&classification=retained_ambiguous&limit=1000&marker=opaque",
))
.expect("bounded filtered query should parse");
assert_eq!(filtered.protocol, IlmRecoveryProtocol::TierDeleteJournal);
assert_eq!(filtered.classification, Some(IlmRecoveryClassification::RetainedAmbiguous));
assert_eq!(filtered.limit, 1000);
assert!(parse_recovery_control_list_query(None).is_err());
assert!(parse_recovery_control_list_query(Some("protocol=transition_transaction&limit=0")).is_err());
assert!(parse_recovery_control_list_query(Some("protocol=transition_transaction&limit=1001")).is_err());
assert!(parse_recovery_control_list_query(Some("protocol=unknown")).is_err());
assert!(parse_recovery_control_list_query(Some("protocol=transition_transaction&extra=true")).is_err());
}
#[test]
fn recovery_control_id_is_canonical_lowercase_sha256() {
let id = "ab".repeat(32);
with_recovery_control_params(&format!("/rustfs/admin/v3/ilm/recovery/records/{id}"), |params| {
assert_eq!(recovery_control_id_from_params(params).expect("control id should parse"), id);
});
let uppercase = "AB".repeat(32);
with_recovery_control_params(&format!("/rustfs/admin/v3/ilm/recovery/records/{uppercase}"), |params| {
assert!(recovery_control_id_from_params(params).is_err())
});
}
fn manual_transition_job_request(method: Method, path: &'static str) -> S3Request<Body> {
S3Request {
input: Body::empty(),
+3
View File
@@ -232,6 +232,9 @@ pub(crate) mod lifecycle {
pub(crate) type ManualTransitionRunOptions =
super::ecstore_bucket::lifecycle::bucket_lifecycle_ops::ManualTransitionRunOptions;
pub(crate) type ManualTransitionRunReport = super::ecstore_bucket::lifecycle::bucket_lifecycle_ops::ManualTransitionRunReport;
pub(crate) use super::ecstore_bucket::lifecycle::recovery_control::{
IlmRecoveryClassification, IlmRecoveryProtocol, inspect_recovery_control, list_recovery_controls,
};
pub(crate) use super::ecstore_bucket::lifecycle::transition_transaction::{
TransitionOperatorDeleteResult, TransitionOperatorError, delete_transition_candidate_for_operator,
finalize_missing_transition_transaction_for_operator, inspect_transition_transaction_for_operator,