From 6a59c0a474b47783132e066e60daadeef690e108 Mon Sep 17 00:00:00 2001 From: weisd Date: Fri, 24 Oct 2025 18:23:32 +0800 Subject: [PATCH] fix: multipart upload checksum validation (#712) * fix multipart upload checksum --- crates/ecstore/src/set_disk.rs | 212 ++++++++++++++++++++++++-------- crates/filemeta/src/filemeta.rs | 2 +- crates/rio/src/checksum.rs | 6 + crates/rio/src/hash_reader.rs | 9 +- rustfs/src/storage/ecfs.rs | 31 +++-- 5 files changed, 190 insertions(+), 70 deletions(-) diff --git a/crates/ecstore/src/set_disk.rs b/crates/ecstore/src/set_disk.rs index 4803f52a2..71c05df32 100644 --- a/crates/ecstore/src/set_disk.rs +++ b/crates/ecstore/src/set_disk.rs @@ -73,9 +73,9 @@ use rustfs_filemeta::{ use rustfs_lock::fast_lock::types::LockResult; use rustfs_madmin::heal_commands::{HealDriveInfo, HealResultItem}; use rustfs_rio::{EtagResolvable, HashReader, HashReaderMut, TryGetIndex as _, WarpReader}; -use rustfs_utils::http::headers::AMZ_OBJECT_TAGGING; +use rustfs_utils::http::RUSTFS_BUCKET_REPLICATION_SSEC_CHECKSUM; use rustfs_utils::http::headers::AMZ_STORAGE_CLASS; -use rustfs_utils::http::headers::RESERVED_METADATA_PREFIX_LOWER; +use rustfs_utils::http::headers::{AMZ_OBJECT_TAGGING, RESERVED_METADATA_PREFIX, RESERVED_METADATA_PREFIX_LOWER}; use rustfs_utils::{ HashAlgorithm, crypto::hex, @@ -4953,6 +4953,8 @@ impl StorageAPI for SetDisks { } } + let checksums = data.as_hash_reader().content_crc(); + let part_info = ObjectPartInfo { etag: etag.clone(), number: part_id, @@ -4960,13 +4962,10 @@ impl StorageAPI for SetDisks { mod_time: Some(OffsetDateTime::now_utc()), actual_size, index: index_op, + checksums: if checksums.is_empty() { None } else { Some(checksums) }, ..Default::default() }; - // debug!("put_object_part part_info {:?}", part_info); - - // fi.parts = vec![part_info.clone()]; - let part_info_buff = part_info.marshal_msg()?; drop(writers); // drop writers to close all files @@ -5317,7 +5316,13 @@ impl StorageAPI for SetDisks { } fi.data_dir = Some(Uuid::new_v4()); - fi.fresh = true; + + if let Some(cssum) = user_defined.get(RUSTFS_BUCKET_REPLICATION_SSEC_CHECKSUM) + && !cssum.is_empty() + { + fi.checksum = base64_simd::STANDARD.decode_to_vec(cssum).ok().map(Bytes::from); + user_defined.remove(RUSTFS_BUCKET_REPLICATION_SSEC_CHECKSUM); + } let parts_metadata = vec![fi.clone(); disks.len()]; @@ -5343,10 +5348,10 @@ impl StorageAPI for SetDisks { let mod_time = opts.mod_time.unwrap_or(OffsetDateTime::now_utc()); - for fi in parts_metadatas.iter_mut() { - fi.metadata = user_defined.clone(); - fi.mod_time = Some(mod_time); - fi.fresh = true; + for f in parts_metadatas.iter_mut() { + f.metadata = user_defined.clone(); + f.mod_time = Some(mod_time); + f.fresh = true; } // fi.mod_time = Some(now); @@ -5463,23 +5468,27 @@ impl StorageAPI for SetDisks { return Err(Error::other("part result number err")); } + let mut checksum_type = rustfs_rio::ChecksumType::NONE; + if let Some(cs) = fi.metadata.get(rustfs_rio::RUSTFS_MULTIPART_CHECKSUM) { - let Some(checksum_type) = fi.metadata.get(rustfs_rio::RUSTFS_MULTIPART_CHECKSUM_TYPE) else { + let Some(ct) = fi.metadata.get(rustfs_rio::RUSTFS_MULTIPART_CHECKSUM_TYPE) else { return Err(Error::other("checksum type not found")); }; if opts.want_checksum.is_some() && !opts.want_checksum.as_ref().is_some_and(|v| { v.checksum_type - .is(rustfs_rio::ChecksumType::from_string_with_obj_type(cs, checksum_type)) + .is(rustfs_rio::ChecksumType::from_string_with_obj_type(cs, ct)) }) { return Err(Error::other(format!( "checksum type mismatch, got {:?}, want {:?}", opts.want_checksum.as_ref().unwrap(), - rustfs_rio::ChecksumType::from_string_with_obj_type(cs, checksum_type) + rustfs_rio::ChecksumType::from_string_with_obj_type(cs, ct) ))); } + + checksum_type = rustfs_rio::ChecksumType::from_string_with_obj_type(cs, ct); } for (i, part) in object_parts.iter().enumerate() { @@ -5515,6 +5524,12 @@ impl StorageAPI for SetDisks { let mut object_size: usize = 0; let mut object_actual_size: i64 = 0; + let mut checksum_combined = bytes::BytesMut::new(); + let mut checksum = rustfs_rio::Checksum { + checksum_type, + ..Default::default() + }; + for (i, p) in uploaded_parts.iter().enumerate() { let has_part = curr_fi.parts.iter().find(|v| v.number == p.part_num); if has_part.is_none() { @@ -5555,6 +5570,75 @@ impl StorageAPI for SetDisks { )); } + if checksum_type.is_set() { + let Some(crc) = ext_part + .checksums + .as_ref() + .and_then(|f| f.get(checksum_type.to_string().as_str())) + .cloned() + else { + error!( + "complete_multipart_upload fi.checksum not found type={checksum_type}, part_id={}, bucket={}, object={}", + p.part_num, bucket, object + ); + return Err(Error::InvalidPart(p.part_num, ext_part.etag.clone(), p.etag.clone().unwrap_or_default())); + }; + + let part_crc = match checksum_type { + rustfs_rio::ChecksumType::SHA256 => p.checksum_sha256.clone(), + rustfs_rio::ChecksumType::SHA1 => p.checksum_sha1.clone(), + rustfs_rio::ChecksumType::CRC32 => p.checksum_crc32.clone(), + rustfs_rio::ChecksumType::CRC32C => p.checksum_crc32c.clone(), + rustfs_rio::ChecksumType::CRC64_NVME => p.checksum_crc64nvme.clone(), + _ => { + error!( + "complete_multipart_upload checksum type={checksum_type}, part_id={}, bucket={}, object={}", + p.part_num, bucket, object + ); + return Err(Error::InvalidPart(p.part_num, ext_part.etag.clone(), p.etag.clone().unwrap_or_default())); + } + }; + + if part_crc.clone().unwrap_or_default() != crc { + error!("complete_multipart_upload checksum_type={checksum_type:?}, part_crc={part_crc:?}, crc={crc:?}"); + error!( + "complete_multipart_upload checksum mismatch part_id={}, bucket={}, object={}", + p.part_num, bucket, object + ); + return Err(Error::InvalidPart(p.part_num, ext_part.etag.clone(), p.etag.clone().unwrap_or_default())); + } + + let Some(cs) = rustfs_rio::Checksum::new_with_type(checksum_type, &crc) else { + error!( + "complete_multipart_upload checksum new_with_type failed part_id={}, bucket={}, object={}", + p.part_num, bucket, object + ); + return Err(Error::InvalidPart(p.part_num, ext_part.etag.clone(), p.etag.clone().unwrap_or_default())); + }; + + if !cs.valid() { + error!( + "complete_multipart_upload checksum valid failed part_id={}, bucket={}, object={}", + p.part_num, bucket, object + ); + return Err(Error::InvalidPart(p.part_num, ext_part.etag.clone(), p.etag.clone().unwrap_or_default())); + } + + if checksum_type.full_object_requested() { + if let Err(err) = checksum.add_part(&cs, ext_part.actual_size) { + error!( + "complete_multipart_upload checksum add_part failed part_id={}, bucket={}, object={}", + p.part_num, bucket, object + ); + return Err(Error::InvalidPart(p.part_num, ext_part.etag.clone(), p.etag.clone().unwrap_or_default())); + } + } + + checksum_combined.extend_from_slice(cs.raw.as_slice()); + } + + // TODO: check min part size + object_size += ext_part.size; object_actual_size += ext_part.actual_size; @@ -5569,6 +5653,52 @@ impl StorageAPI for SetDisks { }); } + if let Some(wtcs) = opts.want_checksum.as_ref() { + if checksum_type.full_object_requested() { + if wtcs.encoded != checksum.encoded { + error!( + "complete_multipart_upload checksum mismatch want={}, got={}", + wtcs.encoded, checksum.encoded + ); + return Err(Error::other(format!( + "complete_multipart_upload checksum mismatch want={}, got={}", + wtcs.encoded, checksum.encoded + ))); + } + } else if let Err(err) = wtcs.matches(&checksum_combined, uploaded_parts.len() as i32) { + error!( + "complete_multipart_upload checksum matches failed want={}, got={}", + wtcs.encoded, checksum.encoded + ); + return Err(Error::other(format!( + "complete_multipart_upload checksum matches failed want={}, got={}", + wtcs.encoded, checksum.encoded + ))); + } + } + + if let Some(rc_crc) = opts.user_defined.get(RUSTFS_BUCKET_REPLICATION_SSEC_CHECKSUM) { + if let Ok(rc_crc_bytes) = base64_simd::STANDARD.decode_to_vec(rc_crc) { + fi.checksum = Some(Bytes::from(rc_crc_bytes)); + } else { + error!("complete_multipart_upload decode rc_crc failed rc_crc={}", rc_crc); + } + } + + if checksum_type.is_set() { + checksum_type + .merge(rustfs_rio::ChecksumType::MULTIPART) + .merge(rustfs_rio::ChecksumType::INCLUDES_MULTIPART); + if !checksum_type.full_object_requested() { + checksum = rustfs_rio::Checksum::new_from_data(checksum_type, &checksum_combined) + .ok_or_else(|| Error::other("checksum new_from_data failed"))?; + } + fi.checksum = Some(checksum.to_bytes(&checksum_combined)); + } + + fi.metadata.remove(rustfs_rio::RUSTFS_MULTIPART_CHECKSUM); + fi.metadata.remove(rustfs_rio::RUSTFS_MULTIPART_CHECKSUM_TYPE); + fi.size = object_size as i64; fi.mod_time = opts.mod_time; if fi.mod_time.is_none() { @@ -5586,11 +5716,22 @@ impl StorageAPI for SetDisks { fi.metadata.insert("etag".to_owned(), etag); - fi.metadata - .insert(format!("{RESERVED_METADATA_PREFIX_LOWER}actual-size"), object_actual_size.to_string()); - - fi.metadata - .insert("x-rustfs-encryption-original-size".to_string(), object_actual_size.to_string()); + if opts.replication_request { + if let Some(actual_size) = opts + .user_defined + .get(format!("{RESERVED_METADATA_PREFIX_LOWER}Actual-Object-Size").as_str()) + { + fi.metadata + .insert(format!("{RESERVED_METADATA_PREFIX}actual-size"), actual_size.clone()); + fi.metadata + .insert("x-rustfs-encryption-original-size".to_string(), actual_size.to_string()); + } + } else { + fi.metadata + .insert(format!("{RESERVED_METADATA_PREFIX}actual-size"), object_actual_size.to_string()); + fi.metadata + .insert("x-rustfs-encryption-original-size".to_string(), object_actual_size.to_string()); + } if fi.is_compressed() { fi.metadata @@ -5601,9 +5742,6 @@ impl StorageAPI for SetDisks { fi.set_data_moved(); } - // TODO: object_actual_size - let _ = object_actual_size; - for meta in parts_metadatas.iter_mut() { if meta.is_valid() { meta.size = fi.size; @@ -5611,13 +5749,12 @@ impl StorageAPI for SetDisks { meta.parts.clone_from(&fi.parts); meta.metadata = fi.metadata.clone(); meta.versioned = opts.versioned || opts.version_suspended; - - // TODO: Checksum + meta.checksum = fi.checksum.clone(); } } let mut parts = Vec::with_capacity(curr_fi.parts.len()); - // TODO: 优化 cleanupMultipartPath + for p in curr_fi.parts.iter() { parts.push(path_join_buf(&[ &upload_id_path, @@ -5632,28 +5769,6 @@ impl StorageAPI for SetDisks { format!("part.{}", p.number).as_str(), ])); } - - // let _ = self - // .remove_part_meta( - // bucket, - // object, - // upload_id, - // curr_fi.data_dir.unwrap_or(Uuid::nil()).to_string().as_str(), - // p.number, - // ) - // .await; - - // if !fi.parts.iter().any(|v| v.number == p.number) { - // let _ = self - // .remove_object_part( - // bucket, - // object, - // upload_id, - // curr_fi.data_dir.unwrap_or(Uuid::nil()).to_string().as_str(), - // p.number, - // ) - // .await; - // } } { @@ -5672,9 +5787,6 @@ impl StorageAPI for SetDisks { ) .await?; - // debug!("complete fileinfo {:?}", &fi); - - // TODO: reduce_common_data_dir if let Some(old_dir) = op_old_dir { self.commit_rename_data_dir(&shuffle_disks, bucket, object, &old_dir.to_string(), write_quorum) .await?; diff --git a/crates/filemeta/src/filemeta.rs b/crates/filemeta/src/filemeta.rs index a3a1cd48d..070182e0c 100644 --- a/crates/filemeta/src/filemeta.rs +++ b/crates/filemeta/src/filemeta.rs @@ -1382,7 +1382,7 @@ impl From for FileMetaVersion { FileMetaVersion { version_type: VersionType::Object, delete_marker: None, - object: Some(value.into()), + object: Some(MetaObject::from(value)), write_version: 0, } } diff --git a/crates/rio/src/checksum.rs b/crates/rio/src/checksum.rs index 94445a94c..15596282f 100644 --- a/crates/rio/src/checksum.rs +++ b/crates/rio/src/checksum.rs @@ -78,6 +78,12 @@ impl ChecksumType { (self.0 & t.0) == t.0 } + /// Merge another checksum type into this one + pub fn merge(&mut self, other: ChecksumType) -> &mut Self { + self.0 |= other.0; + self + } + /// Get the base checksum type (without flags) pub fn base(self) -> ChecksumType { ChecksumType(self.0 & Self::BASE_TYPE_MASK) diff --git a/crates/rio/src/hash_reader.rs b/crates/rio/src/hash_reader.rs index b2a3c6bc6..03f13864f 100644 --- a/crates/rio/src/hash_reader.rs +++ b/crates/rio/src/hash_reader.rs @@ -331,6 +331,8 @@ impl HashReader { } else { return Err(std::io::Error::new(std::io::ErrorKind::InvalidData, "Invalid checksum type")); } + + tracing::debug!("add_non_trailing_checksum checksum={checksum:?}"); } Ok(()) } @@ -359,9 +361,10 @@ impl HashReader { if checksum.checksum_type.trailing() { if let Some(trailer) = self.trailer_s3s.as_ref() { if let Some(Some(checksum_str)) = trailer.read(|headers| { - headers - .get(checksum.checksum_type.to_string()) - .and_then(|value| value.to_str().ok().map(|s| s.to_string())) + checksum + .checksum_type + .key() + .and_then(|key| headers.get(key).and_then(|value| value.to_str().ok().map(|s| s.to_string()))) }) { map.insert(checksum.checksum_type.to_string(), checksum_str); } diff --git a/rustfs/src/storage/ecfs.rs b/rustfs/src/storage/ecfs.rs index 7453ea65a..1b125c628 100644 --- a/rustfs/src/storage/ecfs.rs +++ b/rustfs/src/storage/ecfs.rs @@ -1934,7 +1934,7 @@ impl S3 for FS { .decrypt_checksums(opts.part_number.unwrap_or(0), &req.headers) .map_err(ApiError::from)?; - warn!("get object metadata checksums: {:?}", checksums); + debug!("get object metadata checksums: {:?}", checksums); for (key, checksum) in checksums { if key == AMZ_CHECKSUM_TYPE { checksum_type = Some(ChecksumType::from(checksum)); @@ -3319,26 +3319,25 @@ impl S3 for FS { &self, req: S3Request, ) -> S3Result> { + let input = req.input; let CompleteMultipartUploadInput { multipart_upload, bucket, key, upload_id, .. - } = req.input; - - // error!("complete_multipart_upload {:?}", multipart_upload); - // mc cp step 5 + } = input; let Some(multipart_upload) = multipart_upload else { return Err(s3_error!(InvalidPart)) }; let opts = &get_complete_multipart_upload_opts(&req.headers).map_err(ApiError::from)?; - let mut uploaded_parts = Vec::new(); - - for part in multipart_upload.parts.into_iter().flatten() { - uploaded_parts.push(CompletePart::from(part)); - } + let uploaded_parts = multipart_upload + .parts + .unwrap_or_default() + .into_iter() + .map(CompletePart::from) + .collect::>(); // is part number sorted? if !uploaded_parts.is_sorted_by_key(|p| p.part_num) { @@ -3397,12 +3396,12 @@ impl S3 for FS { server_side_encryption, ssekms_key_id ); - let mut checksum_crc32 = None; - let mut checksum_crc32c = None; - let mut checksum_sha1 = None; - let mut checksum_sha256 = None; - let mut checksum_crc64nvme = None; - let mut checksum_type = None; + let mut checksum_crc32 = input.checksum_crc32; + let mut checksum_crc32c = input.checksum_crc32c; + let mut checksum_sha1 = input.checksum_sha1; + let mut checksum_sha256 = input.checksum_sha256; + let mut checksum_crc64nvme = input.checksum_crc64nvme; + let mut checksum_type = input.checksum_type; // checksum let (checksums, _is_multipart) = obj_info