feat(ilm): execute legacy recovery dispositions

This commit is contained in:
cxymds
2026-09-06 19:43:36 +08:00
parent c3e6d90c6f
commit c353d987a3
10 changed files with 3163 additions and 200 deletions
+2 -4
View File
@@ -78,10 +78,8 @@ pub mod bucket {
pub mod recovery_disposition {
pub use crate::bucket::lifecycle::recovery_disposition::{
CreatedIlmRecoveryDisposition, IlmRecoveryDisposition, IlmRecoveryDispositionAction, IlmRecoveryDispositionError,
IlmRecoveryDispositionIdentity, IlmRecoveryDispositionOwnerLease, IlmRecoveryDispositionReasonCode,
IlmRecoveryDispositionState, ObservedIlmRecoveryDisposition, create_recovery_disposition_if_absent,
load_recovery_disposition, recovery_disposition_id, save_recovery_disposition_if_current,
IlmRecoveryDispositionExecutionOutcome, IlmRecoveryDispositionReasonCode, IlmRecoveryDispositionState,
dry_run_recovery_disposition, execute_recovery_disposition,
};
}
@@ -32,6 +32,7 @@ use crate::bucket::lifecycle::manual_transition_job::{
record_manual_transition_worker_result_with_reason, renew_manual_transition_job_lease_if_owned,
save_manual_transition_job_record_if_current, save_manual_transition_task_if_absent, update_manual_transition_job_record,
};
use crate::bucket::lifecycle::recovery_disposition_runtime::run_recovery_disposition_maintenance_loop;
use crate::bucket::lifecycle::replication_sink;
use crate::bucket::lifecycle::replication_sink::{
DeleteReplicationConfigSnapshot, ReplicationObjectBridge, ReplicationStatusType, replication_state_to_filemeta,
@@ -149,6 +150,7 @@ pub type ExpiryOpType = Box<dyn ExpiryOp + Send + Sync + 'static>;
static XXHASH_SEED: u64 = 0;
static TIER_FREE_VERSION_RECOVERY_STARTED: OnceLock<()> = OnceLock::new();
static MANUAL_TRANSITION_JOB_RECOVERY_STARTED: OnceLock<()> = OnceLock::new();
static RECOVERY_DISPOSITION_MAINTENANCE_STARTED: OnceLock<()> = OnceLock::new();
#[cfg(test)]
#[derive(Default)]
@@ -2398,9 +2400,20 @@ pub async fn init_background_expiry(api: Arc<ECStore>) {
let _ = spawn_tier_free_version_recovery_once(api.clone(), &TIER_FREE_VERSION_RECOVERY_STARTED);
spawn_tier_delete_journal_recovery_once(api.clone());
spawn_transition_transaction_recovery_once(api.clone());
spawn_recovery_disposition_maintenance_once(api.clone());
spawn_manual_transition_job_recovery_once(api);
}
fn spawn_recovery_disposition_maintenance_once(api: Arc<ECStore>) -> Option<JoinHandle<()>> {
let cancel_token = api.ctx.background_cancel_token()?;
if RECOVERY_DISPOSITION_MAINTENANCE_STARTED.set(()).is_err() {
return None;
}
Some(tokio::spawn(async move {
run_recovery_disposition_maintenance_loop(api, cancel_token).await;
}))
}
fn spawn_manual_transition_job_recovery_once(api: Arc<ECStore>) -> Option<JoinHandle<()>> {
if MANUAL_TRANSITION_JOB_RECOVERY_STARTED.set(()).is_err() {
return None;
@@ -161,20 +161,30 @@ where
DeletedObject = DeletedObject,
>,
{
match api
.delete_object(
RUSTFS_META_BUCKET,
file,
ObjectOptions {
http_preconditions: Some(HTTPPreconditions {
if_match: Some(etag.to_string()),
..Default::default()
}),
..Default::default()
},
)
.await
{
delete_config_if_match_with_opts(api, file, etag, ObjectOptions::default()).await
}
pub(crate) async fn delete_config_if_match_with_opts<S>(
api: Arc<S>,
file: &str,
etag: &str,
mut options: ObjectOptions,
) -> Result<()>
where
S: ObjectOperations<
Error = Error,
ObjectInfo = ObjectInfo,
ObjectOptions = ObjectOptions,
FileInfo = FileInfo,
ObjectToDelete = ObjectToDelete,
DeletedObject = DeletedObject,
>,
{
options.http_preconditions = Some(HTTPPreconditions {
if_match: Some(etag.to_string()),
..Default::default()
});
match api.delete_object(RUSTFS_META_BUCKET, file, options).await {
Ok(_) => Ok(()),
Err(err) => {
if err == Error::FileNotFound || matches!(err, Error::ObjectNotFound(_, _)) {
@@ -26,6 +26,7 @@ mod object_lock_boundary;
pub use self::core as lifecycle;
pub mod recovery_control;
pub mod recovery_disposition;
pub(crate) mod recovery_disposition_runtime;
pub mod recovery_export;
mod replication_sink;
pub mod rule;
@@ -485,6 +485,20 @@ impl IlmRecoveryControl {
self.validate()
}
pub fn abandon_for_operator(&mut self, expected_source_generation: &IlmRecoverySourceGeneration) -> Result<()> {
if self.owner.is_some()
|| self.classification != IlmRecoveryClassification::RetainedAmbiguous
|| &self.observed_source_generation != expected_source_generation
{
return Err(IlmRecoveryControlError::InvalidSuccessor(
"operator abandonment requires the exact ownerless retained source generation",
));
}
self.bump_revision()?;
self.classification = IlmRecoveryClassification::Abandoned;
self.validate()
}
pub fn validate_successor(&self, next: &Self) -> Result<()> {
self.validate()?;
next.validate()?;
@@ -508,12 +522,34 @@ impl IlmRecoveryControl {
self.validate_failure_successor(next)
}
(Some(_), None) => self.validate_finish_successor(next),
(None, None)
if self.classification == IlmRecoveryClassification::RetainedAmbiguous
&& next.classification == IlmRecoveryClassification::Abandoned =>
{
self.validate_operator_abandon_successor(next)
}
(None, None) => Err(IlmRecoveryControlError::InvalidSuccessor(
"ownerless control cannot advance without a claim",
)),
}
}
fn validate_operator_abandon_successor(&self, next: &Self) -> Result<()> {
if next.observed_source_generation != self.observed_source_generation
|| next.attempt_count != self.attempt_count
|| next.consecutive_failure_count != self.consecutive_failure_count
|| next.first_failure_at_unix_nanos != self.first_failure_at_unix_nanos
|| next.last_failure_at_unix_nanos != self.last_failure_at_unix_nanos
|| next.next_attempt_at_unix_nanos != self.next_attempt_at_unix_nanos
|| next.last_error_code != self.last_error_code
{
return Err(IlmRecoveryControlError::InvalidSuccessor(
"operator abandonment changed recovery history or source generation",
));
}
Ok(())
}
fn validate_claim_successor(&self, next: &Self) -> Result<()> {
if self.classification != IlmRecoveryClassification::Retrying
|| next.classification != IlmRecoveryClassification::Retrying
@@ -827,6 +863,23 @@ pub async fn observe_recovery_source(
api: Arc<ECStore>,
canonical_path: &str,
source_schema: &str,
) -> EcstoreResult<ObservedIlmRecoverySource> {
observe_recovery_source_with_options(api, canonical_path, source_schema, false).await
}
pub(crate) async fn observe_recovery_source_no_lock(
api: Arc<ECStore>,
canonical_path: &str,
source_schema: &str,
) -> EcstoreResult<ObservedIlmRecoverySource> {
observe_recovery_source_with_options(api, canonical_path, source_schema, true).await
}
async fn observe_recovery_source_with_options(
api: Arc<ECStore>,
canonical_path: &str,
source_schema: &str,
no_lock: bool,
) -> EcstoreResult<ObservedIlmRecoverySource> {
validate_canonical_source_path(canonical_path).map_err(recovery_control_store_error)?;
if source_schema.trim().is_empty() {
@@ -837,7 +890,16 @@ pub async fn observe_recovery_source(
let mut observations = Vec::new();
for set in api.all_set_disks() {
let authority = format!("pool-{}/set-{}", set.pool_index, set.set_index);
match config_boundary::read_config_with_metadata(set, canonical_path, &ObjectOptions::default()).await {
match config_boundary::read_config_with_metadata(
set,
canonical_path,
&ObjectOptions {
no_lock,
..Default::default()
},
)
.await
{
Ok((data, metadata)) => {
let etag = metadata
.etag
@@ -1255,6 +1317,36 @@ mod tests {
));
}
#[test]
fn operator_abandonment_is_an_exact_ownerless_retained_successor() {
let mut retained = IlmRecoveryControl::new(
control().identity,
generation(),
IlmRecoveryClassification::RetainedAmbiguous,
1_000_000_000,
IlmRecoveryErrorCode::OperatorDispositionRequired,
)
.expect("retained control should build");
let previous = retained.clone();
retained
.abandon_for_operator(&previous.observed_source_generation)
.expect("exact retained generation should be abandonable");
previous
.validate_successor(&retained)
.expect("operator abandonment should be a valid successor");
assert_eq!(retained.classification, IlmRecoveryClassification::Abandoned);
assert_eq!(retained.revision, previous.revision + 1);
let mut wrong_generation = previous.clone();
let mut generation = previous.observed_source_generation.clone();
generation.source_etag = "different".to_string();
assert!(wrong_generation.abandon_for_operator(&generation).is_err());
let mut mutated_history = retained.clone();
mutated_history.attempt_count += 1;
assert!(previous.validate_successor(&mutated_history).is_err());
}
#[test]
fn recovery_control_view_redacts_source_and_owner_details() {
let mut control = control();
File diff suppressed because it is too large Load Diff
File diff suppressed because it is too large Load Diff
+302 -1
View File
@@ -828,8 +828,14 @@ mod tests {
recovery_control::{
IlmRecoveryClassification, IlmRecoveryControl, IlmRecoveryControlIdentity, IlmRecoveryErrorCode,
IlmRecoveryProtocol, MAX_RECOVERY_ATTEMPTS, list_recovery_controls, load_recovery_control,
observe_recovery_source, save_recovery_control_if_absent,
observe_recovery_source, recovery_control_record_object_name, save_recovery_control_if_absent,
},
recovery_disposition::{
IlmRecoveryDispositionExecutionOutcome, IlmRecoveryDispositionState, RecoveryDispositionCrashStage,
dry_run_recovery_disposition, execute_recovery_disposition, inject_recovery_disposition_crash_once,
load_recovery_disposition,
},
recovery_disposition_runtime::garbage_collect_completed_recovery_disposition,
recovery_export::{
create_recovery_export, inspect_recovery_export_observation, load_recovery_export,
recovery_export_record_object_name,
@@ -16936,6 +16942,301 @@ mod tests {
);
}
#[cfg(feature = "test-util")]
#[tokio::test]
#[serial_test::serial(storage_class_env)]
async fn legacy_recovery_disposition_removes_only_local_journals_and_replays() {
Box::pin(legacy_recovery_disposition_removes_only_local_journals_and_replays_case()).await;
}
#[cfg(feature = "test-util")]
async fn legacy_recovery_disposition_removes_only_local_journals_and_replays_case() {
let temp_dir = tempfile::tempdir().expect("create legacy disposition store dir");
let (ctx, store, _shutdown) =
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "legacy-recovery-disposition", &[4])).await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
let tier_name = "LEGACY-DISPOSITION";
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("legacy disposition tier lease should resolve")
.backend_identity();
let fixtures = [
serde_json::json!({
"version": 1,
"obj_name": "legacy/disposition-v1",
"version_id": "opaque-disposition-v1",
"tier_name": tier_name,
}),
serde_json::json!({
"version": 2,
"obj_name": "legacy/disposition-v2",
"version_id": "opaque-disposition-v2",
"tier_name": tier_name,
"backend_identity": backend_identity,
}),
];
let mut journal_paths = Vec::new();
for fixture in &fixtures {
let data = serde_json::to_vec(fixture).expect("legacy disposition fixture should encode");
let entry = crate::bucket::lifecycle::tier_delete_journal::decode_tier_delete_journal_entry(&data)
.expect("legacy disposition fixture should decode");
let path = tier_delete_journal_object_name(&entry);
com::save_config(store.clone(), &path, data)
.await
.expect("legacy disposition fixture should persist");
journal_paths.push(path);
}
let recovered = recover_tier_delete_journal_entries(store.clone(), 100, None)
.await
.expect("legacy disposition recovery scan should finish");
assert_eq!((recovered.scanned, recovered.deleted, recovered.failed), (2, 0, 0));
assert_eq!(tier_delete_journal_count(store.clone()).await, 2);
let mut controls = list_recovery_controls(
store.clone(),
IlmRecoveryProtocol::TierDeleteJournal,
Some(IlmRecoveryClassification::RetainedAmbiguous),
100,
None,
)
.await
.expect("legacy disposition controls should be listable")
.records;
controls.sort_by(|left, right| left.control_id.cmp(&right.control_id));
assert_eq!(controls.len(), 2, "both legacy schemas must support disposition");
let actor_sha256 = rustfs_utils::crypto::hex_sha256(b"legacy-disposition-actor", ToOwned::to_owned);
let wrong_actor_sha256 = rustfs_utils::crypto::hex_sha256(b"different-disposition-actor", ToOwned::to_owned);
let wrong_export_sha256 = "ff".repeat(32);
for (index, control) in controls.iter().enumerate() {
let observation = inspect_recovery_export_observation(store.clone(), &control.control_id)
.await
.expect("legacy disposition source should be observable");
let export = create_recovery_export(store.clone(), &observation, &actor_sha256)
.await
.expect("legacy disposition export should persist");
let confirmed_at_unix_nanos = i64::try_from(OffsetDateTime::now_utc().unix_timestamp_nanos())
.expect("legacy disposition timestamp should fit i64");
if index == 0 {
let wrong_hash = Box::pin(execute_recovery_disposition(
store.clone(),
&observation,
&export.export_id,
&wrong_export_sha256,
&actor_sha256,
confirmed_at_unix_nanos,
))
.await
.expect_err("a mismatched export checksum must fail before local deletion");
assert_eq!(wrong_hash, Error::PreconditionFailed);
assert_eq!(tier_delete_journal_count(store.clone()).await, 2);
}
let dry_run = dry_run_recovery_disposition(
store.clone(),
&observation,
&export.export_id,
&export.content_sha256,
&actor_sha256,
confirmed_at_unix_nanos,
)
.await
.expect("legacy disposition dry-run should validate exact local state");
assert_eq!(dry_run.source_copy_count, observation.source_generation.copies.len());
assert_eq!(
tier_delete_journal_count(store.clone()).await,
fixtures.len() - index,
"dry-run must not delete a legacy journal"
);
assert!(
matches!(
load_recovery_disposition(store.clone(), IlmRecoveryProtocol::TierDeleteJournal, &dry_run.disposition_id,)
.await,
Err(Error::ConfigNotFound)
),
"dry-run must not persist a disposition record"
);
if index == 0 {
inject_recovery_disposition_crash_once(RecoveryDispositionCrashStage::AfterLocalDelete);
Box::pin(execute_recovery_disposition(
store.clone(),
&observation,
&export.export_id,
&export.content_sha256,
&actor_sha256,
confirmed_at_unix_nanos,
))
.await
.expect_err("the injected crash must stop after local delete commits");
assert_eq!(tier_delete_journal_count(store.clone()).await, 1);
let interrupted =
load_recovery_disposition(store.clone(), IlmRecoveryProtocol::TierDeleteJournal, &dry_run.disposition_id)
.await
.expect("the applying disposition must survive the post-delete crash");
assert_eq!(interrupted.disposition.state, IlmRecoveryDispositionState::Applying);
assert!(
interrupted.disposition.confirmed_absent.is_empty(),
"the crash must occur before absence progress is persisted"
);
} else {
inject_recovery_disposition_crash_once(RecoveryDispositionCrashStage::AfterControlAbandon);
Box::pin(execute_recovery_disposition(
store.clone(),
&observation,
&export.export_id,
&export.content_sha256,
&actor_sha256,
confirmed_at_unix_nanos,
))
.await
.expect_err("the injected crash must stop after control abandonment commits");
let interrupted =
load_recovery_disposition(store.clone(), IlmRecoveryProtocol::TierDeleteJournal, &dry_run.disposition_id)
.await
.expect("the applying disposition must survive the post-control crash");
assert_eq!(interrupted.disposition.state, IlmRecoveryDispositionState::Applying);
assert_eq!(
interrupted.disposition.confirmed_absent.len(),
interrupted.disposition.identity.source_generation.copies.len()
);
let abandoned =
load_recovery_control(store.clone(), IlmRecoveryProtocol::TierDeleteJournal, &observation.control_id)
.await
.expect("the abandoned control must survive the injected crash");
let exact_abandoned = abandoned.control.encode().expect("the exact abandoned control should encode");
let mut wrong_history = abandoned.control;
wrong_history.last_error_code = IlmRecoveryErrorCode::CleanupFailed;
let control_path = recovery_control_record_object_name(observation.protocol, &observation.control_id)
.expect("control path should remain canonical");
com::save_config(
store.clone(),
&control_path,
wrong_history
.encode()
.expect("the alternate valid control history should encode"),
)
.await
.expect("the alternate control history fixture should persist");
let wrong_history_err = Box::pin(execute_recovery_disposition(
store.clone(),
&observation,
&export.export_id,
&export.content_sha256,
&actor_sha256,
confirmed_at_unix_nanos + 1,
))
.await
.expect_err("a different abandoned control history must not bridge to completion");
assert_eq!(wrong_history_err, Error::PreconditionFailed);
com::save_config(store.clone(), &control_path, exact_abandoned)
.await
.expect("the exact abandoned control fixture should be restored");
}
let replay_confirmed_at_unix_nanos = confirmed_at_unix_nanos + 2;
let executed = Box::pin(execute_recovery_disposition(
store.clone(),
&observation,
&export.export_id,
&export.content_sha256,
&actor_sha256,
replay_confirmed_at_unix_nanos,
))
.await
.expect("a later request must resume and complete the interrupted disposition");
assert_eq!(executed.state, IlmRecoveryDispositionState::Completed);
assert_eq!(executed.outcome, IlmRecoveryDispositionExecutionOutcome::Completed);
assert_eq!(executed.confirmed_absent_copy_count, executed.source_copy_count);
assert_eq!(tier_delete_journal_count(store.clone()).await, fixtures.len() - index - 1);
assert!(matches!(
com::read_config(store.clone(), &observation.canonical_source_path).await,
Err(Error::ConfigNotFound)
));
if index == 0 {
let untouched = journal_paths
.iter()
.find(|path| *path != &observation.canonical_source_path)
.expect("the other legacy journal should remain");
com::read_config(store.clone(), untouched)
.await
.expect("disposition must not remove a different legacy journal");
}
let abandoned = load_recovery_control(store.clone(), IlmRecoveryProtocol::TierDeleteJournal, &observation.control_id)
.await
.expect("abandoned recovery control should remain inspectable");
assert_eq!(abandoned.control.classification, IlmRecoveryClassification::Abandoned);
assert_eq!(abandoned.control.revision, observation.control_revision + 1);
assert_eq!(abandoned.control.observed_source_generation, observation.source_generation);
let persisted =
load_recovery_disposition(store.clone(), IlmRecoveryProtocol::TierDeleteJournal, &executed.disposition_id)
.await
.expect("completed disposition should remain durable");
assert_eq!(persisted.disposition.state, IlmRecoveryDispositionState::Completed);
let replayed = Box::pin(execute_recovery_disposition(
store.clone(),
&observation,
&export.export_id,
&export.content_sha256,
&actor_sha256,
replay_confirmed_at_unix_nanos + 1,
))
.await
.expect("same actor should replay the completed disposition");
assert_eq!(replayed.state, IlmRecoveryDispositionState::Completed);
assert_eq!(replayed.outcome, IlmRecoveryDispositionExecutionOutcome::Replayed);
let wrong_actor = Box::pin(execute_recovery_disposition(
store.clone(),
&observation,
&export.export_id,
&export.content_sha256,
&wrong_actor_sha256,
replay_confirmed_at_unix_nanos + 2,
))
.await
.expect_err("a different actor must not replay a completed disposition");
assert_eq!(wrong_actor, Error::PreconditionFailed);
assert!(
!Box::pin(garbage_collect_completed_recovery_disposition(
store.clone(),
&persisted,
persisted.disposition.retain_until_unix_nanos - 1,
))
.await
.expect("completed disposition should remain before retention expires")
);
assert!(
Box::pin(garbage_collect_completed_recovery_disposition(
store.clone(),
&persisted,
persisted.disposition.retain_until_unix_nanos,
))
.await
.expect("expired completed disposition should be garbage collected")
);
assert!(matches!(
load_recovery_disposition(store.clone(), IlmRecoveryProtocol::TierDeleteJournal, &executed.disposition_id).await,
Err(Error::ConfigNotFound)
));
assert_eq!(backend.remove_count().await, 0, "legacy disposition must not call the remote tier");
assert_eq!(backend.exact_remove_count(), 0, "legacy disposition must not issue exact remote DELETE");
assert!(
backend.op_log().await.is_empty(),
"legacy disposition must not invoke any backend operation"
);
}
}
#[cfg(feature = "test-util")]
#[tokio::test]
#[serial_test::serial(storage_class_env)]
+341 -107
View File
@@ -18,12 +18,13 @@ use crate::admin::runtime_sources::{current_action_credentials, object_store_fro
use crate::admin::storage_api::bucket::is_reserved_or_invalid_bucket;
use crate::admin::storage_api::error::StorageError;
use crate::admin::storage_api::lifecycle::{
IlmRecoveryClassification, IlmRecoveryControlView, IlmRecoveryExportObservation, IlmRecoveryProtocol,
ManualTransitionCancelCheck, ManualTransitionJobRecord, ManualTransitionJobState, ManualTransitionProgressSink,
ManualTransitionQueueSnapshot, ManualTransitionRunOptions, ManualTransitionRunReport, ManualTransitionScopeAdmission,
ManualTransitionScopeAdmissionClaim, TransitionOperatorDeleteResult, TransitionOperatorError,
claim_manual_transition_scope_admission, create_recovery_export, delete_manual_transition_scope_admission_if_current,
delete_transition_candidate_for_operator, enqueue_transition_for_existing_objects_scoped,
IlmRecoveryClassification, IlmRecoveryControlView, IlmRecoveryDispositionExecutionOutcome, IlmRecoveryDispositionReasonCode,
IlmRecoveryDispositionState, IlmRecoveryExportObservation, IlmRecoveryProtocol, ManualTransitionCancelCheck,
ManualTransitionJobRecord, ManualTransitionJobState, ManualTransitionProgressSink, ManualTransitionQueueSnapshot,
ManualTransitionRunOptions, ManualTransitionRunReport, ManualTransitionScopeAdmission, ManualTransitionScopeAdmissionClaim,
TransitionOperatorDeleteResult, TransitionOperatorError, claim_manual_transition_scope_admission, create_recovery_export,
delete_manual_transition_scope_admission_if_current, delete_transition_candidate_for_operator, dry_run_recovery_disposition,
enqueue_transition_for_existing_objects_scoped, execute_recovery_disposition,
finalize_missing_transition_transaction_for_operator, inspect_recovery_control, inspect_recovery_export_observation,
inspect_transition_transaction_for_operator, list_recovery_controls, load_manual_transition_job_record,
load_manual_transition_scope_admission, load_recovery_export, manual_transition_job_lease_expired,
@@ -257,7 +258,7 @@ pub fn register_ilm_transition_route(r: &mut S3Router<AdminOperation>) -> std::i
r.insert(
Method::POST,
format!("{ADMIN_PREFIX}/v3/ilm/recovery/records/{{control_id}}").as_str(),
AdminOperation(&IlmRecoveryExportCreateHandler {}),
AdminOperation(&IlmRecoveryRecordMutationHandler {}),
)?;
r.insert(
Method::GET,
@@ -537,6 +538,16 @@ fn map_recovery_export_error(err: StorageError) -> S3Error {
}
}
fn map_recovery_disposition_error(err: StorageError) -> S3Error {
if err == StorageError::ConfigNotFound {
admin_s3_error(AdminS3ErrorCode::NoSuchKey, "ILM recovery export or disposition not found")
} else if err == StorageError::SlowDown {
admin_s3_error(AdminS3ErrorCode::SlowDown, "ILM recovery disposition admission capacity is exhausted")
} else {
admin_s3_error(AdminS3ErrorCode::OperationAborted, "ILM recovery disposition request cannot proceed")
}
}
fn recovery_export_download_headers(export_id: &str, encoded_len: usize) -> S3Result<HeaderMap> {
let mut headers = HeaderMap::new();
headers.insert(header::CONTENT_TYPE, HeaderValue::from_static("application/json"));
@@ -600,12 +611,10 @@ struct IlmRecoveryControlInspectResponse {
observation_receipt: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
observation_receipt_expires_at_unix_nanos: Option<i64>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
enum IlmRecoveryDispositionReasonCode {
LegacyRemoteCleanupAbandoned,
#[serde(skip_serializing_if = "Option::is_none")]
disposition_dry_run_receipt: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
disposition_dry_run_receipt_expires_at_unix_nanos: Option<i64>,
}
#[derive(Debug, Deserialize)]
@@ -650,7 +659,6 @@ fn parse_recovery_record_mutation_request(body: &[u8]) -> S3Result<IlmRecoveryRe
Ok(request)
}
#[allow(dead_code)]
enum ValidatedIlmRecoveryRecordMutation<'a> {
Export {
observation_receipt: &'a str,
@@ -733,33 +741,30 @@ struct IlmRecoveryExportCreateResponse {
outcome: &'static str,
}
// These response envelopes pin the future disposition wire contract before
// its storage state machine is connected to this handler.
#[allow(dead_code)]
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
enum IlmRecoveryDispositionDryRunStatus {
Ready,
}
#[allow(dead_code)]
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
enum IlmRecoveryDispositionState {
enum IlmRecoveryDispositionResponseState {
Applying,
Completed,
}
#[allow(dead_code)]
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
enum IlmRecoveryDispositionOutcome {
AcceptedForRecovery,
Completed,
Replayed,
fn recovery_disposition_response_state(state: IlmRecoveryDispositionState) -> S3Result<IlmRecoveryDispositionResponseState> {
match state {
IlmRecoveryDispositionState::Applying => Ok(IlmRecoveryDispositionResponseState::Applying),
IlmRecoveryDispositionState::Completed => Ok(IlmRecoveryDispositionResponseState::Completed),
IlmRecoveryDispositionState::Prepared => Err(admin_s3_error(
AdminS3ErrorCode::OperationAborted,
"ILM recovery disposition is accepted but not yet applying",
)),
}
}
#[allow(dead_code)]
#[derive(Debug, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
struct IlmRecoveryDispositionDryRunResponse {
@@ -776,15 +781,14 @@ struct IlmRecoveryDispositionDryRunResponse {
observation_receipt_expires_at_unix_nanos: i64,
}
#[allow(dead_code)]
#[derive(Debug, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
struct IlmRecoveryDispositionExecuteResponse {
action: IlmRecoveryReceiptAction,
mode: IlmRecoveryReceiptMode,
disposition_id: String,
state: IlmRecoveryDispositionState,
outcome: IlmRecoveryDispositionOutcome,
state: IlmRecoveryDispositionResponseState,
outcome: IlmRecoveryDispositionExecutionOutcome,
confirmed_absent_copy_count: usize,
source_copy_count: usize,
}
@@ -1546,20 +1550,48 @@ impl Operation for IlmRecoveryControlInspectHandler {
.await
.map_err(map_recovery_control_error)?;
let now = OffsetDateTime::now_utc();
let (export_ready, export_not_ready_reason, observation_receipt, expires_at) =
match inspect_recovery_export_observation(store, &control_id).await {
Ok(observation) => match issue_recovery_observation_receipt(
observation,
actor_sha256,
let (
export_ready,
export_not_ready_reason,
observation_receipt,
observation_receipt_expires_at_unix_nanos,
disposition_dry_run_receipt,
disposition_dry_run_receipt_expires_at_unix_nanos,
) = match inspect_recovery_export_observation(store, &control_id).await {
Ok(observation) => {
let export_receipt = issue_recovery_observation_receipt(
observation.clone(),
actor_sha256.clone(),
IlmRecoveryReceiptAction::Export,
IlmRecoveryReceiptMode::Execute,
now,
) {
);
let disposition_receipt = issue_recovery_observation_receipt(
observation,
actor_sha256,
IlmRecoveryReceiptAction::AbandonRemoteCleanup,
IlmRecoveryReceiptMode::DryRun,
now,
);
let (export_ready, export_not_ready_reason, observation_receipt, export_expires_at) = match export_receipt {
Ok((token, expires_at)) => (true, None, Some(token), Some(expires_at)),
Err(_) => (false, Some("receipt_key_unavailable"), None, None),
},
Err(_) => (false, Some("fleet_or_source_not_ready"), None, None),
};
};
let (disposition_receipt, disposition_expires_at) = match disposition_receipt {
Ok((token, expires_at)) => (Some(token), Some(expires_at)),
Err(_) => (None, None),
};
(
export_ready,
export_not_ready_reason,
observation_receipt,
export_expires_at,
disposition_receipt,
disposition_expires_at,
)
}
Err(_) => (false, Some("fleet_or_source_not_ready"), None, None, None, None),
};
json_response(
StatusCode::OK,
&IlmRecoveryControlInspectResponse {
@@ -1567,16 +1599,18 @@ impl Operation for IlmRecoveryControlInspectHandler {
export_ready,
export_not_ready_reason,
observation_receipt,
observation_receipt_expires_at_unix_nanos: expires_at,
observation_receipt_expires_at_unix_nanos,
disposition_dry_run_receipt,
disposition_dry_run_receipt_expires_at_unix_nanos,
},
)
}
}
pub struct IlmRecoveryExportCreateHandler {}
pub struct IlmRecoveryRecordMutationHandler {}
#[async_trait::async_trait]
impl Operation for IlmRecoveryExportCreateHandler {
impl Operation for IlmRecoveryRecordMutationHandler {
async fn call(&self, mut req: S3Request<Body>, params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
let actor_sha256 = authorize_recovery_admin_request(&req, AdminAction::SetTierAction).await?;
let control_id = recovery_control_id_from_params(&params)?;
@@ -1584,35 +1618,115 @@ impl Operation for IlmRecoveryExportCreateHandler {
return Err(admin_s3_error(AdminS3ErrorCode::InternalError, "object store is not initialized"));
};
let body = req.input.store_all_limited(MAX_ADMIN_REQUEST_BODY_SIZE).await.map_err(|_| {
admin_s3_error(AdminS3ErrorCode::InvalidRequest, "ILM recovery export body is too large or unreadable")
admin_s3_error(AdminS3ErrorCode::InvalidRequest, "ILM recovery request body is too large or unreadable")
})?;
let request = parse_recovery_record_mutation_request(&body)?;
let ValidatedIlmRecoveryRecordMutation::Export { observation_receipt } =
validate_recovery_record_mutation_request(&request)?
else {
return Err(admin_s3_error(AdminS3ErrorCode::InvalidArgument, "unsupported ILM recovery action"));
};
let receipt = decode_recovery_receipt(observation_receipt, &recovery_receipt_credentials()?)?;
let now = i64::try_from(OffsetDateTime::now_utc().unix_timestamp_nanos())
let mutation = validate_recovery_record_mutation_request(&request)?;
let now = OffsetDateTime::now_utc();
let now_unix_nanos = i64::try_from(now.unix_timestamp_nanos())
.map_err(|_| admin_s3_error(AdminS3ErrorCode::InternalError, "ILM recovery receipt timestamp is invalid"))?;
let observation = validate_recovery_observation_receipt(
receipt,
&actor_sha256,
&control_id,
IlmRecoveryReceiptAction::Export,
IlmRecoveryReceiptMode::Execute,
now,
)?;
let created = create_recovery_export(store, &observation, &actor_sha256)
.await
.map_err(map_recovery_export_error)?;
let response = IlmRecoveryExportCreateResponse {
download_url: format!("{ADMIN_PREFIX}/v3/ilm/recovery/exports/{}", created.export_id),
outcome: if created.replayed { "replayed" } else { "created" },
export_id: created.export_id,
export_sha256: created.content_sha256,
};
json_response(StatusCode::OK, &response)
let receipt_credentials = recovery_receipt_credentials()?;
match mutation {
ValidatedIlmRecoveryRecordMutation::Export { observation_receipt } => {
let receipt = decode_recovery_receipt(observation_receipt, &receipt_credentials)?;
let observation = validate_recovery_observation_receipt(
receipt,
&actor_sha256,
&control_id,
IlmRecoveryReceiptAction::Export,
IlmRecoveryReceiptMode::Execute,
now_unix_nanos,
)?;
let created = create_recovery_export(store, &observation, &actor_sha256)
.await
.map_err(map_recovery_export_error)?;
let response = IlmRecoveryExportCreateResponse {
download_url: format!("{ADMIN_PREFIX}/v3/ilm/recovery/exports/{}", created.export_id),
outcome: if created.replayed { "replayed" } else { "created" },
export_id: created.export_id,
export_sha256: created.content_sha256,
};
json_response(StatusCode::OK, &response)
}
ValidatedIlmRecoveryRecordMutation::AbandonDryRun {
observation_receipt,
export_id,
export_sha256,
reason_code: IlmRecoveryDispositionReasonCode::LegacyRemoteCleanupAbandoned,
} => {
let receipt = decode_recovery_receipt(observation_receipt, &receipt_credentials)?;
let observation = validate_recovery_observation_receipt(
receipt,
&actor_sha256,
&control_id,
IlmRecoveryReceiptAction::AbandonRemoteCleanup,
IlmRecoveryReceiptMode::DryRun,
now_unix_nanos,
)?;
let dry_run =
dry_run_recovery_disposition(store, &observation, export_id, export_sha256, &actor_sha256, now_unix_nanos)
.await
.map_err(map_recovery_disposition_error)?;
let execute_receipt_now = OffsetDateTime::now_utc();
let (execute_receipt, execute_receipt_expires_at_unix_nanos) = issue_recovery_observation_receipt(
observation,
actor_sha256,
IlmRecoveryReceiptAction::AbandonRemoteCleanup,
IlmRecoveryReceiptMode::Execute,
execute_receipt_now,
)?;
json_response(
StatusCode::OK,
&IlmRecoveryDispositionDryRunResponse {
action: IlmRecoveryReceiptAction::AbandonRemoteCleanup,
mode: IlmRecoveryReceiptMode::DryRun,
status: IlmRecoveryDispositionDryRunStatus::Ready,
disposition_id: dry_run.disposition_id,
export_id: dry_run.export_id,
export_sha256: dry_run.export_content_sha256,
source_generation_sha256: dry_run.source_generation_sha256,
copy_set_sha256: dry_run.copy_set_sha256,
source_copy_count: dry_run.source_copy_count,
observation_receipt: execute_receipt,
observation_receipt_expires_at_unix_nanos: execute_receipt_expires_at_unix_nanos,
},
)
}
ValidatedIlmRecoveryRecordMutation::AbandonExecute {
observation_receipt,
export_id,
export_sha256,
reason_code: IlmRecoveryDispositionReasonCode::LegacyRemoteCleanupAbandoned,
} => {
let receipt = decode_recovery_receipt(observation_receipt, &receipt_credentials)?;
let observation = validate_recovery_observation_receipt(
receipt,
&actor_sha256,
&control_id,
IlmRecoveryReceiptAction::AbandonRemoteCleanup,
IlmRecoveryReceiptMode::Execute,
now_unix_nanos,
)?;
let execution =
execute_recovery_disposition(store, &observation, export_id, export_sha256, &actor_sha256, now_unix_nanos)
.await
.map_err(map_recovery_disposition_error)?;
let state = recovery_disposition_response_state(execution.state)?;
json_response(
StatusCode::OK,
&IlmRecoveryDispositionExecuteResponse {
action: IlmRecoveryReceiptAction::AbandonRemoteCleanup,
mode: IlmRecoveryReceiptMode::Execute,
disposition_id: execution.disposition_id,
state,
outcome: execution.outcome,
confirmed_absent_copy_count: execution.confirmed_absent_copy_count,
source_copy_count: execution.source_copy_count,
},
)
}
}
}
}
@@ -1809,17 +1923,74 @@ mod tests {
assert!(!token.contains("actor-a"));
assert!(!token.contains("ilm/tier-delete-journal"));
assert_eq!(decode_recovery_receipt(&token, &credentials).unwrap(), payload);
assert!(
validate_recovery_observation_receipt(
payload.clone(),
&payload.actor_sha256,
&payload.observation.control_id,
IlmRecoveryReceiptAction::Export,
let receipt_classes = [
(payload.clone(), IlmRecoveryReceiptAction::Export, IlmRecoveryReceiptMode::Execute),
(
IlmRecoveryObservationReceipt {
action: IlmRecoveryReceiptAction::AbandonRemoteCleanup,
mode: IlmRecoveryReceiptMode::DryRun,
..payload.clone()
},
IlmRecoveryReceiptAction::AbandonRemoteCleanup,
IlmRecoveryReceiptMode::DryRun,
),
(
IlmRecoveryObservationReceipt {
action: IlmRecoveryReceiptAction::AbandonRemoteCleanup,
mode: IlmRecoveryReceiptMode::Execute,
..payload.clone()
},
IlmRecoveryReceiptAction::AbandonRemoteCleanup,
IlmRecoveryReceiptMode::Execute,
payload.issued_at_unix_nanos,
),
];
let expected_classes = [
(IlmRecoveryReceiptAction::Export, IlmRecoveryReceiptMode::Execute),
(IlmRecoveryReceiptAction::AbandonRemoteCleanup, IlmRecoveryReceiptMode::DryRun),
(IlmRecoveryReceiptAction::AbandonRemoteCleanup, IlmRecoveryReceiptMode::Execute),
];
for (receipt, actual_action, actual_mode) in &receipt_classes {
for (expected_action, expected_mode) in expected_classes {
let result = validate_recovery_observation_receipt(
receipt.clone(),
&receipt.actor_sha256,
&receipt.observation.control_id,
expected_action,
expected_mode,
receipt.issued_at_unix_nanos,
);
if (*actual_action, *actual_mode) == (expected_action, expected_mode) {
assert!(result.is_ok(), "the matching receipt class must validate");
} else {
assert_eq!(
result.expect_err("receipts must not cross action or mode boundaries").code(),
&S3ErrorCode::AccessDenied
);
}
}
let actor_mismatch = validate_recovery_observation_receipt(
receipt.clone(),
&hex_sha256(b"actor-b", ToOwned::to_owned),
&receipt.observation.control_id,
*actual_action,
*actual_mode,
receipt.issued_at_unix_nanos,
)
.is_ok()
);
.expect_err("receipts must remain bound to the authenticated actor");
assert_eq!(actor_mismatch.code(), &S3ErrorCode::AccessDenied);
let expired = validate_recovery_observation_receipt(
receipt.clone(),
&receipt.actor_sha256,
&receipt.observation.control_id,
*actual_action,
*actual_mode,
receipt.expires_at_unix_nanos,
)
.expect_err("expired receipts must fail closed");
assert_eq!(expired.code(), &S3ErrorCode::AccessDenied);
}
let assert_denied = |receipt: IlmRecoveryObservationReceipt, actor: &str, control: &str, now: i64| {
let err = validate_recovery_observation_receipt(
receipt,
@@ -1832,19 +2003,7 @@ mod tests {
.expect_err("invalid observation receipt must be denied");
assert_eq!(err.code(), &S3ErrorCode::AccessDenied);
};
assert_denied(
payload.clone(),
&hex_sha256(b"actor-b", ToOwned::to_owned),
&payload.observation.control_id,
payload.issued_at_unix_nanos,
);
assert_denied(payload.clone(), &payload.actor_sha256, &"cd".repeat(32), payload.issued_at_unix_nanos);
assert_denied(
payload.clone(),
&payload.actor_sha256,
&payload.observation.control_id,
payload.expires_at_unix_nanos,
);
let mut invalid = payload.clone();
invalid.schema = "rustfs-ilm-recovery-observation-receipt-v2".to_string();
@@ -1994,10 +2153,15 @@ mod tests {
for invalid in [
execute_json.replace(r#""confirm":true,"#, ""),
execute_json.replace(r#""confirm":true"#, r#""confirm":false"#),
execute_json.replace(r#""confirm":true"#, r#""confirm":null"#),
execute_json.replace(
r#""acknowledge_remote_cleanup_abandoned":true"#,
r#""acknowledge_remote_cleanup_abandoned":false"#,
),
execute_json.replace(
r#""acknowledge_remote_cleanup_abandoned":true"#,
r#""acknowledge_remote_cleanup_abandoned":null"#,
),
execute_json.replace(export_id.as_str(), uppercase_export_id.as_str()),
execute_json.replace(export_sha256.as_str(), "too-short"),
execute_json.replace("opaque-execute", ""),
@@ -2010,8 +2174,19 @@ mod tests {
}
}
assert!(parse_recovery_record_mutation_request(execute_json.replace(r#""mode":"execute","#, "").as_bytes()).is_err());
assert!(
parse_recovery_record_mutation_request(
dry_run_json
.replace("legacy_remote_cleanup_abandoned", "operator_override")
.as_bytes()
)
.is_err()
);
for invalid in [
br#"{"action":"export","observation_receipt":"opaque","extra":true}"#.as_slice(),
br#"{"action":"abandon_remote_cleanup","mode":null}"#.as_slice(),
br#"{"action":"abandon_remote_cleanup","mode":"preview"}"#.as_slice(),
br#"{"action":"unknown","observation_receipt":"opaque"}"#.as_slice(),
] {
@@ -2051,11 +2226,24 @@ mod tests {
observation_receipt_expires_at_unix_nanos: 900_000_000_001,
};
let dry_run_json = serde_json::to_value(&dry_run).unwrap();
assert_eq!(dry_run_json["action"], "abandon_remote_cleanup");
assert_eq!(dry_run_json["mode"], "dry_run");
assert_eq!(dry_run_json["status"], "ready");
assert_eq!(
serde_json::from_value::<IlmRecoveryDispositionDryRunResponse>(dry_run_json).unwrap(),
dry_run_json,
serde_json::json!({
"action": "abandon_remote_cleanup",
"mode": "dry_run",
"status": "ready",
"disposition_id": "ab".repeat(32),
"export_id": "cd".repeat(32),
"export_sha256": "ef".repeat(32),
"source_generation_sha256": "12".repeat(32),
"copy_set_sha256": "34".repeat(32),
"source_copy_count": 2,
"observation_receipt": "opaque-execute",
"observation_receipt_expires_at_unix_nanos": 900_000_000_001_i64,
})
);
assert_eq!(
serde_json::from_value::<IlmRecoveryDispositionDryRunResponse>(dry_run_json.clone()).unwrap(),
dry_run
);
@@ -2063,22 +2251,68 @@ mod tests {
action: IlmRecoveryReceiptAction::AbandonRemoteCleanup,
mode: IlmRecoveryReceiptMode::Execute,
disposition_id: "ab".repeat(32),
state: IlmRecoveryDispositionState::Applying,
outcome: IlmRecoveryDispositionOutcome::AcceptedForRecovery,
state: IlmRecoveryDispositionResponseState::Applying,
outcome: IlmRecoveryDispositionExecutionOutcome::AcceptedForRecovery,
confirmed_absent_copy_count: 1,
source_copy_count: 2,
};
let execute_json = serde_json::to_value(&execute).unwrap();
assert_eq!(execute_json["state"], "applying");
assert_eq!(execute_json["outcome"], "accepted_for_recovery");
assert_eq!(
serde_json::from_value::<IlmRecoveryDispositionExecuteResponse>(execute_json).unwrap(),
execute_json,
serde_json::json!({
"action": "abandon_remote_cleanup",
"mode": "execute",
"disposition_id": "ab".repeat(32),
"state": "applying",
"outcome": "accepted_for_recovery",
"confirmed_absent_copy_count": 1,
"source_copy_count": 2,
})
);
assert_eq!(
serde_json::from_value::<IlmRecoveryDispositionExecuteResponse>(execute_json.clone()).unwrap(),
execute
);
let mut unknown = serde_json::to_value(&dry_run).unwrap();
unknown["unexpected"] = serde_json::json!(true);
assert!(serde_json::from_value::<IlmRecoveryDispositionDryRunResponse>(unknown).is_err());
let mut unknown_dry_run = dry_run_json;
unknown_dry_run["unexpected"] = serde_json::json!(true);
assert!(serde_json::from_value::<IlmRecoveryDispositionDryRunResponse>(unknown_dry_run).is_err());
let mut unknown_execute = execute_json;
unknown_execute["unexpected"] = serde_json::json!(true);
assert!(serde_json::from_value::<IlmRecoveryDispositionExecuteResponse>(unknown_execute).is_err());
}
#[test]
fn prepared_recovery_disposition_is_operation_aborted_and_not_a_wire_state() {
let err = recovery_disposition_response_state(IlmRecoveryDispositionState::Prepared)
.expect_err("Prepared must not escape through the closed execute response");
assert_eq!(err.code(), &S3ErrorCode::OperationAborted);
assert_eq!(err.message(), Some("ILM recovery disposition is accepted but not yet applying"));
assert!(serde_json::from_str::<IlmRecoveryDispositionResponseState>(r#""prepared""#).is_err());
assert_eq!(
serde_json::to_string(&IlmRecoveryDispositionResponseState::Applying).unwrap(),
r#""applying""#
);
assert_eq!(
serde_json::to_string(&IlmRecoveryDispositionResponseState::Completed).unwrap(),
r#""completed""#
);
}
#[test]
fn recovery_disposition_error_mapping_is_stable_and_fail_closed() {
let not_found = map_recovery_disposition_error(StorageError::ConfigNotFound);
assert_eq!(not_found.code(), &S3ErrorCode::NoSuchKey);
assert_eq!(not_found.message(), Some("ILM recovery export or disposition not found"));
let overloaded = map_recovery_disposition_error(StorageError::SlowDown);
assert_eq!(overloaded.code(), &S3ErrorCode::SlowDown);
assert_eq!(overloaded.message(), Some("ILM recovery disposition admission capacity is exhausted"));
let stale = map_recovery_disposition_error(StorageError::PreconditionFailed);
assert_eq!(stale.code(), &S3ErrorCode::OperationAborted);
assert_eq!(stale.message(), Some("ILM recovery disposition request cannot proceed"));
}
#[test]
@@ -2094,7 +2328,7 @@ mod tests {
);
}
fn manual_transition_job_request(method: Method, path: &'static str) -> S3Request<Body> {
fn credential_less_admin_request(method: Method, path: &'static str) -> S3Request<Body> {
S3Request {
input: Body::empty(),
method,
@@ -2470,7 +2704,7 @@ mod tests {
#[tokio::test]
async fn transition_admin_gate_keeps_its_missing_credentials_response() {
let err = authorize_transition_admin_request(
&manual_transition_job_request(Method::GET, "/rustfs/admin/v3/ilm/transition/jobs/job-123"),
&credential_less_admin_request(Method::GET, "/rustfs/admin/v3/ilm/transition/jobs/job-123"),
AdminAction::ListTierAction,
)
.await
@@ -2483,8 +2717,8 @@ mod tests {
#[tokio::test]
async fn recovery_admin_gate_keeps_its_missing_credentials_response() {
let err = authorize_recovery_admin_request(
&manual_transition_job_request(Method::GET, "/rustfs/admin/v3/ilm/recovery/controls/control-123"),
AdminAction::ListTierAction,
&credential_less_admin_request(Method::POST, "/rustfs/admin/v3/ilm/recovery/records/control-id"),
AdminAction::SetTierAction,
)
.await
.expect_err("a recovery admin request without credentials must fail");
@@ -2616,7 +2850,7 @@ mod tests {
async fn manual_transition_job_handlers_reject_missing_credentials_before_status_contract() {
let status_err = ManualTransitionJobStatusHandler {}
.call(
manual_transition_job_request(Method::GET, "/rustfs/admin/v3/ilm/transition/jobs/job-123"),
credential_less_admin_request(Method::GET, "/rustfs/admin/v3/ilm/transition/jobs/job-123"),
Params::new(),
)
.await
@@ -2626,7 +2860,7 @@ mod tests {
let cancel_err = ManualTransitionJobCancelHandler {}
.call(
manual_transition_job_request(Method::DELETE, "/rustfs/admin/v3/ilm/transition/jobs/job-123"),
credential_less_admin_request(Method::DELETE, "/rustfs/admin/v3/ilm/transition/jobs/job-123"),
Params::new(),
)
.await
+4
View File
@@ -235,6 +235,10 @@ pub(crate) mod lifecycle {
pub(crate) use super::ecstore_bucket::lifecycle::recovery_control::{
IlmRecoveryClassification, IlmRecoveryControlView, IlmRecoveryProtocol, inspect_recovery_control, list_recovery_controls,
};
pub(crate) use super::ecstore_bucket::lifecycle::recovery_disposition::{
IlmRecoveryDispositionExecutionOutcome, IlmRecoveryDispositionReasonCode, IlmRecoveryDispositionState,
dry_run_recovery_disposition, execute_recovery_disposition,
};
pub(crate) use super::ecstore_bucket::lifecycle::recovery_export::{
IlmRecoveryExportObservation, create_recovery_export, inspect_recovery_export_observation, load_recovery_export,
};