fix(heal): restore bucket metadata before replacement completion (#8042)

This commit is contained in:
Chris
2026-09-21 12:18:06 +08:00
committed by GitHub
parent ae6bbaef27
commit 4047ed6c8d
9 changed files with 451 additions and 7 deletions
@@ -2141,7 +2141,7 @@ impl RemoteDisk {
let file_info_bin = encode_file_info_msgpack(fi)?;
let mut client = self.get_client().await?;
let mut request = Request::new(RenameDataRequest {
bucket_incarnation_id: crate::store::bucket_heal_scope(dst_volume)
bucket_incarnation_id: crate::store::bucket_heal_scope_for_object(dst_volume, dst_path)
.map(|scope| scope.incarnation.as_bytes().to_vec().into())
.unwrap_or_default(),
disk: self.endpoint.to_string(),
+1 -1
View File
@@ -825,7 +825,7 @@ impl Disk {
dst_volume: &str,
dst_path: &str,
) -> Result<RenameDataResp> {
let Some(scope) = crate::store::bucket_heal_scope(dst_volume) else {
let Some(scope) = crate::store::bucket_heal_scope_for_object(dst_volume, dst_path) else {
return self
.rename_data_borrowed_with_fence(src_volume, src_path, fi, dst_volume, dst_path, None)
.await;
+1 -1
View File
@@ -854,7 +854,7 @@ impl SetDisks {
);
let disks = self.get_disks_internal().await;
let bucket_heal_scope = crate::store::bucket_heal_scope(bucket);
let bucket_heal_scope = crate::store::bucket_heal_scope_for_object(bucket, object);
if let Some(scope) = &bucket_heal_scope {
scope.check()?;
}
+309
View File
@@ -58,6 +58,26 @@ pub(crate) fn bucket_heal_scope(bucket: &str) -> Option<Arc<BucketHealScope>> {
.flatten()
}
/// Resolve only the two canonical bucket records; other internal paths carry no bucket authority.
pub(super) fn bucket_metadata_owner(object: &str) -> Option<&str> {
use crate::bucket::metadata::{BUCKET_INCARNATION_FILE, BUCKET_METADATA_FILE};
let (prefix, rest) = object.split_once('/')?;
let (bucket, file) = rest.split_once('/')?;
(prefix == disk::BUCKET_META_PREFIX
&& matches!(file, BUCKET_METADATA_FILE | BUCKET_INCARNATION_FILE)
&& check_valid_bucket_name_strict(bucket).is_ok())
.then_some(bucket)
}
pub(crate) fn bucket_heal_scope_for_object(bucket: &str, object: &str) -> Option<Arc<BucketHealScope>> {
let owner = if bucket == RUSTFS_META_BUCKET {
bucket_metadata_owner(object)?
} else {
bucket
};
bucket_heal_scope(owner)
}
/// Storage-owned proof for the exact version and every selected erasure location.
/// This is an in-process result, never reconstructed from admin drive telemetry.
#[derive(Debug)]
@@ -422,6 +442,97 @@ impl ECStore {
Ok(Arc::ptr_eq(&target_set, &pool.get_disks_by_key(POOL_META_NAME)))
}
/// Restore the bucket configuration and incarnation before retiring a replacement marker.
pub async fn heal_replacement_bucket_metadata(
self: &Arc<Self>,
bucket: &str,
opts: &HealOpts,
targets: &[String],
) -> Result<()> {
use crate::bucket::metadata::{BUCKET_INCARNATION_FILE, BUCKET_METADATA_FILE};
use crate::bucket::metadata_sys::acquire_bucket_metadata_transaction_read_lock_in;
let (Some(pool_index), Some(set_index)) = (opts.pool, opts.set) else {
return Err(Error::PreconditionFailed);
};
if targets.is_empty() || opts.dry_run || opts.no_lock || check_valid_bucket_name_strict(bucket).is_err() {
return Err(Error::PreconditionFailed);
}
let pool = self
.pools
.get(pool_index)
.ok_or_else(|| invalid_heal_pool_index(pool_index, self.pools.len()))?;
let target_set = pool.get_disks_for_heal_object(BUCKET_METADATA_FILE, opts)?;
if targets.iter().any(|target| {
target_set
.set_endpoints
.iter()
.filter(|endpoint| endpoint.to_string() == *target)
.count()
!= 1
}) {
return Err(Error::PreconditionFailed);
}
let objects: Vec<_> = [BUCKET_METADATA_FILE, BUCKET_INCARNATION_FILE]
.into_iter()
.map(|file| format!("{}/{bucket}/{file}", disk::BUCKET_META_PREFIX))
.filter(|object| Arc::ptr_eq(&target_set, &pool.get_disks_by_key(object)))
.collect();
if objects.is_empty() {
return Ok(());
}
// Read the persisted identity without lazy migration: a replacement
// must never turn lost configuration into newly fabricated defaults.
let incarnation = self.bucket_incarnation_id_from_disk(bucket).await?;
let targets = targets.to_vec();
self.run_bucket_heal_at_incarnation(bucket, incarnation, opts, move |store, bucket, opts| async move {
// Match config writers: lifecycle, metadata transaction, object locks.
// Keep the transaction stable through target readback, including when
// the caller is cancelled and the storage owner finishes its write.
let transaction = acquire_bucket_metadata_transaction_read_lock_in(&store.ctx, &bucket).await?;
for object in objects {
let scope = bucket_heal_scope(&bucket).ok_or(Error::PreconditionFailed)?;
scope.check()?;
if transaction.is_lock_lost() {
return Err(Error::PreconditionFailed);
}
// Config objects occupy one pool, unlike pool.bin. Require a
// successful all-pool lookup before accepting another pool's ownership.
if !store.single_pool() {
let (owner, _) = store
.get_pool_info_for_delete_marker(RUSTFS_META_BUCKET, &object, &ObjectOptions::default())
.await?;
if owner.index != pool_index {
continue;
}
}
let metadata_opts = HealOpts { remove: false, ..opts };
let (result, error) = store.heal_object(RUSTFS_META_BUCKET, &object, "", &metadata_opts).await?;
if let Some(error) = error {
return Err(error);
}
for target in &targets {
let mut drives = result.after.drives.iter().filter(|drive| &drive.endpoint == target);
if !drives.next().is_some_and(|drive| drive.state == "ok") || drives.next().is_some() {
return Err(Error::PreconditionFailed);
}
}
if !store
.replacement_targets_have_version(RUSTFS_META_BUCKET, &object, "", pool_index, set_index, &targets)
.await?
{
return Err(Error::PreconditionFailed);
}
}
if transaction.is_lock_lost() {
return Err(Error::PreconditionFailed);
}
Ok(())
})
.await
}
/// Return every live erasure set selected by an object-heal scope.
pub async fn heal_erasure_set_scopes(&self, opts: &HealOpts) -> Result<Vec<(usize, usize)>> {
let pools = self.get_pools_for_heal_object(opts)?;
@@ -1345,6 +1456,16 @@ mod tests {
.await
.is_err()
);
for file in [".metadata.bin", ".bucket-incarnation"] {
let metadata = format!("buckets/{bucket}/{file}");
assert!(
store
.rename_local_data_at_incarnation(&disk_ref, (&bucket, object), &fi, (RUSTFS_META_BUCKET, &metadata), old)
.await
.is_err(),
"stale remote heal must not replace the successor's {file}"
);
}
assert!(store.heal_bucket_at_incarnation(&bucket, old, &opts).await.is_err());
assert!(
store
@@ -1743,6 +1864,194 @@ mod tests {
(temp_dir, store, shutdown)
}
#[test]
fn bucket_metadata_heal_owner_rejects_path_aliases_and_unrelated_records() {
for file in [".metadata.bin", ".bucket-incarnation"] {
assert_eq!(bucket_metadata_owner(&format!("buckets/example/{file}")), Some("example"));
}
for path in [
"/buckets/example/.metadata.bin",
"buckets/../.metadata.bin",
"buckets/example/../.metadata.bin",
"buckets/example/.metadata.bin/extra",
"buckets/example//.metadata.bin",
"buckets/example/.metadata.bin.old",
"config/iam/example/.metadata.bin",
"buckets/example/usage.json",
] {
assert_eq!(bucket_metadata_owner(path), None, "unexpected bucket authority for {path}");
}
}
#[tokio::test]
async fn replacement_bucket_metadata_respects_pool_ownership_and_required_records() {
let (root, store, shutdown) = multi_pool_heal_store().await;
let bucket = "replacement-bucket-metadata";
for invalid_bucket in ["MixedCase", "../example", RUSTFS_META_BUCKET] {
assert!(matches!(
store
.heal_replacement_bucket_metadata(
invalid_bucket,
&HealOpts {
pool: Some(0),
set: Some(0),
..Default::default()
},
&[root.path().join("pool0-disk0").to_string_lossy().into_owned()],
)
.await,
Err(Error::PreconditionFailed)
));
}
store
.make_bucket(bucket, &MakeBucketOptions::default())
.await
.expect("create bucket and its persisted records");
let mut owned = 0;
for pool in 0..2 {
let target = root.path().join(format!("pool{pool}-disk0"));
let records: Vec<_> = [".metadata.bin", ".bucket-incarnation"]
.into_iter()
.map(|file| {
let path = target.join(RUSTFS_META_BUCKET).join("buckets").join(bucket).join(file);
let exists = path.join("xl.meta").exists();
if exists {
std::fs::remove_dir_all(&path).expect("remove replacement metadata shard");
owned += 1;
}
(path, exists)
})
.collect();
store
.heal_replacement_bucket_metadata(
bucket,
&HealOpts {
pool: Some(pool),
set: Some(0),
recreate: true,
..Default::default()
},
&[target.to_string_lossy().into_owned()],
)
.await
.expect("repair owned records and accept authoritative ownership in the other pool");
for (path, existed) in records {
assert_eq!(path.join("xl.meta").exists(), existed, "wrong placement for {}", path.display());
}
}
assert_eq!(owned, 2, "each required record must have exactly one owning pool");
let object = format!("buckets/{bucket}/.metadata.bin");
let owner = store
.get_pool_idx_existing_with_opts(RUSTFS_META_BUCKET, &object, &ObjectOptions::default())
.await
.expect("persisted metadata owner");
let disk = store.pools[owner].disk_set[0].disks.read().await[0]
.clone()
.expect("metadata target disk");
let fi = disk
.read_version(
"",
RUSTFS_META_BUCKET,
&object,
"",
&disk::ReadOptions {
read_data: true,
..Default::default()
},
)
.await
.expect("inline metadata to repair through target-side admission");
assert!(fi.data.is_some(), "fixture must contain an inline metadata shard");
let target_path = root.path().join(format!("pool{owner}-disk0"));
std::fs::remove_dir_all(target_path.join(RUSTFS_META_BUCKET).join(&object)).expect("remove remote target shard");
let incarnation = store
.bucket_incarnation_id_from_disk(bucket)
.await
.expect("current incarnation");
store
.rename_local_data_at_incarnation(
&disk.endpoint().to_string(),
(disk::RUSTFS_META_TMP_BUCKET, "metadata-heal-control"),
&fi,
(RUSTFS_META_BUCKET, &object),
incarnation,
)
.await
.expect("the target-side RPC admission must accept current bucket metadata repair");
assert!(target_path.join(RUSTFS_META_BUCKET).join(&object).join("xl.meta").exists());
let (configuration, _) = read_config_no_lock_preserve_empty_with_metadata(store.clone(), &object)
.await
.expect("read the current bucket configuration");
save_config(store.clone(), &object, configuration)
.await
.expect("publish a newer config revision in the same bucket incarnation");
let committed =
std::fs::read(target_path.join(RUSTFS_META_BUCKET).join(&object).join("xl.meta")).expect("new target configuration");
assert!(
store
.rename_local_data_at_incarnation(
&disk.endpoint().to_string(),
(disk::RUSTFS_META_TMP_BUCKET, "delayed-metadata-heal"),
&fi,
(RUSTFS_META_BUCKET, &object),
incarnation,
)
.await
.is_err(),
"a delayed repair must not overwrite a newer config in the same incarnation"
);
assert_eq!(
std::fs::read(target_path.join(RUSTFS_META_BUCKET).join(&object).join("xl.meta")).expect("retained configuration"),
committed,
"rejected repair must leave the newer configuration unchanged"
);
for pool in 0..2 {
for disk in 0..4 {
let path = root
.path()
.join(format!("pool{pool}-disk{disk}"))
.join(RUSTFS_META_BUCKET)
.join("buckets")
.join(bucket)
.join(".metadata.bin");
if path.exists() {
std::fs::remove_dir_all(path).expect("remove every copy of required bucket configuration");
}
}
}
metadata_sys::remove_bucket_metadata_in(&store.ctx, bucket)
.await
.expect("remove the cached configuration so repair cannot rely on it");
assert!(
store
.heal_replacement_bucket_metadata(
bucket,
&HealOpts {
pool: Some(0),
set: Some(0),
recreate: true,
..Default::default()
},
&[root.path().join("pool0-disk0").to_string_lossy().into_owned()],
)
.await
.is_err(),
"total metadata loss must not be accepted as ownership in another pool"
);
for pool in 0..2 {
assert!(
!root
.path()
.join(format!("pool{pool}-disk0"))
.join(RUSTFS_META_BUCKET)
.join(&object)
.exists(),
"missing required configuration must not be recreated with defaults"
);
}
shutdown.cancel();
}
#[tokio::test]
#[serial_test::serial]
async fn absence_proof_requires_every_selected_pool() {
+1 -1
View File
@@ -426,8 +426,8 @@ mod bucket_fence;
pub(crate) use bucket::await_bucket_namespace_operation;
pub use bucket_fence::BucketIncarnationFenceGuard;
mod heal;
pub(crate) use heal::bucket_heal_scope;
pub use heal::{HealObjectAbsenceProof, HealObjectStorageResult};
pub(crate) use heal::{bucket_heal_scope, bucket_heal_scope_for_object};
mod heal_walk;
pub use heal_walk::HealWalkVersion;
mod init;
+40 -2
View File
@@ -161,14 +161,52 @@ impl ECStore {
"incarnation-bound rename requires namespace locking and a non-nil identity",
));
}
let bucket = if destination.0 == RUSTFS_META_BUCKET {
super::heal::bucket_metadata_owner(destination.1).ok_or(DiskError::FileAccessDenied)?
} else {
destination.0
};
let guard = Arc::new(
self.acquire_bucket_incarnation_fence(destination.0, expected)
self.acquire_bucket_incarnation_fence(bucket, expected)
.await
.map_err(|error| DiskError::other(error.to_string()))?,
);
if guard.is_lock_lost() {
return Err(DiskError::other("bucket heal incarnation fence was lost"));
}
let transaction = if destination.0 == RUSTFS_META_BUCKET {
// The receiver must retain the config fence after an RPC timeout.
// A delayed request must also reject a newer config in the same
// bucket incarnation before publishing its older repair payload.
let transaction = crate::bucket::metadata_sys::acquire_bucket_metadata_transaction_read_lock_in(&self.ctx, bucket)
.await
.map_err(|error| DiskError::other(error.to_string()))?;
let (current, _) = self
.get_pool_info_for_delete_marker(
destination.0,
destination.1,
&ObjectOptions {
no_lock: true,
..Default::default()
},
)
.await
.map_err(|error| DiskError::other(error.to_string()))?;
let current = current.object_info;
if transaction.is_lock_lost()
|| guard.is_lock_lost()
|| fi.mod_time.is_none()
|| current.mod_time != fi.mod_time
|| current.data_dir != fi.data_dir
|| current.size != fi.size
|| current.etag != fi.get_etag()
{
return Err(DiskError::FileAccessDenied);
}
Some(transaction)
} else {
None
};
rename_local_data_with_ctx(
&self.ctx,
disk_ref,
@@ -176,7 +214,7 @@ impl ECStore {
fi,
destination,
RenameDataGuards {
external_guard: Some(guard),
external_guard: Some(Arc::new((guard, transaction))),
..Default::default()
},
)
+41
View File
@@ -907,6 +907,21 @@ impl ErasureSetHealer {
}
if failed_objects == 0 && skipped_objects == 0 && failed_buckets == 0 {
let targets = if self.pool_metadata_target_endpoints.is_empty() {
self.target_endpoints.as_ref()
} else {
self.pool_metadata_target_endpoints.as_ref()
};
if self.replacement_task_id.is_some() || (!self.heal_opts.dry_run && self.heal_opts.recreate && !targets.is_empty()) {
self.verify_replacement_identity_fence("bucket metadata").await?;
// Recheck even resumed buckets: their user-object cursor does not
// prove that the replacement holds the internal bucket records.
for bucket in buckets {
self.storage
.heal_replacement_bucket_metadata(bucket, &self.heal_opts, targets)
.await?;
}
}
self.heal_replacement_pool_metadata(
set_disk_id,
&mut ErasureSetPassCounters {
@@ -2267,6 +2282,8 @@ mod resume_loop_tests {
list_include_lifecycle_object_info: Mutex<Vec<bool>>,
replacement_target_identity_sequences: Mutex<VecDeque<Vec<ReplacementTargetIdentity>>>,
pool_metadata_placement: Mutex<Option<ReplacementCommitEvidence>>,
bucket_metadata_calls: Mutex<Vec<String>>,
bucket_metadata_failure: AtomicBool,
fail_listing: AtomicBool,
fail_listing_buckets: Mutex<HashSet<String>>,
}
@@ -2462,6 +2479,13 @@ mod resume_loop_tests {
None => Ok(true),
}
}
async fn heal_replacement_bucket_metadata(&self, bucket: &str, _opts: &HealOpts, _targets: &[String]) -> Result<()> {
self.bucket_metadata_calls.lock().unwrap().push(bucket.to_owned());
if self.bucket_metadata_failure.load(Ordering::SeqCst) {
return Err(Error::Storage(EcstoreError::PreconditionFailed));
}
Ok(())
}
async fn replacement_targets_have_version(
&self,
_bucket: &str,
@@ -3058,6 +3082,18 @@ mod resume_loop_tests {
)
.with_replacement_targets(vec!["replacement-a".to_string()], Some(replacement_task_id.clone()));
env.storage.bucket_metadata_failure.store(true, Ordering::SeqCst);
assert!(healer.heal_erasure_set(&["b".to_string()], "pool_0_set_0").await.is_err());
let incomplete = ResumeManager::load_replacement_intent(env.healer.disk.clone(), &replacement_task_id)
.await
.expect("failed metadata repair must retain replacement intent")
.get_state()
.await;
assert!(!incomplete.completed);
assert_ne!(incomplete.replacement_phase, crate::heal::resume::ReplacementPhase::Verified);
assert!(env.storage.calls().is_empty(), "metadata failure must stop completion before pool.bin");
env.storage.bucket_metadata_failure.store(false, Ordering::SeqCst);
let error = healer
.heal_erasure_set(&["b".to_string()], "pool_0_set_0")
.await
@@ -3073,6 +3109,11 @@ mod resume_loop_tests {
assert_eq!(state.replacement_phase, crate::heal::resume::ReplacementPhase::Intent);
assert_eq!(state.retry_count, 1);
assert_eq!(env.storage.calls(), vec![(POOL_META_NAME.to_string(), None)]);
assert_eq!(
*env.storage.bucket_metadata_calls.lock().unwrap(),
vec!["b", "b"],
"resuming a completed user scan must retry bucket metadata"
);
}
#[tokio::test]
+12
View File
@@ -591,6 +591,11 @@ pub trait HealStorageAPI: Send + Sync {
Err(Error::other("replacement pool metadata placement is unsupported"))
}
/// Heal and physically verify bucket configuration on its replacement targets.
async fn heal_replacement_bucket_metadata(&self, _bucket: &str, _opts: &HealOpts, _targets: &[String]) -> Result<()> {
Err(Error::Storage(StorageError::PreconditionFailed))
}
/// Read target-specific physical evidence for one replacement version.
///
/// This is only used by automatic replacement healing after the normal
@@ -1595,6 +1600,13 @@ impl HealStorageAPI for ECStoreHealStorage {
.map_err(Error::Storage)
}
async fn heal_replacement_bucket_metadata(&self, bucket: &str, opts: &HealOpts, targets: &[String]) -> Result<()> {
self.ecstore
.heal_replacement_bucket_metadata(bucket, opts, targets)
.await
.map_err(Error::Storage)
}
async fn replacement_targets_have_version(
&self,
bucket: &str,
@@ -12,7 +12,7 @@
// See the License for the specific language governing permissions and
// limitations under the License.
//! Regression for replacement healing requiring pool.bin on a non-owning set.
//! Replacement healing must restore internal records only on their owning sets.
#![recursion_limit = "256"]
use http::HeaderMap;
@@ -101,6 +101,31 @@ async fn replacement_pool_metadata_follows_real_two_set_placement() {
storage.replacement_pool_metadata_required(&opts).is_err(),
"invalid scope must fail closed"
);
assert!(
storage
.heal_replacement_bucket_metadata(bucket, &opts, &[paths[0].to_string_lossy().into_owned()])
.await
.is_err(),
"bucket metadata repair must reject an incomplete or invalid replacement scope"
);
}
for targets in [vec![], vec![paths[4].to_string_lossy().into_owned()]] {
assert!(
storage
.heal_replacement_bucket_metadata(
bucket,
&HealOpts {
pool: Some(0),
set: Some(0),
..Default::default()
},
&targets,
)
.await
.is_err(),
"empty or wrong-set targets must be rejected before repairing metadata"
);
}
let mut owning_set = None;
@@ -120,6 +145,17 @@ async fn replacement_pool_metadata_follows_real_two_set_placement() {
assert!(owning_set.replace(set).is_none(), "only one set owns pool.bin");
std::fs::remove_dir_all(&metadata_path).unwrap();
}
let bucket_records: Vec<_> = [".metadata.bin", ".bucket-incarnation"]
.into_iter()
.map(|file| {
let path = target_path.join(RUSTFS_META_BUCKET).join("buckets").join(bucket).join(file);
let existed = path.join("xl.meta").exists();
if existed {
std::fs::remove_dir_all(&path).expect("remove bucket metadata shard from replacement");
}
(path, existed)
})
.collect();
// Select a user key using the real placement algorithm, independently
// of the metadata-scope decision under test.
let key = (0..1000)
@@ -191,6 +227,14 @@ async fn replacement_pool_metadata_follows_real_two_set_placement() {
"completed repair must clear its healing marker"
);
assert_eq!(metadata_path.join("xl.meta").exists(), owns_metadata);
for (path, owned) in bucket_records {
assert_eq!(
path.join("xl.meta").exists(),
owned,
"replacement must restore each owned bucket record without creating a wrong-set copy: {}",
path.display()
);
}
for (relative, bytes) in original_parts {
assert_eq!(std::fs::read(object_path.join(relative)).unwrap(), bytes, "reconstructed shard differs");
}