diff --git a/crates/ecstore/src/api/mod.rs b/crates/ecstore/src/api/mod.rs index e885be5b0..db4a082c7 100644 --- a/crates/ecstore/src/api/mod.rs +++ b/crates/ecstore/src/api/mod.rs @@ -43,12 +43,12 @@ pub mod bucket { pub mod bucket_lifecycle_ops { pub use crate::bucket::lifecycle::bucket_lifecycle_ops::{ - ExpiryState, LifecycleOps, ManualTransitionRunOptions, ManualTransitionRunReport, RestoreRequestOps, - TransitionState, TransitionedObject, apply_expiry_rule, apply_transition_rule, + ExpiryState, LifecycleOps, 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_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, + 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, }; } diff --git a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs index 682739472..9deaec6d3 100644 --- a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs +++ b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs @@ -2445,7 +2445,7 @@ pub struct ManualTransitionRunOptions { pub max_duration: Option, } -#[derive(Debug, Clone, Default, PartialEq, Eq, serde::Serialize)] +#[derive(Debug, Clone, Default, PartialEq, Eq, serde::Serialize, serde::Deserialize)] pub struct ManualTransitionRunReport { pub bucket: String, pub prefix: String, @@ -2474,6 +2474,12 @@ pub struct ManualTransitionRunReport { pub next_version_idmarker: Option, } +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct ManualTransitionRunExecution { + pub report: ManualTransitionRunReport, + pub cancelled: bool, +} + impl ManualTransitionRunReport { fn new(bucket: &str, options: &ManualTransitionRunOptions) -> Self { Self { @@ -2522,11 +2528,25 @@ pub async fn enqueue_transition_for_existing_objects_scoped( 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 Some(lc) = runtime_sources::bucket_lifecycle_config(bucket).await else { - return Ok(report); + return Ok(ManualTransitionRunExecution { + report, + cancelled: false, + }); }; report.lifecycle_config_found = true; let mut marker = options.marker.clone(); @@ -2537,24 +2557,39 @@ pub async fn enqueue_transition_for_existing_objects_scoped( 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()) { + report.next_marker.clone_from(&previous_marker); + report.next_version_idmarker.clone_from(&previous_version_marker); + return Ok(ManualTransitionRunExecution { report, cancelled: true }); + } 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(report); + return Ok(ManualTransitionRunExecution { + report, + cancelled: false, + }); } 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(report); + return Ok(ManualTransitionRunExecution { + report, + cancelled: false, + }); } 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) { @@ -2562,20 +2597,29 @@ pub async fn enqueue_transition_for_existing_objects_scoped( report.next_marker = Some(object.name.clone()); report.next_version_idmarker = Some(manual_transition_version_marker(object)); } - return Ok(report); + return Ok(ManualTransitionRunExecution { + report, + cancelled: false, + }); } previous_marker = Some(object.name.clone()); previous_version_marker = Some(manual_transition_version_marker(object)); } if !page.is_truncated { - return Ok(report); + return Ok(ManualTransitionRunExecution { + report, + cancelled: false, + }); } 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(report); + return Ok(ManualTransitionRunExecution { + report, + cancelled: false, + }); } marker = page.next_marker; @@ -2779,6 +2823,10 @@ 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) +} + fn manual_transition_version_marker(oi: &ObjectInfo) -> String { oi.version_id .map(|version| version.to_string()) @@ -3694,7 +3742,7 @@ 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, + lifecycle_rule_has_date_expiration, lifecycle_version_purge_state_from_completed_targets, manual_transition_cancelled, 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, resolve_transition_queue_capacity, resolve_transition_queue_send_timeout, resolve_transition_worker_count, @@ -6446,6 +6494,18 @@ mod tests { assert!(report.was_truncated()); } + #[test] + fn manual_transition_cancel_token_detects_cancelled_state() { + let cancel_token = CancellationToken::new(); + + assert!(!manual_transition_cancelled(None)); + assert!(!manual_transition_cancelled(Some(&cancel_token))); + + cancel_token.cancel(); + + assert!(manual_transition_cancelled(Some(&cancel_token))); + } + #[tokio::test] async fn existing_object_lifecycle_allows_expired_marker_after_replication_completed() { let lc = expired_delete_marker_lifecycle(); diff --git a/rustfs/src/admin/handlers/ilm_transition.rs b/rustfs/src/admin/handlers/ilm_transition.rs index fdbc2b41b..2ea82fc7d 100644 --- a/rustfs/src/admin/handlers/ilm_transition.rs +++ b/rustfs/src/admin/handlers/ilm_transition.rs @@ -16,9 +16,13 @@ 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, }; +use crate::admin::storage_api::runtime::ECStore; use crate::auth::{check_key_valid, get_session_token}; use crate::server::{ADMIN_PREFIX, RemoteAddr}; use http::{HeaderMap, HeaderValue}; @@ -32,8 +36,13 @@ use rustfs_utils::{ use s3s::header::CONTENT_TYPE; use s3s::{Body, S3Error, S3ErrorCode, S3Request, S3Response, S3Result, s3_error}; use serde::{Deserialize, Serialize}; -use std::sync::{Mutex, MutexGuard, OnceLock}; +use std::collections::HashMap; +use std::sync::{Arc, Mutex, MutexGuard, OnceLock}; +use time::OffsetDateTime; +use time::format_description::well_known::Rfc3339; +use tokio_util::sync::CancellationToken; use tracing::{error, info, warn}; +use uuid::Uuid; const JSON_CONTENT_TYPE: &str = "application/json"; const DEFAULT_MANUAL_TRANSITION_MAX_OBJECTS: u64 = 10_000; @@ -42,8 +51,11 @@ 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(); #[derive(Debug, Clone, PartialEq, Eq)] struct ManualTransitionRunScope { @@ -139,6 +151,21 @@ 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, @@ -148,16 +175,71 @@ struct ManualTransitionRunResponse { 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, + job_id: String, + status_endpoint: String, + status: ManualTransitionJobStatus, + 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, + report: ManualTransitionRunReport, +} + pub fn register_ilm_transition_route(r: &mut S3Router) -> std::io::Result<()> { r.insert( Method::POST, format!("{ADMIN_PREFIX}/v3/ilm/transition/run").as_str(), AdminOperation(&ManualTransitionRunHandler {}), )?; + r.insert( + Method::GET, + format!("{ADMIN_PREFIX}/v3/ilm/transition/jobs/{{job_id}}").as_str(), + AdminOperation(&ManualTransitionJobStatusHandler {}), + )?; + r.insert( + Method::DELETE, + format!("{ADMIN_PREFIX}/v3/ilm/transition/jobs/{{job_id}}").as_str(), + AdminOperation(&ManualTransitionJobCancelHandler {}), + )?; Ok(()) } -fn parse_manual_transition_query(query: Option<&str>) -> S3Result<(String, ManualTransitionRunOptions)> { +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()) .map_err(|_| s3_error!(InvalidArgument, "invalid manual transition query"))?, @@ -184,12 +266,11 @@ fn parse_manual_transition_query(query: Option<&str>) -> S3Result<(String, Manua if mode.is_some_and(|mode| mode != "enqueue_only" && mode != "async") { return Err(s3_error!(InvalidArgument, "unsupported manual transition mode")); } - if query.async_mode == Some(true) || mode == Some("async") { - return Err(S3Error::with_message( - S3ErrorCode::NotImplemented, - "durable manual transition jobs are not implemented; omit async/mode for enqueue_only", - )); - } + let run_mode = if query.async_mode == Some(true) || mode == Some("async") { + ManualTransitionRunMode::Async + } else { + ManualTransitionRunMode::EnqueueOnly + }; let max_objects = query.max_objects.unwrap_or(DEFAULT_MANUAL_TRANSITION_MAX_OBJECTS); if max_objects == 0 || max_objects > MAX_MANUAL_TRANSITION_OBJECTS { @@ -213,6 +294,7 @@ fn parse_manual_transition_query(query: Option<&str>) -> S3Result<(String, Manua max_objects: Some(max_objects), max_duration: query.max_duration_seconds.map(std::time::Duration::from_secs), }, + run_mode, )) } @@ -341,7 +423,7 @@ fn response_state(report: &ManualTransitionRunReport) -> &'static str { } } -fn json_response(response: &ManualTransitionRunResponse) -> S3Result> { +fn json_response(status: StatusCode, response: &T) -> S3Result> { let body = serde_json::to_vec(response).map_err(|err| { S3Error::with_message(S3ErrorCode::InternalError, format!("failed to encode manual transition response: {err}")) })?; @@ -349,7 +431,223 @@ fn json_response(response: &ManualTransitionRunResponse) -> 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(), + } +} + +fn manual_transition_job_status_endpoint(job_id: &str) -> String { + format!("{ADMIN_PREFIX}/v3/ilm/transition/jobs/{job_id}") +} + +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 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}")) + }) +} + +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}"), + )), + } +} + +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(mut record) = read_manual_transition_job_record(store, job_id).await? else { + return Err(s3_error!(NoSuchKey, "manual transition job not found")); + }; + 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()); + } + Ok(record) +} + +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 +} + +fn spawn_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; + } + + 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); + } + drop(admission_guard); + }); } pub struct ManualTransitionRunHandler {} @@ -360,7 +658,7 @@ impl Operation for ManualTransitionRunHandler { let request_id = admin_request_id(&req.headers).unwrap_or_default().to_string(); let remote_addr = admin_remote_addr(&req).unwrap_or_default(); let actor = authorize_manual_transition_request(&req).await?; - let (bucket, options) = match parse_manual_transition_query(req.uri.query()) { + let (bucket, options, run_mode) = match parse_manual_transition_query(req.uri.query()) { Ok(parsed) => parsed, Err(err) => { log_manual_transition_rejected("invalid_query_parameters", &request_id, &actor, &remote_addr); @@ -374,7 +672,7 @@ impl Operation for ManualTransitionRunHandler { 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); @@ -382,6 +680,24 @@ 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) => { @@ -402,7 +718,59 @@ impl Operation for ManualTransitionRunHandler { report, }; - json_response(&response) + json_response(StatusCode::OK, &response) + } +} + +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"))?, + )?; + 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) + } +} + +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"))?, + )?; + 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.status = ManualTransitionJobStatus::Unknown; + record.cancel_requested = true; + 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()); + save_manual_transition_job_record(store, &record).await?; + } + record + }; + json_response(StatusCode::OK, &record) } } @@ -412,11 +780,12 @@ mod tests { #[test] fn manual_transition_query_defaults_to_bounded_run() { - let (bucket, options) = + let (bucket, options, run_mode) = parse_manual_transition_query(Some("bucket=data&prefix=logs/&marker=logs/a&versionMarker=v1&tier=warm")) .expect("valid query should parse"); assert_eq!(bucket, "data"); + assert_eq!(run_mode, ManualTransitionRunMode::EnqueueOnly); assert_eq!(options.prefix, "logs/"); assert_eq!(options.marker.as_deref(), Some("logs/a")); assert_eq!(options.version_marker.as_deref(), Some("v1")); @@ -428,32 +797,34 @@ mod tests { #[test] fn manual_transition_query_accepts_duration_budget() { - let (_bucket, options) = + let (_bucket, options, run_mode) = parse_manual_transition_query(Some("bucket=data&maxDurationSeconds=30")).expect("valid query should parse"); + assert_eq!(run_mode, ManualTransitionRunMode::EnqueueOnly); assert_eq!(options.max_duration, Some(std::time::Duration::from_secs(30))); } #[test] fn manual_transition_query_accepts_explicit_enqueue_only_mode() { - let (_bucket, options) = parse_manual_transition_query(Some("bucket=data&mode=enqueue_only&async=false")) + let (_bucket, options, run_mode) = parse_manual_transition_query(Some("bucket=data&mode=enqueue_only&async=false")) .expect("explicit enqueue_only mode should remain compatible"); + assert_eq!(run_mode, ManualTransitionRunMode::EnqueueOnly); assert!(!options.dry_run); assert_eq!(options.max_objects, Some(DEFAULT_MANUAL_TRANSITION_MAX_OBJECTS)); } #[test] - fn manual_transition_query_rejects_unimplemented_durable_mode() { - let err = parse_manual_transition_query(Some("bucket=data&async=true")) - .expect_err("async durable jobs must not be silently accepted"); + 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"); - assert_eq!(err.code(), &S3ErrorCode::NotImplemented); + assert_eq!(run_mode, ManualTransitionRunMode::Async); - let err = parse_manual_transition_query(Some("bucket=data&mode=async")) - .expect_err("mode=async must not fall back to enqueue_only"); + let (_bucket, _options, run_mode) = + parse_manual_transition_query(Some("bucket=data&mode=async")).expect("mode=async should parse"); - assert_eq!(err.code(), &S3ErrorCode::NotImplemented); + assert_eq!(run_mode, ManualTransitionRunMode::Async); } #[test] @@ -479,11 +850,11 @@ mod tests { #[test] fn manual_transition_scope_ignores_resume_and_budget_parameters() { - let (bucket, first) = parse_manual_transition_query(Some( + let (bucket, first, _run_mode) = parse_manual_transition_query(Some( "bucket=data&prefix=logs/&tier=warm&marker=logs/a&versionMarker=v1&maxObjects=10", )) .expect("first query should parse"); - let (_, second) = parse_manual_transition_query(Some( + let (_, second, _run_mode) = parse_manual_transition_query(Some( "bucket=data&prefix=logs/&tier=WARM&marker=logs/z&versionMarker=v9&maxObjects=20", )) .expect("second query should parse"); @@ -496,9 +867,9 @@ mod tests { #[test] fn manual_transition_scope_distinguishes_dry_run_mode() { - let (bucket, real) = + let (bucket, real, _run_mode) = parse_manual_transition_query(Some("bucket=data&prefix=logs/&tier=warm")).expect("real query should parse"); - let (_, dry_run) = parse_manual_transition_query(Some("bucket=data&prefix=logs/&tier=warm&dryRun=true")) + let (_, dry_run, _run_mode) = parse_manual_transition_query(Some("bucket=data&prefix=logs/&tier=warm&dryRun=true")) .expect("dry-run query should parse"); assert_ne!( @@ -509,7 +880,7 @@ mod tests { #[test] fn manual_transition_admission_rejects_same_scope_until_guard_drops() { - let (bucket, options) = + let (bucket, options, _run_mode) = parse_manual_transition_query(Some("bucket=admission-test&prefix=logs/&tier=warm")).expect("query should parse"); let scope = ManualTransitionRunScope::new(&bucket, &options); let first = acquire_manual_transition_admission(scope.clone()).expect("first admission should succeed"); @@ -536,7 +907,7 @@ mod tests { #[test] fn manual_transition_admission_rejects_overlapping_prefix_or_tier() { - let (bucket, options) = + let (bucket, options, _run_mode) = parse_manual_transition_query(Some("bucket=admission-overlap-test&prefix=logs/")).expect("query should parse"); let scope = ManualTransitionRunScope::new(&bucket, &options); let active = acquire_manual_transition_admission(scope).expect("first admission should succeed"); @@ -647,6 +1018,53 @@ mod tests { 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_handler_requires_set_tier_action() { let src = include_str!("ilm_transition.rs"); @@ -679,6 +1097,16 @@ 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/handlers/mod.rs b/rustfs/src/admin/handlers/mod.rs index 17874bf81..b18f541a8 100644 --- a/rustfs/src/admin/handlers/mod.rs +++ b/rustfs/src/admin/handlers/mod.rs @@ -121,6 +121,8 @@ mod tests { let _remove_remote_target_handler = replication::RemoveRemoteTargetHandler {}; let _scanner_status_handler = scanner::ScannerStatusHandler {}; let _manual_transition_handler = ilm_transition::ManualTransitionRunHandler {}; + let _manual_transition_status_handler = ilm_transition::ManualTransitionJobStatusHandler {}; + let _manual_transition_cancel_handler = ilm_transition::ManualTransitionJobCancelHandler {}; let _site_replication_add_handler = site_replication::SiteReplicationAddHandler {}; let _site_replication_info_handler = site_replication::SiteReplicationInfoHandler {}; let _site_replication_status_handler = site_replication::SiteReplicationStatusHandler {}; diff --git a/rustfs/src/admin/route_policy.rs b/rustfs/src/admin/route_policy.rs index e1ef8bf97..3df6d2ebe 100644 --- a/rustfs/src/admin/route_policy.rs +++ b/rustfs/src/admin/route_policy.rs @@ -414,6 +414,18 @@ pub const ADMIN_ROUTE_POLICY_SPECS: &[AdminRouteSpec] = &[ admin(HttpMethod::Put, "/rustfs/admin/v3/config", CONFIG_UPDATE, RouteRiskLevel::High), admin(HttpMethod::Get, "/rustfs/admin/v3/scanner/status", SERVER_INFO, RouteRiskLevel::Sensitive), admin(HttpMethod::Post, "/rustfs/admin/v3/ilm/transition/run", SET_TIER, RouteRiskLevel::High), + admin( + HttpMethod::Get, + "/rustfs/admin/v3/ilm/transition/jobs/{job_id}", + SET_TIER, + RouteRiskLevel::High, + ), + admin( + HttpMethod::Delete, + "/rustfs/admin/v3/ilm/transition/jobs/{job_id}", + SET_TIER, + RouteRiskLevel::High, + ), admin( HttpMethod::Get, "/rustfs/admin/v3/audit/target/list", @@ -1803,7 +1815,11 @@ mod tests { #[test] fn route_policy_requires_set_tier_for_manual_transition_run() { assert_action(HttpMethod::Post, "/rustfs/admin/v3/ilm/transition/run", SET_TIER); + assert_action(HttpMethod::Get, "/rustfs/admin/v3/ilm/transition/jobs/{job_id}", SET_TIER); + assert_action(HttpMethod::Delete, "/rustfs/admin/v3/ilm/transition/jobs/{job_id}", SET_TIER); assert_not_action(HttpMethod::Post, "/rustfs/admin/v3/ilm/transition/run", SERVER_INFO); + assert_not_action(HttpMethod::Get, "/rustfs/admin/v3/ilm/transition/jobs/{job_id}", SERVER_INFO); + assert_not_action(HttpMethod::Delete, "/rustfs/admin/v3/ilm/transition/jobs/{job_id}", SERVER_INFO); } #[test] diff --git a/rustfs/src/admin/route_registration_test.rs b/rustfs/src/admin/route_registration_test.rs index 5b02b9f9e..aef6c049e 100644 --- a/rustfs/src/admin/route_registration_test.rs +++ b/rustfs/src/admin/route_registration_test.rs @@ -196,6 +196,16 @@ fn expected_admin_route_matrix() -> Vec { admin_route_sample(Method::POST, "/v3/tier/{tiername}", "/v3/tier/HOT"), admin_route(Method::POST, "/v3/tier/clear"), admin_route(Method::POST, "/v3/ilm/transition/run"), + admin_route_sample( + Method::GET, + "/v3/ilm/transition/jobs/{job_id}", + "/v3/ilm/transition/jobs/11111111-1111-4111-8111-111111111111", + ), + admin_route_sample( + Method::DELETE, + "/v3/ilm/transition/jobs/{job_id}", + "/v3/ilm/transition/jobs/11111111-1111-4111-8111-111111111111", + ), admin_route(Method::PUT, "/v3/set-bucket-quota"), admin_route(Method::GET, "/v3/get-bucket-quota"), admin_route_sample(Method::PUT, "/v3/quota/{bucket}", "/v3/quota/test-bucket"), @@ -834,6 +844,16 @@ fn test_register_routes_cover_representative_admin_paths() { assert_route(&router, Method::PUT, &admin_path("/v3/config")); assert_route(&router, Method::GET, &admin_path("/v3/scanner/status")); assert_route(&router, Method::POST, &admin_path("/v3/ilm/transition/run")); + assert_route( + &router, + Method::GET, + &admin_path("/v3/ilm/transition/jobs/11111111-1111-4111-8111-111111111111"), + ); + assert_route( + &router, + Method::DELETE, + &admin_path("/v3/ilm/transition/jobs/11111111-1111-4111-8111-111111111111"), + ); assert_route(&router, Method::GET, &table_catalog_path("/config")); assert_route(&router, Method::PUT, &table_catalog_path("/buckets/analytics")); diff --git a/rustfs/src/admin/storage_api.rs b/rustfs/src/admin/storage_api.rs index fa21b756e..dabd5724f 100644 --- a/rustfs/src/admin/storage_api.rs +++ b/rustfs/src/admin/storage_api.rs @@ -195,6 +195,8 @@ pub(crate) mod bucket_target_sys { pub(crate) mod lifecycle { 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( @@ -208,6 +210,21 @@ 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) mod tier_last_day_stats { #[cfg(test)] pub(crate) type LastDayTierStats = super::super::ecstore_bucket::lifecycle::tier_last_day_stats::LastDayTierStats;