diff --git a/crates/e2e_test/src/reliant/tiering.rs b/crates/e2e_test/src/reliant/tiering.rs index 97136c50f..38cd30fc9 100644 --- a/crates/e2e_test/src/reliant/tiering.rs +++ b/crates/e2e_test/src/reliant/tiering.rs @@ -73,9 +73,12 @@ 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_ASYNC_CONFLICT_BUCKET: &str = "ilm7-manual-async-conflict"; 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 MANUAL_ASYNC_CONFLICT_PREFIX: &str = "manual-async-conflict/"; +const MANUAL_ASYNC_CONFLICT_NESTED_PREFIX: &str = "manual-async-conflict/nested/"; const OBJECT_KEY: &str = "tier/鲁A12345/report.bin"; const MANUAL_DUE_KEY: &str = "manual-due/report.bin"; const MANUAL_DRY_RUN_KEY: &str = "manual-dry-run/report.bin"; @@ -367,6 +370,16 @@ struct ManualTransitionJobStatusResponse { report: ManualTransitionRunReport, } +#[derive(Debug, Deserialize)] +struct ManualTransitionJobConflictResponse { + state: String, + mode: String, + active_job_id: String, + status_endpoint: String, + cancel_endpoint: String, + scope_key: String, +} + async fn manual_transition_run( hot: &RustFSTestEnvironment, bucket: &str, @@ -418,15 +431,25 @@ async fn manual_transition_async_run( dry_run: bool, max_objects: u64, ) -> Result> { + let (status, body) = manual_transition_async_run_raw(hot, bucket, prefix, dry_run, max_objects).await?; + assert_eq!(status, reqwest::StatusCode::ACCEPTED, "async manual transition response: {body}"); + Ok(serde_json::from_str(&body)?) +} + +async fn manual_transition_async_run_raw( + hot: &RustFSTestEnvironment, + bucket: &str, + prefix: &str, + dry_run: bool, + max_objects: u64, +) -> Result<(reqwest::StatusCode, String), Box> { let bucket = urlencoding::encode(bucket); let prefix = urlencoding::encode(prefix); let tier = urlencoding::encode(TIER_NAME); let path = format!( "/rustfs/admin/v3/ilm/transition/run?bucket={bucket}&prefix={prefix}&tier={tier}&dryRun={dry_run}&maxObjects={max_objects}&mode=async" ); - let (status, body) = signed_admin_request(&hot.url, Method::POST, &path, None, &hot.access_key, &hot.secret_key).await?; - assert_eq!(status, reqwest::StatusCode::ACCEPTED, "async manual transition response: {body}"); - Ok(serde_json::from_str(&body)?) + signed_admin_request(&hot.url, Method::POST, &path, None, &hot.access_key, &hot.secret_key).await } async fn manual_transition_job_status( @@ -908,6 +931,91 @@ async fn test_manual_transition_async_limit_reports_terminal_partial() -> TestRe Ok(()) } +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn test_manual_transition_async_overlapping_scope_conflict_reports_active_job() -> TestResult { + let mut cold = RustFSTestEnvironment::new().await?; + cold.access_key = "manualasyncconflictcoldtieradmin".to_string(); + cold.secret_key = "manualasyncconflictcoldtiersecret".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_ASYNC_CONFLICT_BUCKET).send().await?; + for idx in 0..50 { + let key = format!("{MANUAL_ASYNC_CONFLICT_NESTED_PREFIX}obj-{idx:02}"); + put_single_part_object(&hot_client, MANUAL_ASYNC_CONFLICT_BUCKET, &key, b"async conflict payload").await?; + } + put_lifecycle_transition_rule( + &hot_client, + MANUAL_ASYNC_CONFLICT_BUCKET, + "manual-async-conflict", + MANUAL_ASYNC_CONFLICT_PREFIX, + 0, + ) + .await?; + + let (first, second) = tokio::join!( + manual_transition_async_run_raw(&hot, MANUAL_ASYNC_CONFLICT_BUCKET, MANUAL_ASYNC_CONFLICT_PREFIX, false, 50), + manual_transition_async_run_raw(&hot, MANUAL_ASYNC_CONFLICT_BUCKET, MANUAL_ASYNC_CONFLICT_NESTED_PREFIX, false, 50) + ); + let responses = [first?, second?]; + let accepted = responses + .iter() + .find(|(status, _)| *status == reqwest::StatusCode::ACCEPTED) + .ok_or("one concurrent async run must be accepted")?; + let conflict = responses + .iter() + .find(|(status, _)| *status == reqwest::StatusCode::CONFLICT) + .ok_or("one concurrent async run must report conflict")?; + + let accepted: ManualTransitionRunResponse = serde_json::from_str(&accepted.1)?; + let conflict: ManualTransitionJobConflictResponse = serde_json::from_str(&conflict.1)?; + let job_id = accepted + .job_id + .as_deref() + .ok_or("accepted async response must include job_id")?; + let status_endpoint = accepted + .status_endpoint + .as_deref() + .ok_or("accepted async response must include status_endpoint")?; + + assert_eq!(accepted.state, "accepted"); + assert_eq!(accepted.mode, "durable_job"); + assert_eq!(conflict.state, "conflict"); + assert_eq!(conflict.mode, "durable_job"); + assert_eq!(conflict.active_job_id, job_id); + assert_eq!(conflict.status_endpoint, status_endpoint); + assert_eq!(conflict.cancel_endpoint, status_endpoint); + assert!(!conflict.scope_key.is_empty()); + + let terminal = wait_for_manual_transition_job_terminal(&hot, status_endpoint, StdDuration::from_secs(30)).await?; + assert_eq!(terminal.job_id, job_id); + assert_eq!(terminal.status, "completed", "terminal conflict winner response: {terminal:#?}"); + assert!(!terminal.report.dry_run); + assert_eq!(terminal.report.bucket, MANUAL_ASYNC_CONFLICT_BUCKET); + assert!( + terminal.report.prefix == MANUAL_ASYNC_CONFLICT_PREFIX || terminal.report.prefix == MANUAL_ASYNC_CONFLICT_NESTED_PREFIX, + "terminal conflict winner response: {terminal:#?}" + ); + assert_eq!(terminal.report.scanned, 50, "terminal conflict winner response: {terminal:#?}"); + assert_eq!(terminal.report.eligible, 50, "terminal conflict winner response: {terminal:#?}"); + assert_eq!(terminal.report.dry_run_eligible, 0, "terminal conflict winner response: {terminal:#?}"); + assert_eq!( + terminal.report.enqueued + terminal.report.skipped_already_in_flight, + 50, + "terminal conflict winner response: {terminal:#?}" + ); + assert!(cold_tier_object_count(&cold_client).await? <= 50); + + Ok(()) +} + #[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn test_manual_transition_run_contract_no_status_cancel_fields() -> TestResult { let mut cold = RustFSTestEnvironment::new().await?; diff --git a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs index 95dd32de8..6f23b663d 100644 --- a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs +++ b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs @@ -4158,10 +4158,11 @@ pub async fn apply_lifecycle_action(event: &lifecycle::Event, src: &LcEventSrc, mod tests { use super::{ DATE_EXPIRY_EXISTING_OBJECTS_GRACE_SECS, DEFAULT_TRANSITION_QUEUE_CAPACITY, DEFAULT_TRANSITION_WORKERS_ABSOLUTE_MAX, - DEFAULT_TRANSITION_WORKERS_CAP, ExpiryState, FreeVersionTask, ManualTransitionRunOptions, ManualTransitionRunReport, - StaleMultipartUploadCandidate, TIER_FREE_VERSION_RECOVERY_BASE_INTERVAL, TIER_FREE_VERSION_RECOVERY_MAX_IDLE_INTERVAL, - TRANSITION_COMPLETE, TierFreeVersionRecoverySchedule, TransitionEnqueueOutcome, TransitionState, TransitionedObject, - VersionReplicationScan, cleanup_empty_multipart_sha_dirs_on_local_disks, cleanup_stale_multipart_uploads_once_at, + DEFAULT_TRANSITION_WORKERS_CAP, ExpiryState, FreeVersionTask, ManualTransitionQueueSnapshot, ManualTransitionRunOptions, + ManualTransitionRunReport, StaleMultipartUploadCandidate, TIER_FREE_VERSION_RECOVERY_BASE_INTERVAL, + TIER_FREE_VERSION_RECOVERY_MAX_IDLE_INTERVAL, TRANSITION_COMPLETE, TierFreeVersionRecoverySchedule, + TransitionEnqueueOutcome, TransitionState, TransitionedObject, VersionReplicationScan, + cleanup_empty_multipart_sha_dirs_on_local_disks, cleanup_stale_multipart_uploads_once_at, 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, @@ -4182,8 +4183,9 @@ mod tests { 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, + ManualTransitionJobRecord, ManualTransitionJobState, ManualTransitionScopeAdmission, ManualTransitionScopeAdmissionClaim, + claim_manual_transition_scope_admission, legacy_manual_transition_scope_key, load_manual_transition_job_record, + load_manual_transition_scope_admission, request_manual_transition_job_cancel, save_manual_transition_job_record, save_manual_transition_scope_admission_if_absent, }; use crate::bucket::lifecycle::replication_sink::{ @@ -7031,7 +7033,7 @@ mod tests { .await .expect("manual transition recovery should run"); - assert_eq!(stats.scanned, 1); + assert!(stats.scanned >= 1); assert_eq!(stats.resumed, 1); assert_eq!(stats.failed, 0); let recovered = load_manual_transition_job_record(ecstore.clone(), job_id) @@ -7064,6 +7066,106 @@ mod tests { assert!(!report.has_partial_enqueue()); } + #[tokio::test] + #[serial] + async fn manual_transition_job_cancel_marks_running_record_only() { + let (_paths, ecstore) = setup_test_env().await; + let running_id = Uuid::new_v4(); + let running = ManualTransitionJobRecord::new( + running_id, + "manual-cancel-running-bucket", + &ManualTransitionRunOptions::default(), + "owner-a", + ); + save_manual_transition_job_record(ecstore.clone(), &running) + .await + .expect("running job record should save"); + + let cancelled = request_manual_transition_job_cancel(ecstore.clone(), running_id) + .await + .expect("running job cancel request should persist"); + + assert_eq!(cancelled.state, ManualTransitionJobState::Running); + assert!(cancelled.cancel_requested); + let loaded = load_manual_transition_job_record(ecstore.clone(), running_id) + .await + .expect("cancelled running job should reload"); + assert!(loaded.cancel_requested); + + let terminal_id = Uuid::new_v4(); + let mut terminal = ManualTransitionJobRecord::new( + terminal_id, + "manual-cancel-terminal-bucket", + &ManualTransitionRunOptions::default(), + "owner-a", + ); + terminal.complete( + ManualTransitionRunReport { + bucket: "manual-cancel-terminal-bucket".to_string(), + ..Default::default() + }, + ManualTransitionQueueSnapshot::default(), + ); + save_manual_transition_job_record(ecstore.clone(), &terminal) + .await + .expect("terminal job record should save"); + + let after_terminal_cancel = request_manual_transition_job_cancel(ecstore, terminal_id) + .await + .expect("terminal job cancel should be idempotent"); + + assert_eq!(after_terminal_cancel.state, ManualTransitionJobState::Completed); + assert!(!after_terminal_cancel.cancel_requested); + } + + #[tokio::test] + #[serial] + async fn manual_transition_admission_blocks_active_legacy_scope_record() { + let (_paths, ecstore) = setup_test_env().await; + let legacy_options = ManualTransitionRunOptions { + prefix: "logs/".to_string(), + tier: Some("warm".to_string()), + ..Default::default() + }; + let legacy_id = Uuid::new_v4(); + let mut legacy = ManualTransitionJobRecord::new(legacy_id, "manual-legacy-scope-bucket", &legacy_options, "old-owner"); + legacy.scope_key = legacy_manual_transition_scope_key(&legacy.bucket, &legacy_options); + save_manual_transition_job_record(ecstore.clone(), &legacy) + .await + .expect("legacy job record should save"); + save_manual_transition_scope_admission_if_absent(ecstore.clone(), &ManualTransitionScopeAdmission::from_job(&legacy)) + .await + .expect("legacy scope admission should save"); + + let new_options = ManualTransitionRunOptions { + prefix: "archive/".to_string(), + tier: Some("cold".to_string()), + ..Default::default() + }; + let new_record = ManualTransitionJobRecord::new(Uuid::new_v4(), "manual-legacy-scope-bucket", &new_options, "new-owner"); + save_manual_transition_job_record(ecstore.clone(), &new_record) + .await + .expect("new job record should save"); + + let claim = + claim_manual_transition_scope_admission(ecstore.clone(), &ManualTransitionScopeAdmission::from_job(&new_record)) + .await + .expect("new bucket-level admission claim should resolve"); + + let ManualTransitionScopeAdmissionClaim::Conflict(active) = claim else { + panic!("active legacy job must block new bucket-level admission"); + }; + assert_eq!(active.job_id, legacy_id); + assert_eq!(active.scope_key, legacy.scope_key); + assert!( + matches!( + load_manual_transition_scope_admission(ecstore, &new_record.scope_key).await, + Err(Error::ConfigNotFound) + ), + "conflicted bucket-level admission must be released" + ); + } + #[tokio::test] async fn existing_object_lifecycle_allows_expired_marker_after_replication_completed() { let lc = expired_delete_marker_lifecycle(); diff --git a/crates/ecstore/src/bucket/lifecycle/manual_transition_job.rs b/crates/ecstore/src/bucket/lifecycle/manual_transition_job.rs index 45917ef06..2c35e2594 100644 --- a/crates/ecstore/src/bucket/lifecycle/manual_transition_job.rs +++ b/crates/ecstore/src/bucket/lifecycle/manual_transition_job.rs @@ -23,8 +23,10 @@ use crate::bucket::lifecycle::bucket_lifecycle_ops::{ ManualTransitionQueueSnapshot, ManualTransitionRunOptions, ManualTransitionRunReport, }; use crate::bucket::lifecycle::config_boundary; +use crate::disk::RUSTFS_META_BUCKET; use crate::error::{Error, Result as EcstoreResult}; use crate::object_api::ObjectOptions; +use crate::storage_api_contracts::list::ListOperations as _; use crate::storage_api_contracts::object::HTTPPreconditions; use crate::store::ECStore; @@ -33,6 +35,7 @@ pub const MANUAL_TRANSITION_JOB_RECORD_PREFIX: &str = "ilm/manual-transition/job 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; +const MANUAL_TRANSITION_LEGACY_SCOPE_SCAN_LIMIT: i32 = 1000; #[derive(Debug, thiserror::Error)] pub enum ManualTransitionJobError { @@ -351,6 +354,16 @@ pub enum ManualTransitionScopeAdmissionClaim { } pub fn manual_transition_scope_key(bucket: &str, options: &ManualTransitionRunOptions) -> String { + let mut scope = String::new(); + scope.push_str(bucket); + scope.push('\0'); + // Durable admission v1 is bucket-level: a single CAS alias must cover prefix + // overlap and wildcard-tier conflicts without a non-atomic scope scan. + scope.push_str(if options.dry_run { "dry_run" } else { "run" }); + hex_sha256(scope.as_bytes(), ToOwned::to_owned) +} + +pub(crate) fn legacy_manual_transition_scope_key(bucket: &str, options: &ManualTransitionRunOptions) -> String { let mut scope = String::new(); scope.push_str(bucket); scope.push('\0'); @@ -553,7 +566,7 @@ pub async fn claim_manual_transition_scope_admission( admission: &ManualTransitionScopeAdmission, ) -> EcstoreResult { match save_manual_transition_scope_admission_if_absent(api.clone(), admission).await { - Ok(()) => return Ok(ManualTransitionScopeAdmissionClaim::Claimed), + Ok(()) => return finish_manual_transition_scope_admission_claim(api, admission).await, Err(Error::PreconditionFailed) => {} Err(err) => return Err(err), } @@ -572,8 +585,8 @@ pub async fn claim_manual_transition_scope_admission( } }; if active_job_reclaimable { - return match save_manual_transition_scope_admission_if_current(api, admission, &etag).await { - Ok(()) => Ok(ManualTransitionScopeAdmissionClaim::Claimed), + return match save_manual_transition_scope_admission_if_current(api.clone(), admission, &etag).await { + Ok(()) => finish_manual_transition_scope_admission_claim(api, admission).await, Err(Error::PreconditionFailed) => Ok(ManualTransitionScopeAdmissionClaim::Conflict(Box::new(active))), Err(err) => Err(err), }; @@ -582,6 +595,78 @@ pub async fn claim_manual_transition_scope_admission( Ok(ManualTransitionScopeAdmissionClaim::Conflict(Box::new(active))) } +async fn finish_manual_transition_scope_admission_claim( + api: Arc, + admission: &ManualTransitionScopeAdmission, +) -> EcstoreResult { + if let Some(active) = find_active_legacy_manual_transition_scope_conflict(api.clone(), admission).await? { + delete_manual_transition_scope_admission_if_current(api, &admission.scope_key, admission.job_id, admission.lease_id) + .await?; + return Ok(ManualTransitionScopeAdmissionClaim::Conflict(Box::new(active))); + } + Ok(ManualTransitionScopeAdmissionClaim::Claimed) +} + +async fn find_active_legacy_manual_transition_scope_conflict( + api: Arc, + admission: &ManualTransitionScopeAdmission, +) -> EcstoreResult> { + let mut marker = None; + loop { + let page = api + .clone() + .list_objects_v2( + RUSTFS_META_BUCKET, + MANUAL_TRANSITION_JOB_RECORD_PREFIX, + marker, + None, + MANUAL_TRANSITION_LEGACY_SCOPE_SCAN_LIMIT, + false, + None, + false, + ) + .await?; + for object in page.objects { + let job_id = + manual_transition_job_id_from_record_object_name(&object.name).map_err(manual_transition_job_store_error)?; + if job_id == admission.job_id { + continue; + } + let record = match load_manual_transition_job_record(api.clone(), job_id).await { + Ok(record) => record, + Err(Error::ConfigNotFound) => continue, + Err(err) => return Err(err), + }; + if record.scope_key == admission.scope_key { + continue; + } + let legacy_scope_key = legacy_manual_transition_scope_key( + &record.bucket, + &ManualTransitionRunOptions { + prefix: record.prefix.clone(), + tier: record.tier.clone(), + dry_run: record.dry_run, + ..Default::default() + }, + ); + if record.scope_key != legacy_scope_key { + continue; + } + if record.bucket == admission.bucket + && record.dry_run == admission.dry_run + && !record.is_terminal() + && !manual_transition_job_lease_expired(&record) + { + return Ok(Some(ManualTransitionScopeAdmission::from_job(&record))); + } + } + if !page.is_truncated { + return Ok(None); + } + marker = page.next_continuation_token; + } +} + 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?; @@ -876,6 +961,46 @@ mod tests { ); } + #[test] + fn manual_transition_scope_key_uses_bucket_level_admission_for_durable_v1() { + let broad = manual_transition_scope_key( + "bucket", + &ManualTransitionRunOptions { + prefix: "logs/".to_string(), + tier: None, + ..Default::default() + }, + ); + let nested = manual_transition_scope_key( + "bucket", + &ManualTransitionRunOptions { + prefix: "logs/2026/".to_string(), + tier: Some("warm".to_string()), + ..Default::default() + }, + ); + let disjoint = manual_transition_scope_key( + "bucket", + &ManualTransitionRunOptions { + prefix: "archive/".to_string(), + tier: Some("cold".to_string()), + ..Default::default() + }, + ); + let dry_run = manual_transition_scope_key( + "bucket", + &ManualTransitionRunOptions { + prefix: "logs/".to_string(), + dry_run: true, + ..Default::default() + }, + ); + + assert_eq!(broad, nested); + assert_eq!(broad, disjoint); + assert_ne!(broad, dry_run); + } + #[test] fn manual_transition_scope_admission_carries_job_lease_fence() { let options = ManualTransitionRunOptions::default();