From cda76889097ebbd3e85ab61daa1129c615189537 Mon Sep 17 00:00:00 2001 From: Zhengchao An Date: Wed, 8 Jul 2026 15:01:37 +0800 Subject: [PATCH] fix(multipart): clean temp part data on failure (#4412) fix(multipart): clean failed part temp data --- crates/ecstore/src/set_disk/ops/multipart.rs | 485 +++++++++++-------- crates/ecstore/src/set_disk/ops/object.rs | 9 +- 2 files changed, 282 insertions(+), 212 deletions(-) diff --git a/crates/ecstore/src/set_disk/ops/multipart.rs b/crates/ecstore/src/set_disk/ops/multipart.rs index 36bab1982..8bc8dc9f2 100644 --- a/crates/ecstore/src/set_disk/ops/multipart.rs +++ b/crates/ecstore/src/set_disk/ops/multipart.rs @@ -298,225 +298,238 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks { let tmp_part = format!("{}x{}", Uuid::new_v4(), OffsetDateTime::now_utc().unix_timestamp()); let tmp_part_path = Arc::new(format!("{tmp_part}/{part_suffix}")); - let erasure = coding::Erasure::new(fi.erasure.data_blocks, fi.erasure.parity_blocks, fi.erasure.block_size); - let writer_setup_stage_start = rustfs_io_metrics::put_stage_metrics_enabled().then(Instant::now); + let result: Result = async { + let erasure = coding::Erasure::new(fi.erasure.data_blocks, fi.erasure.parity_blocks, fi.erasure.block_size); + let writer_setup_stage_start = rustfs_io_metrics::put_stage_metrics_enabled().then(Instant::now); - let mut writers = Vec::with_capacity(shuffle_disks.len()); - let mut errors = Vec::with_capacity(shuffle_disks.len()); - for disk_op in shuffle_disks.iter() { - if let Some(disk) = disk_op { - let writer = match create_bitrot_writer( - false, - Some(disk), + let mut writers = Vec::with_capacity(shuffle_disks.len()); + let mut errors = Vec::with_capacity(shuffle_disks.len()); + for disk_op in shuffle_disks.iter() { + if let Some(disk) = disk_op { + let writer = match create_bitrot_writer( + false, + Some(disk), + RUSTFS_META_TMP_BUCKET, + &tmp_part_path, + erasure.shard_file_size(data.size()), + erasure.shard_size(), + HashAlgorithm::HighwayHash256S, + ) + .await + { + Ok(writer) => writer, + Err(err) => { + warn!( + event = EVENT_SET_DISK_MULTIPART, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_SET_DISK, + disk = ?disk, + state = "bitrot_writer_skipped", + error = ?err, + "Set disk multipart bitrot writer skipped" + ); + errors.push(Some(err)); + writers.push(None); + continue; + } + }; + + writers.push(Some(writer)); + errors.push(None); + } else { + errors.push(Some(DiskError::DiskNotFound)); + writers.push(None); + } + } + + if let Some(stage_start) = writer_setup_stage_start { + rustfs_io_metrics::record_put_object_stage_duration( + "multipart_set_disk_writer_setup", + stage_start.elapsed().as_secs_f64() * 1000.0, + ); + } + + 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, + ); + Err(to_object_err(write_err.into(), vec![bucket, object]))?; + } + + Err(Error::other(format!("not enough disks to write: {errors:?}")))?; + } + + // Capture the original part size before swapping the stream out for encoding. + let multipart_part_size = data.size(); + let stream = mem::replace( + &mut data.stream, + HashReader::from_stream(Cursor::new(Vec::new()), 0, 0, None, None, false)?, + ); + + let write_path = classify_multipart_part_write_path(multipart_part_size, fi.erasure.block_size); + rustfs_io_metrics::record_put_object_path(write_path.multipart_metric_label()); + let encode_stage_start = rustfs_io_metrics::put_stage_metrics_enabled().then(Instant::now); + + let (reader, w_size) = match write_path { + SmallWritePath::SingleBlockNonInline => { + Arc::new(erasure) + .encode_single_block_non_inline(stream, &mut writers, write_quorum) + .await? + } + SmallWritePath::PipelineBatchedLarge => { + Arc::new(erasure).encode_batched(stream, &mut writers, write_quorum).await? + } + SmallWritePath::Inline | SmallWritePath::Pipeline => { + Arc::new(erasure).encode(stream, &mut writers, write_quorum).await? + } + }; + + if let Some(stage_start) = encode_stage_start { + rustfs_io_metrics::record_put_object_stage_duration( + "multipart_set_disk_encode", + stage_start.elapsed().as_secs_f64() * 1000.0, + ); + } + + let _ = mem::replace(&mut data.stream, reader); + + if (w_size as i64) < data.size() { + warn!( + event = EVENT_SET_DISK_MULTIPART, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_SET_DISK, + bucket, + object, + part_number = part_id, + written_size = w_size, + expected_size = data.size(), + state = "short_write", + "Set disk multipart write produced fewer bytes than expected" + ); + Err(Error::other(format!( + "put_object_part write size < data.size(), w_size={}, data.size={}", + w_size, + data.size() + )))?; + } + + let index_op = data + .stream + .try_get_index() + .map(crate::io_support::rio::compression_index_storage_bytes); + + let mut etag = data.stream.try_resolve_etag().unwrap_or_default(); + + if let Some(ref tag) = opts.preserve_etag { + etag = tag.clone(); + } + + let mut actual_size = data.actual_size(); + if actual_size < 0 { + let is_compressed = fi.is_compressed(); + if !is_compressed { + actual_size = w_size as i64; + } + } + + if fi.is_compressed() { + record_compression_total_memory(actual_size as u64, w_size as u64).await; + } + let checksums = data.as_hash_reader().content_crc(); + + let part_info = ObjectPartInfo { + etag: etag.clone(), + number: part_id, + size: w_size, + mod_time: Some(OffsetDateTime::now_utc()), + actual_size, + index: index_op, + checksums: if checksums.is_empty() { None } else { Some(checksums) }, + ..Default::default() + }; + + let part_info_buff = part_info.marshal_msg()?; + + drop(writers); // drop writers to close all files + + let part_path = format!("{}/{}/{}", upload_id_path, fi.data_dir.unwrap_or_default(), part_suffix); + + // Serialize only the commit (rename_part), not the whole upload. Each + // concurrent stream writes to its own unique temp dir (see `tmp_part` + // above), so the encode/stream phase never conflicts and must stay + // lock-free — holding a lock across it would serialize slow re-transmits + // of the same part and defeat the S3 "last finisher wins" semantics + // (it also caused UploadPart lock-acquire timeouts). The mixed-generation + // hazard is confined to rename_part, where two temp parts are moved + // cross-disk onto the SAME final part_path: interleaving there can leave + // shards from two generations, each individually bitrot-valid, that only + // surface as silent corruption at read time (backlog#853). A write lock + // scoped to the uploadId namespace makes each commit atomic across disks, + // so the last committer wins consistently. Mirrors MinIO's per-uploadID + // NS lock; distinct from the object lock held by complete_multipart_upload + // (disjoint namespaces, no lock-ordering cycle). + let _upload_commit_guard = if opts.no_lock { + None + } else { + Some( + self.acquire_write_lock_diag("put_object_part_commit", RUSTFS_META_MULTIPART_BUCKET, &upload_id_path) + .await?, + ) + }; + + let _ = self + .rename_part( + &disks, RUSTFS_META_TMP_BUCKET, &tmp_part_path, - erasure.shard_file_size(data.size()), - erasure.shard_size(), - HashAlgorithm::HighwayHash256S, - ) - .await - { - Ok(writer) => writer, - Err(err) => { - warn!( - event = EVENT_SET_DISK_MULTIPART, - component = LOG_COMPONENT_ECSTORE, - subsystem = LOG_SUBSYSTEM_SET_DISK, - disk = ?disk, - state = "bitrot_writer_skipped", - error = ?err, - "Set disk multipart bitrot writer skipped" - ); - errors.push(Some(err)); - writers.push(None); - continue; - } - }; - - writers.push(Some(writer)); - errors.push(None); - } else { - errors.push(Some(DiskError::DiskNotFound)); - writers.push(None); - } - } - - if let Some(stage_start) = writer_setup_stage_start { - rustfs_io_metrics::record_put_object_stage_duration( - "multipart_set_disk_writer_setup", - stage_start.elapsed().as_secs_f64() * 1000.0, - ); - } - - 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, + RUSTFS_META_MULTIPART_BUCKET, + &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), - }, - &errors, - write_quorum, - &write_err, - ); - return Err(to_object_err(write_err.into(), vec![bucket, object])); - } + }), + ) + .await?; - return Err(Error::other(format!("not enough disks to write: {errors:?}"))); + drop(_upload_commit_guard); + + let ret: PartInfo = PartInfo { + etag: Some(etag.clone()), + part_num: part_id, + last_mod: Some(OffsetDateTime::now_utc()), + size: w_size, + actual_size, + }; + + // error!("put_object_part ret {:?}", &ret); + + Ok(ret) + } + .await; + + if result.is_err() + && let Err(err) = self.delete_all(RUSTFS_META_TMP_BUCKET, &tmp_part).await + { + warn!(tmp_part = %tmp_part, error = ?err, "failed to cleanup multipart temporary data"); } - // Capture the original part size before swapping the stream out for encoding. - let multipart_part_size = data.size(); - let stream = mem::replace( - &mut data.stream, - HashReader::from_stream(Cursor::new(Vec::new()), 0, 0, None, None, false)?, - ); - - let write_path = classify_multipart_part_write_path(multipart_part_size, fi.erasure.block_size); - rustfs_io_metrics::record_put_object_path(write_path.multipart_metric_label()); - let encode_stage_start = rustfs_io_metrics::put_stage_metrics_enabled().then(Instant::now); - - let (reader, w_size) = match write_path { - SmallWritePath::SingleBlockNonInline => { - Arc::new(erasure) - .encode_single_block_non_inline(stream, &mut writers, write_quorum) - .await? - } - SmallWritePath::PipelineBatchedLarge => Arc::new(erasure).encode_batched(stream, &mut writers, write_quorum).await?, - SmallWritePath::Inline | SmallWritePath::Pipeline => { - Arc::new(erasure).encode(stream, &mut writers, write_quorum).await? - } - }; // TODO: delete temporary directory on error - - if let Some(stage_start) = encode_stage_start { - rustfs_io_metrics::record_put_object_stage_duration( - "multipart_set_disk_encode", - stage_start.elapsed().as_secs_f64() * 1000.0, - ); - } - - let _ = mem::replace(&mut data.stream, reader); - - if (w_size as i64) < data.size() { - warn!( - event = EVENT_SET_DISK_MULTIPART, - component = LOG_COMPONENT_ECSTORE, - subsystem = LOG_SUBSYSTEM_SET_DISK, - bucket, - object, - part_number = part_id, - written_size = w_size, - expected_size = data.size(), - state = "short_write", - "Set disk multipart write produced fewer bytes than expected" - ); - return Err(Error::other(format!( - "put_object_part write size < data.size(), w_size={}, data.size={}", - w_size, - data.size() - ))); - } - - let index_op = data - .stream - .try_get_index() - .map(crate::io_support::rio::compression_index_storage_bytes); - - let mut etag = data.stream.try_resolve_etag().unwrap_or_default(); - - if let Some(ref tag) = opts.preserve_etag { - etag = tag.clone(); - } - - let mut actual_size = data.actual_size(); - if actual_size < 0 { - let is_compressed = fi.is_compressed(); - if !is_compressed { - actual_size = w_size as i64; - } - } - - if fi.is_compressed() { - record_compression_total_memory(actual_size as u64, w_size as u64).await; - } - let checksums = data.as_hash_reader().content_crc(); - - let part_info = ObjectPartInfo { - etag: etag.clone(), - number: part_id, - size: w_size, - mod_time: Some(OffsetDateTime::now_utc()), - actual_size, - index: index_op, - checksums: if checksums.is_empty() { None } else { Some(checksums) }, - ..Default::default() - }; - - let part_info_buff = part_info.marshal_msg()?; - - drop(writers); // drop writers to close all files - - let part_path = format!("{}/{}/{}", upload_id_path, fi.data_dir.unwrap_or_default(), part_suffix); - - // Serialize only the commit (rename_part), not the whole upload. Each - // concurrent stream writes to its own unique temp dir (see `tmp_part` - // above), so the encode/stream phase never conflicts and must stay - // lock-free — holding a lock across it would serialize slow re-transmits - // of the same part and defeat the S3 "last finisher wins" semantics - // (it also caused UploadPart lock-acquire timeouts). The mixed-generation - // hazard is confined to rename_part, where two temp parts are moved - // cross-disk onto the SAME final part_path: interleaving there can leave - // shards from two generations, each individually bitrot-valid, that only - // surface as silent corruption at read time (backlog#853). A write lock - // scoped to the uploadId namespace makes each commit atomic across disks, - // so the last committer wins consistently. Mirrors MinIO's per-uploadID - // NS lock; distinct from the object lock held by complete_multipart_upload - // (disjoint namespaces, no lock-ordering cycle). - let _upload_commit_guard = if opts.no_lock { - None - } else { - Some( - self.acquire_write_lock_diag("put_object_part_commit", RUSTFS_META_MULTIPART_BUCKET, &upload_id_path) - .await?, - ) - }; - - let _ = self - .rename_part( - &disks, - RUSTFS_META_TMP_BUCKET, - &tmp_part_path, - RUSTFS_META_MULTIPART_BUCKET, - &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?; - - drop(_upload_commit_guard); - - let ret: PartInfo = PartInfo { - etag: Some(etag.clone()), - part_num: part_id, - last_mod: Some(OffsetDateTime::now_utc()), - size: w_size, - actual_size, - }; - - // error!("put_object_part ret {:?}", &ret); - - Ok(ret) + result } #[tracing::instrument(skip(self))] @@ -1434,6 +1447,28 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks { #[cfg(test)] mod tests { use super::*; + use crate::disk::DiskAPI as _; + use crate::set_disk::ops::object::hermetic_set_disks_support::hermetic_set_disks; + use tempfile::TempDir; + + async fn non_trash_tmp_entries(temp_dirs: &[TempDir]) -> Vec { + let mut leftovers = Vec::new(); + for temp_dir in temp_dirs { + let tmp_path = temp_dir.path().join(RUSTFS_META_TMP_BUCKET); + let mut read_dir = match tokio::fs::read_dir(&tmp_path).await { + Ok(read_dir) => read_dir, + Err(err) if err.kind() == std::io::ErrorKind::NotFound => continue, + Err(err) => panic!("tmp dir {tmp_path:?} should be listable: {err}"), + }; + while let Some(entry) = read_dir.next_entry().await.expect("tmp dir entry should be readable") { + let name = entry.file_name().to_string_lossy().to_string(); + if name != ".trash" { + leftovers.push(format!("{}/{name}", tmp_path.display())); + } + } + } + leftovers + } #[tokio::test] async fn collect_list_parts_results_fails_early_when_quorum_is_impossible() { @@ -1522,6 +1557,38 @@ mod tests { assert!(started.elapsed() < Duration::from_millis(120)); } + #[tokio::test] + async fn put_object_part_failure_cleans_tmp_workspace_inline() { + let (temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await; + let bucket = "multipart-tmp-clean-bucket"; + let object = "object"; + for disk in &disk_stores { + disk.make_volume(bucket).await.expect("bucket volume should be created"); + } + + let upload = set_disks + .new_multipart_upload(bucket, object, &ObjectOptions::default()) + .await + .expect("multipart upload should be created"); + let declared_size = 1024 * 1024; + let short_stream = Cursor::new(vec![7u8; 512]); + let mut reader = PutObjReader::new( + HashReader::from_stream(short_stream, declared_size, declared_size, None, None, false) + .expect("hash reader should be constructed"), + ); + + let err = set_disks + .put_object_part(bucket, object, &upload.upload_id, 1, &mut reader, &ObjectOptions::default()) + .await + .expect_err("short multipart stream should fail"); + + let leftovers = non_trash_tmp_entries(&temp_dirs).await; + assert!( + leftovers.is_empty(), + "failed multipart upload part must not leave tmp shards behind, leftovers: {leftovers:?}, err: {err}" + ); + } + #[test] fn reduce_quorum_part_numbers_only_keeps_parts_present_on_quorum_of_drives() { let object_parts = vec![ diff --git a/crates/ecstore/src/set_disk/ops/object.rs b/crates/ecstore/src/set_disk/ops/object.rs index 064f6f51d..76c686968 100644 --- a/crates/ecstore/src/set_disk/ops/object.rs +++ b/crates/ecstore/src/set_disk/ops/object.rs @@ -2322,7 +2322,7 @@ mod b3_write_quorum_tests { } #[cfg(test)] -mod hermetic_set_disks_support { +pub(in crate::set_disk::ops) mod hermetic_set_disks_support { //! Shared hermetic `SetDisks` construction for the ops tests below: the //! `SetDisks` under test is built directly on formatted local disks (same //! pattern as the `ops/locking.rs` tests) so the tests stay hermetic — no @@ -2334,7 +2334,10 @@ mod hermetic_set_disks_support { use tempfile::TempDir; use tokio::sync::RwLock; - pub(super) async fn make_formatted_local_disk(disk_idx: usize, format: &FormatV3) -> (TempDir, Endpoint, DiskStore) { + pub(in crate::set_disk::ops) async fn make_formatted_local_disk( + disk_idx: usize, + format: &FormatV3, + ) -> (TempDir, Endpoint, DiskStore) { let dir = tempfile::tempdir().expect("tempdir should be created"); let mut endpoint = Endpoint::try_from(dir.path().to_str().expect("tempdir path should be utf8")).expect("endpoint should parse"); @@ -2361,7 +2364,7 @@ mod hermetic_set_disks_support { (dir, endpoint, disk) } - pub(super) async fn hermetic_set_disks(disk_count: usize) -> (Vec, Vec, Arc) { + pub(in crate::set_disk::ops) async fn hermetic_set_disks(disk_count: usize) -> (Vec, Vec, Arc) { let format = FormatV3::new(1, disk_count); let mut temp_dirs = Vec::with_capacity(disk_count);