fix(ecstore): require write quorum for metadata early stop (#4300)

This commit is contained in:
GatewayJ
2026-07-10 15:57:09 +08:00
committed by GitHub
parent 40cf94f243
commit a269f8df05
7 changed files with 1241 additions and 68 deletions
+1
View File
@@ -812,6 +812,7 @@ impl LocalDiskWrapper {
recursive: false,
immediate: false,
undo_write: false,
undo_delete: false,
old_data_dir: None,
},
)
File diff suppressed because it is too large Load Diff
+4
View File
@@ -883,6 +883,8 @@ pub struct DeleteOptions {
pub recursive: bool,
pub immediate: bool,
pub undo_write: bool,
#[serde(default)]
pub undo_delete: bool,
pub old_data_dir: Option<Uuid>,
}
@@ -1122,12 +1124,14 @@ mod tests {
recursive: true,
immediate: false,
undo_write: true,
undo_delete: false,
old_data_dir: Some(Uuid::new_v4()),
};
assert!(opts.recursive);
assert!(!opts.immediate);
assert!(opts.undo_write);
assert!(!opts.undo_delete);
assert!(opts.old_data_dir.is_some());
}
@@ -357,7 +357,7 @@ impl MetadataQuorumAccumulator {
if !self.allow_early_stop {
return None;
}
if self.delete_marker_votes >= self.missing_response_quorum() {
if self.delete_marker_votes >= self.default_write_quorum() {
return Some(MetadataEarlyStopDecision {
reason: GET_METADATA_EARLY_STOP_REASON_DELETE_MARKER,
});
@@ -373,8 +373,8 @@ impl MetadataQuorumAccumulator {
if self
.candidate
.as_ref()
.and_then(|candidate| self.candidate_read_quorum(candidate))
.is_some_and(|read_quorum| self.candidate_votes >= read_quorum)
.and_then(|candidate| self.candidate_latest_quorum(candidate))
.is_some_and(|latest_quorum| self.candidate_votes >= latest_quorum)
{
return Some(MetadataEarlyStopDecision {
reason: GET_METADATA_EARLY_STOP_REASON_VALID_QUORUM,
@@ -441,14 +441,26 @@ impl MetadataQuorumAccumulator {
GET_METADATA_EARLY_STOP_REASON_INSUFFICIENT_QUORUM
}
pub(in crate::set_disk) fn candidate_read_quorum(&self, candidate: &FileInfo) -> Option<usize> {
pub(in crate::set_disk) fn candidate_latest_quorum(&self, candidate: &FileInfo) -> Option<usize> {
if self.default_parity_count == 0 {
return Some(self.total_disks);
}
if candidate.deleted || candidate.size == 0 || candidate.erasure.parity_blocks >= self.total_disks {
return None;
}
Some(self.total_disks.saturating_sub(candidate.erasure.parity_blocks))
Some(candidate.write_quorum(self.default_write_quorum()))
}
pub(in crate::set_disk) fn default_write_quorum(&self) -> usize {
if self.default_parity_count == 0 {
return self.total_disks;
}
let data_blocks = self.total_disks.saturating_sub(self.default_parity_count);
if data_blocks == self.default_parity_count {
data_blocks.saturating_add(1)
} else {
data_blocks
}
}
pub(in crate::set_disk) fn missing_response_quorum(&self) -> usize {
@@ -1772,7 +1784,7 @@ pub(in crate::set_disk) fn should_allow_metadata_early_stop(
}
(is_get_metadata_early_stop_enabled() && version_id.is_empty() && !healing && !incl_free_versions)
|| (is_version_early_stop_enabled() && !version_id.is_empty() && !healing)
|| (is_version_early_stop_enabled() && !version_id.is_empty() && !healing && !incl_free_versions)
}
/// Final gate for the metadata early-stop fast path.
@@ -1885,7 +1897,7 @@ impl SetDisks {
read_data: bool,
healing: bool,
incl_free_versions: bool,
allow_early_stop: bool,
caller_allows_early_stop: bool,
default_parity_count: usize,
) -> disk::error::Result<(Vec<FileInfo>, Vec<Option<DiskError>>, MetadataFanoutDiagnostics)> {
Self::read_all_fileinfo_inner(
@@ -1898,7 +1910,7 @@ impl SetDisks {
healing,
incl_free_versions,
true,
allow_early_stop,
caller_allows_early_stop,
default_parity_count,
)
.await
@@ -4064,24 +4076,24 @@ mod tests {
}
#[test]
fn metadata_quorum_accumulator_candidate_quorum_handles_zero_parity_and_invalid_candidates() {
fn metadata_quorum_accumulator_candidate_latest_quorum_handles_zero_parity_and_invalid_candidates() {
let accumulator = MetadataQuorumAccumulator::new(4, 0, true);
let candidate = metadata_test_fileinfo("object");
assert_eq!(accumulator.candidate_read_quorum(&candidate), Some(4));
assert_eq!(accumulator.candidate_latest_quorum(&candidate), Some(4));
assert_eq!(accumulator.missing_response_quorum(), 4);
let accumulator = MetadataQuorumAccumulator::new(4, 2, true);
let mut deleted = candidate.clone();
deleted.deleted = true;
assert_eq!(accumulator.candidate_read_quorum(&deleted), None);
assert_eq!(accumulator.candidate_latest_quorum(&deleted), None);
let mut empty = candidate.clone();
empty.size = 0;
assert_eq!(accumulator.candidate_read_quorum(&empty), None);
assert_eq!(accumulator.candidate_latest_quorum(&empty), None);
let mut impossible_parity = candidate;
impossible_parity.erasure.parity_blocks = 4;
assert_eq!(accumulator.candidate_read_quorum(&impossible_parity), None);
assert_eq!(accumulator.candidate_latest_quorum(&impossible_parity), None);
}
#[test]
+63
View File
@@ -3903,6 +3903,69 @@ mod tests {
);
}
#[tokio::test]
async fn test_rename_data_inline_quorum_failure_rolls_back_destination_object() {
let dir = tempfile::tempdir().expect("tempdir should be created");
let disk_root = dir.path().join("disk0");
fs::create_dir_all(&disk_root).await.expect("disk root should be created");
let endpoint = Endpoint::try_from(disk_root.to_str().expect("disk path should be utf8")).expect("endpoint should parse");
let disk = new_disk(
&endpoint,
&DiskOption {
cleanup: false,
health_check: false,
},
)
.await
.expect("disk should be created");
let bucket = "bucket";
let object = "inline-object";
let tmp_object = "tmp-inline-object";
let version_id = Uuid::parse_str("aaaaaaaa-aaaa-aaaa-aaaa-aaaaaaaaaaaa").expect("version id should parse");
match disk.make_volume(bucket).await {
Ok(()) | Err(DiskError::VolumeExists) => {}
Err(err) => panic!("bucket should be available: {err:?}"),
}
match disk.make_volume(RUSTFS_META_TMP_BUCKET).await {
Ok(()) | Err(DiskError::VolumeExists) => {}
Err(err) => panic!("tmp bucket should be available: {err:?}"),
}
let object_dir = disk_root.join(bucket).join(object);
fs::create_dir_all(&object_dir).await.expect("object dir should be created");
let mut old_fi = FileInfo::new(&format!("{bucket}/{object}"), 1, 1);
old_fi.name = object.to_string();
old_fi.version_id = Some(version_id);
old_fi.data = Some(Bytes::from_static(b"old-inline"));
old_fi.size = 10;
old_fi.mod_time = Some(OffsetDateTime::now_utc());
let mut old_meta = FileMeta::default();
old_meta.add_version(old_fi).expect("old metadata should accept file info");
let old_meta_buf = old_meta.marshal_msg().expect("old metadata should encode");
fs::write(object_dir.join(STORAGE_FORMAT_FILE), old_meta_buf.clone())
.await
.expect("old metadata should be written");
let mut new_fi = FileInfo::new(&format!("{bucket}/{object}"), 1, 1);
new_fi.name = object.to_string();
new_fi.version_id = Some(version_id);
new_fi.data = Some(Bytes::from_static(b"new-inline"));
new_fi.size = 10;
new_fi.mod_time = Some(OffsetDateTime::now_utc());
let disks = vec![Some(disk), None];
let file_infos = vec![new_fi.clone(), new_fi];
let result = SetDisks::rename_data(&disks, RUSTFS_META_TMP_BUCKET, tmp_object, &file_infos, bucket, object, 2).await;
assert!(result.is_err());
let restored_meta = fs::read(object_dir.join(STORAGE_FORMAT_FILE))
.await
.expect("destination metadata should remain readable");
assert_eq!(restored_meta, old_meta_buf);
}
#[test]
fn disk_health_entry_returns_cached_value_within_ttl() {
let entry = DiskHealthEntry {
+169 -8
View File
@@ -78,7 +78,7 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks {
};
let metadata_stage_start = Instant::now();
let (fi, files, disks) = match self.get_object_fileinfo(bucket, object, opts, true).await {
let (fi, files, disks) = match self.get_object_fileinfo(bucket, object, opts, true, true).await {
Ok(result) => result,
Err(err) => {
rustfs_io_metrics::record_get_object_metadata_phase_duration(metadata_stage_start.elapsed().as_secs_f64());
@@ -1312,6 +1312,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
async fn delete_object_version(&self, bucket: &str, object: &str, fi: &FileInfo, force_del_marker: bool) -> Result<()> {
let disks = self.disk_inventory().await;
let write_quorum = disks.len() / 2 + 1;
let rollback_dir = Uuid::new_v4();
let mut futures = Vec::with_capacity(disks.len());
let mut errs = Vec::with_capacity(disks.len());
@@ -1320,7 +1321,16 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
futures.push(async move {
if let Some(disk) = disk {
match disk
.delete_version(bucket, object, fi.clone(), force_del_marker, DeleteOptions::default())
.delete_version(
bucket,
object,
fi.clone(),
force_del_marker,
DeleteOptions {
old_data_dir: Some(rollback_dir),
..Default::default()
},
)
.await
{
Ok(r) => Ok(r),
@@ -1344,7 +1354,77 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
}
}
resolve_tiered_decommission_write_quorum_result(&errs, write_quorum, bucket, object)
let quorum_result = resolve_tiered_decommission_write_quorum_result(&errs, write_quorum, bucket, object);
let should_rollback = quorum_result.is_err();
let mut rollback_futures = Vec::new();
for (index, err) in errs.iter().enumerate() {
if err.is_some() {
continue;
}
let Some(disk) = disks[index].as_ref() else {
continue;
};
let disk = disk.clone();
let bucket = bucket.to_string();
let object = object.to_string();
let fi = fi.clone();
rollback_futures.push(async move {
if should_rollback {
if let Err(err) = disk
.delete_version(
&bucket,
&object,
fi,
force_del_marker,
DeleteOptions {
undo_write: true,
undo_delete: true,
old_data_dir: Some(rollback_dir),
..Default::default()
},
)
.await
{
warn!(
bucket = %bucket,
object = %object,
rollback_dir = %rollback_dir,
error = ?err,
"failed to roll back delete after write quorum failure"
);
}
} else {
let rollback_path = format!("{object}/{rollback_dir}");
if let Err(err) = disk
.delete(
&bucket,
&rollback_path,
DeleteOptions {
recursive: true,
immediate: true,
..Default::default()
},
)
.await
&& err != DiskError::FileNotFound
&& err != DiskError::VolumeNotFound
{
warn!(
bucket = %bucket,
object = %object,
rollback_dir = %rollback_dir,
error = ?err,
"failed to clean delete rollback state after quorum success"
);
}
}
});
}
join_all(rollback_futures).await;
quorum_result
}
#[tracing::instrument(skip(self))]
@@ -1542,6 +1622,8 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
vers.push(fi_vers);
}
let rollback_dir = Uuid::new_v4();
let disks = self.disks.read().await;
let disks = disks.clone();
@@ -1554,7 +1636,15 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
let vers = vers.clone();
futures.push(async move {
if let Some(disk) = disk {
disk.delete_versions(bucket, vers, DeleteOptions::default()).await
disk.delete_versions(
bucket,
vers,
DeleteOptions {
old_data_dir: Some(rollback_dir),
..Default::default()
},
)
.await
} else {
let mut errs = Vec::with_capacity(vers.len());
for _ in 0..vers.len() {
@@ -1618,6 +1708,77 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
record_capacity_scope_if_needed(opts.capacity_scope_token, &disks);
let mut rollback_futures = Vec::new();
for fi_vers in &vers {
// delete_versions commits one xl.meta per object group, so rollback must use the same boundary.
let should_rollback = fi_vers.versions.iter().any(|fi| del_errs[fi.idx].is_some());
for (disk_idx, disk) in disks.iter().enumerate() {
if fi_vers.versions.iter().any(|fi| del_obj_errs[disk_idx][fi.idx].is_some()) {
continue;
}
let Some(disk) = disk.as_ref() else {
continue;
};
let disk = disk.clone();
let bucket = bucket.to_string();
let object = fi_vers.name.clone();
let versions = fi_vers.clone();
rollback_futures.push(async move {
if should_rollback {
let errs = disk
.delete_versions(
&bucket,
vec![versions],
DeleteOptions {
undo_write: true,
undo_delete: true,
old_data_dir: Some(rollback_dir),
..Default::default()
},
)
.await;
if let Some(err) = errs.into_iter().flatten().next() {
warn!(
bucket = %bucket,
object = %object,
rollback_dir = %rollback_dir,
error = ?err,
"failed to roll back batch delete after write quorum failure"
);
}
} else {
let rollback_path = format!("{object}/{rollback_dir}");
if let Err(err) = disk
.delete(
&bucket,
&rollback_path,
DeleteOptions {
recursive: true,
immediate: true,
..Default::default()
},
)
.await
&& err != DiskError::FileNotFound
&& err != DiskError::VolumeNotFound
{
warn!(
bucket = %bucket,
object = %object,
rollback_dir = %rollback_dir,
error = ?err,
"failed to clean batch delete rollback state after quorum success"
);
}
}
});
}
}
join_all(rollback_futures).await;
// TODO: add_partial
if dist_erasure {
@@ -1789,7 +1950,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
// This avoids HEAD/GetObject metadata visibility skew immediately after
// PutObject/CompleteMultipartUpload.
let (fi, _, _) = self
.get_object_fileinfo(bucket, object, opts, true)
.get_object_fileinfo(bucket, object, opts, true, false)
.await
.map_err(|e| to_object_err(e, vec![bucket, object]))?;
@@ -1935,7 +2096,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
// _lock_guard = guard_opt;
// }
let (mut fi, meta_arr, online_disks) = self.get_object_fileinfo(bucket, object, opts, true).await?;
let (mut fi, meta_arr, online_disks) = self.get_object_fileinfo(bucket, object, opts, true, false).await?;
/*if err != nil {
return Err(to_object_err(err, vec![bucket, object]));
}*/
@@ -1969,7 +2130,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
if let Err(err) = self.heal_object(bucket, object, "", &HealOpts {no_lock: true, ..Default::default()}) {
return err.expect("err");
}
(fi, meta_arr, online_disks) = self.get_object_fileinfo(&bucket, &object, &opts, true);
(fi, meta_arr, online_disks) = self.get_object_fileinfo(&bucket, &object, &opts, true, false);
if err != nil {
return to_object_err(err, vec![bucket, object]);
}
@@ -2116,7 +2277,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
Err(rerr.unwrap())
};
let mut oi = ObjectInfo::default();
let fi = self_.clone().get_object_fileinfo(bucket, object, opts, true).await;
let fi = self_.clone().get_object_fileinfo(bucket, object, opts, true, false).await;
if let Err(err) = fi {
return set_restore_header_fn(&mut oi, Some(to_object_err(err, vec![bucket, object]))).await;
}
+34 -8
View File
@@ -170,10 +170,10 @@ impl SetDisks {
object: &str,
opts: &ObjectOptions,
read_data: bool,
caller_allows_early_stop: bool,
) -> Result<(FileInfo, Vec<FileInfo>, Vec<Option<DiskStore>>)> {
// Read-only callers (GET/HEAD/tag read) may use the metadata early-stop
// fast path.
self.get_object_fileinfo_gated(bucket, object, opts, read_data, true).await
self.get_object_fileinfo_gated(bucket, object, opts, read_data, caller_allows_early_stop)
.await
}
/// Like `get_object_fileinfo`, but `allow_early_stop=false` forces the full
@@ -329,7 +329,7 @@ impl SetDisks {
object: &str,
opts: &ObjectOptions,
) -> (ObjectInfo, usize, Option<StorageError>) {
let fi = match self.get_object_fileinfo(bucket, object, opts, false).await {
let fi = match self.get_object_fileinfo(bucket, object, opts, false, false).await {
Ok((fi, _, _)) => fi,
Err(e) => return (ObjectInfo::default(), 0, Some(e)),
};
@@ -2783,6 +2783,8 @@ mod tests {
accumulator.observe_file_info(&metadata_early_stop_candidate("object", 1));
assert!(accumulator.early_stop_decision().is_none());
accumulator.observe_file_info(&metadata_early_stop_candidate("object", 2));
assert!(accumulator.early_stop_decision().is_none());
accumulator.observe_file_info(&metadata_early_stop_candidate("object", 3));
assert_eq!(
accumulator.early_stop_decision(),
@@ -2790,7 +2792,7 @@ mod tests {
reason: GET_METADATA_EARLY_STOP_REASON_VALID_QUORUM
})
);
assert_eq!(accumulator.valid_responses, 2);
assert_eq!(accumulator.valid_responses, 3);
}
#[test]
@@ -2799,25 +2801,31 @@ mod tests {
accumulator.observe_file_info(&metadata_early_stop_candidate("object", 1));
accumulator.observe_file_info(&metadata_early_stop_candidate("object", 2));
assert!(accumulator.early_stop_decision().is_none());
accumulator.observe_file_info(&metadata_early_stop_candidate("object", 3));
assert!(accumulator.early_stop_decision().is_some());
assert_eq!(accumulator.candidate_votes, 2);
assert_eq!(accumulator.candidate_votes, 3);
}
#[test]
fn metadata_quorum_accumulator_partial_result_remains_quorum_compatible() {
let first = metadata_early_stop_candidate("object", 1);
let second = metadata_early_stop_candidate("object", 2);
let parts_metadata = vec![first.clone(), second, FileInfo::default(), FileInfo::default()];
let third = metadata_early_stop_candidate("object", 3);
let parts_metadata = vec![first.clone(), second, third, FileInfo::default()];
let errs = vec![None, None, None, None];
let (read_quorum, _) = SetDisks::object_quorum_from_meta(&parts_metadata, &errs, 2)
let (read_quorum, write_quorum) = SetDisks::object_quorum_from_meta(&parts_metadata, &errs, 2)
.expect("partial early-stop metadata should preserve read quorum");
let read_quorum = usize::try_from(read_quorum).expect("read quorum should be non-negative");
let write_quorum = usize::try_from(write_quorum).expect("write quorum should be non-negative");
let selection_quorum = SetDisks::latest_fileinfo_selection_quorum("", &parts_metadata, &errs, read_quorum, write_quorum);
let selected = SetDisks::pick_valid_fileinfo(&parts_metadata, None, Some("etag-1".to_string()), read_quorum)
.expect("partial early-stop metadata should preserve selected FileInfo");
assert_eq!(read_quorum, 2);
assert_eq!(selection_quorum, 3);
assert_eq!(selected.name, first.name);
assert_eq!(selected.get_etag(), first.get_etag());
}
@@ -2870,6 +2878,8 @@ mod tests {
accumulator.observe_file_info(&deleted);
accumulator.observe_file_info(&deleted);
assert!(accumulator.early_stop_decision().is_none());
accumulator.observe_file_info(&deleted);
assert_eq!(
accumulator.early_stop_decision(),
@@ -3008,6 +3018,22 @@ mod tests {
);
}
#[test]
fn metadata_early_stop_rejects_healing_and_free_version_requests() {
temp_env::with_vars(
[
(ENV_RUSTFS_GET_METADATA_EARLY_STOP_ENABLE, Some("true")),
(ENV_RUSTFS_GET_METADATA_VERSION_EARLY_STOP_ENABLE, Some("true")),
],
|| {
assert!(!should_allow_metadata_early_stop(false, "", true, false));
assert!(!should_allow_metadata_early_stop(false, "", false, true));
assert!(!should_allow_metadata_early_stop(false, "version-id", true, false));
assert!(!should_allow_metadata_early_stop(false, "version-id", false, true));
},
);
}
#[test]
fn version_early_stop_gate_defaults_to_disabled() {
temp_env::with_var(ENV_RUSTFS_GET_METADATA_VERSION_EARLY_STOP_ENABLE, None::<&str>, || {