diff --git a/crates/ecstore/src/erasure_coding/encode.rs b/crates/ecstore/src/erasure_coding/encode.rs index f39ac2960..e2ffa81cc 100644 --- a/crates/ecstore/src/erasure_coding/encode.rs +++ b/crates/ecstore/src/erasure_coding/encode.rs @@ -191,7 +191,18 @@ impl<'a> MultiWriter<'a> { let summary = build_write_quorum_failure_summary(&self.errs, OBJECT_OP_IGNORED_ERRS, self.write_quorum); let summary_text = format_write_quorum_failure(&summary); runtime_sources::record_erasure_write_quorum_failure("write", quorum_dominant_error_metric_label(&summary)); - error!("reduce_write_quorum_errs: {:?}, {}, errs={:?}", write_err, summary_text, self.errs); + error!( + required = summary.required, + achieved = summary.achieved, + failed = summary.failed, + total = summary.total, + offline_disks = summary.offline_disks, + retryable_failures = summary.retryable_failures, + dominant_error = summary.dominant_error_label, + returned_error = %write_err, + errs = ?self.errs, + "Erasure encode write quorum unavailable: {summary_text}" + ); return Err(std::io::Error::other(format!("Failed to write data: {summary_text}"))); } @@ -246,8 +257,16 @@ impl<'a> MultiWriter<'a> { let summary_text = format_write_quorum_failure(&summary); runtime_sources::record_erasure_write_quorum_failure("shutdown", quorum_dominant_error_metric_label(&summary)); error!( - "reduce_write_quorum_errs during shutdown: {:?}, {}, errs={:?}", - write_err, summary_text, self.errs + required = summary.required, + achieved = summary.achieved, + failed = summary.failed, + total = summary.total, + offline_disks = summary.offline_disks, + retryable_failures = summary.retryable_failures, + dominant_error = summary.dominant_error_label, + returned_error = %write_err, + errs = ?self.errs, + "Erasure encode shutdown quorum unavailable: {summary_text}" ); return Err(std::io::Error::other(format!("Failed to shutdown writers: {summary_text}"))); } diff --git a/crates/ecstore/src/set_disk/mod.rs b/crates/ecstore/src/set_disk/mod.rs index 1cdbcd77f..1be841f40 100644 --- a/crates/ecstore/src/set_disk/mod.rs +++ b/crates/ecstore/src/set_disk/mod.rs @@ -25,7 +25,8 @@ use crate::bucket::versioning::VersioningApi; use crate::bucket::versioning_sys::BucketVersioningSys; use crate::client::{object_api_utils::get_raw_etag, transition_api::ReaderImpl}; use crate::disk::error_reduce::{ - BUCKET_OP_IGNORED_ERRS, OBJECT_OP_IGNORED_ERRS, count_errs, reduce_read_quorum_errs, reduce_write_quorum_errs, + BUCKET_OP_IGNORED_ERRS, OBJECT_OP_IGNORED_ERRS, build_write_quorum_failure_summary, count_errs, reduce_read_quorum_errs, + reduce_write_quorum_errs, }; use crate::disk::{ self, CHECK_PART_DISK_NOT_FOUND, CHECK_PART_FILE_CORRUPT, CHECK_PART_FILE_NOT_FOUND, CHECK_PART_SUCCESS, CHECK_PART_UNKNOWN, @@ -152,6 +153,12 @@ type WalkOptions = StorageWalkOptions bool>; const LOG_COMPONENT_ECSTORE: &str = "ecstore"; const LOG_SUBSYSTEM_SET_DISK: &str = "set_disk"; const EVENT_SET_DISK_MULTIPART: &str = "set_disk_multipart"; +const COMPLETE_MULTIPART_PART_MISSING: &str = "part_missing"; +const COMPLETE_MULTIPART_PART_READ_QUORUM_UNAVAILABLE: &str = "read_quorum_unavailable"; +const COMPLETE_MULTIPART_PART_ERROR: &str = "part_error"; +const MULTIPART_WRITE_QUORUM_UPLOAD_METADATA: &str = "upload_metadata"; +const MULTIPART_WRITE_QUORUM_WRITER_SETUP: &str = "writer_setup"; +const MULTIPART_WRITE_QUORUM_RENAME_PART: &str = "rename_part"; const EVENT_SET_DISK_WRITE: &str = "set_disk_write"; const EVENT_SET_DISK_HEAL: &str = "set_disk_heal"; const EVENT_SET_DISK_COMMIT_TAIL_SLOW: &str = "set_disk_commit_tail_slow"; @@ -449,6 +456,28 @@ fn get_codec_streaming_reader_decision( GetCodecStreamingDecision::Use } +fn is_confirmed_complete_part_missing(err: &str) -> bool { + err.contains("file not found") + || err.contains("Specified part could not be found") + || (err.starts_with("part.") && err.ends_with(" not found")) +} + +fn complete_multipart_part_error(part_number: usize, err: &str, bucket: &str, object: &str) -> Error { + if is_confirmed_complete_part_missing(err) { + return Error::InvalidPart(part_number, bucket.to_owned(), object.to_owned()); + } + + to_object_err(Error::ErasureReadQuorum, vec![bucket, object]) +} + +fn complete_multipart_part_error_result(err: &Error) -> &'static str { + match err { + Error::InvalidPart(_, _, _) => COMPLETE_MULTIPART_PART_MISSING, + Error::ErasureReadQuorum | Error::InsufficientReadQuorum(_, _) => COMPLETE_MULTIPART_PART_READ_QUORUM_UNAVAILABLE, + _ => COMPLETE_MULTIPART_PART_ERROR, + } +} + /// Record a lock acquisition for deadlock detection. /// This records detailed lock information for deadlock analysis. /// Returns the lock_id for later release tracking. @@ -492,6 +521,47 @@ fn record_lock_release(bucket: &str, object: &str, lock_id: &str, lock_type: &st ); } +#[derive(Clone, Copy, Debug)] +pub(super) struct MultipartWriteQuorumContext<'a> { + stage: &'static str, + bucket: &'a str, + object: &'a str, + upload_id: &'a str, + part_number: Option, +} + +fn log_multipart_write_quorum_failure( + context: MultipartWriteQuorumContext<'_>, + errs: &[Option], + write_quorum: usize, + returned_error: &DiskError, +) { + let summary = build_write_quorum_failure_summary(errs, OBJECT_OP_IGNORED_ERRS, write_quorum); + runtime_sources::record_erasure_write_quorum_failure(context.stage, summary.dominant_error_label); + warn!( + target: "rustfs_ecstore::set_disk", + event = EVENT_SET_DISK_MULTIPART, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_SET_DISK, + op = "upload_part", + state = "write_quorum_unavailable", + stage = context.stage, + bucket = %context.bucket, + object = %context.object, + upload_id = %context.upload_id, + part_number = context.part_number, + required = summary.required, + achieved = summary.achieved, + failed = summary.failed, + total = summary.total, + offline_disks = summary.offline_disks, + retryable_failures = summary.retryable_failures, + dominant_error = summary.dominant_error_label, + returned_error = %returned_error, + "Set disk multipart write quorum unavailable" + ); +} + fn issue3031_diag_enabled() -> bool { rustfs_utils::get_env_bool(ENV_ISSUE3031_DIAG_ENABLE, false) } @@ -3702,6 +3772,18 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks { let nil_count = errors.iter().filter(|&e| e.is_none()).count(); if nil_count < write_quorum { if let Some(write_err) = reduce_write_quorum_errs(&errors, OBJECT_OP_IGNORED_ERRS, write_quorum) { + log_multipart_write_quorum_failure( + MultipartWriteQuorumContext { + stage: MULTIPART_WRITE_QUORUM_WRITER_SETUP, + bucket, + object, + upload_id, + part_number: Some(part_id), + }, + &errors, + write_quorum, + &write_err, + ); return Err(to_object_err(write_err.into(), vec![bucket, object])); } @@ -3803,6 +3885,13 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks { &part_path, part_info_buff.into(), write_quorum, + Some(MultipartWriteQuorumContext { + stage: MULTIPART_WRITE_QUORUM_RENAME_PART, + bucket, + object, + upload_id, + part_number: Some(part_id), + }), ) .await?; @@ -4293,8 +4382,9 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks { let part_numbers = uploaded_parts.iter().map(|v| v.part_num).collect::>(); - let object_parts = - Self::read_parts(&disks, RUSTFS_META_MULTIPART_BUCKET, &part_meta_paths, &part_numbers, read_quorum).await?; + let object_parts = Self::read_parts(&disks, RUSTFS_META_MULTIPART_BUCKET, &part_meta_paths, &part_numbers, read_quorum) + .await + .map_err(|err| to_object_err(err.into(), vec![bucket, object]))?; if object_parts.len() != uploaded_parts.len() { return Err(Error::other("part result number err")); @@ -4317,11 +4407,16 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks { for (i, part) in object_parts.iter().enumerate() { if let Some(err) = &part.error { - error!("complete_multipart_upload part error: {:?}", &err); - if issue3031_diag_enabled() { - warn!( + let mapped_err = complete_multipart_part_error(uploaded_parts[i].part_num, err, bucket, object); + let result = complete_multipart_part_error_result(&mapped_err); + if matches!(mapped_err, Error::InvalidPart(_, _, _)) { + debug!( target: "rustfs_ecstore::set_disk", + event = EVENT_SET_DISK_MULTIPART, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_SET_DISK, op = "complete_multipart_upload", + result = result, bucket = %bucket, object = %object, upload_id = %upload_id, @@ -4330,9 +4425,28 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks { read_quorum = read_quorum, write_quorum = write_quorum, error = %err, - "issue3031_complete_part_error" + "Set disk multipart part missing" + ); + } else { + warn!( + target: "rustfs_ecstore::set_disk", + event = EVENT_SET_DISK_MULTIPART, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_SET_DISK, + op = "complete_multipart_upload", + result = result, + bucket = %bucket, + object = %object, + upload_id = %upload_id, + uploaded_part_num = uploaded_parts[i].part_num, + observed_part_num = part.number, + read_quorum = read_quorum, + write_quorum = write_quorum, + error = %err, + "Set disk multipart part resolution failed" ); } + return Err(mapped_err); } if uploaded_parts[i].part_num != part.number { @@ -5655,6 +5769,41 @@ mod tests { use tempfile::TempDir; use time::OffsetDateTime; + #[test] + fn complete_part_error_maps_confirmed_missing_to_invalid_part() { + for err in ["file not found", "Specified part could not be found", "part.7 not found"] { + let mapped = complete_multipart_part_error(7, err, "bucket", "object"); + + assert!(matches!( + mapped, + Error::InvalidPart(7, ref bucket, ref object) if bucket == "bucket" && object == "object" + )); + assert_eq!(complete_multipart_part_error_result(&mapped), COMPLETE_MULTIPART_PART_MISSING); + } + } + + #[test] + fn complete_part_error_maps_read_quorum_to_retryable_server_error() { + let mapped = complete_multipart_part_error(1, "erasure read quorum", "bucket", "object"); + + assert!(matches!( + mapped, + Error::InsufficientReadQuorum(ref bucket, ref object) if bucket == "bucket" && object == "object" + )); + assert_eq!( + complete_multipart_part_error_result(&mapped), + COMPLETE_MULTIPART_PART_READ_QUORUM_UNAVAILABLE + ); + } + + #[test] + fn complete_part_error_maps_unknown_part_error_to_retryable_server_error() { + let mapped = complete_multipart_part_error(1, "metadata decode failed", "bucket", "object"); + + assert!(matches!(mapped, Error::InsufficientReadQuorum(_, _))); + assert_ne!(complete_multipart_part_error_result(&mapped), COMPLETE_MULTIPART_PART_MISSING); + } + #[derive(Debug, Default)] struct FailingClient; diff --git a/crates/ecstore/src/set_disk/multipart.rs b/crates/ecstore/src/set_disk/multipart.rs index 47402c576..14bdbb559 100644 --- a/crates/ecstore/src/set_disk/multipart.rs +++ b/crates/ecstore/src/set_disk/multipart.rs @@ -17,6 +17,13 @@ use std::future::Future; use std::time::Duration; use tokio::task::JoinSet; +fn map_upload_id_metadata_error(bucket: &str, object: &str, upload_id: &str, err: DiskError) -> Error { + if err == DiskError::FileNotFound { + return StorageError::InvalidUploadID(bucket.to_owned(), object.to_owned(), upload_id.to_owned()); + } + err.into() +} + fn empty_upload_fallback_possible(successful_responses: usize, errs: &[Option]) -> bool { successful_responses == 0 && errs.iter().any(|err| matches!(err, Some(DiskError::FileNotFound))) @@ -165,15 +172,8 @@ impl SetDisks { Self::read_all_fileinfo(&disks, bucket, RUSTFS_META_MULTIPART_BUCKET, &upload_id_path, "", false, false, false) .await?; - let map_err_notfound = |err: DiskError| { - if err == DiskError::FileNotFound { - return StorageError::InvalidUploadID(bucket.to_owned(), object.to_owned(), upload_id.to_owned()); - } - err.into() - }; - - let (read_quorum, write_quorum) = - Self::object_quorum_from_meta(&parts_metadata, &errs, self.default_parity_count).map_err(map_err_notfound)?; + let (read_quorum, write_quorum) = Self::object_quorum_from_meta(&parts_metadata, &errs, self.default_parity_count) + .map_err(|err| map_upload_id_metadata_error(bucket, object, upload_id, err))?; if read_quorum < 0 { error!("check_upload_id_exists: read_quorum < 0, errs={:?}", errs); @@ -189,10 +189,22 @@ impl SetDisks { quorum = write_quorum as usize; if let Some(err) = reduce_write_quorum_errs(&errs, OBJECT_OP_IGNORED_ERRS, quorum) { - return Err(map_err_notfound(err)); + log_multipart_write_quorum_failure( + MultipartWriteQuorumContext { + stage: MULTIPART_WRITE_QUORUM_UPLOAD_METADATA, + bucket, + object, + upload_id, + part_number: None, + }, + &errs, + quorum, + &err, + ); + return Err(map_upload_id_metadata_error(bucket, object, upload_id, err)); } } else if let Some(err) = reduce_read_quorum_errs(&errs, OBJECT_OP_IGNORED_ERRS, quorum) { - return Err(map_err_notfound(err)); + return Err(map_upload_id_metadata_error(bucket, object, upload_id, err)); } let (_, mod_time, etag) = Self::list_online_disks(&disks, &parts_metadata, &errs, quorum); @@ -320,4 +332,58 @@ mod tests { let parts = reduce_quorum_part_numbers(object_parts, 2); assert_eq!(parts, vec![1, 2]); } + + fn test_multipart_fileinfo(object: &str, data_blocks: usize, parity_blocks: usize, index: usize) -> FileInfo { + let mut file_info = FileInfo::new(object, data_blocks, parity_blocks); + file_info.erasure.index = index; + file_info.data_dir = Some(Uuid::new_v4()); + file_info + } + + #[test] + fn upload_id_write_quorum_fails_when_only_read_quorum_metadata_is_visible() { + let parts_metadata = vec![ + test_multipart_fileinfo("bucket/object", 2, 2, 1), + test_multipart_fileinfo("bucket/object", 2, 2, 2), + FileInfo::default(), + FileInfo::default(), + ]; + let errs = vec![None, None, Some(DiskError::DiskNotFound), Some(DiskError::DiskNotFound)]; + + let (read_quorum, write_quorum) = + SetDisks::object_quorum_from_meta(&parts_metadata, &errs, 2).expect("read quorum should resolve metadata geometry"); + + assert_eq!(read_quorum, 2); + assert_eq!(write_quorum, 3); + assert!(reduce_read_quorum_errs(&errs, OBJECT_OP_IGNORED_ERRS, read_quorum as usize).is_none()); + + let err = reduce_write_quorum_errs(&errs, OBJECT_OP_IGNORED_ERRS, write_quorum as usize) + .expect("write quorum should fail with only two writable metadata copies"); + + assert_eq!(err, DiskError::ErasureWriteQuorum); + } + + #[test] + fn upload_id_all_not_found_maps_to_invalid_upload_id() { + let parts_metadata = vec![ + FileInfo::default(), + FileInfo::default(), + FileInfo::default(), + FileInfo::default(), + ]; + let errs = vec![ + Some(DiskError::FileNotFound), + Some(DiskError::FileNotFound), + Some(DiskError::FileNotFound), + Some(DiskError::FileNotFound), + ]; + + let err = SetDisks::object_quorum_from_meta(&parts_metadata, &errs, 2) + .map(|_| ()) + .map_err(|err| map_upload_id_metadata_error("bucket", "object", "upload-id", err)) + .expect_err("all missing upload metadata should remain an invalid upload id"); + + assert!(matches!(err, StorageError::InvalidUploadID(bucket, object, upload_id) + if bucket == "bucket" && object == "object" && upload_id == "upload-id")); + } } diff --git a/crates/ecstore/src/set_disk/read.rs b/crates/ecstore/src/set_disk/read.rs index 1694f3b09..0e3199f0e 100644 --- a/crates/ecstore/src/set_disk/read.rs +++ b/crates/ecstore/src/set_disk/read.rs @@ -362,6 +362,112 @@ fn is_metadata_fanout_ignored_error(err: &DiskError) -> bool { OBJECT_OP_IGNORED_ERRS.iter().any(|ignored| ignored == err) } +fn is_confirmed_missing_part_error(err: Option<&str>) -> bool { + let Some(err) = err else { + return false; + }; + + err.contains("file not found") + || err.contains("No such file or directory") + || err.contains("Specified part could not be found") + || (err.starts_with("part.") && err.ends_with(" not found")) +} + +fn resolve_read_part_from_responses( + bucket: &str, + part_meta_path: &str, + part_number: usize, + part_idx: usize, + expected_part_count: usize, + responses: &[Option>], + read_quorum: usize, +) -> disk::error::Result { + let mut etag_quorum = HashMap::new(); + let mut part_infos = Vec::new(); + let mut present_count = 0usize; + let mut missing_count = 0usize; + let mut transient_error_count = 0usize; + let mut mismatched_response_count = 0usize; + for response in responses.iter() { + let Some(parts) = response else { + transient_error_count += 1; + continue; + }; + + if parts.len() != expected_part_count { + mismatched_response_count += 1; + continue; + } + + if !parts[part_idx].etag.is_empty() { + present_count += 1; + *etag_quorum.entry(parts[part_idx].etag.clone()).or_insert(0) += 1; + part_infos.push(parts[part_idx].clone()); + continue; + } + + if is_confirmed_missing_part_error(parts[part_idx].error.as_deref()) { + missing_count += 1; + } else { + transient_error_count += 1; + } + } + + let mut max_etag_quorum = 0; + let mut max_etag = None; + for (etag, quorum) in etag_quorum.iter() { + if quorum > &max_etag_quorum { + max_etag_quorum = *quorum; + max_etag = Some(etag); + } + } + let max_quorum = max_etag_quorum.max(missing_count); + + let mut found = None; + for info in part_infos.iter() { + if let Some(etag) = max_etag + && info.etag == *etag + { + found = Some(info.clone()); + break; + } + } + + if let (Some(found), Some(max_etag)) = (found, max_etag) + && !found.etag.is_empty() + && etag_quorum.get(max_etag).unwrap_or(&0) >= &read_quorum + { + return Ok(found); + } + + if missing_count >= read_quorum { + return Ok(ObjectPartInfo { + number: part_number, + error: Some(format!("part.{part_number} not found")), + ..Default::default() + }); + } + + if issue3031_diag_enabled() { + warn!( + target: "rustfs_ecstore::set_disk", + bucket = %bucket, + part_meta_path = %part_meta_path, + part_id = part_number, + read_quorum = read_quorum, + max_quorum = max_quorum, + disk_response_count = responses.len(), + present_count = present_count, + missing_count = missing_count, + transient_error_count = transient_error_count, + mismatched_response_count = mismatched_response_count, + "issue3031_read_parts_part_quorum" + ); + } + + Err(DiskError::ErasureReadQuorum) +} + fn shard_read_cost_for_disk(disk: Option<&DiskStore>) -> ShardReadCost { match disk { Some(disk) if disk.is_local() => ShardReadCost::Local, @@ -723,8 +829,6 @@ impl SetDisks { part_numbers: &[usize], read_quorum: usize, ) -> disk::error::Result> { - let mut errs = Vec::with_capacity(disks.len()); - let mut object_parts = Vec::with_capacity(disks.len()); let bucket = bucket.to_string(); let part_meta_paths = part_meta_paths.to_vec(); @@ -750,96 +854,22 @@ impl SetDisks { Err(()) => return Err(DiskError::ErasureReadQuorum), }; - errs.extend(collected_errors); - object_parts.extend(responses.into_iter().map(|resp| resp.unwrap_or_default())); - - if let Some(err) = reduce_read_quorum_errs(&errs, OBJECT_OP_IGNORED_ERRS, read_quorum) { + if let Some(err) = reduce_read_quorum_errs(&collected_errors, OBJECT_OP_IGNORED_ERRS, read_quorum) { return Err(err); } let mut ret = vec![ObjectPartInfo::default(); part_meta_paths.len()]; for (part_idx, part_info) in part_meta_paths.iter().enumerate() { - let mut part_meta_quorum = HashMap::new(); - let mut part_infos = Vec::new(); - let mut present_count = 0usize; - let mut missing_or_empty_count = 0usize; - let mut mismatched_response_count = 0usize; - for parts in object_parts.iter() { - if parts.len() != part_meta_paths.len() { - mismatched_response_count += 1; - *part_meta_quorum.entry(part_info.clone()).or_insert(0) += 1; - continue; - } - - if !parts[part_idx].etag.is_empty() { - present_count += 1; - *part_meta_quorum.entry(parts[part_idx].etag.clone()).or_insert(0) += 1; - part_infos.push(parts[part_idx].clone()); - continue; - } - - missing_or_empty_count += 1; - *part_meta_quorum.entry(part_info.clone()).or_insert(0) += 1; - } - - let mut max_quorum = 0; - let mut max_etag = None; - let mut max_part_meta = None; - for (etag, quorum) in part_meta_quorum.iter() { - if quorum > &max_quorum { - max_quorum = *quorum; - max_etag = Some(etag); - max_part_meta = Some(etag); - } - } - - let mut found = None; - for info in part_infos.iter() { - if let Some(etag) = max_etag - && info.etag == *etag - { - found = Some(info.clone()); - break; - } - - if let Some(part_meta) = max_part_meta - && info.etag.is_empty() - && part_meta.ends_with(format!("part.{0}.meta", info.number).as_str()) - { - found = Some(info.clone()); - break; - } - } - - if let (Some(found), Some(max_etag)) = (found, max_etag) - && !found.etag.is_empty() - && part_meta_quorum.get(max_etag).unwrap_or(&0) >= &read_quorum - { - ret[part_idx] = found.clone(); - } else { - if issue3031_diag_enabled() { - warn!( - target: "rustfs_ecstore::set_disk", - bucket = %bucket, - part_meta_path = %part_info, - part_id = part_numbers[part_idx], - read_quorum = read_quorum, - max_quorum = max_quorum, - disk_response_count = object_parts.len(), - present_count = present_count, - missing_or_empty_count = missing_or_empty_count, - mismatched_response_count = mismatched_response_count, - max_vote_is_missing_marker = max_etag.map(|etag| etag == part_info).unwrap_or(false), - "issue3031_read_parts_part_quorum" - ); - } - ret[part_idx] = ObjectPartInfo { - number: part_numbers[part_idx], - error: Some(format!("part.{} not found", part_numbers[part_idx])), - ..Default::default() - }; - } + ret[part_idx] = resolve_read_part_from_responses( + &bucket, + part_info, + part_numbers[part_idx], + part_idx, + part_meta_paths.len(), + &responses, + read_quorum, + )?; } Ok(ret) @@ -2471,6 +2501,118 @@ mod tests { fi } + fn read_part_test_part(number: usize, etag: &str) -> ObjectPartInfo { + ObjectPartInfo { + number, + etag: etag.to_string(), + ..Default::default() + } + } + + fn read_part_test_missing(number: usize) -> ObjectPartInfo { + ObjectPartInfo { + number, + error: Some("file not found".to_string()), + ..Default::default() + } + } + + #[test] + fn resolve_read_part_returns_part_when_etag_reaches_quorum() { + let part = read_part_test_part(1, "etag-1"); + let responses = vec![ + Some(vec![part.clone()]), + Some(vec![part]), + Some(vec![read_part_test_missing(1)]), + None, + ]; + + let resolved = resolve_read_part_from_responses("bucket", "upload/part.1.meta", 1, 0, 1, &responses, 2) + .expect("etag quorum should resolve part"); + + assert_eq!(resolved.etag, "etag-1"); + assert!(resolved.error.is_none()); + } + + #[test] + fn resolve_read_part_returns_missing_only_when_missing_reaches_quorum() { + let responses = vec![ + Some(vec![read_part_test_missing(1)]), + Some(vec![read_part_test_missing(1)]), + None, + ]; + + let resolved = resolve_read_part_from_responses("bucket", "upload/part.1.meta", 1, 0, 1, &responses, 2) + .expect("missing quorum should resolve as a confirmed missing part"); + + assert_eq!(resolved.number, 1); + assert_eq!(resolved.error.as_deref(), Some("part.1 not found")); + } + + #[test] + fn resolve_read_part_treats_os_not_found_as_confirmed_missing() { + let responses = vec![ + Some(vec![ObjectPartInfo { + number: 9999, + error: Some("No such file or directory (os error 2)".to_string()), + ..Default::default() + }]), + Some(vec![ObjectPartInfo { + number: 9999, + error: Some("No such file or directory (os error 2)".to_string()), + ..Default::default() + }]), + Some(vec![read_part_test_part(1, "stale-etag")]), + None, + ]; + + let resolved = resolve_read_part_from_responses("bucket", "upload/part.9999.meta", 9999, 0, 1, &responses, 2) + .expect("OS not-found quorum should resolve as a missing part"); + + assert_eq!(resolved.number, 9999); + assert_eq!(resolved.error.as_deref(), Some("part.9999 not found")); + } + + #[test] + fn resolve_read_part_returns_missing_when_missing_quorum_beats_stale_present_part() { + let responses = vec![ + Some(vec![read_part_test_missing(1)]), + Some(vec![read_part_test_missing(1)]), + Some(vec![read_part_test_part(1, "stale-etag")]), + None, + ]; + + let resolved = resolve_read_part_from_responses("bucket", "upload/part.1.meta", 1, 0, 1, &responses, 2) + .expect("confirmed missing quorum should resolve as a missing part despite stale metadata"); + + assert_eq!(resolved.number, 1); + assert_eq!(resolved.error.as_deref(), Some("part.1 not found")); + } + + #[test] + fn resolve_read_part_preserves_read_quorum_when_present_part_lacks_quorum() { + let responses = vec![ + Some(vec![read_part_test_part(1, "etag-1")]), + Some(vec![read_part_test_missing(1)]), + None, + ]; + + let err = resolve_read_part_from_responses("bucket", "upload/part.1.meta", 1, 0, 1, &responses, 2) + .expect_err("mixed present and missing observations should not become InvalidPart"); + + assert_eq!(err, DiskError::ErasureReadQuorum); + } + + #[test] + fn resolve_read_part_does_not_count_mismatched_response_as_missing() { + let responses = vec![Some(Vec::new()), Some(vec![read_part_test_missing(1)]), None]; + + let err = resolve_read_part_from_responses("bucket", "upload/part.1.meta", 1, 0, 1, &responses, 2) + .expect_err("invalid disk responses must not vote for a missing part"); + + assert_eq!(err, DiskError::ErasureReadQuorum); + } + #[test] fn metadata_fanout_diagnostics_classifies_response_outcomes() { let valid = metadata_fanout_test_fileinfo("object"); diff --git a/crates/ecstore/src/set_disk/write.rs b/crates/ecstore/src/set_disk/write.rs index 2ff88fd4f..4b60a20a4 100644 --- a/crates/ecstore/src/set_disk/write.rs +++ b/crates/ecstore/src/set_disk/write.rs @@ -316,6 +316,7 @@ impl SetDisks { dst_object: &str, meta: Bytes, write_quorum: usize, + quorum_context: Option>, ) -> disk::error::Result>> { let src_bucket = Arc::new(src_bucket.to_string()); let src_object = Arc::new(src_object.to_string()); @@ -380,7 +381,11 @@ impl SetDisks { } if let Some(err) = reduce_write_quorum_errs(&errs, OBJECT_OP_IGNORED_ERRS, write_quorum) { - warn!("rename_part errs {:?}", &errs); + if let Some(context) = quorum_context { + log_multipart_write_quorum_failure(context, &errs, write_quorum, &err); + } else { + warn!("rename_part errs {:?}", &errs); + } self.cleanup_multipart_path(&[dst_object.to_string(), format!("{dst_object}.meta")]) .await; return Err(err); diff --git a/rustfs/src/error.rs b/rustfs/src/error.rs index 46b4b9ef0..cffa6739e 100644 --- a/rustfs/src/error.rs +++ b/rustfs/src/error.rs @@ -242,6 +242,10 @@ impl From for ApiError { StorageError::BucketExists(_) => S3ErrorCode::BucketAlreadyOwnedByYou, StorageError::StorageFull => S3ErrorCode::ServiceUnavailable, StorageError::SlowDown => S3ErrorCode::SlowDown, + StorageError::ErasureReadQuorum + | StorageError::InsufficientReadQuorum(_, _) + | StorageError::ErasureWriteQuorum + | StorageError::InsufficientWriteQuorum(_, _) => S3ErrorCode::SlowDown, StorageError::NamespaceLockQuorumUnavailable { .. } => S3ErrorCode::ServiceUnavailable, StorageError::Lock(_) => S3ErrorCode::ServiceUnavailable, StorageError::DecommissionNotStarted => S3ErrorCode::InvalidRequest, @@ -451,6 +455,10 @@ mod tests { (StorageError::BucketExists("test".into()), S3ErrorCode::BucketAlreadyOwnedByYou), (StorageError::StorageFull, S3ErrorCode::ServiceUnavailable), (StorageError::SlowDown, S3ErrorCode::SlowDown), + (StorageError::ErasureReadQuorum, S3ErrorCode::SlowDown), + (StorageError::InsufficientReadQuorum("test".into(), "test".into()), S3ErrorCode::SlowDown), + (StorageError::ErasureWriteQuorum, S3ErrorCode::SlowDown), + (StorageError::InsufficientWriteQuorum("test".into(), "test".into()), S3ErrorCode::SlowDown), ( StorageError::NamespaceLockQuorumUnavailable { mode: "write",