From 65953bfdb3d98c570f6b664b166380867450803a Mon Sep 17 00:00:00 2001 From: Zhengchao An Date: Wed, 8 Jul 2026 06:01:50 +0800 Subject: [PATCH] fix(ecstore): reduce rename data version signatures (#4383) --- crates/ecstore/src/disk/local.rs | 24 +++++-- .../src/set_disk/core/io_primitives.rs | 68 ++++++++++++++++++- crates/ecstore/src/set_disk/mod.rs | 35 ++++++++++ 3 files changed, 120 insertions(+), 7 deletions(-) diff --git a/crates/ecstore/src/disk/local.rs b/crates/ecstore/src/disk/local.rs index 3fe4d432f..250e3d228 100644 --- a/crates/ecstore/src/disk/local.rs +++ b/crates/ecstore/src/disk/local.rs @@ -3103,6 +3103,18 @@ fn local_disk_scan_lock_key(bucket: &str, base_dir: &str, filter_prefix: Option< (bucket.to_owned(), prefix) } +fn rename_data_versions_signature(meta: &FileMeta) -> Option> { + if meta.versions.len() > 10 { + return None; + } + + let mut signature = Vec::with_capacity(meta.versions.len() * 16); + for version in meta.versions.iter() { + signature.extend_from_slice(version.header.version_id.unwrap_or_default().as_bytes()); + } + Some(signature) +} + fn is_root_path(path: impl AsRef) -> bool { path.as_ref().components().count() == 1 && path.as_ref().has_root() } @@ -4121,6 +4133,7 @@ impl DiskAPI for LocalDisk { let _ = xlmeta.data.remove_two(version_id, *old_data_dir); } xlmeta.add_version(fi)?; + let version_signature = rename_data_versions_signature(&xlmeta); let new_dst_buf = xlmeta.marshal_msg()?; let src_file_parent = src_file_path.parent().unwrap_or(src_volume_dir.as_path()); @@ -4249,7 +4262,7 @@ impl DiskAPI for LocalDisk { Ok(RenameDataResp { old_data_dir: has_old_data_dir, - sign: None, + sign: version_signature, }) } else { // Inline: merge read + parse + write + rename into single spawn_blocking @@ -4261,7 +4274,7 @@ impl DiskAPI for LocalDisk { None }; - let (old_data_dir, _dst_buf) = tokio::task::spawn_blocking(move || { + let (old_data_dir, version_signature) = tokio::task::spawn_blocking(move || { // Read existing xl.meta let has_dst_buf = match std::fs::read(&dst) { Ok(buf) => Some(Bytes::from(buf)), @@ -4283,6 +4296,7 @@ impl DiskAPI for LocalDisk { let _ = xlmeta.data.remove_two(version_id, *d); } xlmeta.add_version(fi)?; + let version_signature = rename_data_versions_signature(&xlmeta); let new_buf = xlmeta.marshal_msg()?; // Write new xl.meta + rename @@ -4347,7 +4361,7 @@ impl DiskAPI for LocalDisk { os::fsync_dir_std(dst_parent)?; } - Ok::<(Option, Option), std::io::Error>((old_data_dir, has_dst_buf)) + Ok::<(Option, Option>), std::io::Error>((old_data_dir, version_signature)) }) .await .map_err(DiskError::from)??; @@ -4361,7 +4375,7 @@ impl DiskAPI for LocalDisk { Ok(RenameDataResp { old_data_dir, - sign: None, + sign: version_signature, }) } } @@ -5079,6 +5093,7 @@ mod test { .expect("rename_data should commit"); assert_eq!(resp.old_data_dir, Some(old_data_dir)); + assert_eq!(resp.sign, Some(version_id.as_bytes().to_vec())); assert!( dst_object_dir .join(old_data_dir.to_string()) @@ -5286,6 +5301,7 @@ mod test { .expect("inline rename_data should commit"); assert_eq!(resp.old_data_dir, Some(old_data_dir)); + assert_eq!(resp.sign, Some(version_id.as_bytes().to_vec())); let backup_path = dst_object_dir.join(old_data_dir.to_string()).join(STORAGE_FORMAT_FILE_BACKUP); assert!(backup_path.exists()); // The rollback backup must contain the previous metadata bytes verbatim so diff --git a/crates/ecstore/src/set_disk/core/io_primitives.rs b/crates/ecstore/src/set_disk/core/io_primitives.rs index 407e256c4..b87b6473a 100644 --- a/crates/ecstore/src/set_disk/core/io_primitives.rs +++ b/crates/ecstore/src/set_disk/core/io_primitives.rs @@ -2406,6 +2406,13 @@ enum OrphanDirScan { Missing, } +fn rename_data_versions_key(versions: &[u8]) -> Option<[u8; 8]> { + let prefix = versions.get(..8)?; + let mut key = [0; 8]; + key.copy_from_slice(prefix); + Some(key) +} + impl SetDisks { pub(in crate::set_disk) fn default_read_quorum(&self) -> usize { self.set_drive_count - self.default_parity_count @@ -2577,10 +2584,8 @@ impl SetDisks { return Err(ret_err); } - let versions = None; - // TODO: reduceCommonVersions - let data_dir = Self::reduce_common_data_dir(&data_dirs, write_quorum); + let versions = Self::select_rename_data_versions(&disk_versions, &errs, write_quorum); let online_disks = Self::eval_disks(disks, &errs); let cleanup_disks = if let Some(data_dir) = data_dir { disks @@ -2602,6 +2607,63 @@ impl SetDisks { Ok((online_disks, versions, data_dir, cleanup_disks)) } + pub(in crate::set_disk) fn reduce_common_versions(disk_versions: &[Option>], write_quorum: usize) -> Option> { + let mut versions_count = HashMap::new(); + + for versions in disk_versions.iter().flatten() { + if let Some(key) = rename_data_versions_key(versions) { + *versions_count.entry(key).or_insert(0usize) += 1; + } + } + + let (common_versions, max_count) = versions_count + .into_iter() + .max_by_key(|(_, count)| *count) + .unwrap_or(([0; 8], 0)); + + if max_count < write_quorum { + return None; + } + + disk_versions + .iter() + .flatten() + .find(|versions| rename_data_versions_key(versions).is_some_and(|key| key == common_versions)) + .cloned() + } + + pub(in crate::set_disk) fn select_rename_data_versions( + disk_versions: &[Option>], + errs: &[Option], + write_quorum: usize, + ) -> Option> { + let mut versions = Self::reduce_common_versions(disk_versions, write_quorum); + for (dversions, err) in disk_versions.iter().zip(errs.iter()) { + if err.is_some() { + continue; + } + let Some(dversions) = dversions.as_ref().filter(|versions| !versions.is_empty()) else { + continue; + }; + + match versions.as_ref() { + Some(current_versions) if dversions != current_versions => { + if dversions.len() > current_versions.len() { + versions = Some(dversions.clone()); + } + break; + } + Some(_) => {} + None => { + versions = Some(dversions.clone()); + break; + } + } + } + + versions + } + #[allow(dead_code)] #[tracing::instrument(level = "debug", skip(self, disks))] pub(in crate::set_disk) async fn commit_rename_data_dir( diff --git a/crates/ecstore/src/set_disk/mod.rs b/crates/ecstore/src/set_disk/mod.rs index caa55136f..779f8826e 100644 --- a/crates/ecstore/src/set_disk/mod.rs +++ b/crates/ecstore/src/set_disk/mod.rs @@ -5958,6 +5958,41 @@ mod tests { assert_eq!(result, Some(uuid1)); // Ignore None votes; uuid1 should still meet quorum } + fn rename_versions_signature(first_key: u8, version_count: usize) -> Vec { + let mut signature = vec![0; version_count * 16]; + signature[7] = first_key; + signature + } + + #[test] + fn test_reduce_common_versions_requires_write_quorum() { + let common = rename_versions_signature(1, 1); + let other = rename_versions_signature(2, 1); + + let disk_versions = vec![Some(common.clone()), Some(common.clone()), Some(other)]; + let result = SetDisks::reduce_common_versions(&disk_versions, 2); + assert_eq!(result, Some(common)); + + let split_versions = vec![ + Some(rename_versions_signature(1, 1)), + Some(rename_versions_signature(2, 1)), + None, + ]; + let result = SetDisks::reduce_common_versions(&split_versions, 2); + assert_eq!(result, None); + } + + #[test] + fn test_select_rename_data_versions_keeps_longer_success_disparity() { + let common = rename_versions_signature(1, 1); + let longer = rename_versions_signature(2, 2); + let disk_versions = vec![Some(common.clone()), Some(common), Some(longer.clone())]; + let errs = vec![None, None, None]; + + let result = SetDisks::select_rename_data_versions(&disk_versions, &errs, 2); + assert_eq!(result, Some(longer)); + } + #[test] fn test_object_quorum_from_meta_returns_not_found_when_all_metadata_is_missing() { let errs = vec![