mirror of
https://github.com/rustfs/rustfs.git
synced 2026-07-29 17:48:58 +00:00
fix(tiering): base restore expiry on completion (#5366)
* fix(tiering): base restore expiry on completion * fix(tiering): import restore metadata types * fix(restore): resolve metadata finalization build errors * test(ecstore): import restore expiry helper
This commit is contained in:
@@ -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(),
|
||||
|
||||
@@ -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<ObjectInfo> {
|
||||
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,
|
||||
|
||||
@@ -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())
|
||||
|
||||
Reference in New Issue
Block a user