diff --git a/crates/ecstore/src/set_disk/ops/object.rs b/crates/ecstore/src/set_disk/ops/object.rs index 41775da3e..0236b7cd0 100644 --- a/crates/ecstore/src/set_disk/ops/object.rs +++ b/crates/ecstore/src/set_disk/ops/object.rs @@ -4185,6 +4185,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { let mut p_reader = PutObjReader::new(hash_reader); return match self_.clone().put_object(bucket, object, &mut p_reader, &ropts).await { Ok(restored_info) => { + let restored_info = self_.finalize_restore_metadata(bucket, object, &restored_info, &opts).await?; send_event(EventArgs { event_name: EventName::ObjectRestoreCompleted.as_str().to_string(), bucket_name: bucket.to_string(), @@ -4319,6 +4320,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { return set_restore_header_fn(&mut oi, Some(err)).await; } }; + let restored_info = self_.finalize_restore_metadata(bucket, object, &restored_info, opts).await?; send_event(EventArgs { event_name: EventName::ObjectRestoreCompleted.as_str().to_string(), bucket_name: bucket.to_string(), diff --git a/crates/ecstore/src/set_disk/replication.rs b/crates/ecstore/src/set_disk/replication.rs index 611abeb7a..ad84d4842 100644 --- a/crates/ecstore/src/set_disk/replication.rs +++ b/crates/ecstore/src/set_disk/replication.rs @@ -13,7 +13,10 @@ // limitations under the License. use super::*; +use crate::bucket::lifecycle::lifecycle; +use rustfs_filemeta::RestoreStatusOps; use rustfs_utils::http::headers::{AMZ_RESTORE_EXPIRY_DAYS, AMZ_RESTORE_REQUEST_DATE}; +use s3s::dto::{RestoreStatus, Timestamp}; #[derive(Clone, Copy, Debug, Eq, PartialEq)] struct RestoreCleanupIdentity { @@ -43,6 +46,69 @@ impl RestoreCleanupIdentity { } impl SetDisks { + pub(super) async fn finalize_restore_metadata( + &self, + bucket: &str, + object: &str, + obj_info: &ObjectInfo, + opts: &ObjectOptions, + ) -> Result { + let expected = RestoreCleanupIdentity::from_object_info(obj_info); + let expected_operation_id = restore_operation_id_from_metadata(&opts.user_defined)?; + let expected_etag = obj_info + .etag + .clone() + .unwrap_or_else(|| get_raw_etag(obj_info.user_defined.as_ref())); + let version_id = expected.version_id.map(|v| v.to_string()); + let _lock_guard = if !opts.no_lock { + Some( + self.acquire_write_lock_diag("restore_finalize_metadata", bucket, object) + .await?, + ) + } else { + None + }; + let read_opts = ObjectOptions { + version_id, + versioned: opts.versioned, + version_suspended: opts.version_suspended, + ..Default::default() + }; + let (mut fi, _, disks) = self + .get_object_fileinfo_gated(bucket, object, &read_opts, false, false) + .await?; + if let Some(expected_operation_id) = expected_operation_id { + require_restore_operation_id(&fi.metadata, expected_operation_id)?; + } + if !expected.matches_file_info(&fi, &expected_etag) { + return Err(Error::other("restored object changed before restore metadata finalization")); + } + let restore_expiry = + lifecycle::expected_expiry_time(OffsetDateTime::now_utc(), opts.transition.restore_request.days.unwrap_or(1)); + fi.metadata.insert( + X_AMZ_RESTORE.as_str().to_string(), + RestoreStatus { + is_restore_in_progress: Some(false), + restore_expiry_date: Some(Timestamp::from(restore_expiry)), + } + .to_string(), + ); + self.invalidate_get_object_metadata_cache(bucket, object).await; + self.update_object_meta_with_opts( + bucket, + object, + fi.clone(), + disks.as_slice(), + &UpdateMetadataOpts { + replace_user_metadata: true, + ..Default::default() + }, + ) + .await?; + self.invalidate_get_object_metadata_cache(bucket, object).await; + Ok(ObjectInfo::from_file_info(&fi, bucket, object, opts.versioned || opts.version_suspended)) + } + pub async fn update_restore_metadata( &self, bucket: &str, diff --git a/crates/ecstore/src/set_disk/transition_matrix_tests.rs b/crates/ecstore/src/set_disk/transition_matrix_tests.rs index e77a0e437..a137d55a6 100644 --- a/crates/ecstore/src/set_disk/transition_matrix_tests.rs +++ b/crates/ecstore/src/set_disk/transition_matrix_tests.rs @@ -13,10 +13,11 @@ // limitations under the License. use super::*; -use crate::bucket::lifecycle::lifecycle::{TRANSITION_COMPLETE, TRANSITION_PENDING, TransitionOptions}; +use crate::bucket::lifecycle::lifecycle::{TRANSITION_COMPLETE, TRANSITION_PENDING, TransitionOptions, expected_expiry_time}; use crate::ecstore_validation_blackbox::make_local_set_disks; use crate::services::tier::test_util::register_mock_tier; use crate::storage_api_contracts::object::{ObjectIO as _, ObjectOperations as _}; +use rustfs_filemeta::{RestoreStatusOps as _, parse_restore_obj_status}; use tokio::io::AsyncReadExt; async fn prime_metadata_generation(set_disks: &SetDisks, bucket: &str, object: &str) -> GetObjectMetadataCacheKey { @@ -83,13 +84,57 @@ async fn transition_and_restore_reclaim_prior_metadata_generations() { let transitioned_generation = prime_metadata_generation(&set_disks, bucket, object).await; let mut restore_opts = ObjectOptions::default(); restore_opts.transition.restore_request.days = Some(1); - Arc::clone(&set_disks) - .restore_transitioned_object(bucket, object, &restore_opts) - .await - .expect("restore should succeed"); + let restore_started = OffsetDateTime::now_utc(); + let expiry_from_restore_start = temp_env::async_with_vars( + [ + ("RUSTFS_ILM_DEBUG_DAY_SECS", Some("1")), + ("RUSTFS_ILM_PROCESS_TIME", Some("1")), + ], + async { + let expiry_from_restore_start = expected_expiry_time(restore_started, 1); + let get_barrier = backend.arm_get_barrier().await; + let restore_set = Arc::clone(&set_disks); + let restore = + tokio::spawn(async move { restore_set.restore_transitioned_object(bucket, object, &restore_opts).await }); + get_barrier.wait_until_paused().await; + tokio::time::timeout(Duration::from_secs(5), async { + loop { + if expected_expiry_time(OffsetDateTime::now_utc(), 1) > expiry_from_restore_start { + break; + } + tokio::time::sleep(Duration::from_millis(10)).await; + } + }) + .await + .expect("test clock should cross the next accelerated lifecycle boundary"); + get_barrier.release(); + restore + .await + .expect("restore task should join") + .expect("restore should succeed"); + expiry_from_restore_start + }, + ) + .await; assert_generation_reclaimed(&set_disks, &transitioned_generation).await; assert_eq!(backend.get_count().await, 1, "restore should read the remote candidate exactly once"); + let restored_info = set_disks + .get_object_info(bucket, object, &ObjectOptions::default()) + .await + .expect("restored object metadata should be readable"); + let restore_status = parse_restore_obj_status( + restored_info + .user_defined + .get(s3s::header::X_AMZ_RESTORE.as_str()) + .expect("completed restore header should be present"), + ) + .expect("completed restore header should parse"); + assert!( + restore_status.expiry().expect("completed restore should have an expiry") > expiry_from_restore_start, + "restore expiry must be based on completion, not the time the remote copy started" + ); + let mut restored = Vec::new(); set_disks .get_object_reader(bucket, object, None, HeaderMap::new(), &ObjectOptions::default())