diff --git a/crates/e2e_test/src/admin_auth_test.rs b/crates/e2e_test/src/admin_auth_test.rs index ddcd4da40..8f50426d3 100644 --- a/crates/e2e_test/src/admin_auth_test.rs +++ b/crates/e2e_test/src/admin_auth_test.rs @@ -48,6 +48,7 @@ mod tests { const ADMIN_INFO_PATH: &str = "/rustfs/admin/v3/info"; const ADMIN_MANUAL_TRANSITION_PATH: &str = "/rustfs/admin/v3/ilm/transition/run?bucket=auth-deny-manual-transition&maxObjects=1"; + const ADMIN_MANUAL_TRANSITION_JOB_PATH: &str = "/rustfs/admin/v3/ilm/transition/jobs/11111111-1111-4111-8111-111111111111"; /// Send a SigV4-signed request to `path` (optionally with a JSON `body`) and /// return `(status, body)`. Uses the `UNSIGNED_PAYLOAD` content hash so a @@ -189,6 +190,42 @@ mod tests { root_body.contains("\"mode\":\"enqueue_only\""), "root response should be the manual transition JSON contract, body: {root_body}" ); + let (root_status, root_body) = signed_request( + &env.url, + http::Method::GET, + ADMIN_MANUAL_TRANSITION_JOB_PATH, + None, + &env.access_key, + &env.secret_key, + ) + .await?; + assert_eq!( + root_status, + reqwest::StatusCode::NOT_FOUND, + "root credential must reach the manual transition status handler, body: {root_body}" + ); + assert!( + root_body.contains("NoSuchKey"), + "missing durable job should return NoSuchKey once authorized, body: {root_body}" + ); + let (root_status, root_body) = signed_request( + &env.url, + http::Method::DELETE, + ADMIN_MANUAL_TRANSITION_JOB_PATH, + None, + &env.access_key, + &env.secret_key, + ) + .await?; + assert_eq!( + root_status, + reqwest::StatusCode::NOT_FOUND, + "root credential must reach the manual transition cancel handler, body: {root_body}" + ); + assert!( + root_body.contains("NoSuchKey"), + "missing durable job cancel should return NoSuchKey once authorized, body: {root_body}" + ); let (status, body) = signed_request(&env.url, http::Method::POST, ADMIN_MANUAL_TRANSITION_PATH, None, user_ak, user_sk).await?; @@ -201,6 +238,28 @@ mod tests { body.contains("AccessDenied"), "manual transition rejection must carry the AccessDenied S3 error code, body: {body}" ); + let (status, body) = + signed_request(&env.url, http::Method::GET, ADMIN_MANUAL_TRANSITION_JOB_PATH, None, user_ak, user_sk).await?; + assert_eq!( + status, + reqwest::StatusCode::FORBIDDEN, + "non-admin credential must get 403 on manual transition status, body: {body}" + ); + assert!( + body.contains("AccessDenied"), + "manual transition status rejection must carry the AccessDenied S3 error code, body: {body}" + ); + let (status, body) = + signed_request(&env.url, http::Method::DELETE, ADMIN_MANUAL_TRANSITION_JOB_PATH, None, user_ak, user_sk).await?; + assert_eq!( + status, + reqwest::StatusCode::FORBIDDEN, + "non-admin credential must get 403 on manual transition cancel, body: {body}" + ); + assert!( + body.contains("AccessDenied"), + "manual transition cancel rejection must carry the AccessDenied S3 error code, body: {body}" + ); env.stop_server(); Ok(()) diff --git a/crates/e2e_test/src/reliant/tiering.rs b/crates/e2e_test/src/reliant/tiering.rs index 3a8be7800..97136c50f 100644 --- a/crates/e2e_test/src/reliant/tiering.rs +++ b/crates/e2e_test/src/reliant/tiering.rs @@ -71,8 +71,10 @@ const MANUAL_DRY_RUN_BUCKET: &str = "ilm7-manual-dry-run"; const MANUAL_NOT_DUE_BUCKET: &str = "ilm7-manual-not-due"; const MANUAL_QUEUE_PRESSURE_BUCKET: &str = "ilm7-manual-queue-pressure"; const MANUAL_ASYNC_STATUS_BUCKET: &str = "ilm7-manual-async-status"; +const MANUAL_CONTINUATION_BUCKET: &str = "ilm7-manual-continuation"; const MANUAL_ASYNC_LIMIT_BUCKET: &str = "ilm7-manual-async-limit"; const MANUAL_QUEUE_PRESSURE_PREFIX: &str = "manual-queue-pressure/"; +const MANUAL_CONTINUATION_PREFIX: &str = "manual-continuation/"; const MANUAL_ASYNC_LIMIT_PREFIX: &str = "manual-async-limit/"; const OBJECT_KEY: &str = "tier/鲁A12345/report.bin"; const MANUAL_DUE_KEY: &str = "manual-due/report.bin"; @@ -352,6 +354,7 @@ struct ManualTransitionRunReport { skipped_queue_timeout: u64, truncated_by_limit: bool, truncated_by_duration: bool, + continuation_token: Option, } #[derive(Debug, Deserialize)] @@ -379,13 +382,28 @@ async fn manual_transition_run_with_max( prefix: &str, dry_run: bool, max_objects: u64, +) -> Result> { + manual_transition_run_with_max_and_continuation(hot, bucket, prefix, dry_run, max_objects, None).await +} + +async fn manual_transition_run_with_max_and_continuation( + hot: &RustFSTestEnvironment, + bucket: &str, + prefix: &str, + dry_run: bool, + max_objects: u64, + continuation_token: Option<&str>, ) -> Result> { let bucket = urlencoding::encode(bucket); let prefix = urlencoding::encode(prefix); let tier = urlencoding::encode(TIER_NAME); - let path = format!( + let mut path = format!( "/rustfs/admin/v3/ilm/transition/run?bucket={bucket}&prefix={prefix}&tier={tier}&dryRun={dry_run}&maxObjects={max_objects}" ); + if let Some(token) = continuation_token { + path.push_str("&continuationToken="); + path.push_str(&urlencoding::encode(token)); + } let (status, body) = signed_admin_request(&hot.url, Method::POST, &path, None, &hot.access_key, &hot.secret_key).await?; if !status.is_success() { return Err(format!("manual transition run failed: status={status}, body={body}").into()); @@ -928,6 +946,77 @@ async fn test_manual_transition_run_contract_no_status_cancel_fields() -> TestRe Ok(()) } +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn test_manual_transition_run_continuation_token_resumes_without_raw_markers() -> TestResult { + let mut cold = RustFSTestEnvironment::new().await?; + cold.access_key = "manualcontinuationcoldtieradmin".to_string(); + cold.secret_key = "manualcontinuationcoldtiersecret".to_string(); + cold.start_rustfs_server_without_cleanup(vec![]).await?; + let cold_client = cold.create_s3_client(); + cold_client.create_bucket().bucket(TIER_BUCKET).send().await?; + + let mut hot = RustFSTestEnvironment::new().await?; + hot.start_rustfs_server_with_env(vec![], &[("RUSTFS_SCANNER_ENABLED", "false"), ("RUSTFS_SCANNER_CYCLE", "3600")]) + .await?; + let hot_client = hot.create_s3_client(); + add_rustfs_tier(&hot, &cold).await?; + + hot_client.create_bucket().bucket(MANUAL_CONTINUATION_BUCKET).send().await?; + for idx in 0..2 { + let key = format!("{MANUAL_CONTINUATION_PREFIX}obj-{idx:02}"); + put_single_part_object(&hot_client, MANUAL_CONTINUATION_BUCKET, &key, b"manual continuation payload").await?; + } + put_lifecycle_transition_rule( + &hot_client, + MANUAL_CONTINUATION_BUCKET, + "manual-continuation", + MANUAL_CONTINUATION_PREFIX, + 0, + ) + .await?; + + let first = manual_transition_run_with_max(&hot, MANUAL_CONTINUATION_BUCKET, MANUAL_CONTINUATION_PREFIX, true, 1).await?; + assert_eq!(first.state, "partial", "first continuation page: {first:#?}"); + assert_eq!(first.mode, "enqueue_only"); + assert_eq!(first.report.bucket, MANUAL_CONTINUATION_BUCKET); + assert_eq!(first.report.prefix, MANUAL_CONTINUATION_PREFIX); + assert!(first.report.dry_run); + assert_eq!(first.report.scanned, 1, "first continuation page: {first:#?}"); + assert_eq!(first.report.eligible, 1, "first continuation page: {first:#?}"); + assert_eq!(first.report.dry_run_eligible, 1, "first continuation page: {first:#?}"); + assert!(first.report.truncated_by_limit); + let continuation = first + .report + .continuation_token + .as_deref() + .ok_or("partial manual transition run must return an opaque continuation token")?; + assert!( + !continuation.contains(MANUAL_CONTINUATION_PREFIX), + "continuation token must not expose the raw object prefix: {continuation}" + ); + + hot.restart_server_preserving_data(vec![], &[("RUSTFS_SCANNER_ENABLED", "false"), ("RUSTFS_SCANNER_CYCLE", "3600")]) + .await?; + + let second = manual_transition_run_with_max_and_continuation( + &hot, + MANUAL_CONTINUATION_BUCKET, + MANUAL_CONTINUATION_PREFIX, + true, + 10, + Some(continuation), + ) + .await?; + assert_eq!(second.state, "completed", "second continuation page: {second:#?}"); + assert_eq!(second.report.scanned, 1, "second continuation page: {second:#?}"); + assert_eq!(second.report.eligible, 1, "second continuation page: {second:#?}"); + assert_eq!(second.report.dry_run_eligible, 1, "second continuation page: {second:#?}"); + assert!(!second.report.truncated_by_limit); + assert!(second.report.continuation_token.is_none()); + + Ok(()) +} + #[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn test_manual_transition_run_queue_pressure_partial() -> TestResult { let mut cold = RustFSTestEnvironment::new().await?; diff --git a/crates/ecstore/src/api/mod.rs b/crates/ecstore/src/api/mod.rs index 41395f53e..dc839701a 100644 --- a/crates/ecstore/src/api/mod.rs +++ b/crates/ecstore/src/api/mod.rs @@ -43,12 +43,27 @@ pub mod bucket { pub mod bucket_lifecycle_ops { pub use crate::bucket::lifecycle::bucket_lifecycle_ops::{ - ExpiryState, LifecycleOps, ManualTransitionRunExecution, ManualTransitionRunOptions, ManualTransitionRunReport, - RestoreRequestOps, TransitionState, TransitionedObject, apply_expiry_rule, apply_transition_rule, - enqueue_expiry_for_existing_objects, enqueue_transition_for_existing_objects, + ExpiryState, LifecycleOps, ManualTransitionCancelCheck, ManualTransitionProgressSink, + ManualTransitionQueueSnapshot, ManualTransitionRunExecution, ManualTransitionRunOptions, + ManualTransitionRunReport, RestoreRequestOps, TransitionState, TransitionedObject, apply_expiry_rule, + apply_transition_rule, enqueue_expiry_for_existing_objects, enqueue_transition_for_existing_objects, enqueue_transition_for_existing_objects_scoped, enqueue_transition_for_existing_objects_scoped_with_cancel, enqueue_transition_immediate, expire_transitioned_object, get_global_expiry_state, get_global_transition_state, - init_background_expiry, post_restore_opts, run_stale_multipart_upload_cleanup_once, validate_transition_tier, + init_background_expiry, manual_transition_queue_snapshot, post_restore_opts, + run_stale_multipart_upload_cleanup_once, validate_transition_tier, + }; + } + + pub mod manual_transition_job { + pub use crate::bucket::lifecycle::manual_transition_job::{ + ManualTransitionJobRecord, ManualTransitionJobState, ManualTransitionScopeAdmission, + ManualTransitionScopeAdmissionClaim, claim_manual_transition_scope_admission, + delete_manual_transition_scope_admission_if_current, load_manual_transition_job_record, + load_manual_transition_job_record_with_etag, load_manual_transition_scope_admission, + manual_transition_job_lease_expired, manual_transition_scope_admission_lease_expired, + manual_transition_scope_key, persist_manual_transition_job_progress, renew_manual_transition_job_lease, + request_manual_transition_job_cancel, save_manual_transition_job_record, + save_manual_transition_job_record_if_current, save_manual_transition_scope_admission_if_absent, }; } diff --git a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs index 9deaec6d3..95dd32de8 100644 --- a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs +++ b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs @@ -21,6 +21,13 @@ use crate::bucket::lifecycle::evaluator::Evaluator; use crate::bucket::lifecycle::lifecycle::{ self, Lifecycle, ObjectOpts, TransitionOptions, abort_incomplete_multipart_upload_due, }; +use crate::bucket::lifecycle::manual_transition_job::{ + MANUAL_TRANSITION_JOB_RECORD_PREFIX, ManualTransitionJobRecord, ManualTransitionScopeAdmission, + ManualTransitionScopeAdmissionClaim, claim_manual_transition_scope_admission, + delete_manual_transition_scope_admission_if_current, load_manual_transition_job_record, + load_manual_transition_job_record_with_etag, manual_transition_job_id_from_record_object_name, + manual_transition_job_lease_expired, persist_manual_transition_job_progress, save_manual_transition_job_record_if_current, +}; use crate::bucket::lifecycle::replication_sink; use crate::bucket::lifecycle::replication_sink::{ ReplicateDecision, ReplicationState, ReplicationStatusType, VersionPurgeStatusType, replication_state_to_filemeta, @@ -36,7 +43,7 @@ use crate::bucket::lifecycle::transition_transaction::run_transition_transaction use crate::bucket::versioning_sys::BucketVersioningSys; use crate::client::object_api_utils::new_getobjectreader; use crate::disk::error::DiskError; -use crate::disk::{DeleteOptions, Disk, DiskAPI, RUSTFS_META_MULTIPART_BUCKET, STORAGE_FORMAT_FILE}; +use crate::disk::{DeleteOptions, Disk, DiskAPI, RUSTFS_META_BUCKET, RUSTFS_META_MULTIPART_BUCKET, STORAGE_FORMAT_FILE}; use crate::error::Error; use crate::error::StorageError; use crate::error::{error_resp_to_object_err, is_err_object_not_found, is_err_version_not_found, is_network_or_host_down}; @@ -58,7 +65,6 @@ use crate::storage_api_contracts::{ }; use crate::store::ECStore; use async_channel::{Receiver as A_Receiver, Sender as A_Sender, bounded}; -use futures::Future; use http::HeaderMap; use rand::RngExt as _; use rustfs_common::metrics::{ @@ -83,6 +89,7 @@ use sha2::{Digest, Sha256}; use std::any::Any; use std::collections::{HashMap, HashSet}; use std::env; +use std::future::Future; use std::pin::Pin; use std::sync::atomic::{AtomicI64, Ordering}; use std::sync::{Arc, Mutex, OnceLock, Weak}; @@ -116,6 +123,7 @@ pub type ExpiryOpType = Box; static XXHASH_SEED: u64 = 0; static TIER_FREE_VERSION_RECOVERY_STARTED: OnceLock<()> = OnceLock::new(); +static MANUAL_TRANSITION_JOB_RECOVERY_STARTED: OnceLock<()> = OnceLock::new(); pub const AMZ_OBJECT_TAGGING: &str = "X-Amz-Tagging"; pub const AMZ_TAG_COUNT: &str = "x-amz-tagging-count"; @@ -135,6 +143,7 @@ const TIER_FREE_VERSION_RECOVERY_MAX_IDLE_INTERVAL: StdDuration = StdDuration::f const TIER_FREE_VERSION_RECOVERY_JITTER_PERCENT: u64 = 10; const DATE_EXPIRY_EXISTING_OBJECTS_GRACE_SECS: i64 = 5; const EXPIRY_WORKER_QUEUE_CAPACITY: usize = 1000; +const DEFAULT_MANUAL_TRANSITION_JOB_RECOVERY_LIMIT: usize = 100; // Phase 5 (backlog#939): lifecycle expiry/transition state moved into the // per-instance `InstanceContext`; these owner helpers forward to the current @@ -1165,6 +1174,20 @@ impl TransitionState { global_metrics().record_scanner_lifecycle_transition_state(self.scanner_transition_state_update()); } + pub fn manual_transition_queue_snapshot(&self) -> ManualTransitionQueueSnapshot { + let state = self.scanner_transition_state_update(); + ManualTransitionQueueSnapshot { + queue_capacity: state.queue_capacity, + queued: state.queued, + active: state.active, + workers: state.workers, + queue_full: state.queue_full, + queue_send_timeout: state.queue_send_timeout, + compensation_pending: state.compensation_pending, + compensation_running: state.compensation_running, + } + } + fn handle_immediate_enqueue_failure(self: &Arc, oi: &ObjectInfo, src: &LcEventSrc, failure: ImmediateEnqueueFailure) { Self::inc_counter(&self.missed_immediate_tasks); let scheduled = self.schedule_bucket_compensation(&oi.bucket); @@ -1569,7 +1592,276 @@ pub async fn init_background_expiry(api: Arc) { ExpiryState::resize_workers(workers, api.clone()).await; 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); + spawn_transition_transaction_recovery_once(api.clone()); + spawn_manual_transition_job_recovery_once(api); +} + +fn spawn_manual_transition_job_recovery_once(api: Arc) -> Option> { + if MANUAL_TRANSITION_JOB_RECOVERY_STARTED.set(()).is_err() { + return None; + } + + Some(tokio::spawn(async move { + let cancel_token = runtime_sources::background_services_cancel_token().unwrap_or_default(); + select! { + _ = cancel_token.cancelled() => {} + result = recover_manual_transition_jobs_once(api, DEFAULT_MANUAL_TRANSITION_JOB_RECOVERY_LIMIT, None) => { + match result { + Ok(stats) => { + debug!( + event = EVENT_LIFECYCLE_WORKER_STATE, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_LIFECYCLE, + scanned = stats.scanned, + resumed = stats.resumed, + cancelled = stats.cancelled, + skipped = stats.skipped, + failed = stats.failed, + truncated = stats.truncated, + next_marker = ?stats.next_marker, + state = "manual_transition_recovery_completed", + "Manual transition job recovery completed" + ); + } + Err(err) => { + warn!( + event = EVENT_LIFECYCLE_WORKER_STATE, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_LIFECYCLE, + error = %err, + state = "manual_transition_recovery_failed", + "Manual transition job recovery failed" + ); + } + } + } + } + })) +} + +#[derive(Debug, Clone, Default, PartialEq, Eq)] +pub struct ManualTransitionJobRecoveryStats { + pub scanned: u64, + pub resumed: u64, + pub cancelled: u64, + pub skipped: u64, + pub failed: u64, + pub next_marker: Option, + pub truncated: bool, +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum ManualTransitionJobRecoveryOutcome { + Resumed, + Cancelled, + Skipped, +} + +pub async fn recover_manual_transition_jobs_once( + api: Arc, + limit: usize, + marker: Option, +) -> Result { + if limit == 0 { + return Err(Error::other("manual transition job recovery limit must be greater than zero")); + } + let list_limit = i32::try_from(limit).map_or(i32::MAX, |value| value); + let page = api + .clone() + .list_objects_v2( + RUSTFS_META_BUCKET, + MANUAL_TRANSITION_JOB_RECORD_PREFIX, + marker, + None, + list_limit, + false, + None, + false, + ) + .await?; + let mut stats = ManualTransitionJobRecoveryStats { + next_marker: page.next_continuation_token, + truncated: page.is_truncated, + ..Default::default() + }; + + for object in page.objects { + stats.scanned = stats.scanned.saturating_add(1); + let job_id = match manual_transition_job_id_from_record_object_name(&object.name) { + Ok(job_id) => job_id, + Err(err) => { + stats.failed = stats.failed.saturating_add(1); + warn!( + event = EVENT_LIFECYCLE_WORKER_STATE, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_LIFECYCLE, + object = %object.name, + error = %err, + state = "manual_transition_recovery_skipped", + "Manual transition recovery skipped a corrupt job record path" + ); + continue; + } + }; + match recover_manual_transition_job(api.clone(), job_id).await { + Ok(ManualTransitionJobRecoveryOutcome::Resumed) => stats.resumed = stats.resumed.saturating_add(1), + Ok(ManualTransitionJobRecoveryOutcome::Cancelled) => stats.cancelled = stats.cancelled.saturating_add(1), + Ok(ManualTransitionJobRecoveryOutcome::Skipped) => stats.skipped = stats.skipped.saturating_add(1), + Err(err) => { + stats.failed = stats.failed.saturating_add(1); + warn!( + event = EVENT_LIFECYCLE_WORKER_STATE, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_LIFECYCLE, + job_id = %job_id, + error = %err, + state = "manual_transition_recovery_failed", + "Manual transition recovery failed a job" + ); + } + } + } + + Ok(stats) +} + +async fn recover_manual_transition_job(api: Arc, job_id: Uuid) -> Result { + let (mut record, etag) = match load_manual_transition_job_record_with_etag(api.clone(), job_id).await { + Ok(record) => record, + Err(Error::ConfigNotFound) => return Ok(ManualTransitionJobRecoveryOutcome::Skipped), + Err(err) => return Err(err), + }; + if record.is_terminal() || !manual_transition_job_lease_expired(&record) { + return Ok(ManualTransitionJobRecoveryOutcome::Skipped); + } + + if record.cancel_requested { + let mut report = record.report.clone(); + report.cancelled = true; + record.complete(report, manual_transition_queue_snapshot()); + return match save_manual_transition_job_record_if_current(api.clone(), &record, &etag).await { + Ok(()) => { + release_manual_transition_recovery_admission(api, &record).await; + Ok(ManualTransitionJobRecoveryOutcome::Cancelled) + } + Err(Error::PreconditionFailed) => Ok(ManualTransitionJobRecoveryOutcome::Skipped), + Err(err) => Err(err), + }; + } + + record.claim_recovery_lease(manual_transition_recovery_owner_id(), manual_transition_queue_snapshot()); + let recovery_lease_id = record.lease_id; + match save_manual_transition_job_record_if_current(api.clone(), &record, &etag).await { + Ok(()) => {} + Err(Error::PreconditionFailed) => return Ok(ManualTransitionJobRecoveryOutcome::Skipped), + Err(err) => return Err(err), + } + + match claim_manual_transition_scope_admission(api.clone(), &ManualTransitionScopeAdmission::from_job(&record)).await { + Ok(ManualTransitionScopeAdmissionClaim::Claimed) => {} + Ok(ManualTransitionScopeAdmissionClaim::Conflict(_)) => { + abandon_manual_transition_recovery_lease(api, job_id, recovery_lease_id).await?; + return Ok(ManualTransitionJobRecoveryOutcome::Skipped); + } + Err(err) => { + abandon_manual_transition_recovery_lease(api, job_id, recovery_lease_id).await?; + return Err(err); + } + } + + let mut options = record.resume_options(); + options.cancel_check = Some(manual_transition_recovery_cancel_check(api.clone(), job_id)); + options.progress_sink = Some(manual_transition_recovery_progress_sink(api.clone(), job_id)); + let result = enqueue_transition_for_existing_objects_scoped(api.clone(), &record.bucket, options).await; + let final_record = finalize_recovered_manual_transition_job(api.clone(), job_id, result).await?; + release_manual_transition_recovery_admission(api, &final_record).await; + Ok(ManualTransitionJobRecoveryOutcome::Resumed) +} + +fn manual_transition_recovery_owner_id() -> &'static str { + "ecstore-manual-transition-recovery" +} + +fn manual_transition_recovery_cancel_check(api: Arc, job_id: Uuid) -> ManualTransitionCancelCheck { + Arc::new(move || { + let api = api.clone(); + Box::pin(async move { + match load_manual_transition_job_record(api, job_id).await { + Ok(record) => record.cancel_requested || record.is_terminal(), + Err(_) => true, + } + }) + }) +} + +fn manual_transition_recovery_progress_sink(api: Arc, job_id: Uuid) -> ManualTransitionProgressSink { + Arc::new(move |report| { + let api = api.clone(); + Box::pin(async move { + persist_manual_transition_job_progress(api, job_id, &report, manual_transition_queue_snapshot()) + .await + .map(|_| ()) + }) + }) +} + +async fn finalize_recovered_manual_transition_job( + api: Arc, + job_id: Uuid, + result: Result, +) -> Result { + for _ in 0..4 { + let (mut record, etag) = load_manual_transition_job_record_with_etag(api.clone(), job_id).await?; + if record.is_terminal() { + return Ok(record); + } + match &result { + Ok(report) => record.complete(report.clone(), manual_transition_queue_snapshot()), + Err(err) => record.fail(format!("manual transition recovery failed: {err}")), + } + match save_manual_transition_job_record_if_current(api.clone(), &record, &etag).await { + Ok(()) => return Ok(record), + Err(Error::PreconditionFailed) => continue, + Err(err) => return Err(err), + } + } + Err(Error::PreconditionFailed) +} + +async fn release_manual_transition_recovery_admission(api: Arc, record: &ManualTransitionJobRecord) { + if let Err(err) = + delete_manual_transition_scope_admission_if_current(api, &record.scope_key, record.job_id, record.lease_id).await + { + debug!( + event = EVENT_LIFECYCLE_WORKER_STATE, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_LIFECYCLE, + job_id = %record.job_id, + error = %err, + state = "manual_transition_recovery_admission_release_failed", + "Manual transition recovery failed to release admission" + ); + } +} + +async fn abandon_manual_transition_recovery_lease(api: Arc, job_id: Uuid, lease_id: Uuid) -> Result<(), Error> { + for _ in 0..4 { + let (mut record, etag) = match load_manual_transition_job_record_with_etag(api.clone(), job_id).await { + Ok(record) => record, + Err(Error::ConfigNotFound) => return Ok(()), + Err(err) => return Err(err), + }; + if record.lease_id != lease_id || record.is_terminal() { + return Ok(()); + } + record.abandon_recovery_lease(lease_id); + match save_manual_transition_job_record_if_current(api.clone(), &record, &etag).await { + Ok(()) => return Ok(()), + Err(Error::PreconditionFailed) => continue, + Err(err) => return Err(err), + } + } + Ok(()) } fn spawn_tier_free_version_recovery_once(api: Arc, started: &OnceLock<()>) -> Option> { @@ -2434,18 +2726,63 @@ pub async fn enqueue_immediate_expiry(oi: &ObjectInfo, src: LcEventSrc) { } } -#[derive(Debug, Clone, Default, PartialEq, Eq, serde::Serialize)] +pub type ManualTransitionCancelCheck = Arc Pin + Send>> + Send + Sync + 'static>; +pub type ManualTransitionProgressSink = + Arc Pin> + Send>> + Send + Sync + 'static>; + +#[derive(Clone, Default, serde::Serialize, serde::Deserialize)] pub struct ManualTransitionRunOptions { pub prefix: String, pub marker: Option, pub version_marker: Option, + pub continuation_token: Option, pub tier: Option, pub dry_run: bool, pub max_objects: Option, pub max_duration: Option, + #[serde(skip)] + pub cancel_token: Option, + #[serde(skip)] + pub cancel_check: Option, + #[serde(skip)] + pub progress_sink: Option, } +impl std::fmt::Debug for ManualTransitionRunOptions { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.debug_struct("ManualTransitionRunOptions") + .field("prefix", &self.prefix) + .field("marker", &self.marker) + .field("version_marker", &self.version_marker) + .field("continuation_token", &self.continuation_token) + .field("tier", &self.tier) + .field("dry_run", &self.dry_run) + .field("max_objects", &self.max_objects) + .field("max_duration", &self.max_duration) + .field("cancel_token", &self.cancel_token.is_some()) + .field("cancel_check", &self.cancel_check.is_some()) + .field("progress_sink", &self.progress_sink.is_some()) + .finish() + } +} + +impl PartialEq for ManualTransitionRunOptions { + fn eq(&self, other: &Self) -> bool { + self.prefix == other.prefix + && self.marker == other.marker + && self.version_marker == other.version_marker + && self.continuation_token == other.continuation_token + && self.tier == other.tier + && self.dry_run == other.dry_run + && self.max_objects == other.max_objects + && self.max_duration == other.max_duration + } +} + +impl Eq for ManualTransitionRunOptions {} + #[derive(Debug, Clone, Default, PartialEq, Eq, serde::Serialize, serde::Deserialize)] +#[serde(deny_unknown_fields)] pub struct ManualTransitionRunReport { pub bucket: String, pub prefix: String, @@ -2466,14 +2803,30 @@ pub struct ManualTransitionRunReport { pub skipped_queue_full: u64, pub skipped_queue_closed: u64, pub skipped_queue_timeout: u64, + pub tier_failure: u64, pub truncated_by_limit: bool, pub truncated_by_duration: bool, - #[serde(skip_serializing)] + pub cancelled: bool, + #[serde(skip_serializing_if = "Option::is_none")] + pub continuation_token: Option, + #[serde(skip)] pub next_marker: Option, - #[serde(skip_serializing)] + #[serde(skip)] pub next_version_idmarker: Option, } +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, serde::Serialize, serde::Deserialize)] +pub struct ManualTransitionQueueSnapshot { + pub queue_capacity: u64, + pub queued: u64, + pub active: u64, + pub workers: u64, + pub queue_full: u64, + pub queue_send_timeout: u64, + pub compensation_pending: u64, + pub compensation_running: u64, +} + #[derive(Debug, Clone, PartialEq, Eq)] pub struct ManualTransitionRunExecution { pub report: ManualTransitionRunReport, @@ -2496,7 +2849,7 @@ impl ManualTransitionRunReport { } pub fn was_truncated(&self) -> bool { - self.truncated_by_limit || self.truncated_by_duration + self.truncated_by_limit || self.truncated_by_duration || self.cancelled } fn record_enqueue_outcome(&mut self, outcome: TransitionEnqueueOutcome) { @@ -2518,114 +2871,178 @@ impl ManualTransitionRunReport { } } +#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)] +#[serde(deny_unknown_fields)] +struct ManualTransitionContinuationToken { + marker: Option, + version_marker: Option, +} + +fn encode_manual_transition_continuation_token(marker: Option, version_marker: Option) -> Option { + if marker.is_none() && version_marker.is_none() { + return None; + } + let token = ManualTransitionContinuationToken { marker, version_marker }; + serde_json::to_vec(&token) + .ok() + .map(|encoded| base64_simd::URL_SAFE_NO_PAD.encode_to_string(&encoded)) +} + +pub fn decode_manual_transition_continuation_token(token: &str) -> Result<(Option, Option), Error> { + if token.trim().is_empty() { + return Err(Error::other("manual transition continuation token is empty")); + } + let decoded = base64_simd::URL_SAFE_NO_PAD + .decode_to_vec(token.as_bytes()) + .map_err(|err| Error::other(format!("decode manual transition continuation token failed: {err}")))?; + let token: ManualTransitionContinuationToken = serde_json::from_slice(&decoded) + .map_err(|err| Error::other(format!("parse manual transition continuation token failed: {err}")))?; + Ok(( + token.marker.filter(|marker| !marker.is_empty()), + token.version_marker.filter(|marker| !marker.is_empty()), + )) +} + +async fn persist_manual_transition_progress( + options: &ManualTransitionRunOptions, + report: &ManualTransitionRunReport, +) -> Result<(), Error> { + if let Some(progress_sink) = &options.progress_sink { + progress_sink(report.clone()).await?; + } + Ok(()) +} + +async fn persist_manual_transition_page_checkpoint( + options: &ManualTransitionRunOptions, + report: &ManualTransitionRunReport, + marker: Option, + version_marker: Option, +) -> Result<(), Error> { + let mut checkpoint = report.clone(); + checkpoint.next_marker.clone_from(&marker); + checkpoint.next_version_idmarker.clone_from(&version_marker); + checkpoint.continuation_token = encode_manual_transition_continuation_token(marker, version_marker); + persist_manual_transition_progress(options, &checkpoint).await +} + pub async fn enqueue_transition_for_existing_objects(api: Arc, bucket: &str) -> Result<(), Error> { let _ = enqueue_transition_for_existing_objects_scoped(api, bucket, ManualTransitionRunOptions::default()).await?; Ok(()) } +pub fn manual_transition_queue_snapshot() -> ManualTransitionQueueSnapshot { + runtime_sources::transition_state_handle().manual_transition_queue_snapshot() +} + +pub async fn enqueue_transition_for_existing_objects_scoped_with_cancel( + api: Arc, + bucket: &str, + mut options: ManualTransitionRunOptions, + cancel_token: Option, +) -> Result { + if cancel_token.is_some() { + options.cancel_token = cancel_token; + } + let report = enqueue_transition_for_existing_objects_scoped(api, bucket, options).await?; + Ok(ManualTransitionRunExecution { + cancelled: report.cancelled, + report, + }) +} + pub async fn enqueue_transition_for_existing_objects_scoped( api: Arc, bucket: &str, options: ManualTransitionRunOptions, ) -> Result { - Ok(enqueue_transition_for_existing_objects_scoped_with_cancel(api, bucket, options, None) - .await? - .report) -} - -pub async fn enqueue_transition_for_existing_objects_scoped_with_cancel( - api: Arc, - bucket: &str, - options: ManualTransitionRunOptions, - cancel_token: Option, -) -> Result { const LIST_PAGE_SIZE: i32 = 1000; let mut report = ManualTransitionRunReport::new(bucket, &options); + let (mut marker, mut version_marker) = if let Some(token) = options.continuation_token.as_deref() { + decode_manual_transition_continuation_token(token)? + } else { + (options.marker.clone(), options.version_marker.clone()) + }; let Some(lc) = runtime_sources::bucket_lifecycle_config(bucket).await else { - return Ok(ManualTransitionRunExecution { - report, - cancelled: false, - }); + return Ok(report); }; report.lifecycle_config_found = true; - let mut marker = options.marker.clone(); - let mut version_marker = options.version_marker.clone(); let mut previous_marker = marker.clone(); let mut previous_version_marker = version_marker.clone(); let src = LcEventSrc::Scanner; let deadline = options.max_duration.map(|duration| tokio::time::Instant::now() + duration); loop { - if manual_transition_cancelled(cancel_token.as_ref()) { - return Ok(ManualTransitionRunExecution { report, cancelled: true }); - } - let page = api .clone() .list_object_versions(bucket, &options.prefix, marker.clone(), version_marker.clone(), None, LIST_PAGE_SIZE) .await?; for (index, object) in page.objects.iter().enumerate() { - if manual_transition_cancelled(cancel_token.as_ref()) { + if manual_transition_cancel_requested(&options).await { + report.cancelled = true; report.next_marker.clone_from(&previous_marker); report.next_version_idmarker.clone_from(&previous_version_marker); - return Ok(ManualTransitionRunExecution { report, cancelled: true }); + report.continuation_token = + encode_manual_transition_continuation_token(report.next_marker.clone(), report.next_version_idmarker.clone()); + persist_manual_transition_progress(&options, &report).await?; + return Ok(report); } if manual_transition_duration_elapsed(deadline) { report.truncated_by_duration = true; report.next_marker.clone_from(&previous_marker); report.next_version_idmarker.clone_from(&previous_version_marker); - return Ok(ManualTransitionRunExecution { - report, - cancelled: false, - }); + report.continuation_token = + encode_manual_transition_continuation_token(report.next_marker.clone(), report.next_version_idmarker.clone()); + persist_manual_transition_progress(&options, &report).await?; + return Ok(report); } report.scanned = report.scanned.saturating_add(1); enqueue_transition_with_lifecycle_report(object, &lc, &src, &options, &mut report).await; if report.has_partial_enqueue() { report.next_marker.clone_from(&previous_marker); report.next_version_idmarker.clone_from(&previous_version_marker); - return Ok(ManualTransitionRunExecution { - report, - cancelled: false, - }); + report.continuation_token = + encode_manual_transition_continuation_token(report.next_marker.clone(), report.next_version_idmarker.clone()); + persist_manual_transition_progress(&options, &report).await?; + return Ok(report); } if options.max_objects.is_some_and(|max_objects| report.scanned >= max_objects) { if manual_transition_has_more_after_limit(index, page.objects.len(), page.is_truncated) { report.truncated_by_limit = true; report.next_marker = Some(object.name.clone()); report.next_version_idmarker = Some(manual_transition_version_marker(object)); + report.continuation_token = encode_manual_transition_continuation_token( + report.next_marker.clone(), + report.next_version_idmarker.clone(), + ); } - return Ok(ManualTransitionRunExecution { - report, - cancelled: false, - }); + persist_manual_transition_progress(&options, &report).await?; + return Ok(report); } previous_marker = Some(object.name.clone()); previous_version_marker = Some(manual_transition_version_marker(object)); } if !page.is_truncated { - return Ok(ManualTransitionRunExecution { - report, - cancelled: false, - }); + return Ok(report); } if manual_transition_duration_elapsed(deadline) { report.truncated_by_duration = true; report.next_marker.clone_from(&previous_marker); report.next_version_idmarker.clone_from(&previous_version_marker); - return Ok(ManualTransitionRunExecution { - report, - cancelled: false, - }); + report.continuation_token = + encode_manual_transition_continuation_token(report.next_marker.clone(), report.next_version_idmarker.clone()); + persist_manual_transition_progress(&options, &report).await?; + return Ok(report); } marker = page.next_marker; version_marker = page.next_version_idmarker; previous_marker = marker.clone(); previous_version_marker = version_marker.clone(); + persist_manual_transition_page_checkpoint(&options, &report, marker.clone(), version_marker.clone()).await?; } } @@ -2823,8 +3240,14 @@ fn manual_transition_duration_elapsed(deadline: Option) -> deadline.is_some_and(|deadline| tokio::time::Instant::now() >= deadline) } -fn manual_transition_cancelled(cancel_token: Option<&CancellationToken>) -> bool { - cancel_token.is_some_and(CancellationToken::is_cancelled) +async fn manual_transition_cancel_requested(options: &ManualTransitionRunOptions) -> bool { + if options.cancel_token.as_ref().is_some_and(|token| token.is_cancelled()) { + return true; + } + match options.cancel_check.as_ref() { + Some(cancel_check) => cancel_check().await, + None => false, + } } fn manual_transition_version_marker(oi: &ObjectInfo) -> String { @@ -3742,9 +4165,10 @@ mod tests { enqueue_recovered_free_version_with_state, enqueue_transition_with_lifecycle, enqueue_transition_with_lifecycle_report, eval_action_from_lifecycle, jitter_tier_free_version_recovery_delay, lifecycle_action_blocked_by_replication, lifecycle_delete_all_versions_replication_scan, lifecycle_deleted_object, lifecycle_replication_blocks_action, - lifecycle_rule_has_date_expiration, lifecycle_version_purge_state_from_completed_targets, manual_transition_cancelled, + lifecycle_rule_has_date_expiration, lifecycle_version_purge_state_from_completed_targets, manual_transition_duration_elapsed, manual_transition_has_more_after_limit, manual_transition_version_marker, - mark_delete_opts_skip_decommissioned_on_remote_success, merge_stale_multipart_candidate, replication_state_for_delete, + mark_delete_opts_skip_decommissioned_on_remote_success, merge_stale_multipart_candidate, + persist_manual_transition_page_checkpoint, recover_manual_transition_jobs_once, replication_state_for_delete, resolve_transition_queue_capacity, resolve_transition_queue_send_timeout, resolve_transition_worker_count, resolve_transition_workers_absolute_max, run_tier_free_version_recovery_loop, select_restore_s3_location, set_recovered_free_version_enqueue_observer, should_defer_date_expiry_for_recent_config_update, @@ -3754,6 +4178,14 @@ mod tests { #[cfg(feature = "test-util")] use super::{delete_free_version_remote_object_then, encode_dir_object, get_transitioned_object_reader_with_tier_manager}; use crate::bucket::lifecycle::bucket_lifecycle_audit::LcEventSrc; + use crate::bucket::lifecycle::bucket_lifecycle_ops::{ + decode_manual_transition_continuation_token, encode_manual_transition_continuation_token, + }; + use crate::bucket::lifecycle::manual_transition_job::{ + ManualTransitionJobRecord, ManualTransitionJobState, ManualTransitionScopeAdmission, load_manual_transition_job_record, + load_manual_transition_scope_admission, save_manual_transition_job_record, + save_manual_transition_scope_admission_if_absent, + }; use crate::bucket::lifecycle::replication_sink::{ ReplicateDecision, ReplicateTargetDecision, ReplicationStatusType, VersionPurgeStatusType, }; @@ -3769,7 +4201,7 @@ mod tests { use crate::client::transition_api::ReaderImpl; use crate::disk::endpoint::Endpoint; use crate::disk::{RUSTFS_META_MULTIPART_BUCKET, STORAGE_FORMAT_FILE}; - use crate::error::is_err_invalid_upload_id; + use crate::error::{Error, is_err_invalid_upload_id}; use crate::layout::endpoints::{EndpointServerPools, Endpoints, PoolEndpoints}; use crate::object_api::{ObjectInfo, ObjectOptions, PutObjReader}; #[cfg(feature = "test-util")] @@ -6495,15 +6927,141 @@ mod tests { } #[test] - fn manual_transition_cancel_token_detects_cancelled_state() { - let cancel_token = CancellationToken::new(); + fn manual_transition_continuation_token_round_trips_resume_cursor() { + let token = encode_manual_transition_continuation_token(Some("logs/a".to_string()), Some("null".to_string())) + .expect("non-empty cursor should encode"); - assert!(!manual_transition_cancelled(None)); - assert!(!manual_transition_cancelled(Some(&cancel_token))); + let (marker, version_marker) = + decode_manual_transition_continuation_token(&token).expect("continuation token should decode"); - cancel_token.cancel(); + assert_eq!(marker.as_deref(), Some("logs/a")); + assert_eq!(version_marker.as_deref(), Some("null")); + } - assert!(manual_transition_cancelled(Some(&cancel_token))); + #[test] + fn manual_transition_continuation_token_rejects_malformed_input() { + let err = + decode_manual_transition_continuation_token("not-base64").expect_err("malformed continuation token must fail closed"); + + assert!(err.to_string().contains("decode manual transition continuation token failed")); + } + + #[test] + fn manual_transition_report_serializes_public_continuation_only() { + let report = ManualTransitionRunReport { + continuation_token: Some("opaque".to_string()), + next_marker: Some("logs/a".to_string()), + next_version_idmarker: Some("null".to_string()), + ..Default::default() + }; + + let value = serde_json::to_value(report).expect("report should serialize"); + + assert_eq!(value.get("continuation_token").and_then(|value| value.as_str()), Some("opaque")); + assert!(value.get("next_marker").is_none()); + assert!(value.get("next_version_idmarker").is_none()); + } + + #[tokio::test] + async fn manual_transition_page_checkpoint_persists_resume_cursor() { + let observed = Arc::new(StdMutex::new(Vec::new())); + let sink_observed = Arc::clone(&observed); + let options = ManualTransitionRunOptions { + progress_sink: Some(Arc::new(move |report| { + let sink_observed = Arc::clone(&sink_observed); + Box::pin(async move { + sink_observed.lock().expect("observed reports mutex poisoned").push(report); + Ok(()) + }) + })), + ..Default::default() + }; + let report = ManualTransitionRunReport { + bucket: "bucket".to_string(), + prefix: "logs/".to_string(), + scanned: 1000, + ..Default::default() + }; + + persist_manual_transition_page_checkpoint(&options, &report, Some("logs/page-end".to_string()), Some("null".to_string())) + .await + .expect("page checkpoint should persist"); + + assert!(report.continuation_token.is_none()); + let observed = observed.lock().expect("observed reports mutex poisoned"); + assert_eq!(observed.len(), 1); + assert_eq!(observed[0].scanned, 1000); + assert_eq!(observed[0].next_marker.as_deref(), Some("logs/page-end")); + assert_eq!(observed[0].next_version_idmarker.as_deref(), Some("null")); + let token = observed[0] + .continuation_token + .as_deref() + .expect("checkpoint should carry resume token"); + let (marker, version_marker) = + decode_manual_transition_continuation_token(token).expect("checkpoint token should decode"); + assert_eq!(marker.as_deref(), Some("logs/page-end")); + assert_eq!(version_marker.as_deref(), Some("null")); + } + + #[tokio::test] + #[serial] + async fn manual_transition_recovery_replays_expired_running_record() { + let (_paths, ecstore) = setup_test_env().await; + let job_id = Uuid::new_v4(); + let options = ManualTransitionRunOptions { + prefix: "logs/".to_string(), + continuation_token: Some( + encode_manual_transition_continuation_token(Some("logs/page-end".to_string()), Some("null".to_string())) + .expect("resume token should encode"), + ), + max_objects: Some(7), + ..Default::default() + }; + let mut record = ManualTransitionJobRecord::new(job_id, "manual-recovery-bucket", &options, "old-owner"); + record.report.continuation_token.clone_from(&options.continuation_token); + record.lease_expires_at_unix_nanos = 0; + save_manual_transition_job_record(ecstore.clone(), &record) + .await + .expect("expired job record should save"); + save_manual_transition_scope_admission_if_absent(ecstore.clone(), &ManualTransitionScopeAdmission::from_job(&record)) + .await + .expect("expired scope admission should save"); + + let stats = recover_manual_transition_jobs_once(ecstore.clone(), 10, None) + .await + .expect("manual transition recovery should run"); + + assert_eq!(stats.scanned, 1); + assert_eq!(stats.resumed, 1); + assert_eq!(stats.failed, 0); + let recovered = load_manual_transition_job_record(ecstore.clone(), job_id) + .await + .expect("recovered job should load"); + assert_eq!(recovered.state, ManualTransitionJobState::Completed); + assert_eq!(recovered.max_objects, Some(7)); + assert_eq!(recovered.report.bucket, "manual-recovery-bucket"); + assert_eq!(recovered.report.prefix, "logs/"); + assert!(!recovered.report.lifecycle_config_found); + assert!(recovered.report.continuation_token.is_none()); + assert!( + matches!( + load_manual_transition_scope_admission(ecstore, &recovered.scope_key).await, + Err(Error::ConfigNotFound) + ), + "completed recovery must release the scope admission" + ); + } + + #[test] + fn manual_transition_cancelled_report_is_partial_and_resumable() { + let report = ManualTransitionRunReport { + cancelled: true, + continuation_token: Some("opaque".to_string()), + ..Default::default() + }; + + assert!(report.was_truncated()); + assert!(!report.has_partial_enqueue()); } #[tokio::test] diff --git a/crates/ecstore/src/bucket/lifecycle/config_boundary.rs b/crates/ecstore/src/bucket/lifecycle/config_boundary.rs index 4c89162bb..9f4deab7d 100644 --- a/crates/ecstore/src/bucket/lifecycle/config_boundary.rs +++ b/crates/ecstore/src/bucket/lifecycle/config_boundary.rs @@ -40,6 +40,21 @@ where com::read_config(api, file).await } +pub(crate) async fn read_config_with_metadata(api: Arc, file: &str, opts: &ObjectOptions) -> Result<(Vec, ObjectInfo)> +where + S: ObjectIO< + Error = Error, + RangeSpec = HTTPRangeSpec, + HeaderMap = HeaderMap, + ObjectOptions = ObjectOptions, + ObjectInfo = ObjectInfo, + GetObjectReader = GetObjectReader, + PutObjectReader = PutObjReader, + >, +{ + com::read_config_with_metadata(api, file, opts).await +} + pub(crate) async fn save_config(api: Arc, file: &str, data: Vec) -> Result<()> where S: ObjectIO< @@ -55,6 +70,21 @@ where com::save_config(api, file, data).await } +pub(crate) async fn save_config_with_opts(api: Arc, file: &str, data: Vec, opts: &ObjectOptions) -> Result<()> +where + S: ObjectIO< + Error = Error, + RangeSpec = HTTPRangeSpec, + HeaderMap = HeaderMap, + ObjectOptions = ObjectOptions, + ObjectInfo = ObjectInfo, + GetObjectReader = GetObjectReader, + PutObjectReader = PutObjReader, + >, +{ + com::save_config_with_opts(api, file, data, opts).await +} + pub(crate) async fn delete_config(api: Arc, file: &str) -> Result<()> where S: ObjectOperations< diff --git a/crates/ecstore/src/bucket/lifecycle/manual_transition_job.rs b/crates/ecstore/src/bucket/lifecycle/manual_transition_job.rs new file mode 100644 index 000000000..45917ef06 --- /dev/null +++ b/crates/ecstore/src/bucket/lifecycle/manual_transition_job.rs @@ -0,0 +1,892 @@ +// Copyright 2024 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use std::sync::Arc; + +use rustfs_utils::crypto::{hex_sha256, is_sha256_checksum}; +use serde::{Deserialize, Serialize}; +use time::OffsetDateTime; +use uuid::Uuid; + +use crate::bucket::lifecycle::bucket_lifecycle_ops::{ + ManualTransitionQueueSnapshot, ManualTransitionRunOptions, ManualTransitionRunReport, +}; +use crate::bucket::lifecycle::config_boundary; +use crate::error::{Error, Result as EcstoreResult}; +use crate::object_api::ObjectOptions; +use crate::storage_api_contracts::object::HTTPPreconditions; +use crate::store::ECStore; + +pub const MANUAL_TRANSITION_JOB_SCHEMA: &str = "rustfs-manual-transition-job-v1"; +pub const MANUAL_TRANSITION_JOB_RECORD_PREFIX: &str = "ilm/manual-transition/jobs"; +pub const MANUAL_TRANSITION_SCOPE_RECORD_PREFIX: &str = "ilm/manual-transition/scopes"; +pub const MAX_MANUAL_TRANSITION_JOB_RECORD_SIZE: usize = 64 * 1024; +const MANUAL_TRANSITION_JOB_LEASE_SECONDS: i128 = 60; + +#[derive(Debug, thiserror::Error)] +pub enum ManualTransitionJobError { + #[error("manual transition job is corrupt: {0}")] + Corrupt(&'static str), + #[error("manual transition job schema is unsupported: {0}")] + UnsupportedSchema(String), + #[error("manual transition job checksum mismatch")] + ChecksumMismatch, + #[error("manual transition job json error: {0}")] + Json(#[from] serde_json::Error), +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum ManualTransitionJobState { + Running, + Completed, + Partial, + Failed, + Cancelled, + Unknown, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct ManualTransitionJobRecord { + pub job_id: Uuid, + pub scope_key: String, + pub bucket: String, + pub prefix: String, + pub tier: Option, + pub dry_run: bool, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub max_objects: Option, + pub owner_id: String, + pub lease_id: Uuid, + pub lease_expires_at_unix_nanos: i128, + pub state: ManualTransitionJobState, + pub cancel_requested: bool, + pub created_at_unix_nanos: i128, + pub updated_at_unix_nanos: i128, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub completed_at_unix_nanos: Option, + pub report: ManualTransitionRunReport, + pub queue_snapshot: ManualTransitionQueueSnapshot, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub error: Option, +} + +impl ManualTransitionJobRecord { + pub fn new(job_id: Uuid, bucket: &str, options: &ManualTransitionRunOptions, owner_id: impl Into) -> Self { + let scope_key = manual_transition_scope_key(bucket, options); + let now = OffsetDateTime::now_utc().unix_timestamp_nanos(); + let lease_id = Uuid::new_v4(); + Self { + job_id, + scope_key, + bucket: bucket.to_string(), + prefix: options.prefix.clone(), + tier: options.tier.clone(), + dry_run: options.dry_run, + max_objects: options.max_objects, + owner_id: owner_id.into(), + lease_id, + lease_expires_at_unix_nanos: manual_transition_job_lease_expires_at(now), + state: ManualTransitionJobState::Running, + cancel_requested: false, + created_at_unix_nanos: now, + updated_at_unix_nanos: now, + completed_at_unix_nanos: None, + report: ManualTransitionRunReport { + bucket: bucket.to_string(), + prefix: options.prefix.clone(), + tier: options.tier.clone(), + dry_run: options.dry_run, + ..Default::default() + }, + queue_snapshot: ManualTransitionQueueSnapshot::default(), + error: None, + } + } + + pub fn complete(&mut self, report: ManualTransitionRunReport, queue_snapshot: ManualTransitionQueueSnapshot) { + self.state = if report.cancelled { + ManualTransitionJobState::Cancelled + } else if report.was_truncated() || report.has_partial_enqueue() || report.tier_failure > 0 { + ManualTransitionJobState::Partial + } else { + ManualTransitionJobState::Completed + }; + self.report = report; + self.queue_snapshot = queue_snapshot; + self.error = None; + self.mark_updated_terminal(); + } + + pub fn fail(&mut self, error: impl Into) { + self.state = ManualTransitionJobState::Failed; + self.report.tier_failure = self.report.tier_failure.saturating_add(1); + self.error = Some(error.into()); + self.mark_updated_terminal(); + } + + pub fn mark_cancel_requested(&mut self) { + self.cancel_requested = true; + self.updated_at_unix_nanos = OffsetDateTime::now_utc().unix_timestamp_nanos(); + } + + pub fn claim_recovery_lease(&mut self, owner_id: impl Into, queue_snapshot: ManualTransitionQueueSnapshot) { + if self.state == ManualTransitionJobState::Running { + self.owner_id = owner_id.into(); + self.lease_id = Uuid::new_v4(); + self.renew_lease(queue_snapshot); + } + } + + pub fn abandon_recovery_lease(&mut self, lease_id: Uuid) { + if self.state == ManualTransitionJobState::Running && self.lease_id == lease_id { + self.lease_expires_at_unix_nanos = 0; + self.updated_at_unix_nanos = OffsetDateTime::now_utc().unix_timestamp_nanos(); + } + } + + pub fn resume_options(&self) -> ManualTransitionRunOptions { + ManualTransitionRunOptions { + prefix: self.prefix.clone(), + continuation_token: self.report.continuation_token.clone(), + tier: self.tier.clone(), + dry_run: self.dry_run, + max_objects: self.max_objects, + ..Default::default() + } + } + + pub fn renew_lease(&mut self, queue_snapshot: ManualTransitionQueueSnapshot) { + let now = OffsetDateTime::now_utc().unix_timestamp_nanos(); + self.updated_at_unix_nanos = now; + self.lease_expires_at_unix_nanos = manual_transition_job_lease_expires_at(now); + self.queue_snapshot = queue_snapshot; + } + + pub fn update_running_progress(&mut self, report: ManualTransitionRunReport, queue_snapshot: ManualTransitionQueueSnapshot) { + if self.state == ManualTransitionJobState::Running { + self.report = report; + self.renew_lease(queue_snapshot); + } + } + + pub fn mark_unknown_if_unowned(&mut self) { + if self.state == ManualTransitionJobState::Running { + self.state = ManualTransitionJobState::Unknown; + self.error = Some("manual transition job outcome is unknown after restart or owner loss".to_string()); + self.mark_updated_terminal(); + } + } + + pub fn is_terminal(&self) -> bool { + matches!( + self.state, + ManualTransitionJobState::Completed + | ManualTransitionJobState::Partial + | ManualTransitionJobState::Failed + | ManualTransitionJobState::Cancelled + | ManualTransitionJobState::Unknown + ) + } + + fn mark_updated_terminal(&mut self) { + let now = OffsetDateTime::now_utc().unix_timestamp_nanos(); + self.updated_at_unix_nanos = now; + self.completed_at_unix_nanos = Some(now); + } + + pub fn encode(&self) -> Result, ManualTransitionJobError> { + self.validate()?; + let job_bytes = serde_json::to_vec(self)?; + let content_sha256 = hex_sha256(&job_bytes, ToOwned::to_owned); + let persisted = PersistedManualTransitionJobRecord { + schema: MANUAL_TRANSITION_JOB_SCHEMA.to_string(), + content_sha256, + job: self.clone(), + }; + let encoded = serde_json::to_vec(&persisted)?; + if encoded.len() > MAX_MANUAL_TRANSITION_JOB_RECORD_SIZE { + return Err(ManualTransitionJobError::Corrupt("encoded job exceeds maximum size")); + } + Ok(encoded) + } + + pub fn decode(expected_job_id: Uuid, data: &[u8]) -> Result { + if data.len() > MAX_MANUAL_TRANSITION_JOB_RECORD_SIZE { + return Err(ManualTransitionJobError::Corrupt("encoded job exceeds maximum size")); + } + let persisted: PersistedManualTransitionJobRecord = serde_json::from_slice(data)?; + if persisted.schema != MANUAL_TRANSITION_JOB_SCHEMA { + return Err(ManualTransitionJobError::UnsupportedSchema(persisted.schema)); + } + if !is_sha256_checksum(&persisted.content_sha256) { + return Err(ManualTransitionJobError::Corrupt("content checksum is not a sha256 checksum")); + } + let job_bytes = serde_json::to_vec(&persisted.job)?; + let actual_checksum = hex_sha256(&job_bytes, ToOwned::to_owned); + if persisted.content_sha256 != actual_checksum { + return Err(ManualTransitionJobError::ChecksumMismatch); + } + if persisted.job.job_id != expected_job_id { + return Err(ManualTransitionJobError::Corrupt("job_id does not match record key")); + } + persisted.job.validate()?; + Ok(persisted.job) + } + + fn validate(&self) -> Result<(), ManualTransitionJobError> { + if self.job_id.is_nil() { + return Err(ManualTransitionJobError::Corrupt("job_id is nil")); + } + if self.lease_id.is_nil() { + return Err(ManualTransitionJobError::Corrupt("lease_id is nil")); + } + if self.scope_key.is_empty() { + return Err(ManualTransitionJobError::Corrupt("scope_key is empty")); + } + if self.bucket.is_empty() { + return Err(ManualTransitionJobError::Corrupt("bucket is empty")); + } + if self.owner_id.trim().is_empty() { + return Err(ManualTransitionJobError::Corrupt("owner_id is empty")); + } + if self.completed_at_unix_nanos.is_some() && !self.is_terminal() { + return Err(ManualTransitionJobError::Corrupt("non-terminal job has completed timestamp")); + } + if self.state == ManualTransitionJobState::Cancelled && !self.cancel_requested { + return Err(ManualTransitionJobError::Corrupt("cancelled job is missing cancel request")); + } + Ok(()) + } +} + +fn manual_transition_job_lease_expires_at(now_unix_nanos: i128) -> i128 { + now_unix_nanos.saturating_add(MANUAL_TRANSITION_JOB_LEASE_SECONDS.saturating_mul(1_000_000_000)) +} + +pub fn manual_transition_job_lease_expired(record: &ManualTransitionJobRecord) -> bool { + OffsetDateTime::now_utc().unix_timestamp_nanos() > record.lease_expires_at_unix_nanos +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +struct PersistedManualTransitionJobRecord { + schema: String, + content_sha256: String, + job: ManualTransitionJobRecord, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct ManualTransitionScopeAdmission { + pub schema: String, + pub scope_key: String, + pub job_id: Uuid, + pub lease_id: Uuid, + pub owner_id: String, + pub bucket: String, + pub prefix: String, + pub tier: Option, + pub dry_run: bool, + pub lease_expires_at_unix_nanos: i128, + pub updated_at_unix_nanos: i128, +} + +impl ManualTransitionScopeAdmission { + pub fn from_job(record: &ManualTransitionJobRecord) -> Self { + Self { + schema: MANUAL_TRANSITION_JOB_SCHEMA.to_string(), + scope_key: record.scope_key.clone(), + job_id: record.job_id, + lease_id: record.lease_id, + owner_id: record.owner_id.clone(), + bucket: record.bucket.clone(), + prefix: record.prefix.clone(), + tier: record.tier.clone(), + dry_run: record.dry_run, + lease_expires_at_unix_nanos: record.lease_expires_at_unix_nanos, + updated_at_unix_nanos: record.updated_at_unix_nanos, + } + } + + pub fn validate(&self) -> Result<(), ManualTransitionJobError> { + if self.schema != MANUAL_TRANSITION_JOB_SCHEMA { + return Err(ManualTransitionJobError::UnsupportedSchema(self.schema.clone())); + } + if self.job_id.is_nil() { + return Err(ManualTransitionJobError::Corrupt("job_id is nil")); + } + if self.lease_id.is_nil() { + return Err(ManualTransitionJobError::Corrupt("lease_id is nil")); + } + if self.scope_key.is_empty() { + return Err(ManualTransitionJobError::Corrupt("scope_key is empty")); + } + if self.bucket.is_empty() { + return Err(ManualTransitionJobError::Corrupt("bucket is empty")); + } + if self.owner_id.trim().is_empty() { + return Err(ManualTransitionJobError::Corrupt("owner_id is empty")); + } + Ok(()) + } +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum ManualTransitionScopeAdmissionClaim { + Claimed, + Conflict(Box), +} + +pub fn manual_transition_scope_key(bucket: &str, options: &ManualTransitionRunOptions) -> String { + let mut scope = String::new(); + scope.push_str(bucket); + scope.push('\0'); + scope.push_str(&options.prefix); + scope.push('\0'); + if let Some(tier) = &options.tier { + scope.push_str(&tier.to_ascii_uppercase()); + } + scope.push('\0'); + scope.push_str(if options.dry_run { "dry_run" } else { "run" }); + hex_sha256(scope.as_bytes(), ToOwned::to_owned) +} + +pub fn manual_transition_job_record_object_name(job_id: Uuid) -> Result { + if job_id.is_nil() { + return Err(ManualTransitionJobError::Corrupt("job_id is nil")); + } + let job_key = job_id.simple().to_string(); + Ok(format!( + "{}/{}/{}/{}.json", + MANUAL_TRANSITION_JOB_RECORD_PREFIX, + &job_key[..2], + &job_key[2..4], + job_key + )) +} + +pub fn manual_transition_job_id_from_record_object_name(object_name: &str) -> Result { + let Some(rest) = object_name.strip_prefix(MANUAL_TRANSITION_JOB_RECORD_PREFIX) else { + return Err(ManualTransitionJobError::Corrupt("job record object prefix is invalid")); + }; + let rest = rest + .strip_prefix('/') + .ok_or(ManualTransitionJobError::Corrupt("job record object path is invalid"))?; + let mut parts = rest.split('/'); + let first = parts + .next() + .ok_or(ManualTransitionJobError::Corrupt("job record first shard is missing"))?; + let second = parts + .next() + .ok_or(ManualTransitionJobError::Corrupt("job record second shard is missing"))?; + let file = parts + .next() + .ok_or(ManualTransitionJobError::Corrupt("job record file is missing"))?; + if parts.next().is_some() { + return Err(ManualTransitionJobError::Corrupt("job record object path has extra components")); + } + let job_key = file + .strip_suffix(".json") + .ok_or(ManualTransitionJobError::Corrupt("job record suffix is invalid"))?; + if job_key.len() != 32 || first.len() != 2 || second.len() != 2 || first != &job_key[..2] || second != &job_key[2..4] { + return Err(ManualTransitionJobError::Corrupt("job record shards do not match job id")); + } + Uuid::parse_str(job_key).map_err(|_| ManualTransitionJobError::Corrupt("job record job id is invalid")) +} + +pub fn manual_transition_scope_record_object_name(scope_key: &str) -> Result { + if scope_key.len() != 64 + || !scope_key + .bytes() + .all(|byte| byte.is_ascii_hexdigit() && !byte.is_ascii_uppercase()) + { + return Err(ManualTransitionJobError::Corrupt("scope_key is not a lowercase sha256 hex digest")); + } + Ok(format!( + "{}/{}/{}/{}.json", + MANUAL_TRANSITION_SCOPE_RECORD_PREFIX, + &scope_key[..2], + &scope_key[2..4], + scope_key + )) +} + +pub async fn save_manual_transition_job_record(api: Arc, job: &ManualTransitionJobRecord) -> EcstoreResult<()> { + let object = manual_transition_job_record_object_name(job.job_id).map_err(manual_transition_job_store_error)?; + let data = job.encode().map_err(manual_transition_job_store_error)?; + config_boundary::save_config(api, &object, data).await +} + +pub async fn load_manual_transition_job_record(api: Arc, job_id: Uuid) -> EcstoreResult { + let (record, _) = load_manual_transition_job_record_with_etag(api, job_id).await?; + Ok(record) +} + +pub async fn load_manual_transition_job_record_with_etag( + api: Arc, + job_id: Uuid, +) -> EcstoreResult<(ManualTransitionJobRecord, String)> { + let object = manual_transition_job_record_object_name(job_id).map_err(manual_transition_job_store_error)?; + let (data, object_info) = config_boundary::read_config_with_metadata(api, &object, &ObjectOptions::default()).await?; + let etag = object_info + .etag + .filter(|etag| !etag.trim().is_empty()) + .ok_or_else(|| Error::other("manual transition job record is missing an ETag"))?; + let record = ManualTransitionJobRecord::decode(job_id, &data).map_err(manual_transition_job_store_error)?; + Ok((record, etag)) +} + +pub async fn save_manual_transition_job_record_if_current( + api: Arc, + job: &ManualTransitionJobRecord, + current_etag: &str, +) -> EcstoreResult<()> { + if current_etag.trim().is_empty() { + return Err(Error::other("manual transition job current ETag is empty")); + } + let object = manual_transition_job_record_object_name(job.job_id).map_err(manual_transition_job_store_error)?; + let data = job.encode().map_err(manual_transition_job_store_error)?; + config_boundary::save_config_with_opts( + api, + &object, + data, + &ObjectOptions { + max_parity: true, + http_preconditions: Some(HTTPPreconditions { + if_match: Some(current_etag.to_string()), + ..Default::default() + }), + ..Default::default() + }, + ) + .await +} + +pub async fn save_manual_transition_scope_admission_if_absent( + api: Arc, + admission: &ManualTransitionScopeAdmission, +) -> EcstoreResult<()> { + admission.validate().map_err(manual_transition_job_store_error)?; + let object = manual_transition_scope_record_object_name(&admission.scope_key).map_err(manual_transition_job_store_error)?; + let data = serde_json::to_vec(admission).map_err(Error::other)?; + config_boundary::save_config_with_opts( + api, + &object, + data, + &ObjectOptions { + max_parity: true, + http_preconditions: Some(HTTPPreconditions { + if_none_match: Some("*".to_string()), + ..Default::default() + }), + ..Default::default() + }, + ) + .await +} + +pub async fn load_manual_transition_scope_admission( + api: Arc, + scope_key: &str, +) -> EcstoreResult { + let (admission, _) = load_manual_transition_scope_admission_with_etag(api, scope_key).await?; + Ok(admission) +} + +pub async fn load_manual_transition_scope_admission_with_etag( + api: Arc, + scope_key: &str, +) -> EcstoreResult<(ManualTransitionScopeAdmission, String)> { + let object = manual_transition_scope_record_object_name(scope_key).map_err(manual_transition_job_store_error)?; + let (data, object_info) = config_boundary::read_config_with_metadata(api, &object, &ObjectOptions::default()).await?; + let etag = object_info + .etag + .filter(|etag| !etag.trim().is_empty()) + .ok_or_else(|| Error::other("manual transition scope admission is missing an ETag"))?; + let admission: ManualTransitionScopeAdmission = serde_json::from_slice(&data).map_err(Error::other)?; + admission.validate().map_err(manual_transition_job_store_error)?; + Ok((admission, etag)) +} + +pub async fn save_manual_transition_scope_admission_if_current( + api: Arc, + admission: &ManualTransitionScopeAdmission, + current_etag: &str, +) -> EcstoreResult<()> { + if current_etag.trim().is_empty() { + return Err(Error::other("manual transition scope admission current ETag is empty")); + } + admission.validate().map_err(manual_transition_job_store_error)?; + let object = manual_transition_scope_record_object_name(&admission.scope_key).map_err(manual_transition_job_store_error)?; + let data = serde_json::to_vec(admission).map_err(Error::other)?; + config_boundary::save_config_with_opts( + api, + &object, + data, + &ObjectOptions { + max_parity: true, + http_preconditions: Some(HTTPPreconditions { + if_match: Some(current_etag.to_string()), + ..Default::default() + }), + ..Default::default() + }, + ) + .await +} + +pub async fn claim_manual_transition_scope_admission( + api: Arc, + admission: &ManualTransitionScopeAdmission, +) -> EcstoreResult { + match save_manual_transition_scope_admission_if_absent(api.clone(), admission).await { + Ok(()) => return Ok(ManualTransitionScopeAdmissionClaim::Claimed), + Err(Error::PreconditionFailed) => {} + Err(err) => return Err(err), + } + + let (active, etag) = load_manual_transition_scope_admission_with_etag(api.clone(), &admission.scope_key).await?; + let scope_lease_expired = manual_transition_scope_admission_lease_expired(&active); + let active_job_reclaimable = if active.job_id == admission.job_id { + scope_lease_expired + } else { + match load_manual_transition_job_record(api.clone(), active.job_id).await { + Ok(active_job) => { + active_job.is_terminal() || (scope_lease_expired && manual_transition_job_lease_expired(&active_job)) + } + Err(Error::ConfigNotFound) => true, + Err(err) => return Err(err), + } + }; + if active_job_reclaimable { + return match save_manual_transition_scope_admission_if_current(api, admission, &etag).await { + Ok(()) => Ok(ManualTransitionScopeAdmissionClaim::Claimed), + Err(Error::PreconditionFailed) => Ok(ManualTransitionScopeAdmissionClaim::Conflict(Box::new(active))), + Err(err) => Err(err), + }; + } + + Ok(ManualTransitionScopeAdmissionClaim::Conflict(Box::new(active))) +} + +pub async fn request_manual_transition_job_cancel(api: Arc, job_id: Uuid) -> EcstoreResult { + for _ in 0..4 { + let (mut record, etag) = load_manual_transition_job_record_with_etag(api.clone(), job_id).await?; + if record.is_terminal() || record.cancel_requested { + return Ok(record); + } + record.mark_cancel_requested(); + match save_manual_transition_job_record_if_current(api.clone(), &record, &etag).await { + Ok(()) => return Ok(record), + Err(Error::PreconditionFailed) => continue, + Err(err) => return Err(err), + } + } + Err(Error::PreconditionFailed) +} + +pub async fn persist_manual_transition_job_progress( + api: Arc, + job_id: Uuid, + report: &ManualTransitionRunReport, + queue_snapshot: ManualTransitionQueueSnapshot, +) -> EcstoreResult { + let (mut record, etag) = load_manual_transition_job_record_with_etag(api.clone(), job_id).await?; + record.update_running_progress(report.clone(), queue_snapshot); + save_manual_transition_job_record_if_current(api.clone(), &record, &etag).await?; + renew_manual_transition_scope_admission_from_job(api, &record).await?; + Ok(record) +} + +pub async fn renew_manual_transition_job_lease( + api: Arc, + job_id: Uuid, + queue_snapshot: ManualTransitionQueueSnapshot, +) -> EcstoreResult { + let (mut record, etag) = load_manual_transition_job_record_with_etag(api.clone(), job_id).await?; + if record.state == ManualTransitionJobState::Running { + record.renew_lease(queue_snapshot); + save_manual_transition_job_record_if_current(api.clone(), &record, &etag).await?; + renew_manual_transition_scope_admission_from_job(api, &record).await?; + } + Ok(record) +} + +async fn renew_manual_transition_scope_admission_from_job( + api: Arc, + record: &ManualTransitionJobRecord, +) -> EcstoreResult<()> { + if let Ok((admission, admission_etag)) = + load_manual_transition_scope_admission_with_etag(api.clone(), &record.scope_key).await + && admission.job_id == record.job_id + && admission.lease_id == record.lease_id + { + let renewed_admission = ManualTransitionScopeAdmission::from_job(record); + save_manual_transition_scope_admission_if_current(api, &renewed_admission, &admission_etag).await?; + } + Ok(()) +} + +pub async fn delete_manual_transition_scope_admission_if_current( + api: Arc, + scope_key: &str, + job_id: Uuid, + lease_id: Uuid, +) -> EcstoreResult { + match load_manual_transition_scope_admission(api.clone(), scope_key).await { + Ok(admission) if admission.job_id == job_id && admission.lease_id == lease_id => {} + Ok(_) => return Ok(false), + Err(Error::ConfigNotFound) => return Ok(true), + Err(err) => return Err(err), + } + let object = manual_transition_scope_record_object_name(scope_key).map_err(manual_transition_job_store_error)?; + match config_boundary::delete_config(api, &object).await { + Ok(()) | Err(Error::ConfigNotFound) => Ok(true), + Err(err) => Err(err), + } +} + +fn manual_transition_job_store_error(err: ManualTransitionJobError) -> Error { + Error::other(err) +} + +pub fn manual_transition_scope_admission_lease_expired(admission: &ManualTransitionScopeAdmission) -> bool { + OffsetDateTime::now_utc().unix_timestamp_nanos() > admission.lease_expires_at_unix_nanos +} + +#[cfg(test)] +mod tests { + use super::*; + + const TEST_OWNER: &str = "owner-a"; + + #[test] + fn manual_transition_job_record_rejects_nil_job_id() { + let options = ManualTransitionRunOptions::default(); + let mut record = ManualTransitionJobRecord::new(Uuid::new_v4(), "bucket", &options, TEST_OWNER); + record.job_id = Uuid::nil(); + + let err = record.encode().expect_err("nil job id must fail closed"); + + assert!(err.to_string().contains("job_id is nil")); + } + + #[test] + fn manual_transition_job_record_round_trips_with_checksum() { + let options = ManualTransitionRunOptions { + prefix: "logs/".to_string(), + tier: Some("warm".to_string()), + max_objects: Some(17), + ..Default::default() + }; + let mut record = ManualTransitionJobRecord::new(Uuid::new_v4(), "bucket", &options, TEST_OWNER); + record.report.scanned = 3; + record.report.continuation_token = Some("opaque".to_string()); + record.mark_cancel_requested(); + record.complete( + ManualTransitionRunReport { + bucket: "bucket".to_string(), + prefix: "logs/".to_string(), + tier: Some("warm".to_string()), + cancelled: true, + ..Default::default() + }, + ManualTransitionQueueSnapshot::default(), + ); + + let encoded = record.encode().expect("job record should encode"); + let decoded = ManualTransitionJobRecord::decode(record.job_id, &encoded).expect("job record should decode"); + + assert_eq!(decoded.job_id, record.job_id); + assert_eq!(decoded.state, ManualTransitionJobState::Cancelled); + assert!(decoded.cancel_requested); + assert_eq!(decoded.max_objects, Some(17)); + } + + #[test] + fn manual_transition_job_record_builds_resume_options() { + let options = ManualTransitionRunOptions { + prefix: "logs/".to_string(), + continuation_token: Some("start-token".to_string()), + tier: Some("warm".to_string()), + dry_run: true, + max_objects: Some(3), + ..Default::default() + }; + let mut record = ManualTransitionJobRecord::new(Uuid::new_v4(), "bucket", &options, TEST_OWNER); + record.report.continuation_token = Some("resume-token".to_string()); + + let resume = record.resume_options(); + + assert_eq!(resume.prefix, "logs/"); + assert_eq!(resume.continuation_token.as_deref(), Some("resume-token")); + assert_eq!(resume.tier.as_deref(), Some("warm")); + assert!(resume.dry_run); + assert_eq!(resume.max_objects, Some(3)); + assert!(resume.cancel_token.is_none()); + assert!(resume.cancel_check.is_none()); + assert!(resume.progress_sink.is_none()); + } + + #[test] + fn manual_transition_job_record_object_name_round_trips_job_id() { + let job_id = Uuid::new_v4(); + let object_name = manual_transition_job_record_object_name(job_id).expect("job record path should encode"); + + let decoded = manual_transition_job_id_from_record_object_name(&object_name).expect("job record path should decode"); + + assert_eq!(decoded, job_id); + } + + #[test] + fn manual_transition_job_record_object_name_rejects_shard_mismatch() { + let job_id = Uuid::new_v4(); + let object_name = manual_transition_job_record_object_name(job_id).expect("job record path should encode"); + let job_key = job_id.simple().to_string(); + let bad_first_shard = if &job_key[..2] == "ff" { "00" } else { "ff" }; + let object_name = object_name.replacen( + &format!("{MANUAL_TRANSITION_JOB_RECORD_PREFIX}/{}/", &job_key[..2]), + &format!("{MANUAL_TRANSITION_JOB_RECORD_PREFIX}/{bad_first_shard}/"), + 1, + ); + + let err = manual_transition_job_id_from_record_object_name(&object_name).expect_err("bad shard must fail closed"); + + assert!(err.to_string().contains("shards")); + } + + #[test] + fn manual_transition_job_record_rejects_checksum_drift() { + let options = ManualTransitionRunOptions::default(); + let record = ManualTransitionJobRecord::new(Uuid::new_v4(), "bucket", &options, TEST_OWNER); + let encoded = record.encode().expect("job record should encode"); + let mut value: serde_json::Value = serde_json::from_slice(&encoded).expect("encoded job should be json"); + value["job"]["bucket"] = serde_json::Value::String("other-bucket".to_string()); + let mutated = serde_json::to_vec(&value).expect("mutated job should encode"); + + let err = ManualTransitionJobRecord::decode(record.job_id, &mutated).expect_err("checksum drift must fail closed"); + + assert!(matches!(err, ManualTransitionJobError::ChecksumMismatch)); + } + + #[test] + fn manual_transition_job_record_rejects_unknown_report_fields() { + let options = ManualTransitionRunOptions::default(); + let record = ManualTransitionJobRecord::new(Uuid::new_v4(), "bucket", &options, TEST_OWNER); + let encoded = record.encode().expect("job record should encode"); + let mut value: serde_json::Value = serde_json::from_slice(&encoded).expect("encoded job should be json"); + let checksum = value.get_mut("content_sha256").expect("encoded job has checksum"); + *checksum = serde_json::Value::String("0".repeat(64)); + value["job"]["report"]["unexpected"] = serde_json::Value::Bool(true); + let mutated_job_bytes = serde_json::to_vec(&value["job"]).expect("mutated job should encode"); + value["content_sha256"] = serde_json::Value::String(hex_sha256(&mutated_job_bytes, ToOwned::to_owned)); + let mutated = serde_json::to_vec(&value).expect("mutated envelope should encode"); + + let err = ManualTransitionJobRecord::decode(record.job_id, &mutated).expect_err("unknown report field must fail"); + + assert!(err.to_string().contains("unknown field")); + } + + #[test] + fn manual_transition_job_record_marks_restart_unknown_terminal() { + let options = ManualTransitionRunOptions::default(); + let mut record = ManualTransitionJobRecord::new(Uuid::new_v4(), "bucket", &options, TEST_OWNER); + + record.mark_unknown_if_unowned(); + + assert_eq!(record.state, ManualTransitionJobState::Unknown); + assert!(record.is_terminal()); + assert!(record.completed_at_unix_nanos.is_some()); + assert!( + record + .error + .as_deref() + .is_some_and(|error| error.contains("unknown after restart")) + ); + } + + #[test] + fn manual_transition_job_record_failure_counts_tier_failure() { + let options = ManualTransitionRunOptions::default(); + let mut record = ManualTransitionJobRecord::new(Uuid::new_v4(), "bucket", &options, TEST_OWNER); + + record.fail("missing tier"); + + assert_eq!(record.state, ManualTransitionJobState::Failed); + assert_eq!(record.report.tier_failure, 1); + assert_eq!(record.error.as_deref(), Some("missing tier")); + } + + #[test] + fn manual_transition_job_record_tier_failure_report_is_partial() { + let options = ManualTransitionRunOptions::default(); + let mut record = ManualTransitionJobRecord::new(Uuid::new_v4(), "bucket", &options, TEST_OWNER); + + record.complete( + ManualTransitionRunReport { + tier_failure: 1, + ..Default::default() + }, + ManualTransitionQueueSnapshot::default(), + ); + + assert_eq!(record.state, ManualTransitionJobState::Partial); + assert_eq!(record.report.tier_failure, 1); + assert!(record.error.is_none()); + } + + #[test] + fn manual_transition_scope_key_is_stable_and_sanitized() { + let first = manual_transition_scope_key( + "bucket", + &ManualTransitionRunOptions { + prefix: "logs/".to_string(), + tier: Some("warm".to_string()), + ..Default::default() + }, + ); + let second = manual_transition_scope_key( + "bucket", + &ManualTransitionRunOptions { + prefix: "logs/".to_string(), + tier: Some("WARM".to_string()), + marker: Some("ignored".to_string()), + ..Default::default() + }, + ); + + assert_eq!(first, second); + assert!( + manual_transition_scope_record_object_name(&first) + .expect("scope path should encode") + .starts_with(MANUAL_TRANSITION_SCOPE_RECORD_PREFIX) + ); + } + + #[test] + fn manual_transition_scope_admission_carries_job_lease_fence() { + let options = ManualTransitionRunOptions::default(); + let record = ManualTransitionJobRecord::new(Uuid::new_v4(), "bucket", &options, TEST_OWNER); + + let admission = ManualTransitionScopeAdmission::from_job(&record); + + assert_eq!(admission.job_id, record.job_id); + assert_eq!(admission.lease_id, record.lease_id); + assert_eq!(admission.owner_id, TEST_OWNER); + assert_eq!(admission.scope_key, record.scope_key); + assert!(admission.validate().is_ok()); + } +} diff --git a/crates/ecstore/src/bucket/lifecycle/mod.rs b/crates/ecstore/src/bucket/lifecycle/mod.rs index 311b63006..a84973d28 100644 --- a/crates/ecstore/src/bucket/lifecycle/mod.rs +++ b/crates/ecstore/src/bucket/lifecycle/mod.rs @@ -17,6 +17,7 @@ pub mod bucket_lifecycle_ops; mod config_boundary; pub mod core; pub mod evaluator; +pub mod manual_transition_job; mod metadata_boundary; mod object_lock_boundary; pub use self::core as lifecycle; diff --git a/rustfs/src/admin/handlers/ilm_transition.rs b/rustfs/src/admin/handlers/ilm_transition.rs index 27497bcd1..bbf7e5708 100644 --- a/rustfs/src/admin/handlers/ilm_transition.rs +++ b/rustfs/src/admin/handlers/ilm_transition.rs @@ -16,11 +16,16 @@ use crate::admin::auth::validate_admin_request; use crate::admin::router::{AdminOperation, Operation, S3Router}; 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::config::{read_admin_config, save_admin_config}; use crate::admin::storage_api::error::StorageError; use crate::admin::storage_api::lifecycle::{ - ManualTransitionRunOptions, ManualTransitionRunReport, enqueue_transition_for_existing_objects_scoped, - enqueue_transition_for_existing_objects_scoped_with_cancel, + ManualTransitionCancelCheck, ManualTransitionJobRecord, ManualTransitionJobState, ManualTransitionProgressSink, + ManualTransitionQueueSnapshot, ManualTransitionRunOptions, ManualTransitionRunReport, ManualTransitionScopeAdmission, + ManualTransitionScopeAdmissionClaim, claim_manual_transition_scope_admission, + delete_manual_transition_scope_admission_if_current, enqueue_transition_for_existing_objects_scoped, + load_manual_transition_job_record, load_manual_transition_job_record_with_etag, 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, renew_manual_transition_job_lease, request_manual_transition_job_cancel, + save_manual_transition_job_record, save_manual_transition_job_record_if_current, }; use crate::admin::storage_api::runtime::ECStore; use crate::auth::{check_key_valid, get_session_token}; @@ -37,9 +42,9 @@ use s3s::header::CONTENT_TYPE; use s3s::{Body, S3Error, S3ErrorCode, S3Request, S3Response, S3Result, s3_error}; use serde::{Deserialize, Serialize}; use std::collections::HashMap; +use std::pin::Pin; use std::sync::{Arc, Mutex, MutexGuard, OnceLock}; -use time::OffsetDateTime; -use time::format_description::well_known::Rfc3339; +use std::time::Duration as StdDuration; use tokio_util::sync::CancellationToken; use tracing::{error, info, warn}; use uuid::Uuid; @@ -51,11 +56,10 @@ const MAX_MANUAL_TRANSITION_DURATION_SECONDS: u64 = 3600; const LOG_COMPONENT_ADMIN: &str = "admin"; const LOG_SUBSYSTEM_ILM_TRANSITION: &str = "ilm_transition"; const EVENT_ADMIN_ILM_TRANSITION_STATE: &str = "admin_ilm_transition_state"; -const MANUAL_TRANSITION_JOB_SCHEMA_VERSION: u8 = 1; -const MANUAL_TRANSITION_JOB_CONFIG_PREFIX: &str = "ilm/transition/jobs"; static ACTIVE_MANUAL_TRANSITION_SCOPES: OnceLock>> = OnceLock::new(); -static MANUAL_TRANSITION_JOBS: OnceLock>> = OnceLock::new(); +static ACTIVE_MANUAL_TRANSITION_JOBS: OnceLock>> = OnceLock::new(); +static MANUAL_TRANSITION_OWNER_ID: OnceLock = OnceLock::new(); #[derive(Debug, Clone, PartialEq, Eq)] struct ManualTransitionRunScope { @@ -139,6 +143,8 @@ struct ManualTransitionRunQuery { marker: Option, #[serde(rename = "versionMarker")] version_marker: Option, + #[serde(rename = "continuationToken")] + continuation_token: Option, tier: Option, #[serde(rename = "async")] async_mode: Option, @@ -151,73 +157,47 @@ struct ManualTransitionRunQuery { max_duration_seconds: Option, } -#[derive(Debug, Clone, Copy, PartialEq, Eq)] -enum ManualTransitionRunMode { - EnqueueOnly, - Async, -} - -impl ManualTransitionRunMode { - fn response_mode(self) -> &'static str { - match self { - Self::EnqueueOnly => "enqueue_only", - Self::Async => "durable_job", - } - } -} - #[derive(Debug, Serialize)] struct ManualTransitionRunResponse { state: &'static str, mode: &'static str, + #[serde(skip_serializing_if = "Option::is_none")] job_id: Option, + #[serde(skip_serializing_if = "Option::is_none")] status_endpoint: Option, + #[serde(skip_serializing_if = "Option::is_none")] + cancel_endpoint: Option, report: ManualTransitionRunReport, } -#[derive(Debug, Clone)] -struct ManualTransitionJobRuntime { - record: ManualTransitionJobRecord, - cancel_token: CancellationToken, -} - -#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] -#[serde(rename_all = "snake_case")] -enum ManualTransitionJobStatus { - Queued, - Running, - Completed, - Partial, - Cancelled, - Failed, - Unknown, -} - -impl ManualTransitionJobStatus { - fn is_terminal(self) -> bool { - matches!(self, Self::Completed | Self::Partial | Self::Cancelled | Self::Failed | Self::Unknown) - } -} - -#[derive(Debug, Clone, Serialize, Deserialize)] -#[serde(deny_unknown_fields)] -struct ManualTransitionJobRecord { - schema_version: u8, +#[derive(Debug, Serialize)] +struct ManualTransitionJobResponse { + status: ManualTransitionJobState, + mode: &'static str, job_id: String, status_endpoint: String, - status: ManualTransitionJobStatus, + cancel_endpoint: String, + cancel_requested: bool, bucket: String, prefix: String, tier: Option, dry_run: bool, - max_objects: Option, - max_duration_seconds: Option, - created_at: String, - started_at: Option, - finished_at: Option, - cancel_requested: bool, - failure_reason: Option, + created_at_unix_nanos: i128, + updated_at_unix_nanos: i128, + completed_at_unix_nanos: Option, report: ManualTransitionRunReport, + queue_snapshot: ManualTransitionQueueSnapshot, + failure_reason: Option, +} + +#[derive(Debug, Serialize)] +struct ManualTransitionJobConflictResponse { + state: &'static str, + mode: &'static str, + active_job_id: String, + status_endpoint: String, + cancel_endpoint: String, + scope_key: String, } pub fn register_ilm_transition_route(r: &mut S3Router) -> std::io::Result<()> { @@ -239,6 +219,12 @@ pub fn register_ilm_transition_route(r: &mut S3Router) -> std::i Ok(()) } +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum ManualTransitionRunMode { + EnqueueOnly, + Async, +} + fn parse_manual_transition_query(query: Option<&str>) -> S3Result<(String, ManualTransitionRunOptions, ManualTransitionRunMode)> { let query: ManualTransitionRunQuery = match query { Some(query) => serde_urlencoded::from_bytes(query.as_bytes()) @@ -271,6 +257,12 @@ fn parse_manual_transition_query(query: Option<&str>) -> S3Result<(String, Manua } else { ManualTransitionRunMode::EnqueueOnly }; + if query.continuation_token.is_some() && (query.marker.is_some() || query.version_marker.is_some()) { + return Err(s3_error!( + InvalidArgument, + "continuationToken cannot be combined with marker or versionMarker" + )); + } let max_objects = query.max_objects.unwrap_or(DEFAULT_MANUAL_TRANSITION_MAX_OBJECTS); if max_objects == 0 || max_objects > MAX_MANUAL_TRANSITION_OBJECTS { @@ -289,10 +281,14 @@ fn parse_manual_transition_query(query: Option<&str>) -> S3Result<(String, Manua prefix: query.prefix.unwrap_or_default(), marker: query.marker.filter(|marker| !marker.is_empty()), version_marker: query.version_marker.filter(|version_marker| !version_marker.is_empty()), + continuation_token: query.continuation_token.filter(|token| !token.is_empty()), tier: query.tier.map(|tier| tier.trim().to_string()).filter(|tier| !tier.is_empty()), dry_run: query.dry_run.unwrap_or(false), max_objects: Some(max_objects), max_duration: query.max_duration_seconds.map(std::time::Duration::from_secs), + cancel_token: None, + cancel_check: None, + progress_sink: None, }, run_mode, )) @@ -423,7 +419,98 @@ fn response_state(report: &ManualTransitionRunReport) -> &'static str { } } -fn json_response(status: StatusCode, response: &T) -> S3Result> { +fn validate_manual_transition_job_id(params: &Params<'_, '_>) -> S3Result<()> { + let job_id = params.get("job_id").unwrap_or(""); + if job_id.is_empty() { + return Err(s3_error!(InvalidRequest, "manual transition job id is required")); + } + Ok(()) +} + +fn manual_transition_job_id_from_params(params: &Params<'_, '_>) -> S3Result { + validate_manual_transition_job_id(params)?; + Uuid::parse_str(params.get("job_id").unwrap_or("")) + .map_err(|_| s3_error!(InvalidArgument, "invalid manual transition job id")) +} + +fn active_manual_transition_jobs() -> &'static Mutex> { + ACTIVE_MANUAL_TRANSITION_JOBS.get_or_init(|| Mutex::new(HashMap::new())) +} + +fn insert_active_manual_transition_job(job_id: Uuid, cancel_token: CancellationToken) { + let mut jobs = active_manual_transition_jobs() + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()); + jobs.insert(job_id, cancel_token); +} + +fn remove_active_manual_transition_job(job_id: Uuid) { + let mut jobs = active_manual_transition_jobs() + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()); + jobs.remove(&job_id); +} + +fn active_manual_transition_cancel_token(job_id: Uuid) -> Option { + let jobs = active_manual_transition_jobs() + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()); + jobs.get(&job_id).cloned() +} + +fn manual_transition_status_endpoint(job_id: Uuid) -> String { + format!("{ADMIN_PREFIX}/v3/ilm/transition/jobs/{job_id}") +} + +fn manual_transition_owner_id() -> &'static str { + MANUAL_TRANSITION_OWNER_ID.get_or_init(|| Uuid::new_v4().to_string()).as_str() +} + +fn manual_transition_job_response(record: ManualTransitionJobRecord) -> ManualTransitionJobResponse { + let status_endpoint = manual_transition_status_endpoint(record.job_id); + ManualTransitionJobResponse { + status: record.state, + mode: "durable_job", + job_id: record.job_id.to_string(), + status_endpoint: status_endpoint.clone(), + cancel_endpoint: status_endpoint, + cancel_requested: record.cancel_requested, + bucket: record.bucket, + prefix: record.prefix, + tier: record.tier, + dry_run: record.dry_run, + created_at_unix_nanos: record.created_at_unix_nanos, + updated_at_unix_nanos: record.updated_at_unix_nanos, + completed_at_unix_nanos: record.completed_at_unix_nanos, + report: record.report, + queue_snapshot: record.queue_snapshot, + failure_reason: record.error, + } +} + +fn manual_transition_job_conflict_response(admission: ManualTransitionScopeAdmission) -> ManualTransitionJobConflictResponse { + let status_endpoint = manual_transition_status_endpoint(admission.job_id); + ManualTransitionJobConflictResponse { + state: "conflict", + mode: "durable_job", + active_job_id: admission.job_id.to_string(), + status_endpoint: status_endpoint.clone(), + cancel_endpoint: status_endpoint, + scope_key: admission.scope_key, + } +} + +fn map_manual_transition_job_load_error(err: StorageError, job_id: Uuid) -> S3Error { + if err == StorageError::ConfigNotFound { + s3_error!(NoSuchKey, "manual transition job not found: {}", job_id) + } else if err == StorageError::PreconditionFailed { + s3_error!(OperationAborted, "manual transition job record changed concurrently; retry the request") + } else { + S3Error::with_message(S3ErrorCode::InternalError, format!("manual transition job store failed: {err}")) + } +} + +fn json_response(response: &T, status: StatusCode) -> S3Result> { let body = serde_json::to_vec(response).map_err(|err| { S3Error::with_message(S3ErrorCode::InternalError, format!("failed to encode manual transition response: {err}")) })?; @@ -434,224 +521,226 @@ fn json_response(status: StatusCode, response: &T) -> S3Result &'static Mutex> { - MANUAL_TRANSITION_JOBS.get_or_init(|| Mutex::new(HashMap::new())) -} - -fn lock_manual_transition_jobs() -> MutexGuard<'static, HashMap> { - match manual_transition_job_registry().lock() { - Ok(jobs) => jobs, - Err(poisoned) => poisoned.into_inner(), +async fn update_manual_transition_job_record_cas( + store: Arc, + job_id: Uuid, + mut update: impl FnMut(&mut ManualTransitionJobRecord), +) -> S3Result { + for _ in 0..4 { + let (mut record, etag) = load_manual_transition_job_record_with_etag(store.clone(), job_id) + .await + .map_err(|err| map_manual_transition_job_load_error(err, job_id))?; + update(&mut record); + match save_manual_transition_job_record_if_current(store.clone(), &record, &etag).await { + Ok(()) => return Ok(record), + Err(StorageError::PreconditionFailed) => continue, + Err(err) => { + return Err(S3Error::with_message( + S3ErrorCode::InternalError, + format!("manual transition job store failed: {err}"), + )); + } + } } + Err(s3_error!( + OperationAborted, + "manual transition job record changed concurrently; retry the request" + )) } -fn manual_transition_job_status_endpoint(job_id: &str) -> String { - format!("{ADMIN_PREFIX}/v3/ilm/transition/jobs/{job_id}") +fn manual_transition_durable_cancel_check(store: Arc, job_id: Uuid) -> ManualTransitionCancelCheck { + let last_cancelled = Arc::new(std::sync::atomic::AtomicBool::new(false)); + let next_poll_at = Arc::new(Mutex::new(std::time::Instant::now())); + Arc::new(move || { + let store = store.clone(); + let last_cancelled = last_cancelled.clone(); + let next_poll_at = next_poll_at.clone(); + Box::pin(async move { + if last_cancelled.load(std::sync::atomic::Ordering::SeqCst) { + return true; + } + let should_poll = { + let mut next_poll_at = next_poll_at.lock().unwrap_or_else(|poisoned| poisoned.into_inner()); + let now = std::time::Instant::now(); + if now < *next_poll_at { + false + } else { + *next_poll_at = now + StdDuration::from_secs(1); + true + } + }; + if !should_poll { + return false; + } + match load_manual_transition_job_record(store, job_id).await { + Ok(record) if record.cancel_requested => { + last_cancelled.store(true, std::sync::atomic::Ordering::SeqCst); + true + } + _ => false, + } + }) as Pin + Send>> + }) } -fn manual_transition_job_config_path(job_id: &str) -> String { - format!("{MANUAL_TRANSITION_JOB_CONFIG_PREFIX}/{job_id}.json") -} - -fn parse_manual_transition_job_id(job_id: &str) -> S3Result { - Uuid::parse_str(job_id).map_err(|_| s3_error!(InvalidArgument, "invalid manual transition job id"))?; - Ok(job_id.to_string()) -} - -fn manual_transition_timestamp(now: OffsetDateTime) -> String { - now.format(&Rfc3339).unwrap_or_else(|_| now.unix_timestamp().to_string()) -} - -fn initial_manual_transition_report(bucket: &str, options: &ManualTransitionRunOptions) -> ManualTransitionRunReport { - ManualTransitionRunReport { - bucket: bucket.to_string(), - prefix: options.prefix.clone(), - tier: options.tier.clone(), - dry_run: options.dry_run, - ..Default::default() - } -} - -fn manual_transition_status_from_report(report: &ManualTransitionRunReport) -> ManualTransitionJobStatus { - if report.was_truncated() || report.has_partial_enqueue() { - ManualTransitionJobStatus::Partial - } else { - ManualTransitionJobStatus::Completed - } -} - -fn new_manual_transition_job_record( - job_id: String, - bucket: &str, - options: &ManualTransitionRunOptions, - created_at: OffsetDateTime, -) -> ManualTransitionJobRecord { - ManualTransitionJobRecord { - schema_version: MANUAL_TRANSITION_JOB_SCHEMA_VERSION, - status_endpoint: manual_transition_job_status_endpoint(&job_id), - job_id, - status: ManualTransitionJobStatus::Queued, - bucket: bucket.to_string(), - prefix: options.prefix.clone(), - tier: options.tier.clone(), - dry_run: options.dry_run, - max_objects: options.max_objects, - max_duration_seconds: options.max_duration.map(|duration| duration.as_secs()), - created_at: manual_transition_timestamp(created_at), - started_at: None, - finished_at: None, - cancel_requested: false, - failure_reason: None, - report: initial_manual_transition_report(bucket, options), - } -} - -fn insert_manual_transition_job(record: ManualTransitionJobRecord, cancel_token: CancellationToken) { - let mut jobs = lock_manual_transition_jobs(); - jobs.insert(record.job_id.clone(), ManualTransitionJobRuntime { record, cancel_token }); -} - -fn update_manual_transition_job_record( - job_id: &str, - update: impl FnOnce(&mut ManualTransitionJobRecord), -) -> Option { - let mut jobs = lock_manual_transition_jobs(); - let runtime = jobs.get_mut(job_id)?; - update(&mut runtime.record); - Some(runtime.record.clone()) -} - -fn in_memory_manual_transition_job_record(job_id: &str) -> Option { - lock_manual_transition_jobs() - .get(job_id) - .map(|runtime| runtime.record.clone()) -} - -fn request_manual_transition_job_cancel(job_id: &str) -> Option { - let mut jobs = lock_manual_transition_jobs(); - let runtime = jobs.get_mut(job_id)?; - if !runtime.record.status.is_terminal() { - runtime.record.cancel_requested = true; - runtime.cancel_token.cancel(); - } - Some(runtime.record.clone()) -} - -fn mark_manual_transition_job_owner_unknown(mut record: ManualTransitionJobRecord) -> ManualTransitionJobRecord { - if !record.status.is_terminal() { - record.status = ManualTransitionJobStatus::Unknown; - record.finished_at = Some(manual_transition_timestamp(OffsetDateTime::now_utc())); - record.failure_reason = Some("manual transition job owner is not active on this node".to_string()); - } - record -} - -fn remove_manual_transition_job(job_id: &str) { - let mut jobs = lock_manual_transition_jobs(); - jobs.remove(job_id); -} - -async fn save_manual_transition_job_record(store: Arc, record: &ManualTransitionJobRecord) -> S3Result<()> { - let data = serde_json::to_vec(record).map_err(|err| { - S3Error::with_message(S3ErrorCode::InternalError, format!("failed to encode manual transition job: {err}")) - })?; - save_admin_config(store, &manual_transition_job_config_path(&record.job_id), data) - .await - .map_err(|err| { - S3Error::with_message(S3ErrorCode::InternalError, format!("failed to persist manual transition job: {err}")) +fn manual_transition_progress_sink(store: Arc, job_id: Uuid) -> ManualTransitionProgressSink { + Arc::new(move |report| { + let store = store.clone(); + Box::pin(async move { + persist_manual_transition_job_progress(store, job_id, &report, manual_transition_queue_snapshot()) + .await + .map(|_| ()) }) + }) } -async fn read_manual_transition_job_record(store: Arc, job_id: &str) -> S3Result> { - match read_admin_config(store, &manual_transition_job_config_path(job_id)).await { - Ok(data) => serde_json::from_slice(&data).map(Some).map_err(|err| { - S3Error::with_message(S3ErrorCode::InternalError, format!("failed to decode manual transition job: {err}")) - }), - Err(StorageError::ConfigNotFound) => Ok(None), - Err(err) => Err(S3Error::with_message( - S3ErrorCode::InternalError, - format!("failed to read manual transition job: {err}"), - )), +fn release_manual_transition_admission(store: Arc, record: &ManualTransitionJobRecord) { + let scope_key = record.scope_key.clone(); + let job_id = record.job_id; + let lease_id = record.lease_id; + tokio::spawn(async move { + if let Err(err) = delete_manual_transition_scope_admission_if_current(store, &scope_key, job_id, lease_id).await { + warn!( + event = EVENT_ADMIN_ILM_TRANSITION_STATE, + component = LOG_COMPONENT_ADMIN, + subsystem = LOG_SUBSYSTEM_ILM_TRANSITION, + operation = "manual_transition_job", + result = "failed", + job_id = %job_id, + error = %err, + "failed to release manual transition admission" + ); + } + }); +} + +async fn finalize_manual_transition_job( + store: Arc, + job_id: Uuid, + result: Result, +) -> Option { + let updated = update_manual_transition_job_record_cas(store.clone(), job_id, |record| { + let cancel_requested = record.cancel_requested; + match &result { + Ok(report) => { + let mut report = report.clone(); + if cancel_requested { + report.cancelled = true; + } + record.complete(report, manual_transition_queue_snapshot()); + if cancel_requested { + record.mark_cancel_requested(); + } + } + Err(err) => { + record.fail(err.to_string()); + if cancel_requested { + record.mark_cancel_requested(); + } + } + } + }) + .await; + match updated { + Ok(record) => Some(record), + Err(err) => { + error!( + event = EVENT_ADMIN_ILM_TRANSITION_STATE, + component = LOG_COMPONENT_ADMIN, + subsystem = LOG_SUBSYSTEM_ILM_TRANSITION, + operation = "manual_transition_job", + result = "failed", + job_id = %job_id, + error = %err, + "failed to persist manual transition job terminal state" + ); + None + } } } -async fn load_manual_transition_job_record(store: Arc, job_id: &str) -> S3Result { - if let Some(record) = in_memory_manual_transition_job_record(job_id) { - return Ok(record); - } - let Some(record) = read_manual_transition_job_record(store, job_id).await? else { - return Err(s3_error!(NoSuchKey, "manual transition job not found")); - }; - Ok(mark_manual_transition_job_owner_unknown(record)) +fn spawn_manual_transition_job_heartbeat(store: Arc, job_id: Uuid, cancel_token: CancellationToken) { + tokio::spawn(async move { + let mut interval = tokio::time::interval(StdDuration::from_secs(5)); + loop { + tokio::select! { + _ = cancel_token.cancelled() => return, + _ = interval.tick() => { + match renew_manual_transition_job_lease(store.clone(), job_id, manual_transition_queue_snapshot()).await { + Ok(record) if record.cancel_requested => cancel_token.cancel(), + Ok(_) => {} + Err(err) => { + warn!( + event = EVENT_ADMIN_ILM_TRANSITION_STATE, + component = LOG_COMPONENT_ADMIN, + subsystem = LOG_SUBSYSTEM_ILM_TRANSITION, + operation = "manual_transition_job", + result = "failed", + job_id = %job_id, + error = %err, + "failed to renew manual transition job lease" + ); + } + } + } + } + } + }); } -async fn persist_manual_transition_job_update(store: Arc, record: &ManualTransitionJobRecord) -> bool { - if let Err(err) = save_manual_transition_job_record(store, record).await { - error!( - event = EVENT_ADMIN_ILM_TRANSITION_STATE, - component = LOG_COMPONENT_ADMIN, - subsystem = LOG_SUBSYSTEM_ILM_TRANSITION, - operation = "manual_transition_job", - job_id = %record.job_id, - error = %err, - "failed to persist manual transition job update" - ); - return false; - } - true +enum StartManualTransitionJobResult { + Started(Box), + Conflict(ManualTransitionJobConflictResponse), } -fn spawn_manual_transition_job( +async fn start_manual_transition_job( store: Arc, bucket: String, options: ManualTransitionRunOptions, - job_id: String, - cancel_token: CancellationToken, - admission_guard: ManualTransitionRunAdmission, -) { - tokio::spawn(async move { - let started_record = update_manual_transition_job_record(&job_id, |record| { - record.status = ManualTransitionJobStatus::Running; - record.started_at = Some(manual_transition_timestamp(OffsetDateTime::now_utc())); - }); - if let Some(record) = started_record.as_ref() { - persist_manual_transition_job_update(store.clone(), record).await; +) -> S3Result { + let job_id = Uuid::new_v4(); + let record = ManualTransitionJobRecord::new(job_id, &bucket, &options, manual_transition_owner_id()); + save_manual_transition_job_record(store.clone(), &record) + .await + .map_err(|err| S3Error::with_message(S3ErrorCode::InternalError, format!("manual transition job store failed: {err}")))?; + match claim_manual_transition_scope_admission(store.clone(), &ManualTransitionScopeAdmission::from_job(&record)).await { + Ok(ManualTransitionScopeAdmissionClaim::Claimed) => {} + Ok(ManualTransitionScopeAdmissionClaim::Conflict(active)) => { + let _ = update_manual_transition_job_record_cas(store.clone(), job_id, |record| { + record.fail("manual transition admission conflict"); + }) + .await; + return Ok(StartManualTransitionJobResult::Conflict(manual_transition_job_conflict_response(*active))); } + Err(err) => { + let _ = update_manual_transition_job_record_cas(store.clone(), job_id, |record| { + record.fail(format!("manual transition admission failed: {err}")); + }) + .await; + return Err(S3Error::with_message( + S3ErrorCode::InternalError, + format!("manual transition admission failed: {err}"), + )); + } + } - let result = enqueue_transition_for_existing_objects_scoped_with_cancel( - store.clone(), - &bucket, - options, - Some(cancel_token.clone()), - ) - .await; - let finished_at = manual_transition_timestamp(OffsetDateTime::now_utc()); - let finished_record = update_manual_transition_job_record(&job_id, |record| { - record.finished_at = Some(finished_at); - match result { - Ok(execution) => { - record.report = execution.report; - record.status = if execution.cancelled || record.cancel_requested || cancel_token.is_cancelled() { - ManualTransitionJobStatus::Cancelled - } else { - manual_transition_status_from_report(&record.report) - }; - record.failure_reason = None; - } - Err(err) => { - record.status = ManualTransitionJobStatus::Failed; - record.failure_reason = Some(err.to_string()); - } - } - }); - if let Some(record) = finished_record.as_ref() - && persist_manual_transition_job_update(store, record).await - && record.status.is_terminal() - { - remove_manual_transition_job(&job_id); + let cancel_token = CancellationToken::new(); + insert_active_manual_transition_job(job_id, cancel_token.clone()); + let mut run_options = options; + run_options.cancel_token = Some(cancel_token.clone()); + run_options.cancel_check = Some(manual_transition_durable_cancel_check(store.clone(), job_id)); + run_options.progress_sink = Some(manual_transition_progress_sink(store.clone(), job_id)); + let run_store = store.clone(); + spawn_manual_transition_job_heartbeat(store, job_id, cancel_token); + tokio::spawn(async move { + let result = enqueue_transition_for_existing_objects_scoped(run_store.clone(), &bucket, run_options).await; + if let Some(final_record) = finalize_manual_transition_job(run_store.clone(), job_id, result).await { + release_manual_transition_admission(run_store, &final_record); } - drop(admission_guard); + remove_active_manual_transition_job(job_id); }); + + Ok(StartManualTransitionJobResult::Started(Box::new(record))) } pub struct ManualTransitionRunHandler {} @@ -673,10 +762,30 @@ impl Operation for ManualTransitionRunHandler { log_manual_transition_rejected("object_store_not_initialized", &request_id, &actor, &remote_addr); return Err(s3_error!(InternalError, "object store is not initialized")); }; + if run_mode == ManualTransitionRunMode::Async { + match start_manual_transition_job(store, bucket, options).await? { + StartManualTransitionJobResult::Started(record) => { + let record = *record; + let status_endpoint = manual_transition_status_endpoint(record.job_id); + let response = ManualTransitionRunResponse { + state: "accepted", + mode: "durable_job", + job_id: Some(record.job_id.to_string()), + status_endpoint: Some(status_endpoint.clone()), + cancel_endpoint: Some(status_endpoint), + report: record.report, + }; + return json_response(&response, StatusCode::ACCEPTED); + } + StartManualTransitionJobResult::Conflict(response) => { + return json_response(&response, StatusCode::CONFLICT); + } + } + } let max_objects = options.max_objects; let max_duration_seconds = options.max_duration.map(|duration| duration.as_secs()); let scope = ManualTransitionRunScope::new(&bucket, &options); - let admission = match acquire_manual_transition_admission(scope) { + let _admission = match acquire_manual_transition_admission(scope) { Ok(admission) => admission, Err(err) => { log_manual_transition_rejected("already_running", &request_id, &actor, &remote_addr); @@ -684,24 +793,6 @@ impl Operation for ManualTransitionRunHandler { } }; - if run_mode == ManualTransitionRunMode::Async { - let job_id = Uuid::new_v4().to_string(); - let record = new_manual_transition_job_record(job_id.clone(), &bucket, &options, OffsetDateTime::now_utc()); - save_manual_transition_job_record(store.clone(), &record).await?; - let cancel_token = CancellationToken::new(); - insert_manual_transition_job(record.clone(), cancel_token.clone()); - spawn_manual_transition_job(store, bucket, options, job_id.clone(), cancel_token, admission); - let response = ManualTransitionRunResponse { - state: "accepted", - mode: run_mode.response_mode(), - job_id: Some(job_id), - status_endpoint: Some(record.status_endpoint), - report: record.report, - }; - - return json_response(StatusCode::ACCEPTED, &response); - } - let report = match enqueue_transition_for_existing_objects_scoped(store, &bucket, options).await { Ok(report) => report, Err(err) => { @@ -719,10 +810,11 @@ impl Operation for ManualTransitionRunHandler { mode: "enqueue_only", job_id: None, status_endpoint: None, + cancel_endpoint: None, report, }; - json_response(StatusCode::OK, &response) + json_response(&response, StatusCode::OK) } } @@ -731,17 +823,35 @@ pub struct ManualTransitionJobStatusHandler {} #[async_trait::async_trait] impl Operation for ManualTransitionJobStatusHandler { async fn call(&self, req: S3Request, params: Params<'_, '_>) -> S3Result> { - let _actor = authorize_manual_transition_request(&req).await?; - let job_id = parse_manual_transition_job_id( - params - .get("job_id") - .ok_or_else(|| s3_error!(InvalidRequest, "manual transition job id is required"))?, - )?; + authorize_manual_transition_request(&req).await?; + let job_id = manual_transition_job_id_from_params(¶ms)?; let Some(store) = object_store_from_extensions(&req.extensions) else { return Err(s3_error!(InternalError, "object store is not initialized")); }; - let record = load_manual_transition_job_record(store, &job_id).await?; - json_response(StatusCode::OK, &record) + let mut record = load_manual_transition_job_record(store.clone(), job_id) + .await + .map_err(|err| map_manual_transition_job_load_error(err, job_id))?; + if record.state == ManualTransitionJobState::Running { + let local_active = active_manual_transition_cancel_token(job_id).is_some(); + let leased_elsewhere = load_manual_transition_scope_admission(store.clone(), &record.scope_key) + .await + .ok() + .is_some_and(|admission| { + admission.job_id == record.job_id + && admission.lease_id == record.lease_id + && !manual_transition_scope_admission_lease_expired(&admission) + }); + if !local_active && !leased_elsewhere && manual_transition_job_lease_expired(&record) { + record = update_manual_transition_job_record_cas(store.clone(), job_id, |record| { + if record.state == ManualTransitionJobState::Running && manual_transition_job_lease_expired(record) { + record.mark_unknown_if_unowned(); + } + }) + .await?; + release_manual_transition_admission(store, &record); + } + } + json_response(&manual_transition_job_response(record), StatusCode::OK) } } @@ -750,35 +860,51 @@ pub struct ManualTransitionJobCancelHandler {} #[async_trait::async_trait] impl Operation for ManualTransitionJobCancelHandler { async fn call(&self, req: S3Request, params: Params<'_, '_>) -> S3Result> { - let _actor = authorize_manual_transition_request(&req).await?; - let job_id = parse_manual_transition_job_id( - params - .get("job_id") - .ok_or_else(|| s3_error!(InvalidRequest, "manual transition job id is required"))?, - )?; + authorize_manual_transition_request(&req).await?; + let job_id = manual_transition_job_id_from_params(¶ms)?; let Some(store) = object_store_from_extensions(&req.extensions) else { return Err(s3_error!(InternalError, "object store is not initialized")); }; - let record = if let Some(record) = request_manual_transition_job_cancel(&job_id) { - record - } else { - let Some(mut record) = read_manual_transition_job_record(store.clone(), &job_id).await? else { - return Err(s3_error!(NoSuchKey, "manual transition job not found")); - }; - if !record.status.is_terminal() { - record.cancel_requested = true; - record = mark_manual_transition_job_owner_unknown(record); - save_manual_transition_job_record(store, &record).await?; - } - record - }; - json_response(StatusCode::OK, &record) + let record = request_manual_transition_job_cancel(store, job_id) + .await + .map_err(|err| map_manual_transition_job_load_error(err, job_id))?; + if !record.is_terminal() + && let Some(cancel_token) = active_manual_transition_cancel_token(job_id) + { + cancel_token.cancel(); + } + json_response(&manual_transition_job_response(record), StatusCode::OK) } } #[cfg(test)] mod tests { use super::*; + use matchit::Router; + + fn with_manual_transition_job_params(path: &str, f: impl FnOnce(&Params<'_, '_>) -> T) -> T { + let mut router = Router::new(); + router + .insert("/rustfs/admin/v3/ilm/transition/jobs/{job_id}", ()) + .expect("route should insert"); + + let matched = router.at(path).expect("route should match"); + f(&matched.params) + } + + fn manual_transition_job_request(method: Method, path: &'static str) -> S3Request { + S3Request { + input: Body::empty(), + method, + uri: path.parse().expect("valid route"), + headers: HeaderMap::new(), + extensions: http::Extensions::new(), + credentials: None, + region: None, + service: None, + trailing_headers: None, + } + } #[test] fn manual_transition_query_defaults_to_bounded_run() { @@ -817,11 +943,12 @@ mod tests { } #[test] - fn manual_transition_query_accepts_explicit_async_mode() { - let (_bucket, _options, run_mode) = - parse_manual_transition_query(Some("bucket=data&async=true")).expect("async mode should parse"); + fn manual_transition_query_accepts_durable_async_mode() { + let (_bucket, options, run_mode) = + parse_manual_transition_query(Some("bucket=data&async=true")).expect("async durable jobs should parse"); assert_eq!(run_mode, ManualTransitionRunMode::Async); + assert_eq!(options.max_objects, Some(DEFAULT_MANUAL_TRANSITION_MAX_OBJECTS)); let (_bucket, _options, run_mode) = parse_manual_transition_query(Some("bucket=data&mode=async")).expect("mode=async should parse"); @@ -842,6 +969,19 @@ mod tests { assert_eq!(err.code(), &S3ErrorCode::InvalidArgument); } + #[test] + fn manual_transition_query_rejects_continuation_with_raw_markers() { + let err = parse_manual_transition_query(Some("bucket=data&continuationToken=opaque&marker=logs/a")) + .expect_err("continuation token and raw marker must not be mixed"); + + assert_eq!(err.code(), &S3ErrorCode::InvalidArgument); + + let err = parse_manual_transition_query(Some("bucket=data&continuationToken=opaque&versionMarker=v1")) + .expect_err("continuation token and raw version marker must not be mixed"); + + assert_eq!(err.code(), &S3ErrorCode::InvalidArgument); + } + #[test] fn manual_transition_query_rejects_unknown_mode() { let err = parse_manual_transition_query(Some("bucket=data&mode=background")) @@ -1012,102 +1152,18 @@ mod tests { mode: "enqueue_only", job_id: None, status_endpoint: None, + cancel_endpoint: None, report, }; let value = serde_json::to_value(response).expect("response should serialize"); + assert!(value.get("job_id").is_none()); + assert!(value.get("status_endpoint").is_none()); + assert!(value.get("cancel_endpoint").is_none()); assert!(value.pointer("/report/next_marker").is_none()); assert!(value.pointer("/report/next_version_idmarker").is_none()); } - #[test] - fn manual_transition_job_record_omits_raw_resume_markers() { - let (bucket, options, run_mode) = - parse_manual_transition_query(Some("bucket=data&prefix=logs/&async=true")).expect("async query should parse"); - let mut record = new_manual_transition_job_record( - "11111111-1111-4111-8111-111111111111".to_string(), - &bucket, - &options, - OffsetDateTime::UNIX_EPOCH, - ); - record.status = ManualTransitionJobStatus::Partial; - record.report.next_marker = Some("private/object".to_string()); - record.report.next_version_idmarker = Some("null".to_string()); - - assert_eq!(run_mode, ManualTransitionRunMode::Async); - assert_eq!( - record.status_endpoint, - "/rustfs/admin/v3/ilm/transition/jobs/11111111-1111-4111-8111-111111111111" - ); - - let value = serde_json::to_value(record).expect("job record should serialize"); - assert!(value.pointer("/report/next_marker").is_none()); - assert!(value.pointer("/report/next_version_idmarker").is_none()); - } - - #[test] - fn manual_transition_job_id_rejects_non_uuid_path_segment() { - let err = parse_manual_transition_job_id("../config").expect_err("path-like job ids must be rejected"); - - assert_eq!(err.code(), &S3ErrorCode::InvalidArgument); - } - - #[test] - fn manual_transition_cancel_marks_active_job_and_token() { - let (bucket, options, _run_mode) = - parse_manual_transition_query(Some("bucket=data&prefix=logs/&async=true")).expect("async query should parse"); - let job_id = Uuid::new_v4().to_string(); - let cancel_token = CancellationToken::new(); - let record = new_manual_transition_job_record(job_id.clone(), &bucket, &options, OffsetDateTime::UNIX_EPOCH); - insert_manual_transition_job(record, cancel_token.clone()); - - let cancelled = request_manual_transition_job_cancel(&job_id).expect("active job should be found"); - - assert!(cancelled.cancel_requested); - assert!(cancel_token.is_cancelled()); - } - - #[test] - fn manual_transition_inactive_owner_marks_nonterminal_job_unknown() { - let (bucket, options, _run_mode) = - parse_manual_transition_query(Some("bucket=data&prefix=logs/&async=true")).expect("async query should parse"); - let record = new_manual_transition_job_record( - "11111111-1111-4111-8111-111111111111".to_string(), - &bucket, - &options, - OffsetDateTime::UNIX_EPOCH, - ); - - let unknown = mark_manual_transition_job_owner_unknown(record); - - assert_eq!(unknown.status, ManualTransitionJobStatus::Unknown); - assert!(unknown.finished_at.is_some()); - assert_eq!( - unknown.failure_reason.as_deref(), - Some("manual transition job owner is not active on this node") - ); - } - - #[test] - fn manual_transition_inactive_owner_preserves_terminal_job() { - let (bucket, options, _run_mode) = - parse_manual_transition_query(Some("bucket=data&prefix=logs/&async=true")).expect("async query should parse"); - let mut record = new_manual_transition_job_record( - "11111111-1111-4111-8111-111111111111".to_string(), - &bucket, - &options, - OffsetDateTime::UNIX_EPOCH, - ); - record.status = ManualTransitionJobStatus::Completed; - record.finished_at = Some("1970-01-01T00:00:00Z".to_string()); - - let completed = mark_manual_transition_job_owner_unknown(record); - - assert_eq!(completed.status, ManualTransitionJobStatus::Completed); - assert_eq!(completed.finished_at.as_deref(), Some("1970-01-01T00:00:00Z")); - assert_eq!(completed.failure_reason, None); - } - #[test] fn manual_transition_handler_requires_set_tier_action() { let src = include_str!("ilm_transition.rs"); @@ -1117,6 +1173,100 @@ mod tests { assert!(!auth_block.contains("AdminAction::ServerInfoAdminAction")); } + #[test] + fn manual_transition_job_id_path_param_is_required() { + with_manual_transition_job_params("/rustfs/admin/v3/ilm/transition/jobs/job-123", |params| { + assert_eq!(params.get("job_id"), Some("job-123")); + validate_manual_transition_job_id(params) + }) + .expect("job id should validate"); + + let err = validate_manual_transition_job_id(&Params::new()).expect_err("missing job id must fail"); + + assert_eq!(err.code(), &S3ErrorCode::InvalidRequest); + } + + #[test] + fn manual_transition_job_response_exposes_status_and_cancel_contract() { + let record = ManualTransitionJobRecord::new(Uuid::new_v4(), "bucket", &ManualTransitionRunOptions::default(), "owner-a"); + let response = manual_transition_job_response(record); + + assert_eq!(response.status, ManualTransitionJobState::Running); + assert_eq!(response.mode, "durable_job"); + assert!(response.status_endpoint.ends_with(&response.job_id)); + assert_eq!(response.cancel_endpoint, response.status_endpoint); + } + + #[test] + fn manual_transition_active_job_cancel_token_round_trips() { + let job_id = Uuid::new_v4(); + let cancel_token = CancellationToken::new(); + insert_active_manual_transition_job(job_id, cancel_token.clone()); + + let active_cancel_token = active_manual_transition_cancel_token(job_id).expect("active job token should be registered"); + active_cancel_token.cancel(); + + assert!(cancel_token.is_cancelled()); + + remove_active_manual_transition_job(job_id); + assert!(active_manual_transition_cancel_token(job_id).is_none()); + } + + #[tokio::test] + 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"), + Params::new(), + ) + .await + .expect_err("status handler must reject unsigned requests"); + assert_eq!(status_err.code(), &S3ErrorCode::InvalidRequest); + assert_eq!(status_err.message(), Some("authentication required")); + + let cancel_err = ManualTransitionJobCancelHandler {} + .call( + manual_transition_job_request(Method::DELETE, "/rustfs/admin/v3/ilm/transition/jobs/job-123"), + Params::new(), + ) + .await + .expect_err("cancel handler must reject unsigned requests"); + assert_eq!(cancel_err.code(), &S3ErrorCode::InvalidRequest); + assert_eq!(cancel_err.message(), Some("authentication required")); + } + + #[test] + fn manual_transition_job_handlers_authorize_validate_and_load_store() { + let src = include_str!("ilm_transition.rs"); + let status_block = extract_block_between_markers( + src, + "impl Operation for ManualTransitionJobStatusHandler", + "pub struct ManualTransitionJobCancelHandler", + ); + let cancel_block = + extract_block_between_markers(src, "impl Operation for ManualTransitionJobCancelHandler", "#[cfg(test)]"); + + let status_load = status_block + .find("load_manual_transition_job_record") + .expect("status route must load the persisted job record"); + let cancel_load = cancel_block + .find("request_manual_transition_job_cancel") + .expect("cancel route must update the persisted job record"); + + for (block, load) in [(status_block, status_load), (cancel_block, cancel_load)] { + let auth = block + .find("authorize_manual_transition_request(&req).await?;") + .expect("job route must authorize with SetTierAction"); + let job_id = block + .find("manual_transition_job_id_from_params(¶ms)?;") + .expect("job route must validate the path job id"); + + assert!(auth < job_id); + assert!(job_id < load); + assert!(!block.contains("ServerInfoAdminAction")); + } + } + #[test] fn manual_transition_logs_masked_actor_and_aggregate_counters() { let src = include_str!("ilm_transition.rs"); @@ -1140,16 +1290,6 @@ mod tests { assert!(!log_block.contains("next_version_idmarker")); } - #[test] - fn manual_transition_background_job_uses_cancelable_scanner_entrypoint() { - let src = include_str!("ilm_transition.rs"); - let job_block = - extract_block_between_markers(src, "fn spawn_manual_transition_job", "pub struct ManualTransitionRunHandler"); - - assert!(job_block.contains("enqueue_transition_for_existing_objects_scoped_with_cancel")); - assert!(job_block.contains("CancellationToken")); - } - fn extract_block_between_markers<'a>(src: &'a str, start_marker: &str, end_marker: &str) -> &'a str { let start = src .find(start_marker) diff --git a/rustfs/src/admin/storage_api.rs b/rustfs/src/admin/storage_api.rs index 73611f152..f5fb811e7 100644 --- a/rustfs/src/admin/storage_api.rs +++ b/rustfs/src/admin/storage_api.rs @@ -193,10 +193,22 @@ pub(crate) mod bucket_target_sys { } pub(crate) mod lifecycle { + pub(crate) use super::ecstore_bucket::lifecycle::manual_transition_job::{ + ManualTransitionJobRecord, ManualTransitionJobState, ManualTransitionScopeAdmission, ManualTransitionScopeAdmissionClaim, + claim_manual_transition_scope_admission, delete_manual_transition_scope_admission_if_current, + load_manual_transition_job_record, load_manual_transition_job_record_with_etag, load_manual_transition_scope_admission, + manual_transition_job_lease_expired, manual_transition_scope_admission_lease_expired, + persist_manual_transition_job_progress, renew_manual_transition_job_lease, request_manual_transition_job_cancel, + save_manual_transition_job_record, save_manual_transition_job_record_if_current, + }; + pub(crate) type ManualTransitionCancelCheck = + super::ecstore_bucket::lifecycle::bucket_lifecycle_ops::ManualTransitionCancelCheck; + pub(crate) type ManualTransitionQueueSnapshot = + super::ecstore_bucket::lifecycle::bucket_lifecycle_ops::ManualTransitionQueueSnapshot; + pub(crate) type ManualTransitionProgressSink = + super::ecstore_bucket::lifecycle::bucket_lifecycle_ops::ManualTransitionProgressSink; pub(crate) type ManualTransitionRunOptions = super::ecstore_bucket::lifecycle::bucket_lifecycle_ops::ManualTransitionRunOptions; - pub(crate) type ManualTransitionRunExecution = - super::ecstore_bucket::lifecycle::bucket_lifecycle_ops::ManualTransitionRunExecution; pub(crate) type ManualTransitionRunReport = super::ecstore_bucket::lifecycle::bucket_lifecycle_ops::ManualTransitionRunReport; pub(crate) async fn enqueue_transition_for_existing_objects_scoped( @@ -210,19 +222,8 @@ pub(crate) mod lifecycle { .await } - pub(crate) async fn enqueue_transition_for_existing_objects_scoped_with_cancel( - api: std::sync::Arc, - bucket: &str, - options: ManualTransitionRunOptions, - cancel_token: Option, - ) -> super::Result { - super::ecstore_bucket::lifecycle::bucket_lifecycle_ops::enqueue_transition_for_existing_objects_scoped_with_cancel( - api, - bucket, - options, - cancel_token, - ) - .await + pub(crate) fn manual_transition_queue_snapshot() -> ManualTransitionQueueSnapshot { + super::ecstore_bucket::lifecycle::bucket_lifecycle_ops::manual_transition_queue_snapshot() } pub(crate) mod tier_last_day_stats {