fix(ecstore): guard free-version decommission conflicts

This commit is contained in:
overtrue
2026-08-23 07:41:17 +08:00
parent bef7e5bb68
commit 7985251b67
5 changed files with 450 additions and 12 deletions
+61
View File
@@ -227,6 +227,44 @@ pub(super) fn restore_commit_operation_id_from_metadata(metadata: &HashMap<Strin
restore_operation_id_from_metadata(metadata)
}
async fn check_decommission_tier_free_version_target(
disk: &DiskStore,
bucket: &str,
object: &str,
source: &FileInfo,
) -> Result<()> {
let raw = match disk.read_xl(bucket, object, false).await {
Ok(raw) => raw,
Err(DiskError::FileNotFound | DiskError::FileVersionNotFound | DiskError::VolumeNotFound) => return Ok(()),
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)?;
if !existing.tier_free_version() || !crate::store::object::tiered_data_movement_source_matches(source, &existing)? {
all_matching_versions_equivalent = false;
}
}
if matching_count == 0 || (matching_count == 1 && all_matching_versions_equivalent) {
return Ok(());
}
Err(StorageError::DataMovementOverwriteErr(
bucket.to_owned(),
object.to_owned(),
source_version_id.map(|version_id| version_id.to_string()).unwrap_or_default(),
)
.into())
}
impl SetDisks {
pub(super) async fn require_current_restore_operation_id(
&self,
@@ -4698,6 +4736,9 @@ impl SetDisks {
});
}
self.validate_decommission_tier_free_version_target(bucket, object, fi)
.await?;
let disks = self.disks.read().await.clone();
let write_quorum = self.default_write_quorum();
let futures = disks.into_iter().map(|disk| {
@@ -4722,6 +4763,26 @@ impl SetDisks {
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| check_decommission_tier_free_version_target(disk, bucket, object, fi));
for result in join_all(preflight).await {
result?;
}
Ok(())
}
#[tracing::instrument(skip(self, fi, opts))]
pub(crate) async fn decommission_tiered_object(
&self,