From f96346124e9ea5516007b075d6181ca387a1ff7b Mon Sep 17 00:00:00 2001 From: houseme Date: Sat, 8 Aug 2026 20:22:51 +0800 Subject: [PATCH] fix(multipart): recover part transactions by write quorum (#5844) Co-authored-by: heihutu --- .../src/set_disk/core/io_primitives.rs | 166 ++++++++++++++---- crates/ecstore/src/set_disk/ops/multipart.rs | 57 ++++++ rustfs/src/error.rs | 10 +- 3 files changed, 199 insertions(+), 34 deletions(-) diff --git a/crates/ecstore/src/set_disk/core/io_primitives.rs b/crates/ecstore/src/set_disk/core/io_primitives.rs index 416f7c65d..5c67c072c 100644 --- a/crates/ecstore/src/set_disk/core/io_primitives.rs +++ b/crates/ecstore/src/set_disk/core/io_primitives.rs @@ -3412,6 +3412,18 @@ impl SetDisks { } async fn recover_part_transaction(&self, dst_object: &str, write_quorum: usize) -> disk::error::Result { + struct PartTransactionObservation { + transaction_meta: Option, + current_meta: Option, + rollback: bool, + err: Option, + } + + enum PartTransactionOutcome { + Commit, + Rollback, + } + let disks = self.get_disks_internal().await; let transaction_path = part_transaction_path(dst_object); let transaction_meta_path = format!("{transaction_path}/{PART_TRANSACTION_NEW_META}"); @@ -3425,36 +3437,76 @@ impl SetDisks { let current_meta_path = current_meta_path.clone(); async move { let Some(disk) = disk else { - return Ok((None, None, false)); + return PartTransactionObservation { + transaction_meta: None, + current_meta: None, + rollback: false, + err: Some(DiskError::DiskNotFound), + }; }; let transaction_meta = match disk.read_all(RUSTFS_META_MULTIPART_BUCKET, &transaction_meta_path).await { Ok(meta) => Some(meta), Err(DiskError::FileNotFound) => None, - Err(err) => return Err(err), + Err(err) => { + return PartTransactionObservation { + transaction_meta: None, + current_meta: None, + rollback: false, + err: Some(err), + }; + } }; let rollback = match disk.read_all(RUSTFS_META_MULTIPART_BUCKET, &rollback_path).await { Ok(_) => true, Err(DiskError::FileNotFound) => false, - Err(err) => return Err(err), + Err(err) => { + return PartTransactionObservation { + transaction_meta, + current_meta: None, + rollback: false, + err: Some(err), + }; + } }; let current_meta = match disk.read_all(RUSTFS_META_MULTIPART_BUCKET, ¤t_meta_path).await { Ok(meta) => Some(meta), Err(DiskError::FileNotFound | DiskError::DiskNotFound) => None, Err(_) => None, }; - Ok((transaction_meta, current_meta, rollback)) + PartTransactionObservation { + transaction_meta, + current_meta, + rollback, + err: None, + } } }); - let observations = join_all(reads).await.into_iter().collect::>>()?; - if observations.iter().all(|(transaction, _, _)| transaction.is_none()) { + let observations = join_all(reads).await; + let read_errs = observations + .iter() + .map(|observation| observation.err.clone()) + .collect::>(); + if let Some(err) = reduce_write_quorum_errs(&read_errs, OBJECT_OP_IGNORED_ERRS, write_quorum) { + return Err(err); + } + + if observations + .iter() + .filter(|observation| observation.err.is_none()) + .all(|observation| observation.transaction_meta.is_none()) + { return Ok(false); } let mut current_counts: HashMap = HashMap::new(); - for (_, current, _) in &observations { - if let Some(current) = current { + let mut transaction_meta_values = HashSet::new(); + for observation in observations.iter().filter(|observation| observation.err.is_none()) { + if let Some(current) = &observation.current_meta { *current_counts.entry(current.clone()).or_default() += 1; } + if let Some(transaction_meta) = &observation.transaction_meta { + transaction_meta_values.insert(transaction_meta.clone()); + } } let current_quorum = current_counts .into_iter() @@ -3462,11 +3514,29 @@ impl SetDisks { let old_meta_path = format!("{transaction_path}/{PART_TRANSACTION_OLD_META}"); let old_meta_absent_path = format!("{transaction_path}/old.meta.absent"); + let mut outcomes = Vec::with_capacity(observations.len()); + for observation in &observations { + let outcome = if observation.err.is_none() && observation.transaction_meta.is_none() { + match &observation.current_meta { + Some(current_meta) if transaction_meta_values.contains(current_meta) => Some(PartTransactionOutcome::Commit), + _ => Some(PartTransactionOutcome::Rollback), + } + } else { + None + }; + outcomes.push(outcome); + } let decisions = observations .iter() .enumerate() - .filter_map(|(index, (transaction_meta, _, rollback))| { - transaction_meta.as_ref().map(|meta| (index, meta.clone(), *rollback)) + .filter_map(|(index, observation)| { + if observation.err.is_some() { + return None; + } + observation + .transaction_meta + .as_ref() + .map(|meta| (index, meta.clone(), observation.rollback)) }) .map(|(index, transaction_meta, rollback)| { let disk = disks[index].clone(); @@ -3475,38 +3545,68 @@ impl SetDisks { let current_quorum = current_quorum.clone(); async move { let Some(disk) = disk else { - return Err(DiskError::DiskNotFound); + return (index, Err(DiskError::DiskNotFound)); }; - let action = if rollback { - PartTransactionAction::Rollback - } else if current_quorum.as_ref() == Some(&transaction_meta) { - PartTransactionAction::Commit - } else if let Some(current_quorum) = current_quorum { - match disk.read_all(RUSTFS_META_MULTIPART_BUCKET, &old_meta_path).await { - Ok(old_meta) if old_meta == current_quorum => PartTransactionAction::Rollback, - Ok(_) => PartTransactionAction::Commit, - Err(DiskError::FileNotFound) => { - match disk.read_all(RUSTFS_META_MULTIPART_BUCKET, &old_meta_absent_path).await { - Ok(_) => PartTransactionAction::Commit, - Err(_) => return Err(DiskError::FileCorrupt), + let result = async { + let action = if rollback { + PartTransactionAction::Rollback + } else if current_quorum.as_ref() == Some(&transaction_meta) { + PartTransactionAction::Commit + } else if let Some(current_quorum) = current_quorum { + match disk.read_all(RUSTFS_META_MULTIPART_BUCKET, &old_meta_path).await { + Ok(old_meta) if old_meta == current_quorum => PartTransactionAction::Rollback, + Ok(_) => PartTransactionAction::Commit, + Err(DiskError::FileNotFound) => { + match disk.read_all(RUSTFS_META_MULTIPART_BUCKET, &old_meta_absent_path).await { + Ok(_) => PartTransactionAction::Commit, + Err(_) => return Err(DiskError::FileCorrupt), + } } + Err(err) => return Err(err), } - Err(err) => return Err(err), - } - } else { - PartTransactionAction::Rollback + } else { + PartTransactionAction::Rollback + }; + disk.settle_part_transaction(RUSTFS_META_MULTIPART_BUCKET, dst_object, action) + .await?; + let outcome = match action { + PartTransactionAction::Commit => PartTransactionOutcome::Commit, + PartTransactionAction::Rollback => PartTransactionOutcome::Rollback, + }; + Ok(outcome) }; - disk.settle_part_transaction(RUSTFS_META_MULTIPART_BUCKET, dst_object, action) - .await?; - Ok(action == PartTransactionAction::Commit) + (index, result.await) } }); let results = join_all(decisions).await; - if let Some(err) = results.iter().find_map(|result| result.as_ref().err()) { - return Err(err.clone()); + let mut settle_errs = read_errs; + for result in results { + match result { + (index, Ok(outcome)) => outcomes[index] = Some(outcome), + (index, Err(err)) => settle_errs[index] = Some(err), + } } - Ok(results.iter().any(|result| matches!(result, Ok(true)))) + if let Some(err) = reduce_write_quorum_errs(&settle_errs, OBJECT_OP_IGNORED_ERRS, write_quorum) { + return Err(err); + } + + let commit_count = outcomes + .iter() + .filter(|outcome| matches!(outcome, Some(PartTransactionOutcome::Commit))) + .count(); + if commit_count >= write_quorum { + return Ok(true); + } + let rollback_count = outcomes + .iter() + .filter(|outcome| matches!(outcome, Some(PartTransactionOutcome::Rollback))) + .count(); + if rollback_count >= write_quorum { + return Ok(false); + } + + Err(DiskError::ErasureWriteQuorum) } pub(in crate::set_disk) async fn recover_part_transactions( diff --git a/crates/ecstore/src/set_disk/ops/multipart.rs b/crates/ecstore/src/set_disk/ops/multipart.rs index 6c8dda73d..b5b9f136b 100644 --- a/crates/ecstore/src/set_disk/ops/multipart.rs +++ b/crates/ecstore/src/set_disk/ops/multipart.rs @@ -3177,6 +3177,63 @@ mod tests { ); } + #[tokio::test] + async fn put_object_part_recovers_transaction_with_one_faulty_disk_at_write_quorum() { + use tokio::io::AsyncReadExt as _; + + let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks_for_pool_with_default_parity(4, 0, 2).await; + assert_eq!(set_disks.default_read_quorum(), 2); + assert_eq!(set_disks.default_write_quorum(), 3); + + let bucket = "multipart-degraded-upload-part-bucket"; + let object = "object"; + make_bucket_on_all(&disk_stores, bucket).await; + + let upload = set_disks + .new_multipart_upload(bucket, object, &ObjectOptions::default()) + .await + .expect("multipart upload should be created before the disk fault"); + disk_stores[0] + .set_disk_id_state(Some(Uuid::new_v4())) + .await + .expect("test should mark one disk stale"); + + let payload = vec![0x5b; 4096]; + let mut reader = PutObjReader::from_vec(payload.clone()); + let part = set_disks + .put_object_part(bucket, object, &upload.upload_id, 1, &mut reader, &ObjectOptions::default()) + .await + .expect("upload part should commit with exactly write quorum healthy disks"); + + set_disks + .clone() + .complete_multipart_upload( + bucket, + object, + &upload.upload_id, + vec![CompletePart { + part_num: part.part_num, + etag: part.etag, + ..Default::default() + }], + &ObjectOptions::default(), + ) + .await + .expect("completion should settle the write-quorum part"); + + let mut object_reader = set_disks + .get_object_reader(bucket, object, None, HeaderMap::new(), &ObjectOptions::default()) + .await + .expect("completed object should be readable through read quorum"); + let mut restored = Vec::new(); + object_reader + .stream + .read_to_end(&mut restored) + .await + .expect("completed object should stream fully"); + assert_eq!(restored, payload); + } + #[tokio::test(flavor = "multi_thread")] #[serial] async fn put_object_part_rechecks_upload_after_commit_lock() { diff --git a/rustfs/src/error.rs b/rustfs/src/error.rs index 93f69cd50..7d27a5241 100644 --- a/rustfs/src/error.rs +++ b/rustfs/src/error.rs @@ -288,7 +288,11 @@ impl From for ApiError { StorageError::ObjectNameInvalid(_, _) => S3ErrorCode::InvalidArgument, StorageError::BucketExists(_) => S3ErrorCode::BucketAlreadyOwnedByYou, StorageError::StorageFull => S3ErrorCode::ServiceUnavailable, - StorageError::SlowDown => S3ErrorCode::SlowDown, + StorageError::SlowDown + | StorageError::FaultyDisk + | StorageError::FaultyRemoteDisk + | StorageError::DiskNotFound + | StorageError::TooManyOpenFiles => S3ErrorCode::SlowDown, StorageError::ErasureReadQuorum | StorageError::InsufficientReadQuorum(_, _) | StorageError::ErasureWriteQuorum @@ -598,6 +602,10 @@ mod tests { (StorageError::BucketExists("test".into()), S3ErrorCode::BucketAlreadyOwnedByYou), (StorageError::StorageFull, S3ErrorCode::ServiceUnavailable), (StorageError::SlowDown, S3ErrorCode::SlowDown), + (StorageError::FaultyDisk, S3ErrorCode::SlowDown), + (StorageError::FaultyRemoteDisk, S3ErrorCode::SlowDown), + (StorageError::DiskNotFound, S3ErrorCode::SlowDown), + (StorageError::TooManyOpenFiles, S3ErrorCode::SlowDown), (StorageError::ErasureReadQuorum, S3ErrorCode::SlowDown), (StorageError::InsufficientReadQuorum("test".into(), "test".into()), S3ErrorCode::SlowDown), (StorageError::ErasureWriteQuorum, S3ErrorCode::SlowDown),