Compare commits

..

6 Commits

Author SHA1 Message Date
houseme 566877d3ba Merge branch 'main' into overtrue/activate-group-e2e 2026-08-23 12:37:24 +08:00
overtrue 3294c64fcc test(e2e): correct group regression fixtures 2026-08-23 11:11:01 +08:00
overtrue a8be0f81f5 test(e2e): bind group selection to Linux listing 2026-08-23 08:55:07 +08:00
overtrue 3179b7acb8 test(e2e): pin group deletion errors 2026-08-23 06:42:55 +08:00
overtrue d25b84a793 Merge remote-tracking branch 'origin/main' into overtrue/activate-group-e2e 2026-08-23 05:43:17 +08:00
overtrue a79b806fb4 test(e2e): activate group management regressions 2026-08-23 02:37:15 +08:00
12 changed files with 151 additions and 1328 deletions
+2 -2
View File
@@ -1,2 +1,2 @@
sha256-darwin=9f767b37ed8b1c82da62ea441462d75487785c8086e56f08fb6f6cd89c6e2e52 sha256-darwin=f832043fcca8c0b616c5d820a3a652da7544298ef5812a8668a3a9a3e4607b8b
sha256-linux=fbdaf42b220958d4b1e8880e0f8b5a7992d38e21051bb60596dd4538424757d6 sha256-linux=93b94adb110b86a41d0b7313909e0bf53cb1515e2d08e8f105652b29b249990f
+118 -36
View File
@@ -14,7 +14,7 @@
//! E2E tests for group management (fixes #2028). //! E2E tests for group management (fixes #2028).
use crate::common::{RustFSTestEnvironment, admin_request, awscurl_delete, awscurl_get, awscurl_put, init_logging}; use crate::common::{RustFSTestEnvironment, admin_ok, admin_request, init_logging};
use aws_sdk_s3::config::{Credentials, Region}; use aws_sdk_s3::config::{Credentials, Region};
use aws_sdk_s3::{Client, Config}; use aws_sdk_s3::{Client, Config};
use tracing::info; use tracing::info;
@@ -83,7 +83,6 @@ async fn update_group_members_rejects_invalid_new_group_names() -> Result<(), Bo
/// Test that deleting a group with members fails, and deleting an empty group succeeds. /// Test that deleting a group with members fails, and deleting an empty group succeeds.
#[tokio::test(flavor = "multi_thread")] #[tokio::test(flavor = "multi_thread")]
#[ignore = "requires awscurl and spawns a real RustFS server"]
async fn test_delete_group_requires_empty_membership() -> Result<(), Box<dyn std::error::Error + Send + Sync>> { async fn test_delete_group_requires_empty_membership() -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
init_logging(); init_logging();
@@ -91,29 +90,58 @@ async fn test_delete_group_requires_empty_membership() -> Result<(), Box<dyn std
env.start_rustfs_server(vec![]).await?; env.start_rustfs_server(vec![]).await?;
// 1. Create a user // 1. Create a user
let add_user_url = format!("{}/rustfs/admin/v3/add-user?accessKey=testuser1", env.url);
let user_body = serde_json::json!({ let user_body = serde_json::json!({
"secretKey": "testuser1secret", "secretKey": "testuser1secret",
"status": "enabled" "status": "enabled"
}); });
awscurl_put(&add_user_url, &user_body.to_string(), &env.access_key, &env.secret_key).await?; admin_ok(
&env,
http::Method::PUT,
"/rustfs/admin/v3/add-user?accessKey=testuser1",
Some(user_body.to_string()),
)
.await?;
info!("Created testuser1"); info!("Created testuser1");
// 2. Create a group with testuser1 as a member // 2. Create a group with testuser1 as a member
let update_members_url = format!("{}/rustfs/admin/v3/update-group-members", env.url);
let add_member_body = serde_json::json!({ let add_member_body = serde_json::json!({
"group": "testgroup", "group": "testgroup",
"members": ["testuser1"], "members": ["testuser1"],
"isRemove": false, "isRemove": false,
"groupStatus": "enabled" "groupStatus": "enabled"
}); });
awscurl_put(&update_members_url, &add_member_body.to_string(), &env.access_key, &env.secret_key).await?; admin_ok(
&env,
http::Method::PUT,
"/rustfs/admin/v3/update-group-members",
Some(add_member_body.to_string()),
)
.await?;
info!("Added testuser1 to testgroup"); info!("Added testuser1 to testgroup");
// 3. Attempt to delete the group while it still has members — should fail // 3. Attempt to delete the group while it still has members — should fail
let delete_group_url = format!("{}/rustfs/admin/v3/group/testgroup", env.url); let (delete_status, delete_body) = admin_request(
let delete_result = awscurl_delete(&delete_group_url, &env.access_key, &env.secret_key).await; &env.url,
assert!(delete_result.is_err(), "deleting a non-empty group should fail"); http::Method::DELETE,
"/rustfs/admin/v3/group/testgroup",
None,
&env.access_key,
&env.secret_key,
)
.await?;
assert_eq!(
delete_status,
reqwest::StatusCode::BAD_REQUEST,
"deleting a non-empty group must return HTTP 400, body: {delete_body}"
);
assert!(
delete_body.contains("<Code>InvalidRequest</Code>"),
"deleting a non-empty group must return InvalidRequest, body: {delete_body}"
);
assert!(
delete_body.contains("<Message>group is not empty</Message>"),
"deleting a non-empty group returned an unexpected message: {delete_body}"
);
info!("Delete of non-empty group correctly rejected"); info!("Delete of non-empty group correctly rejected");
// 4. Remove the member from the group // 4. Remove the member from the group
@@ -123,17 +151,42 @@ async fn test_delete_group_requires_empty_membership() -> Result<(), Box<dyn std
"isRemove": true, "isRemove": true,
"groupStatus": "enabled" "groupStatus": "enabled"
}); });
awscurl_put(&update_members_url, &remove_member_body.to_string(), &env.access_key, &env.secret_key).await?; admin_ok(
&env,
http::Method::PUT,
"/rustfs/admin/v3/update-group-members",
Some(remove_member_body.to_string()),
)
.await?;
info!("Removed testuser1 from testgroup"); info!("Removed testuser1 from testgroup");
// 5. Delete the now-empty group — should succeed // 5. Delete the now-empty group — should succeed
awscurl_delete(&delete_group_url, &env.access_key, &env.secret_key).await?; admin_ok(&env, http::Method::DELETE, "/rustfs/admin/v3/group/testgroup", None).await?;
info!("Deleted empty testgroup successfully"); info!("Deleted empty testgroup successfully");
// 6. Verify the group no longer exists // 6. Verify the group no longer exists
let get_group_url = format!("{}/rustfs/admin/v3/group?group=testgroup", env.url); let (get_status, get_body) = admin_request(
let get_result = awscurl_get(&get_group_url, &env.access_key, &env.secret_key).await; &env.url,
assert!(get_result.is_err(), "group should no longer exist after deletion"); http::Method::GET,
"/rustfs/admin/v3/group?group=testgroup",
None,
&env.access_key,
&env.secret_key,
)
.await?;
assert_eq!(
get_status,
reqwest::StatusCode::NOT_FOUND,
"a deleted group must return HTTP 404, body: {get_body}"
);
assert!(
get_body.contains("<Code>NoSuchResource</Code>"),
"a deleted group must return NoSuchResource, body: {get_body}"
);
assert!(
get_body.contains("<Message>group &apos;testgroup&apos; does not exist</Message>"),
"a deleted group returned an unexpected message: {get_body}"
);
info!("Confirmed testgroup no longer exists"); info!("Confirmed testgroup no longer exists");
Ok(()) Ok(())
@@ -142,7 +195,6 @@ async fn test_delete_group_requires_empty_membership() -> Result<(), Box<dyn std
/// Test that a user with only group membership (no explicit user policy) gets group policies /// Test that a user with only group membership (no explicit user policy) gets group policies
/// and can perform actions allowed by the group (regression test for #2028.1). /// and can perform actions allowed by the group (regression test for #2028.1).
#[tokio::test(flavor = "multi_thread")] #[tokio::test(flavor = "multi_thread")]
#[ignore = "requires awscurl and spawns a real RustFS server"]
async fn test_user_with_only_group_gets_group_policies() -> Result<(), Box<dyn std::error::Error + Send + Sync>> { async fn test_user_with_only_group_gets_group_policies() -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
init_logging(); init_logging();
@@ -160,39 +212,56 @@ async fn test_user_with_only_group_gets_group_policies() -> Result<(), Box<dyn s
"Statement": [{ "Statement": [{
"Effect": "Allow", "Effect": "Allow",
"Action": ["s3:ListAllMyBuckets"], "Action": ["s3:ListAllMyBuckets"],
"Resource": ["*"] "Resource": ["arn:aws:s3:::*"]
}] }]
}); });
let add_policy_url = format!("{}/rustfs/admin/v3/add-canned-policy?name={}", env.url, policy_name); admin_ok(
awscurl_put(&add_policy_url, &policy_doc.to_string(), &env.access_key, &env.secret_key).await?; &env,
http::Method::PUT,
&format!("/rustfs/admin/v3/add-canned-policy?name={policy_name}"),
Some(policy_doc.to_string()),
)
.await?;
info!("Created canned policy {}", policy_name); info!("Created canned policy {}", policy_name);
// 2. Create user with no explicit policy // 2. Create user with no explicit policy
let add_user_url = format!("{}/rustfs/admin/v3/add-user?accessKey={}", env.url, user_name);
let user_body = serde_json::json!({ let user_body = serde_json::json!({
"secretKey": user_secret, "secretKey": user_secret,
"status": "enabled" "status": "enabled"
}); });
awscurl_put(&add_user_url, &user_body.to_string(), &env.access_key, &env.secret_key).await?; admin_ok(
&env,
http::Method::PUT,
&format!("/rustfs/admin/v3/add-user?accessKey={user_name}"),
Some(user_body.to_string()),
)
.await?;
info!("Created user {} with no explicit policy", user_name); info!("Created user {} with no explicit policy", user_name);
// 3. Add user to group (creates group with this member; user_group_memberships must be updated) // 3. Add user to group (creates group with this member; user_group_memberships must be updated)
let update_members_url = format!("{}/rustfs/admin/v3/update-group-members", env.url);
let add_member_body = serde_json::json!({ let add_member_body = serde_json::json!({
"group": group_name, "group": group_name,
"members": [user_name], "members": [user_name],
"isRemove": false, "isRemove": false,
"groupStatus": "enabled" "groupStatus": "enabled"
}); });
awscurl_put(&update_members_url, &add_member_body.to_string(), &env.access_key, &env.secret_key).await?; admin_ok(
&env,
http::Method::PUT,
"/rustfs/admin/v3/update-group-members",
Some(add_member_body.to_string()),
)
.await?;
info!("Added {} to group {}", user_name, group_name); info!("Added {} to group {}", user_name, group_name);
// 4. Attach policy to group // 4. Attach policy to group
let set_policy_url = format!( admin_ok(
"{}/rustfs/admin/v3/set-user-or-group-policy?policyName={}&userOrGroup={}&isGroup=true", &env,
env.url, policy_name, group_name http::Method::PUT,
); &format!("/rustfs/admin/v3/set-user-or-group-policy?policyName={policy_name}&userOrGroup={group_name}&isGroup=true"),
awscurl_put(&set_policy_url, "", &env.access_key, &env.secret_key).await?; Some(String::new()),
)
.await?;
info!("Attached policy {} to group {}", policy_name, group_name); info!("Attached policy {} to group {}", policy_name, group_name);
// 5. User with only group (no user policy) should be able to list buckets // 5. User with only group (no user policy) should be able to list buckets
@@ -209,7 +278,6 @@ async fn test_user_with_only_group_gets_group_policies() -> Result<(), Box<dyn s
/// Test that after deleting a user who was the only member of a group, the group can be deleted /// Test that after deleting a user who was the only member of a group, the group can be deleted
/// (regression test for #2028.2: delete group uses backend membership, not stale cache). /// (regression test for #2028.2: delete group uses backend membership, not stale cache).
#[tokio::test(flavor = "multi_thread")] #[tokio::test(flavor = "multi_thread")]
#[ignore = "requires awscurl and spawns a real RustFS server"]
async fn test_delete_group_after_deleting_user() -> Result<(), Box<dyn std::error::Error + Send + Sync>> { async fn test_delete_group_after_deleting_user() -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
init_logging(); init_logging();
@@ -221,33 +289,47 @@ async fn test_delete_group_after_deleting_user() -> Result<(), Box<dyn std::erro
let group_name = "soledeletegroup"; let group_name = "soledeletegroup";
// 1. Create user // 1. Create user
let add_user_url = format!("{}/rustfs/admin/v3/add-user?accessKey={}", env.url, user_name);
let user_body = serde_json::json!({ let user_body = serde_json::json!({
"secretKey": user_secret, "secretKey": user_secret,
"status": "enabled" "status": "enabled"
}); });
awscurl_put(&add_user_url, &user_body.to_string(), &env.access_key, &env.secret_key).await?; admin_ok(
&env,
http::Method::PUT,
&format!("/rustfs/admin/v3/add-user?accessKey={user_name}"),
Some(user_body.to_string()),
)
.await?;
info!("Created user {}", user_name); info!("Created user {}", user_name);
// 2. Add user to group // 2. Add user to group
let update_members_url = format!("{}/rustfs/admin/v3/update-group-members", env.url);
let add_member_body = serde_json::json!({ let add_member_body = serde_json::json!({
"group": group_name, "group": group_name,
"members": [user_name], "members": [user_name],
"isRemove": false, "isRemove": false,
"groupStatus": "enabled" "groupStatus": "enabled"
}); });
awscurl_put(&update_members_url, &add_member_body.to_string(), &env.access_key, &env.secret_key).await?; admin_ok(
&env,
http::Method::PUT,
"/rustfs/admin/v3/update-group-members",
Some(add_member_body.to_string()),
)
.await?;
info!("Added {} to group {}", user_name, group_name); info!("Added {} to group {}", user_name, group_name);
// 3. Delete the user (backend and cache update so group membership becomes empty) // 3. Delete the user (backend and cache update so group membership becomes empty)
let remove_user_url = format!("{}/rustfs/admin/v3/remove-user?accessKey={}", env.url, user_name); admin_ok(
awscurl_delete(&remove_user_url, &env.access_key, &env.secret_key).await?; &env,
http::Method::DELETE,
&format!("/rustfs/admin/v3/remove-user?accessKey={user_name}"),
None,
)
.await?;
info!("Deleted user {}", user_name); info!("Deleted user {}", user_name);
// 4. Deleting the group should succeed (backend has empty members; no stale cache) // 4. Deleting the group should succeed (backend has empty members; no stale cache)
let delete_group_url = format!("{}/rustfs/admin/v3/group/{}", env.url, group_name); admin_ok(&env, http::Method::DELETE, &format!("/rustfs/admin/v3/group/{group_name}"), None).await?;
awscurl_delete(&delete_group_url, &env.access_key, &env.secret_key).await?;
info!("Deleted group {} after user was removed", group_name); info!("Deleted group {} after user was removed", group_name);
Ok(()) Ok(())
+13 -263
View File
@@ -13,8 +13,6 @@
// limitations under the License. // limitations under the License.
use crate::bucket::replication::replication_state_from_filemeta; use crate::bucket::replication::replication_state_from_filemeta;
#[cfg(test)]
use crate::bucket::utils::is_meta_bucketname;
use crate::bucket::versioning_sys::BucketVersioningSys; use crate::bucket::versioning_sys::BucketVersioningSys;
use crate::bucket::{ use crate::bucket::{
lifecycle::{ lifecycle::{
@@ -1349,53 +1347,6 @@ fn should_cleanup_decommission_source_entry(decommissioned: usize, total_version
decommissioned.saturating_add(expired) == total_versions decommissioned.saturating_add(expired) == total_versions
} }
const DECOMMISSION_FREE_VERSION_MIGRATED_REASON: &str = "tier_free_version_migrated";
const DECOMMISSION_FREE_VERSION_CONSUMED_REASON: &str = "tier_free_version_already_consumed";
const DECOMMISSION_FREE_VERSION_RETAINED_REASON: &str = "tier_free_version_migration_failed";
const DECOMMISSION_FREE_VERSION_SWEEP_REASON: &str = "tier_free_version_unresolved_after_decommission";
const DECOMMISSION_FREE_VERSION_DISPOSITION_REASON: &str = "tier_free_version_disposition_recorded";
#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)]
struct DecommissionFreeVersionDisposition {
migrated: usize,
consumed: usize,
retained: usize,
}
impl DecommissionFreeVersionDisposition {
fn record_migrated(&mut self) {
self.migrated += 1;
}
fn record_consumed(&mut self) {
self.consumed += 1;
}
fn record_retained(&mut self) {
self.retained += 1;
}
fn total(self) -> usize {
self.migrated.saturating_add(self.consumed).saturating_add(self.retained)
}
}
enum DecommissionFreeVersionAttempt {
Migrated,
Consumed,
CapacityFailure(Error),
Retry(Error),
}
fn classify_decommission_free_version_attempt(result: Result<()>) -> DecommissionFreeVersionAttempt {
match result {
Ok(()) => DecommissionFreeVersionAttempt::Migrated,
Err(err) if is_decommission_copy_cleanup_safe_error(&err) => DecommissionFreeVersionAttempt::Consumed,
Err(err) if is_decommission_target_capacity_error(&err) => DecommissionFreeVersionAttempt::CapacityFailure(err),
Err(err) => DecommissionFreeVersionAttempt::Retry(err),
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)] #[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[allow( #[allow(
dead_code, dead_code,
@@ -2664,12 +2615,8 @@ fn determine_decommission_final_state(items_failed: usize, was_cancelled: bool)
} }
} }
fn decommission_remaining_version_count(versions: &[rustfs_filemeta::FileInfo], expired: usize) -> usize { fn decommission_remaining_version_count(total_versions: usize, expired: usize) -> usize {
versions total_versions.saturating_sub(expired)
.iter()
.filter(|version| !version.tier_free_version())
.count()
.saturating_sub(expired)
} }
fn should_skip_decommission_delete_marker( fn should_skip_decommission_delete_marker(
@@ -2732,7 +2679,6 @@ fn decommission_remote_tiered_opts(
user_defined: version.metadata.clone(), user_defined: version.metadata.clone(),
src_pool_idx, src_pool_idx,
data_movement: true, data_movement: true,
incl_free_versions: version.tier_free_version(),
include_part_checksums: true, include_part_checksums: true,
http_preconditions: Some(crate::data_movement::data_movement_target_precondition()), http_preconditions: Some(crate::data_movement::data_movement_target_precondition()),
expected_bucket_incarnation_id, expected_bucket_incarnation_id,
@@ -3876,7 +3822,6 @@ impl ECStore {
let mut decommissioned: usize = 0; let mut decommissioned: usize = 0;
let mut expired: usize = 0; let mut expired: usize = 0;
let mut free_version_disposition = DecommissionFreeVersionDisposition::default();
let mut cleanup_preflight_allowed_missing = Vec::new(); let mut cleanup_preflight_allowed_missing = Vec::new();
for version in fivs.versions.iter() { for version in fivs.versions.iter() {
@@ -3885,115 +3830,6 @@ impl ECStore {
} }
decommission_cancel_signal_result(rx.is_cancelled())?; decommission_cancel_signal_result(rx.is_cancelled())?;
if version.tier_free_version() {
let version_id = version.version_id.map(|v| v.to_string());
let mut migration_error = None;
let mut migrated = false;
let mut consumed = false;
let mut capacity_failure = false;
for _ in 0..3 {
match classify_decommission_free_version_attempt(
run_decommission_side_effect(&rx, &operation_gate, || async {
self.decommission_tiered_object(
bucket.as_str(),
&version.name,
version,
&decommission_remote_tiered_opts(
version,
version_id.clone(),
idx,
expected_bucket_incarnation_id,
),
)
.await
})
.await,
) {
DecommissionFreeVersionAttempt::Migrated => {
migrated = true;
migration_error = None;
break;
}
DecommissionFreeVersionAttempt::Consumed => {
consumed = true;
migration_error = None;
break;
}
DecommissionFreeVersionAttempt::CapacityFailure(err) => {
capacity_failure = true;
migration_error = Some(err);
break;
}
DecommissionFreeVersionAttempt::Retry(err) => migration_error = Some(err),
}
}
{
let mut pool_meta = self.pool_meta.write().await;
ensure_decommission_generation(&pool_meta, idx, generation)?;
if let Err(err) = count_decommission_item(&mut pool_meta, idx, 0, !migrated && !consumed) {
return Err(with_decommission_entry_context(
"count_decommission_item",
bucket.as_str(),
entry.name.as_str(),
err,
));
}
}
if migrated || consumed {
decommissioned += 1;
cleanup_preflight_allowed_missing.push(data_movement::source_cleanup_version_identity(version));
}
if migrated {
free_version_disposition.record_migrated();
} else if consumed {
free_version_disposition.record_consumed();
} else {
free_version_disposition.record_retained();
}
debug!(
event = EVENT_DECOMMISSION_ENTRY,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_POOLS,
pool_index = idx,
bucket = %bucket,
object = %version.name,
version_id = ?version_id,
result = ?migration_error,
reason = if migrated {
DECOMMISSION_FREE_VERSION_MIGRATED_REASON
} else if consumed {
DECOMMISSION_FREE_VERSION_CONSUMED_REASON
} else {
DECOMMISSION_FREE_VERSION_RETAINED_REASON
},
state = if migrated {
"free_version_migrated"
} else if consumed {
"free_version_consumed"
} else {
"free_version_retained"
},
"Decommission free-version disposition recorded"
);
if capacity_failure {
return Err(with_decommission_entry_context(
"decommission_tier_free_version",
bucket.as_str(),
version.name.as_str(),
migration_error.expect("capacity failure must retain its error"),
));
}
if !migrated && !consumed {
break;
}
continue;
}
if run_decommission_side_effect(&rx, &operation_gate, || async { if run_decommission_side_effect(&rx, &operation_gate, || async {
should_skip_lifecycle_for_data_movement( should_skip_lifecycle_for_data_movement(
self.clone(), self.clone(),
@@ -4014,7 +3850,7 @@ impl ECStore {
continue; continue;
} }
let remaining_versions = decommission_remaining_version_count(&fivs.versions, expired); let remaining_versions = decommission_remaining_version_count(fivs.versions.len(), expired);
if should_skip_decommission_delete_marker(version, remaining_versions, replication_config.is_some()) { if should_skip_decommission_delete_marker(version, remaining_versions, replication_config.is_some()) {
// //
decommissioned += 1; decommissioned += 1;
@@ -4293,24 +4129,6 @@ impl ECStore {
} }
} }
if free_version_disposition.total() > 0 {
debug!(
event = EVENT_DECOMMISSION_ENTRY,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_POOLS,
pool_index = idx,
bucket = %bucket,
object = %entry.name,
free_versions_migrated = free_version_disposition.migrated,
free_versions_consumed = free_version_disposition.consumed,
free_versions_retained = free_version_disposition.retained,
free_versions_total = free_version_disposition.total(),
reason = DECOMMISSION_FREE_VERSION_DISPOSITION_REASON,
state = "free_version_disposition",
"Decommission free-version disposition summary"
);
}
if should_cleanup_decommission_source_entry(decommissioned, fivs.versions.len(), expired) { if should_cleanup_decommission_source_entry(decommissioned, fivs.versions.len(), expired) {
if bucket_incarnation_fence.as_ref().is_some_and(|guard| guard.is_lock_lost()) { if bucket_incarnation_fence.as_ref().is_some_and(|guard| guard.is_lock_lost()) {
return Err(Error::other("decommission bucket incarnation fence was lost before source cleanup")); return Err(Error::other("decommission bucket incarnation fence was lost before source cleanup"));
@@ -4480,39 +4298,6 @@ impl ECStore {
.await .await
} }
#[cfg(test)]
pub(crate) async fn decommission_entry_for_test_with_bucket_incarnation(
self: &Arc<Self>,
idx: usize,
entry: MetaCacheEntry,
bucket: String,
set: Arc<SetDisks>,
) -> Result<()> {
let expected_bucket_incarnation_id = if is_meta_bucketname(&bucket) {
None
} else {
Some(self.bucket_incarnation_id_from_disk(&bucket).await?)
};
self.decommission_entry(
CancellationToken::new(),
idx,
OffsetDateTime::now_utc(),
entry,
bucket,
set,
None,
None,
None,
expected_bucket_incarnation_id,
)
.await
}
#[cfg(test)]
pub(crate) async fn check_after_decommission_for_test(self: &Arc<Self>, idx: usize) -> Result<()> {
self.check_after_decommission(idx).await
}
#[tracing::instrument(skip(self, rx))] #[tracing::instrument(skip(self, rx))]
async fn decommission_pool( async fn decommission_pool(
self: &Arc<Self>, self: &Arc<Self>,
@@ -5377,7 +5162,6 @@ impl ECStore {
let lifecycle_config_cb = lifecycle_config.clone(); let lifecycle_config_cb = lifecycle_config.clone();
let object_lock_config_cb = object_lock_config.clone(); let object_lock_config_cb = object_lock_config.clone();
let store = Arc::clone(self); let store = Arc::clone(self);
let set_cb = Arc::clone(set);
let callback_rx_cb = callback_rx.clone(); let callback_rx_cb = callback_rx.clone();
let callback: ListCallback = Arc::new(move |entry: MetaCacheEntry| { let callback: ListCallback = Arc::new(move |entry: MetaCacheEntry| {
@@ -5387,7 +5171,6 @@ impl ECStore {
let lifecycle_config = lifecycle_config_cb.clone(); let lifecycle_config = lifecycle_config_cb.clone();
let object_lock_config = object_lock_config_cb.clone(); let object_lock_config = object_lock_config_cb.clone();
let store = Arc::clone(&store); let store = Arc::clone(&store);
let set = Arc::clone(&set_cb);
let callback_rx = callback_rx_cb.clone(); let callback_rx = callback_rx_cb.clone();
Box::pin(async move { Box::pin(async move {
if callback_rx.is_cancelled() { if callback_rx.is_cancelled() {
@@ -5402,14 +5185,11 @@ impl ECStore {
return; return;
} }
let fivs = match load_decommission_entry_exact_versions( let fivs = match load_decommission_entry_versions(
&set,
&entry, &entry,
&bucket_name, &bucket_name,
"check_after_decommission.file_info_versions", "check_after_decommission.file_info_versions",
) ) {
.await
{
Ok(fivs) => fivs, Ok(fivs) => fivs,
Err(err) => { Err(err) => {
let mut first_err = entry_error.lock().await; let mut first_err = entry_error.lock().await;
@@ -5422,23 +5202,7 @@ impl ECStore {
}; };
let mut remaining = 0; let mut remaining = 0;
for version in fivs.versions.iter().chain(fivs.free_versions.iter()) { for version in &fivs.versions {
if version.tier_free_version() {
remaining += 1;
debug!(
event = EVENT_DECOMMISSION_ENTRY,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_POOLS,
pool_index = idx,
bucket = %bucket_name,
object = %entry.name,
version_id = ?version.version_id,
reason = DECOMMISSION_FREE_VERSION_SWEEP_REASON,
state = "free_version_retained",
"Decommission final sweep retained a free version"
);
continue;
}
if version.deleted { if version.deleted {
continue; continue;
} }
@@ -5553,6 +5317,13 @@ mod tests {
assert_eq!(determine_decommission_final_state(0, true), DecommissionFinalState::Failed); assert_eq!(determine_decommission_final_state(0, true), DecommissionFinalState::Failed);
} }
#[test]
fn decommission_remaining_version_count_excludes_only_expired_versions() {
assert_eq!(decommission_remaining_version_count(1, 0), 1);
assert_eq!(decommission_remaining_version_count(2, 1), 1);
assert_eq!(decommission_remaining_version_count(1, 1), 0);
}
#[test] #[test]
fn lifecycle_action_removes_data_movement_version_rejects_delete_marker_action() { fn lifecycle_action_removes_data_movement_version_rejects_delete_marker_action() {
assert!(!lifecycle_action_removes_data_movement_version(IlmAction::DeleteAction)); assert!(!lifecycle_action_removes_data_movement_version(IlmAction::DeleteAction));
@@ -5609,21 +5380,6 @@ mod tests {
))); )));
} }
#[test]
fn decommission_free_version_attempt_treats_missing_source_as_consumed() {
let attempt =
classify_decommission_free_version_attempt(Err(Error::ObjectNotFound("bucket".to_string(), "object".to_string())));
assert!(matches!(attempt, DecommissionFreeVersionAttempt::Consumed));
}
#[test]
fn decommission_free_version_attempt_preserves_capacity_failure() {
let attempt = classify_decommission_free_version_attempt(Err(Error::DiskFull));
assert!(matches!(attempt, DecommissionFreeVersionAttempt::CapacityFailure(Error::DiskFull)));
}
#[test] #[test]
fn decommission_delete_marker_copy_error_rejects_data_movement_overwrite() { fn decommission_delete_marker_copy_error_rejects_data_movement_overwrite() {
let err = Error::DataMovementOverwriteErr("bucket".to_string(), "object".to_string(), "version".to_string()); let err = Error::DataMovementOverwriteErr("bucket".to_string(), "object".to_string(), "version".to_string());
@@ -5805,12 +5561,6 @@ mod tests {
assert!(opts.include_part_checksums); assert!(opts.include_part_checksums);
assert!(opts.http_preconditions.is_some()); assert!(opts.http_preconditions.is_some());
assert_eq!(opts.expected_bucket_incarnation_id, Some(incarnation)); assert_eq!(opts.expected_bucket_incarnation_id, Some(incarnation));
assert!(!opts.incl_free_versions);
let mut free_version = version;
free_version.set_tier_free_version();
let free_opts = decommission_remote_tiered_opts(&free_version, Some("free-version-id".to_string()), 9, Some(incarnation));
assert!(free_opts.incl_free_versions);
} }
#[test] #[test]
-13
View File
@@ -1950,19 +1950,6 @@ mod tests {
assert!(source_cleanup_versions_match_with_allowed_missing(&expected, &current, &allowed_missing)); assert!(source_cleanup_versions_match_with_allowed_missing(&expected, &current, &allowed_missing));
} }
#[test]
fn test_decommission_cleanup_preflight_accepts_migrated_free_version_consumed_from_source() {
let migrated = cleanup_test_file_info("object.txt", Uuid::from_u128(1), "migrated");
let mut free_version = cleanup_test_file_info("object.txt", Uuid::from_u128(2), "tier-cleanup");
free_version.deleted = true;
free_version.set_tier_free_version();
let expected = cleanup_test_versions(vec![migrated.clone(), free_version.clone()]);
let current = cleanup_test_versions(vec![migrated]);
let allowed_missing = vec![source_cleanup_version_identity(&free_version)];
assert!(source_cleanup_versions_match_with_allowed_missing(&expected, &current, &allowed_missing));
}
#[test] #[test]
fn test_decommission_cleanup_preflight_rejects_unexpected_missing_version() { fn test_decommission_cleanup_preflight_rejects_unexpected_missing_version() {
let migrated = cleanup_test_file_info("object.txt", Uuid::from_u128(1), "migrated"); let migrated = cleanup_test_file_info("object.txt", Uuid::from_u128(1), "migrated");
-297
View File
@@ -227,70 +227,6 @@ pub(super) fn restore_commit_operation_id_from_metadata(metadata: &HashMap<Strin
restore_operation_id_from_metadata(metadata) restore_operation_id_from_metadata(metadata)
} }
async fn inspect_decommission_tier_free_version_target(
disk: &DiskStore,
bucket: &str,
object: &str,
source: &FileInfo,
) -> Result<bool> {
let raw = match disk.read_xl(bucket, object, false).await {
Ok(raw) => raw,
Err(DiskError::FileNotFound | DiskError::FileVersionNotFound | DiskError::VolumeNotFound) => return Ok(false),
Err(err) => return Err(err.into()),
};
let meta = FileMeta::load(&raw.buf)?;
let source_version_id = source.version_id.filter(|version_id| !version_id.is_nil());
let mut matching_count = 0;
let mut all_matching_versions_equivalent = true;
for existing in meta
.versions
.iter()
.filter(|version| version.header.version_id.filter(|version_id| !version_id.is_nil()) == source_version_id)
{
matching_count += 1;
let existing = existing.into_fileinfo(bucket, object, true)?;
existing.validate_for_metadata_read()?;
if !existing.tier_free_version() || !crate::store::tiered_data_movement_source_matches(source, &existing)? {
all_matching_versions_equivalent = false;
}
}
if matching_count == 0 {
return Ok(false);
}
if matching_count == 1 && all_matching_versions_equivalent {
return Ok(true);
}
Err(StorageError::DataMovementOverwriteErr(
bucket.to_owned(),
object.to_owned(),
source_version_id.map(|version_id| version_id.to_string()).unwrap_or_default(),
)
.into())
}
fn ensure_decommission_tier_free_version_commit_fence(bucket: &str, object: &str, opts: &ObjectOptions) -> Result<()> {
if opts
.namespace_lock_fence
.as_ref()
.is_some_and(NamespaceLockFence::is_lock_lost)
|| opts
.bucket_lifecycle_lock_fence
.as_ref()
.is_some_and(NamespaceLockFence::is_lock_lost)
{
return Err(StorageError::NamespaceLockQuorumUnavailable {
mode: "decommission_tier_free_version_commit",
bucket: bucket.to_string(),
object: object.to_string(),
required: 1,
achieved: 0,
});
}
Ok(())
}
impl SetDisks { impl SetDisks {
pub(super) async fn require_current_restore_operation_id( pub(super) async fn require_current_restore_operation_id(
&self, &self,
@@ -4729,98 +4665,6 @@ fn resolve_delete_version_state(opts: &ObjectOptions, goi: &ObjectInfo, version_
} }
impl SetDisks { impl SetDisks {
/// Publish an internal tier free-version record without changing its
/// delete-marker shape or remote-tier identity. The caller holds the
/// source and target object locks; a write quorum is required before the
/// source cleanup may remove the original record.
#[tracing::instrument(skip(self, fi, opts))]
pub(crate) async fn decommission_tier_free_version(
&self,
bucket: &str,
object: &str,
fi: &FileInfo,
opts: &ObjectOptions,
) -> Result<()> {
if !fi.deleted || !fi.tier_free_version() {
return Err(Error::other("decommission tier free-version write requires a free version record"));
}
ensure_decommission_tier_free_version_commit_fence(bucket, object, opts)?;
self.validate_decommission_tier_free_version_target(bucket, object, fi)
.await?;
ensure_decommission_tier_free_version_commit_fence(bucket, object, opts)?;
let disks = self.disks.read().await.clone();
let write_quorum = self.default_write_quorum();
let futures = disks.into_iter().map(|disk| {
let file_info = fi.clone();
async move {
if let Some(disk) = disk {
disk.write_metadata("", bucket, object, file_info).await
} else {
Err(DiskError::DiskNotFound)
}
}
});
let mut errs = Vec::new();
for result in join_all(futures).await {
match result {
Ok(_) => errs.push(None),
Err(err) => errs.push(Some(err)),
}
}
ensure_decommission_tier_free_version_commit_fence(bucket, object, opts)?;
resolve_tiered_decommission_write_quorum_result(&errs, write_quorum, bucket, object)
}
pub(crate) async fn validate_decommission_tier_free_version_target(
&self,
bucket: &str,
object: &str,
fi: &FileInfo,
) -> Result<()> {
// The caller holds the source and target object locks. Inspect every
// target disk before an idempotent return or metadata fan-out so a
// sub-quorum conflict cannot be hidden by a successful quorum.
let disks = self.disks.read().await.clone();
let preflight = disks
.iter()
.flatten()
.map(|disk| inspect_decommission_tier_free_version_target(disk, bucket, object, fi));
for result in join_all(preflight).await {
result?;
}
Ok(())
}
pub(crate) async fn has_decommission_tier_free_version_write_quorum(
&self,
bucket: &str,
object: &str,
fi: &FileInfo,
opts: &ObjectOptions,
) -> Result<bool> {
ensure_decommission_tier_free_version_commit_fence(bucket, object, opts)?;
let disks = self.disks.read().await.clone();
let preflight = disks.iter().map(|disk| async {
match disk {
Some(disk) => inspect_decommission_tier_free_version_target(disk, bucket, object, fi).await,
None => Ok(false),
}
});
let mut equivalent = 0;
for result in join_all(preflight).await {
if result? {
equivalent += 1;
}
}
ensure_decommission_tier_free_version_commit_fence(bucket, object, opts)?;
Ok(equivalent >= self.default_write_quorum())
}
#[tracing::instrument(skip(self, fi, opts))] #[tracing::instrument(skip(self, fi, opts))]
pub(crate) async fn decommission_tiered_object( pub(crate) async fn decommission_tiered_object(
&self, &self,
@@ -10043,147 +9887,6 @@ mod tests {
assert_ne!(updated.erasure.distribution, original.erasure.distribution); assert_ne!(updated.erasure.distribution, original.erasure.distribution);
} }
#[tokio::test]
async fn decommission_tier_free_version_preserves_remote_identity() {
let set_disks = make_local_bucket_test_set_disks().await;
let bucket = "free-version-decommission";
let object = "object.txt";
set_disks
.make_bucket(bucket, &MakeBucketOptions::default())
.await
.expect("target bucket should exist before free-version migration");
let version_id = Uuid::new_v4();
let mut free_version = FileInfo {
name: object.to_string(),
volume: bucket.to_string(),
version_id: Some(version_id),
mod_time: Some(time::OffsetDateTime::now_utc()),
deleted: true,
transition_tier: "WARM-TIER".to_string(),
transitioned_objname: "remote/object".to_string(),
..Default::default()
};
free_version.set_tier_free_version();
// Decoded free versions always carry the on-disk free-version
// suffix alongside the in-memory tier marker; mirror that here so
// the record satisfies delete-marker metadata validation.
rustfs_utils::http::metadata_compat::insert_str(
&mut free_version.metadata,
rustfs_utils::http::metadata_compat::SUFFIX_FREE_VERSION,
String::new(),
);
set_disks
.decommission_tier_free_version(bucket, object, &free_version, &ObjectOptions::default())
.await
.expect("free-version metadata should reach the target quorum");
set_disks
.decommission_tier_free_version(bucket, object, &free_version, &ObjectOptions::default())
.await
.expect("replaying the same free-version metadata should be idempotent");
let versions = set_disks
.load_file_info_versions_exact(bucket, object)
.await
.expect("migrated free-version metadata should decode")
.expect("migrated free-version metadata should exist");
let migrated = versions
.versions
.iter()
.find(|version| version.version_id == Some(version_id))
.expect("free version should be present on the target");
assert_eq!(
versions
.versions
.iter()
.filter(|version| version.version_id == Some(version_id))
.count(),
1
);
assert!(migrated.tier_free_version());
assert_eq!(migrated.transition_tier, "WARM-TIER");
assert_eq!(migrated.transitioned_objname, "remote/object");
}
#[tokio::test]
async fn decommission_tier_free_version_resume_requires_write_quorum() {
let set_disks = make_local_bucket_test_set_disks_with_drive_count(4).await;
let bucket = "free-version-decommission-resume";
let object = "object.txt";
set_disks
.make_bucket(bucket, &MakeBucketOptions::default())
.await
.expect("target bucket should exist before free-version migration");
let mut free_version = FileInfo {
name: object.to_string(),
volume: bucket.to_string(),
version_id: Some(Uuid::new_v4()),
mod_time: Some(time::OffsetDateTime::now_utc()),
deleted: true,
transition_tier: "WARM-TIER".to_string(),
transitioned_objname: "remote/object".to_string(),
..Default::default()
};
free_version.set_tier_free_version();
// Decoded free versions always carry the on-disk free-version
// suffix alongside the in-memory tier marker; mirror that here so
// the record satisfies delete-marker metadata validation.
rustfs_utils::http::metadata_compat::insert_str(
&mut free_version.metadata,
rustfs_utils::http::metadata_compat::SUFFIX_FREE_VERSION,
String::new(),
);
let opts = ObjectOptions::default();
let disks = set_disks.get_disks_internal().await;
for disk in disks.iter().take(2).flatten() {
disk.write_metadata("", bucket, object, free_version.clone())
.await
.expect("partial first attempt should leave equivalent metadata");
}
assert!(
!set_disks
.has_decommission_tier_free_version_write_quorum(bucket, object, &free_version, &opts)
.await
.expect("partial target metadata should remain valid"),
"write-quorum-minus-one must not be accepted as an idempotent migration"
);
disks[2]
.as_ref()
.expect("third target disk should be online")
.write_metadata("", bucket, object, free_version.clone())
.await
.expect("third equivalent target write should complete quorum");
assert!(
set_disks
.has_decommission_tier_free_version_write_quorum(bucket, object, &free_version, &opts)
.await
.expect("write-quorum target metadata should remain valid")
);
}
#[test]
fn decommission_tier_free_version_commit_rejects_lost_fence() {
let opts = ObjectOptions {
namespace_lock_fence: Some(NamespaceLockFence::lost_for_test()),
..Default::default()
};
let err = ensure_decommission_tier_free_version_commit_fence("bucket", "object", &opts)
.expect_err("lost target lock must fail the free-version commit");
assert!(matches!(
err,
Error::NamespaceLockQuorumUnavailable {
mode: "decommission_tier_free_version_commit",
required: 1,
achieved: 0,
..
}
));
}
#[test] #[test]
fn test_resolve_tiered_decommission_write_quorum_result_allows_successful_quorum() { fn test_resolve_tiered_decommission_write_quorum_result_allows_successful_quorum() {
let errs = vec![None, None, Some(DiskError::DiskNotFound), None]; let errs = vec![None, None, Some(DiskError::DiskNotFound), None];
-450
View File
@@ -553,8 +553,6 @@ mod tests {
should_retry_format_load, should_retry_local_decommission_resume, wait_for_local_decommission_resume_delay, should_retry_format_load, should_retry_local_decommission_resume, wait_for_local_decommission_resume_delay,
}; };
#[cfg(feature = "test-util")] #[cfg(feature = "test-util")]
use crate::disk::DiskAPI;
#[cfg(feature = "test-util")]
use crate::{ use crate::{
bucket::lifecycle::{ bucket::lifecycle::{
lifecycle::{TRANSITION_PENDING, TransitionOptions}, lifecycle::{TRANSITION_PENDING, TransitionOptions},
@@ -1309,77 +1307,6 @@ mod tests {
}); });
} }
#[cfg(feature = "test-util")]
async fn seed_transitioned_free_version(
ctx: &Arc<crate::runtime::instance::InstanceContext>,
store: &Arc<crate::store::ECStore>,
bucket: &str,
object: &str,
) -> (uuid::Uuid, uuid::Uuid) {
let tier_name = format!("DECOMFREE{}", uuid::Uuid::new_v4().simple());
register_mock_tier(&ctx.tier_config_mgr(), &tier_name).await;
let mut reader = PutObjReader::from_vec(b"transitioned source bytes".to_vec());
let source = store.pools[0]
.put_object(
bucket,
object,
&mut reader,
&ObjectOptions {
versioned: true,
..Default::default()
},
)
.await
.expect("write transitioned decommission source");
let source_version = source.version_id.expect("transitioned source must be versioned");
store.pools[0]
.transition_object(
bucket,
object,
&ObjectOptions {
versioned: true,
version_id: Some(source_version.to_string()),
transition: TransitionOptions {
status: TRANSITION_PENDING.to_string(),
tier: tier_name,
etag: source.etag.clone().expect("transitioned source must have an ETag"),
..Default::default()
},
mod_time: source.mod_time,
..Default::default()
},
)
.await
.expect("transition source before decommission");
store.pools[0]
.delete_object(
bucket,
object,
ObjectOptions {
versioned: true,
version_id: Some(source_version.to_string()),
..Default::default()
},
)
.await
.expect("delete transitioned source version");
let versions = store.pools[0]
.get_disks_by_key(object)
.load_file_info_versions_exact(bucket, object)
.await
.expect("source versions should decode after transition delete")
.expect("source free version should remain after transition delete");
let free_version = versions
.versions
.iter()
.find(|version| version.tier_free_version())
.and_then(|version| version.version_id)
.expect("transition delete should create a free version");
(source_version, free_version)
}
async fn write_decommission_test_multipart_source( async fn write_decommission_test_multipart_source(
store: &Arc<crate::store::ECStore>, store: &Arc<crate::store::ECStore>,
pool_idx: usize, pool_idx: usize,
@@ -3919,383 +3846,6 @@ mod tests {
shutdown.cancel(); shutdown.cancel();
} }
#[cfg(feature = "test-util")]
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
#[serial_test::serial(storage_class_env)]
async fn decommission_entry_skips_cleanup_only_marker_when_free_version_is_present() {
let temp_dir = tempfile::tempdir().expect("create free-version decommission store dir");
let (ctx, store, shutdown) =
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "decommission-free-marker", &[4, 4])).await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
let bucket = format!("decom-free-marker-{}", uuid::Uuid::new_v4());
let object = "free-marker-object";
store
.make_bucket(&bucket, &MakeBucketOptions::default())
.await
.expect("create free-version decommission bucket");
let (_, free_version) = seed_transitioned_free_version(&ctx, &store, &bucket, object).await;
let source_free = store.pools[0]
.get_disks_by_key(object)
.load_file_info_versions_exact(&bucket, object)
.await
.expect("source free-version metadata should decode")
.and_then(|versions| {
versions
.versions
.into_iter()
.find(|version| version.version_id == Some(free_version) && version.tier_free_version())
})
.expect("source free-version identity should be present before decommission");
let mut target_reader = PutObjReader::from_vec(b"target ordinary bytes".to_vec());
let target_w = store.pools[1]
.put_object(
&bucket,
object,
&mut target_reader,
&ObjectOptions {
versioned: true,
..Default::default()
},
)
.await
.expect("write unrelated target version");
let target_w_version = target_w.version_id.expect("target version should have an id");
let marker = store.pools[0]
.delete_object(
&bucket,
object,
ObjectOptions {
versioned: true,
..Default::default()
},
)
.await
.expect("write cleanup-only delete marker");
assert!(marker.delete_marker);
mark_test_pool_decommissioning(&store, 0).await;
let source_set = store.pools[0].get_disks_by_key(object);
store
.decommission_entry_for_test_with_bucket_incarnation(
0,
MetaCacheEntry {
name: object.to_string(),
..Default::default()
},
bucket.clone(),
source_set.clone(),
)
.await
.expect("real decommission entry should migrate the free version");
let target_versions = store.pools[1]
.get_disks_by_key(object)
.load_file_info_versions_exact(&bucket, object)
.await
.expect("target versions should decode")
.expect("target free version should be present");
assert!(
target_versions
.versions
.iter()
.any(|version| { version.version_id == Some(free_version) && version.tier_free_version() })
);
let migrated_free = target_versions
.versions
.iter()
.find(|version| version.version_id == Some(free_version) && version.tier_free_version())
.expect("migrated free-version identity should remain readable from target disks");
assert!(
crate::store::tiered_data_movement_source_matches(&source_free, migrated_free)
.expect("migrated free-version identity should decode")
);
let retained_w = target_versions
.versions
.iter()
.find(|version| version.version_id == Some(target_w_version))
.expect("unrelated target version should remain");
assert!(!retained_w.deleted && !retained_w.tier_free_version());
assert_eq!(retained_w.size, target_w.size);
assert_eq!(retained_w.get_etag(), target_w.etag);
assert!(
target_versions
.versions
.iter()
.all(|version| { version.tier_free_version() || !version.deleted })
);
assert!(
source_set
.load_file_info_versions_exact(&bucket, object)
.await
.expect("source versions should be readable after cleanup")
.is_none(),
"successful free-version migration should permit source cleanup"
);
let (heal_versions, _, _) = store
.heal_walk_versions_page(1, 0, &bucket, "", None, 2, 16, true)
.await
.expect("heal walk should decode the migrated free version");
let free_version_string = free_version.to_string();
let healed_free = heal_versions
.iter()
.find(|version| version.version_id.as_deref() == Some(free_version_string.as_str()))
.expect("heal walk should surface the migrated free version");
let healed_info = healed_free
.lifecycle_object_info
.as_ref()
.expect("heal walk should retain lifecycle identity for the migrated free version");
assert!(healed_info.transitioned_object.free_version);
assert_eq!(healed_info.transitioned_object.tier, source_free.transition_tier);
assert_eq!(healed_info.transitioned_object.name, source_free.transitioned_objname);
shutdown.cancel();
}
#[cfg(feature = "test-util")]
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
#[serial_test::serial(storage_class_env)]
async fn decommission_entry_allows_free_version_consumed_before_source_lock() {
let temp_dir = tempfile::tempdir().expect("create consumed free-version store dir");
let (ctx, store, shutdown) =
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "decommission-free-consumed", &[4, 4])).await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
let bucket = format!("decom-free-consumed-{}", uuid::Uuid::new_v4());
let object = "free-consumed-object";
store
.make_bucket(&bucket, &MakeBucketOptions::default())
.await
.expect("create consumed free-version bucket");
let (_, free_version) = seed_transitioned_free_version(&ctx, &store, &bucket, object).await;
mark_test_pool_decommissioning(&store, 0).await;
let source_set = store.pools[0].get_disks_by_key(object);
let barrier = crate::store::object::DecommissionFreeVersionSourceRaceBarrier::install(&bucket, object);
let decommission = tokio::spawn({
let store = store.clone();
let bucket = bucket.clone();
let source_set = source_set.clone();
async move {
store
.decommission_entry_for_test_with_bucket_incarnation(
0,
MetaCacheEntry {
name: object.to_string(),
..Default::default()
},
bucket,
source_set,
)
.await
}
});
barrier.wait_until_paused().await;
store.pools[0]
.delete_object(
&bucket,
object,
ObjectOptions {
versioned: true,
version_id: Some(free_version.to_string()),
incl_free_versions: true,
..Default::default()
},
)
.await
.expect("lifecycle should consume the source free version before decommission locks it");
assert!(
source_set
.load_file_info_versions_exact(&bucket, object)
.await
.expect("consumed source metadata should remain readable")
.is_none(),
"the lifecycle delete should remove the source free version"
);
barrier.release();
decommission
.await
.expect("decommission task should join")
.expect("a concurrently consumed free version should not fail source cleanup");
assert!(
store.pools[1]
.get_disks_by_key(object)
.load_file_info_versions_exact(&bucket, object)
.await
.expect("target metadata should remain readable")
.is_none(),
"an already consumed free version should not be recreated on the target"
);
shutdown.cancel();
}
#[cfg(feature = "test-util")]
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
#[serial_test::serial(storage_class_env)]
async fn decommission_entry_rejects_subquorum_free_version_conflict_and_retains_source() {
let temp_dir = tempfile::tempdir().expect("create sub-quorum free-version store dir");
let (ctx, store, shutdown) =
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "decommission-free-conflict", &[4, 4])).await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
let bucket = format!("decom-free-conflict-{}", uuid::Uuid::new_v4());
let object = "free-conflict-object";
store
.make_bucket(&bucket, &MakeBucketOptions::default())
.await
.expect("create sub-quorum conflict bucket");
let (_, free_version) = seed_transitioned_free_version(&ctx, &store, &bucket, object).await;
let source_free = store.pools[0]
.get_disks_by_key(object)
.load_file_info_versions_exact(&bucket, object)
.await
.expect("source free version should decode before crash replay setup")
.and_then(|versions| {
versions
.versions
.into_iter()
.find(|version| version.version_id == Some(free_version))
})
.expect("source free version should be available for crash replay setup");
let mut target_reader = PutObjReader::from_vec(b"target ordinary bytes".to_vec());
let target = store.pools[1]
.put_object(
&bucket,
object,
&mut target_reader,
&ObjectOptions {
versioned: true,
..Default::default()
},
)
.await
.expect("seed ordinary target version");
let target_version = target.version_id.expect("target version must have an ID");
let target_disks = store.pools[1].get_disks_by_key(object).disks.read().await.clone();
for disk in target_disks.iter().skip(1) {
disk.as_ref()
.expect("target crash replay quorum disk should be online")
.write_metadata("", &bucket, object, source_free.clone())
.await
.expect("seed an equivalent free version on the target quorum");
}
let conflict_path = temp_dir
.path()
.join(format!("pool1/set0/disk0/{bucket}/{object}/{STORAGE_FORMAT_FILE}"));
let encoded = tokio::fs::read(&conflict_path)
.await
.expect("target metadata should be readable");
let mut metadata = FileMeta::load(&encoded).expect("target metadata should decode");
let target_index = metadata
.versions
.iter()
.position(|version| version.header.version_id == Some(target_version))
.expect("target version should be present on the conflict disk");
let mut target_meta = metadata.versions[target_index]
.parse_version_meta()
.expect("target version metadata should decode");
target_meta
.object
.as_mut()
.expect("target conflict must remain an ordinary object")
.version_id = Some(free_version);
metadata.versions[target_index] = target_meta.try_into().expect("conflict metadata should encode");
let expected_conflict_meta = metadata.versions[target_index].meta.clone();
let expected_conflict = metadata.versions[target_index]
.into_fileinfo(&bucket, object, true)
.expect("conflict metadata should decode as an ordinary object");
let duplicate_free: rustfs_filemeta::FileMetaShallowVersion = rustfs_filemeta::FileMetaVersion::from(source_free.clone())
.try_into()
.expect("duplicate free metadata should encode");
metadata.versions.insert(target_index, duplicate_free);
tokio::fs::write(&conflict_path, metadata.marshal_msg().expect("conflict metadata should encode"))
.await
.expect("write sub-quorum conflict metadata");
mark_test_pool_decommissioning(&store, 0).await;
let source_set = store.pools[0].get_disks_by_key(object);
store
.decommission_entry_for_test_with_bucket_incarnation(
0,
MetaCacheEntry {
name: object.to_string(),
..Default::default()
},
bucket.clone(),
source_set.clone(),
)
.await
.expect("conflicted decommission entry should retain the source and retry later");
let source_versions = source_set
.load_file_info_versions_exact(&bucket, object)
.await
.expect("retained source versions should decode")
.expect("source free version should be retained after conflict");
assert!(
source_versions
.versions
.iter()
.any(|version| { version.version_id == Some(free_version) && version.tier_free_version() })
);
let post_encoded = tokio::fs::read(&conflict_path)
.await
.expect("conflict metadata should remain readable");
let post_metadata = FileMeta::load(&post_encoded).expect("post-conflict metadata should decode");
let same_id = post_metadata
.versions
.iter()
.filter(|version| version.header.version_id == Some(free_version))
.collect::<Vec<_>>();
assert_eq!(same_id.len(), 2, "conflict metadata should retain both same-ID records");
assert_eq!(same_id.iter().filter(|version| version.header.free_version()).count(), 1);
assert_eq!(same_id.iter().filter(|version| !version.header.free_version()).count(), 1);
let post_conflict = same_id
.into_iter()
.find(|version| !version.header.free_version())
.expect("ordinary conflict version must remain addressable by the source ID");
let post_conflict_info = post_conflict
.into_fileinfo(&bucket, object, true)
.expect("post-conflict ordinary metadata should decode");
assert!(!post_conflict_info.deleted && !post_conflict_info.tier_free_version());
assert_eq!(post_conflict.meta, expected_conflict_meta);
assert_eq!(post_conflict_info.size, expected_conflict.size);
assert_eq!(post_conflict_info.data_dir, expected_conflict.data_dir);
assert_eq!(post_conflict_info.metadata, expected_conflict.metadata);
assert_eq!(post_conflict_info.get_etag(), expected_conflict.get_etag());
for disk_index in 0..4 {
let target_path = temp_dir
.path()
.join(format!("pool1/set0/disk{disk_index}/{bucket}/{object}/{STORAGE_FORMAT_FILE}"));
let target_encoded = tokio::fs::read(&target_path)
.await
.expect("target metadata should remain readable");
let target_meta = FileMeta::load(&target_encoded).expect("target metadata should decode");
let same_id = target_meta
.versions
.iter()
.filter(|version| version.header.version_id == Some(free_version))
.collect::<Vec<_>>();
if disk_index == 0 {
assert_eq!(same_id.len(), 2);
assert!(same_id[0].header.free_version());
assert!(!same_id[1].header.free_version());
} else {
assert_eq!(same_id.len(), 1);
assert!(same_id[0].header.free_version());
}
}
let sweep_err = store
.check_after_decommission_for_test(0)
.await
.expect_err("final sweep must report the retained free version");
assert!(
sweep_err.to_string().contains("version(s) were found"),
"unexpected final sweep error: {sweep_err}"
);
shutdown.cancel();
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)] #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
#[serial_test::serial(storage_class_env)] #[serial_test::serial(storage_class_env)]
async fn versioned_batch_delete_marker_skips_decommission_source() { async fn versioned_batch_delete_marker_skips_decommission_source() {
+1 -1
View File
@@ -151,7 +151,7 @@ pub(crate) mod init_format;
pub(crate) mod list_objects; pub(crate) mod list_objects;
mod multipart; mod multipart;
mod object; mod object;
pub(crate) use object::{ObjectLockDiagGuard, SourceCleanupMutationFence, tiered_data_movement_source_matches}; pub(crate) use object::{ObjectLockDiagGuard, SourceCleanupMutationFence};
pub use object::{ pub use object::{
PrepareSelectObjectSnapshotError, PreparedGetObjectReader, SelectObjectSnapshot, SelectObjectSnapshotReadError, PrepareSelectObjectSnapshotError, PreparedGetObjectReader, SelectObjectSnapshot, SelectObjectSnapshotReadError,
SnapshotConsistencyError, SnapshotConsistencyError,
+15 -143
View File
@@ -490,82 +490,6 @@ fn decommission_mutation_fence_for_test(
.map(|hook| hook.fence.clone()) .map(|hook| hook.fence.clone())
} }
#[cfg(test)]
struct DecommissionFreeVersionSourceRaceState {
bucket: String,
object: String,
arrived: tokio::sync::Notify,
release: tokio::sync::Notify,
}
#[cfg(test)]
pub(crate) struct DecommissionFreeVersionSourceRaceBarrier {
state: Arc<DecommissionFreeVersionSourceRaceState>,
}
#[cfg(test)]
static DECOMMISSION_FREE_VERSION_SOURCE_RACE_BARRIER: std::sync::OnceLock<
std::sync::Mutex<Option<Arc<DecommissionFreeVersionSourceRaceState>>>,
> = std::sync::OnceLock::new();
#[cfg(test)]
impl DecommissionFreeVersionSourceRaceBarrier {
pub(crate) fn install(bucket: &str, object: &str) -> Self {
let state = Arc::new(DecommissionFreeVersionSourceRaceState {
bucket: bucket.to_string(),
object: object.to_string(),
arrived: tokio::sync::Notify::new(),
release: tokio::sync::Notify::new(),
});
let mut slot = DECOMMISSION_FREE_VERSION_SOURCE_RACE_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("decommission free-version source race barrier should not poison");
assert!(slot.is_none(), "decommission free-version source race barrier must be unique");
*slot = Some(Arc::clone(&state));
Self { state }
}
pub(crate) async fn wait_until_paused(&self) {
tokio::time::timeout(Duration::from_secs(30), self.state.arrived.notified())
.await
.expect("decommission should pause before acquiring the free-version source lock");
}
pub(crate) fn release(&self) {
self.state.release.notify_one();
}
}
#[cfg(test)]
impl Drop for DecommissionFreeVersionSourceRaceBarrier {
fn drop(&mut self) {
self.state.release.notify_one();
let mut slot = DECOMMISSION_FREE_VERSION_SOURCE_RACE_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("decommission free-version source race barrier should not poison");
if slot.as_ref().is_some_and(|state| Arc::ptr_eq(state, &self.state)) {
*slot = None;
}
}
}
#[cfg(test)]
async fn pause_decommission_free_version_before_source_lock(bucket: &str, object: &str) {
let state = DECOMMISSION_FREE_VERSION_SOURCE_RACE_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("decommission free-version source race barrier should not poison")
.as_ref()
.filter(|state| state.bucket == bucket && state.object == object)
.cloned();
if let Some(state) = state {
state.arrived.notify_one();
state.release.notified().await;
}
}
pub(crate) struct SourceCleanupMutationFence { pub(crate) struct SourceCleanupMutationFence {
guard: ObjectLockDiagGuard, guard: ObjectLockDiagGuard,
source_lock_covered: bool, source_lock_covered: bool,
@@ -1570,15 +1494,13 @@ fn is_equivalent_data_movement_tiered_object(source: &rustfs_filemeta::FileInfo,
&& source_actual_size == target_actual_size && source_actual_size == target_actual_size
} }
pub(crate) fn tiered_data_movement_source_matches( fn tiered_data_movement_source_matches(
expected: &rustfs_filemeta::FileInfo, expected: &rustfs_filemeta::FileInfo,
current: &rustfs_filemeta::FileInfo, current: &rustfs_filemeta::FileInfo,
) -> Result<bool> { ) -> Result<bool> {
let expected_backend = crate::services::tier::tier::tier_destination_id_from_metadata(&expected.metadata)?; let expected_backend = crate::services::tier::tier::tier_destination_id_from_metadata(&expected.metadata)?;
let current_backend = crate::services::tier::tier::tier_destination_id_from_metadata(&current.metadata)?; let current_backend = crate::services::tier::tier::tier_destination_id_from_metadata(&current.metadata)?;
Ok(expected.version_id == current.version_id Ok(expected.version_id == current.version_id
&& expected.deleted == current.deleted
&& expected.tier_free_version() == current.tier_free_version()
&& expected.data_dir == current.data_dir && expected.data_dir == current.data_dir
&& expected.mod_time == current.mod_time && expected.mod_time == current.mod_time
&& expected.size == current.size && expected.size == current.size
@@ -1592,15 +1514,6 @@ pub(crate) fn tiered_data_movement_source_matches(
&& expected_backend == current_backend) && expected_backend == current_backend)
} }
fn decommission_free_version_overwrite_error(bucket: &str, object: &str, version_id: Option<Uuid>) -> Error {
StorageError::DataMovementOverwriteErr(
bucket.to_owned(),
object.to_owned(),
version_id.map(|id| id.to_string()).unwrap_or_default(),
)
.into()
}
fn should_check_data_movement_resume_target(src_pool_idx: usize, target_pool_idx: usize) -> bool { fn should_check_data_movement_resume_target(src_pool_idx: usize, target_pool_idx: usize) -> bool {
target_pool_idx != src_pool_idx target_pool_idx != src_pool_idx
} }
@@ -2308,23 +2221,6 @@ impl ECStore {
) )
} }
async fn has_equivalent_data_movement_tier_free_version(
&self,
bucket: &str,
object: &str,
source: &rustfs_filemeta::FileInfo,
opts: &ObjectOptions,
target_pool_idx: usize,
) -> Result<bool> {
let pool = self
.pools
.get(target_pool_idx)
.ok_or_else(|| Error::other(format!("invalid tiered data movement target pool {target_pool_idx}")))?;
pool.get_disks_by_key(object)
.has_decommission_tier_free_version_write_quorum(bucket, object, source, opts)
.await
}
fn resolve_decommission_target_pool_idx_result(result: Result<usize>, bucket: &str, object: &str) -> Result<usize> { fn resolve_decommission_target_pool_idx_result(result: Result<usize>, bucket: &str, object: &str) -> Result<usize> {
result.map_err(|err| Error::other(format!("failed to select decommission target pool for {bucket}/{object}: {err}"))) result.map_err(|err| Error::other(format!("failed to select decommission target pool for {bucket}/{object}: {err}")))
} }
@@ -2344,10 +2240,6 @@ impl ECStore {
check_put_object_args(bucket, object)?; check_put_object_args(bucket, object)?;
let mut opts = opts.clone(); let mut opts = opts.clone();
let is_free_version = fi.tier_free_version();
if is_free_version {
opts.incl_free_versions = true;
}
let bucket_incarnation_fence = if is_meta_bucketname(bucket) { let bucket_incarnation_fence = if is_meta_bucketname(bucket) {
None None
} else { } else {
@@ -2385,10 +2277,6 @@ impl ECStore {
&object, &object,
)? )?
}; };
#[cfg(test)]
if is_free_version {
pause_decommission_free_version_before_source_lock(bucket, logical_object).await;
}
let _object_guards = self let _object_guards = self
.acquire_data_movement_object_write_locks(bucket, &object, opts.src_pool_idx, idx, &mut opts) .acquire_data_movement_object_write_locks(bucket, &object, opts.src_pool_idx, idx, &mut opts)
.await?; .await?;
@@ -2406,7 +2294,7 @@ impl ECStore {
versions versions
.versions .versions
.iter() .iter()
.find(|current| current.version_id == fi.version_id && current.tier_free_version() == is_free_version) .find(|current| current.version_id == fi.version_id && !current.tier_free_version())
}) })
.ok_or_else(|| to_object_err(StorageError::FileNotFound, vec![bucket, object.as_str()]))?; .ok_or_else(|| to_object_err(StorageError::FileNotFound, vec![bucket, object.as_str()]))?;
if !tiered_data_movement_source_matches(fi, current_source)? { if !tiered_data_movement_source_matches(fi, current_source)? {
@@ -2421,40 +2309,24 @@ impl ECStore {
.get_available_pool_idx_excluding(bucket, &object, fi.size, opts.src_pool_idx) .get_available_pool_idx_excluding(bucket, &object, fi.size, opts.src_pool_idx)
.await; .await;
let target_pool_idx = resolve_data_movement_resume_target_pool(idx, resume_target_pool_idx, opts.src_pool_idx); let target_pool_idx = resolve_data_movement_resume_target_pool(idx, resume_target_pool_idx, opts.src_pool_idx);
if is_free_version && target_pool_idx == opts.src_pool_idx {
return Err(Error::DiskFull);
}
let equivalent = if is_free_version {
self.has_equivalent_data_movement_tier_free_version(bucket, &object, &fi, &opts, target_pool_idx)
.await?
} else {
self.has_equivalent_data_movement_tiered_object(bucket, &object, &fi, &opts, target_pool_idx)
.await?
};
if equivalent {
return Ok(());
}
return Err(decommission_free_version_overwrite_error(bucket, &object, fi.version_id));
}
let result = if is_free_version {
if self if self
.has_equivalent_data_movement_tier_free_version(bucket, &object, &fi, &opts, idx) .has_equivalent_data_movement_tiered_object(bucket, &object, &fi, &opts, target_pool_idx)
.await? .await?
{ {
return Ok(()); return Ok(());
} }
self.pools[idx]
.get_disks_by_key(&object) return Err(StorageError::DataMovementOverwriteErr(
.decommission_tier_free_version(bucket, &object, &fi, &opts) bucket.to_owned(),
.await object.to_owned(),
} else { opts.version_id.clone().unwrap_or_default(),
self.pools[idx] ));
.get_disks_by_key(&object) }
.decommission_tiered_object(bucket, &object, &fi, &opts)
.await let result = self.pools[idx]
}; .get_disks_by_key(&object)
.decommission_tiered_object(bucket, &object, &fi, &opts)
.await;
if matches!(result, Err(Error::PreconditionFailed)) { if matches!(result, Err(Error::PreconditionFailed)) {
if self if self
.has_equivalent_data_movement_tiered_object(bucket, &object, &fi, &opts, idx) .has_equivalent_data_movement_tiered_object(bucket, &object, &fi, &opts, idx)
-19
View File
@@ -90,21 +90,6 @@ fn legacy_data_key_for_version(version_id: Option<Uuid>) -> Option<String> {
pub const TRANSITION_COMPLETE: &str = "complete"; pub const TRANSITION_COMPLETE: &str = "complete";
pub const TRANSITION_PENDING: &str = "pending"; pub const TRANSITION_PENDING: &str = "pending";
/// xl.meta key marking a tier free-version record.
///
/// A free version is a delete-marker-shaped cleanup hint appended by
/// [`MetaObject::delete_version`] when a version whose remote transition
/// completed is removed from xl.meta; it carries the remote tier identity for
/// an idempotent remote delete and is never a user-visible version
/// (`num_versions` excludes it). While the record exists it is consumed by the
/// lifecycle free-version recovery scan and the usage scanner, which re-enqueue
/// the pending remote delete, and by heal metadata walks. On S3 and lifecycle
/// delete paths the same obligation is also carried by a committed tier-journal
/// entry; deletes without such an entry (for example a removed version whose
/// transition state decodes as unknown) rely on this record alone until the
/// worker removes it after a successful remote delete. Decommission preserves
/// the record and its remote identity on the target pool before source cleanup
/// — see docs/architecture/decommission-compatibility.md.
pub const FREE_VERSION: &str = "free-version"; pub const FREE_VERSION: &str = "free-version";
pub const TRANSITION_STATUS: &str = "transition-status"; pub const TRANSITION_STATUS: &str = "transition-status";
@@ -462,10 +447,6 @@ impl FileMeta {
}; };
if let Some(fidx) = existing_idx { if let Some(fidx) = existing_idx {
let existing = self.versions[fidx].parse_version_meta()?;
if existing.free_version() != version.free_version() {
return Err(Error::other("cannot replace a free version with a non-free version"));
}
return self.set_idx(fidx, version); return self.set_idx(fidx, version);
} }
-9
View File
@@ -2725,15 +2725,6 @@ impl MetaObject {
self.meta_sys.retain(|k, _| !k.starts_with("X-Amz-Restore")); self.meta_sys.retain(|k, _| !k.starts_with("X-Amz-Restore"));
} }
/// Builds the free-version cleanup record appended when a transitioned
/// version is removed from xl.meta. The record keeps the remote tier
/// identity so the lifecycle worker can issue the idempotent remote delete
/// and only then remove the record; until then the recovery scan and the
/// usage scanner keep re-enqueueing it. S3 and lifecycle deletes also
/// persist a committed tier-journal entry for the same remote delete. The
/// decommission path copies this record unchanged before source cleanup,
/// including when the transition state is unknown — see
/// docs/architecture/decommission-compatibility.md.
pub fn init_free_version(&self, fi: &FileInfo) -> Result<(FileMetaVersion, bool)> { pub fn init_free_version(&self, fi: &FileInfo) -> Result<(FileMetaVersion, bool)> {
if fi.skip_tier_free_version() { if fi.skip_tier_free_version() {
return Ok((FileMetaVersion::default(), false)); return Ok((FileMetaVersion::default(), false));
@@ -153,99 +153,6 @@ No migration step is required for these decisions because this note documents th
current RustFS behavior. Changing either decision later requires an operator current RustFS behavior. Changing either decision later requires an operator
compatibility note and updated characterization tests. compatibility note and updated characterization tests.
## Tier Free Versions During Decommission
A tier free version is an internal xl.meta record (`rustfs_filemeta::FREE_VERSION`,
flagged `XL_FLAG_FREE_VERSION`) shaped like a delete marker. It is created by
`MetaObject::init_free_version` when a version whose remote transition completed is
deleted locally: the visible version is removed and the record keeps the remote-tier
identity (tier, object name, version id, state, destination id) needed for an
idempotent remote delete. Free versions are not user-visible versions; `num_versions`
and all listing/GET paths exclude them.
### Lifecycle And Consumers
Creation: any local delete that removes a version whose transition status is
`complete` appends the record via `MetaObject::delete_version`
`init_free_version` (skipped only when `skip_tier_free_version` is set, as on
data-movement copies). The same deletes also persist a durable tier-journal
entry on every user-facing path: S3 single deletes (`execute_delete_object`
`delete_object_with_tier_delete_journal`), S3 batch deletes, lifecycle expiry,
and lifecycle delete-all all prepare and commit a journal entry around the
delete. A journal entry is omitted when the removed version's transition state
decodes as `TransitionVersionState::Unknown`, or on internal journal-less
delete paths that never touch transitioned user objects.
Consumption while the record exists: the background recovery loop started by
`init_background_expiry` (spawned by `spawn_tier_free_version_recovery_once`,
enabled by default) scans disks for pending records and re-enqueues them; the
usage scanner does the same; the lifecycle worker then deletes the remote tier
object idempotently and only afterwards removes the local record. Heal walks
include free-version records in metadata healing. Transition planning,
replication, restore, GET, listings, and usage aggregation never depend on
them.
### Decommission Handling
The exact decommission inventory loader (`load_file_info_versions_exact` via
`get_all_file_info_versions`) keeps free-version records inline in `versions`.
The migration loop handles them before lifecycle expiry and delete-marker
shortcuts. It selects a target pool using the free-version-aware lookup, then
writes the original free record to every target disk with the normal metadata
write quorum. The free-version marker, local version id, transition identity,
transition state, and destination id are preserved at the FileInfo/metadata
boundary.
The source record is physically removed only after the target write quorum has
committed and the source cleanup preflight still matches the exact inventory.
If the lifecycle worker has already completed the remote delete and removed the
source record before decommission acquires the source lock, decommission records
that identity as already consumed and treats the missing source record as safe.
If target capacity, metadata validation, lock fencing, or quorum fails, the
source record remains and the entry records `state = "free_version_retained"`
with reason `tier_free_version_migration_failed`; the worker retries the
operation on a later pass. A target record with the same version id is accepted
only when its free-version identity matches; a conflicting ordinary version or
different free record is an overwrite error. This makes retries idempotent and
prevents a free record from replacing a user-visible version.
### Reference-Audit Result
After migration, user-facing GET/list/transition/replication/restore paths still
exclude the record. Recovery, usage scanning, lifecycle tier cleanup, and heal
continue to see it when they request free versions, so an unresolved remote
delete remains actionable on the target pool. The committed tier journal remains
an independent retry source where one exists; it is not used as a reason to drop
the xl.meta record. In particular, `Unknown` transition state records are
migrated unchanged rather than discarded: the lifecycle worker retains them if
remote identity validation cannot make a delete request.
Each migrated record emits `state = "free_version_migrated"` with reason
`tier_free_version_migrated`. A record consumed before migration emits
`state = "free_version_consumed"` with reason
`tier_free_version_already_consumed`. Each failed record emits the retained state
and failure reason above. The entry also emits a disposition summary with
migrated, consumed, retained, and total counts. The final decommission sweep uses
the exact loader, counts free records still present, and emits one retained
record/reason for each unresolved free version before failing the sweep. This
makes successful migration, completed cleanup, and retained cleanup obligations
visible instead of silently omitting free records.
No new S3-visible version or admin response field is needed: free versions remain
internal and are never counted as user-visible versions. The structured
`decommission_entry` events are the operational status surface for the
free-version disposition; the existing decommission item/failed counters still
report the enclosing object migration result.
Regression guard:
- `decommission_tier_free_version_preserves_remote_identity`
- `decommission_tier_free_version_resume_requires_write_quorum`
- `decommission_tier_free_version_commit_rejects_lost_fence`
- `test_decommission_cleanup_preflight_accepts_migrated_free_version_consumed_from_source`
- `decommission_entry_skips_cleanup_only_marker_when_free_version_is_present`
- `decommission_entry_rejects_subquorum_free_version_conflict_and_retains_source`
## Regression Guard ## Regression Guard
The queued multi-pool contract is guarded by: The queued multi-pool contract is guarded by:
+2 -2
View File
@@ -52,7 +52,7 @@
| fault_proxy | 7 | | | fault_proxy | 7 | |
| get_codec_streaming_compat_test | 1 | | | get_codec_streaming_compat_test | 1 | |
| get_stream_failure_observability_test | 1 | | | get_stream_failure_observability_test | 1 | |
| group_delete_test | 1 | | | group_delete_test | 4 | |
| head_object_consistency_test | 1 | ✅ | | head_object_consistency_test | 1 | ✅ |
| head_object_range_test | 1 | ✅ | | head_object_range_test | 1 | ✅ |
| heal_erasure_disk_rebuild_test | 4 | 🌙 | | heal_erasure_disk_rebuild_test | 4 | 🌙 |
@@ -99,4 +99,4 @@
| tls_hot_reload_test | 1 | ✅ | | tls_hot_reload_test | 1 | ✅ |
| version_id_regression_test | 10 | ✅ | | version_id_regression_test | 10 | ✅ |
**Total listed: 575 tests across 82 modules · PR smoke: 163 tests / 36 modules · merge/main full: 453 tests / 73 modules · nightly replication: 55 tests · nightly cluster faults: 28 tests / 7 modules · nightly protocols: 16 tests** · updated 2026-08-23. **Total listed: 578 tests across 82 modules · PR smoke: 163 tests / 36 modules · merge/main full: 456 tests / 73 modules · nightly replication: 55 tests · nightly cluster faults: 28 tests / 7 modules · nightly protocols: 16 tests** · updated 2026-08-23.