From f9d45e41e159dd42de70137ad4ccf37abd525cbb Mon Sep 17 00:00:00 2001 From: cxymds Date: Sat, 22 Aug 2026 19:39:01 +0800 Subject: [PATCH] fix(quota): account compressed deletes by committed size (#6365) --- crates/ecstore/src/api/mod.rs | 4 +- .../bucket/replication/replication_pool.rs | 2 +- .../replication/replication_resyncer.rs | 2 +- .../replication_target_boundary.rs | 16 +- crates/ecstore/src/core/sets.rs | 123 ++++---- crates/ecstore/src/data_usage/mod.rs | 116 +++++++- crates/ecstore/src/object_api/types.rs | 38 ++- crates/ecstore/src/set_disk/mod.rs | 2 +- crates/ecstore/src/set_disk/ops/object.rs | 163 ++++++++++- .../ecstore/src/storage_api_contracts/mod.rs | 4 +- crates/ecstore/src/store/object.rs | 65 ++++- crates/storage-api/src/lib.rs | 1 + crates/storage-api/src/object.rs | 24 ++ rustfs/src/app/bucket_usecase.rs | 2 +- rustfs/src/app/object_usecase.rs | 271 +++++++++++++++--- rustfs/src/app/storage_api.rs | 10 +- rustfs/src/storage/s3_api/bucket.rs | 44 ++- rustfs/src/storage/storage_api.rs | 2 + 18 files changed, 754 insertions(+), 135 deletions(-) diff --git a/crates/ecstore/src/api/mod.rs b/crates/ecstore/src/api/mod.rs index 17fffac3d..456d38ac7 100644 --- a/crates/ecstore/src/api/mod.rs +++ b/crates/ecstore/src/api/mod.rs @@ -317,8 +317,6 @@ pub mod config { } pub mod data_usage { - #[cfg(feature = "test-util")] - pub use crate::data_usage::seed_bucket_usage_memory_for_test; pub use crate::data_usage::{ DATA_USAGE_CACHE_NAME, apply_bucket_usage_memory_overlay, compute_bucket_usage, init_compression_total_memory_from_backend, invalidate_admin_data_usage_snapshot_cache, @@ -330,6 +328,8 @@ pub mod data_usage { remove_bucket_usage_from_backend, replace_bucket_usage_memory_from_info, store_compression_total_in_backend, store_data_usage_in_backend, }; + #[cfg(feature = "test-util")] + pub use crate::data_usage::{get_bucket_usage_memory, seed_bucket_usage_memory_for_test}; } pub mod disk { diff --git a/crates/ecstore/src/bucket/replication/replication_pool.rs b/crates/ecstore/src/bucket/replication/replication_pool.rs index efe517b99..61dbf2f1e 100644 --- a/crates/ecstore/src/bucket/replication/replication_pool.rs +++ b/crates/ecstore/src/bucket/replication/replication_pool.rs @@ -2855,7 +2855,7 @@ fn replicate_object_info_from_object_info( .map(|v| OffsetDateTime::parse(&v, &Rfc3339).unwrap_or(OffsetDateTime::UNIX_EPOCH)); let mut rstate = oi.replication_state(); rstate.replicate_decision_str = dsc.to_string(); - let asz = oi.get_actual_size().unwrap_or_default(); + let asz = oi.get_actual_size_or_physical(); let ssec = replication_object_is_ssec_encrypted(&oi.user_defined); let checksum = if ssec { oi.checksum.clone() } else { None }; diff --git a/crates/ecstore/src/bucket/replication/replication_resyncer.rs b/crates/ecstore/src/bucket/replication/replication_resyncer.rs index d921ea411..e670a7d6f 100644 --- a/crates/ecstore/src/bucket/replication/replication_resyncer.rs +++ b/crates/ecstore/src/bucket/replication/replication_resyncer.rs @@ -1412,7 +1412,7 @@ pub async fn get_heal_replicate_object_info(oi: &ObjectInfo, rcfg: &ReplicationC }; let mut replication_state = oi.replication_state(); replication_state.replicate_decision_str = dsc.to_string(); - let actual_size = oi.get_actual_size().unwrap_or_default(); + let actual_size = oi.get_actual_size_or_physical(); Ok(ReplicateObjectInfo { name: oi.name.clone(), diff --git a/crates/ecstore/src/bucket/replication/replication_target_boundary.rs b/crates/ecstore/src/bucket/replication/replication_target_boundary.rs index 455a37ed0..4fe89967d 100644 --- a/crates/ecstore/src/bucket/replication/replication_target_boundary.rs +++ b/crates/ecstore/src/bucket/replication/replication_target_boundary.rs @@ -389,7 +389,7 @@ fn replication_source_object(object_info: &ObjectInfo) -> ReplicationSourceObjec .map(|mod_time| OffsetDateTime::from_unix_timestamp(mod_time.unix_timestamp()).unwrap_or(mod_time)), version_id: object_info.version_id.map(|version_id| version_id.to_string()), etag: object_info.etag.as_deref(), - actual_size: object_info.get_actual_size().unwrap_or_default(), + actual_size: object_info.get_actual_size_or_physical(), delete_marker: object_info.delete_marker, content_type: object_info.content_type.as_deref(), content_encoding: object_info.content_encoding.as_deref(), @@ -542,6 +542,20 @@ mod tests { assert!(replication_target_head_is_newer_null_version(&source, &target)); } + #[test] + fn replication_source_uses_physical_size_for_unknown_compressed_object() { + let mut metadata = HashMap::new(); + rustfs_utils::http::insert_str(&mut metadata, rustfs_utils::http::SUFFIX_COMPRESSION, "zstd".to_string()); + let source = ObjectInfo { + size: 128, + actual_size: -1, + user_defined: Arc::new(metadata), + ..Default::default() + }; + + assert_eq!(replication_source_object(&source).actual_size, 128); + } + #[test] fn replication_target_head_content_matches_compare_etag_only() { let source = ObjectInfo { diff --git a/crates/ecstore/src/core/sets.rs b/crates/ecstore/src/core/sets.rs index 3acf1a705..d9b354a08 100644 --- a/crates/ecstore/src/core/sets.rs +++ b/crates/ecstore/src/core/sets.rs @@ -21,7 +21,7 @@ use crate::storage_api_contracts::{ bucket::{BucketInfo, BucketOperations, BucketOptions, DeleteBucketOptions, MakeBucketOptions}, list::{StorageListObjectVersionsInfo, StorageListObjectsV2Info, StorageObjectInfoOrErr, StorageWalkOptions}, multipart::{CompletePart, ListMultipartsInfo, ListPartsInfo, MultipartInfo, MultipartUploadResult, PartInfo}, - object::{DeletedObject, ObjectIO as _, ObjectOperations as _, ObjectToDelete}, + object::{DeleteAccounting, DeletedObject, ObjectIO as _, ObjectOperations as _, ObjectToDelete}, range::HTTPRangeSpec, }; use crate::{ @@ -414,6 +414,66 @@ fn apply_delete_objects_results( } } +fn apply_delete_accounting_results( + accounting: &mut [Option], + set_objects: &[DelObj], + set_accounting: &[Option], +) { + for (obj, value) in set_objects.iter().zip(set_accounting.iter()) { + accounting[obj.orig_idx] = value.clone(); + } +} + +impl Sets { + pub(crate) async fn delete_objects_with_accounting( + &self, + bucket: &str, + objects: Vec, + opts: ObjectOptions, + ) -> (Vec, Vec>, Vec>) { + let mut del_objects = vec![DeletedObject::default(); objects.len()]; + let mut del_errs = vec![None; objects.len()]; + let mut accounting = vec![None; objects.len()]; + let mut set_obj_map = HashMap::new(); + + for (i, obj) in objects.iter().enumerate() { + let idx = self.get_hashed_set_index(obj.object_name.as_str()); + set_obj_map.entry(idx).or_insert_with(Vec::new).push(DelObj { + orig_idx: i, + obj: obj.clone(), + }); + } + + let max_concurrent = set_obj_map.len().min(num_cpus::get()).max(1); + let semaphore = Arc::new(tokio::sync::Semaphore::new(max_concurrent)); + let mut futures = FuturesUnordered::new(); + let bucket = bucket.to_owned(); + + for (set_index, set_objects) in set_obj_map { + let disks = self.get_disks(set_index); + let objects = set_objects.iter().map(|entry| entry.obj.clone()).collect::>(); + let bucket = bucket.clone(); + let opts = opts.clone(); + let semaphore = semaphore.clone(); + futures.push(async move { + let _permit = semaphore + .acquire_owned() + .await + .expect("delete_objects semaphore should remain open"); + let (deleted, errors, accounting) = disks.delete_objects_with_accounting(&bucket, objects, opts).await; + (set_objects, deleted, errors, accounting) + }); + } + + while let Some((set_objects, deleted, errors, set_accounting)) = futures.next().await { + apply_delete_objects_results(&mut del_objects, &mut del_errs, &set_objects, &deleted, errors); + apply_delete_accounting_results(&mut accounting, &set_objects, &set_accounting); + } + + (del_objects, del_errs, accounting) + } +} + #[async_trait::async_trait] impl crate::storage_api_contracts::object::ObjectIO for Sets { type Error = Error; @@ -655,65 +715,8 @@ impl crate::storage_api_contracts::object::ObjectOperations for Sets { objects: Vec, opts: ObjectOptions, ) -> (Vec, Vec>) { - // Default return value - let mut del_objects = vec![DeletedObject::default(); objects.len()]; - - let mut del_errs = Vec::with_capacity(objects.len()); - for _ in 0..objects.len() { - del_errs.push(None) - } - - let mut set_obj_map = HashMap::new(); - - // hash key - for (i, obj) in objects.iter().enumerate() { - let idx = self.get_hashed_set_index(obj.object_name.as_str()); - - if !set_obj_map.contains_key(&idx) { - set_obj_map.insert( - idx, - vec![DelObj { - // set_idx: idx, - orig_idx: i, - obj: obj.clone(), - }], - ); - } else if let Some(val) = set_obj_map.get_mut(&idx) { - val.push(DelObj { - // set_idx: idx, - orig_idx: i, - obj: obj.clone(), - }); - } - } - - let max_concurrent = set_obj_map.len().min(num_cpus::get()).max(1); - let semaphore = Arc::new(tokio::sync::Semaphore::new(max_concurrent)); - let mut futures = FuturesUnordered::new(); - let bucket = bucket.to_string(); - - for (k, v) in set_obj_map { - let disks = self.get_disks(k); - let objs: Vec = v.iter().map(|v| v.obj.clone()).collect(); - let bucket = bucket.clone(); - let opts = opts.clone(); - let semaphore = semaphore.clone(); - - futures.push(async move { - let _permit = semaphore - .acquire_owned() - .await - .expect("delete_objects semaphore should remain open"); - let (dobjects, errs) = disks.delete_objects(&bucket, objs, opts).await; - (v, dobjects, errs) - }); - } - - while let Some((v, dobjects, errs)) = futures.next().await { - apply_delete_objects_results(&mut del_objects, &mut del_errs, &v, &dobjects, errs); - } - - (del_objects, del_errs) + let (deleted, errors, _) = self.delete_objects_with_accounting(bucket, objects, opts).await; + (deleted, errors) } #[tracing::instrument(skip(self))] diff --git a/crates/ecstore/src/data_usage/mod.rs b/crates/ecstore/src/data_usage/mod.rs index 917edc649..3bf9bf500 100644 --- a/crates/ecstore/src/data_usage/mod.rs +++ b/crates/ecstore/src/data_usage/mod.rs @@ -1391,7 +1391,37 @@ impl BucketUsageAccumulator { } pub fn quota_object_size(object: &ObjectInfo) -> Result { - let logical_size = u64::try_from(object.get_actual_size().map_err(Error::other)?).map_err(|_| Error::PartMissingOrCorrupt)?; + // A compressed object may carry -1 while the transformed size is unknown + // (legacy streaming sentinel). In that case the persisted physical size + // is still a valid accounting floor; every other negative value is corrupt. + // An explicit negative `actual-size` metadata value is corrupt, however: + // the sentinel is only valid in the in-memory/object-part field written by + // the legacy streaming path, not as a persisted declared size. + let compressed = object.is_compressed(); + if object.actual_size < -1 || (object.actual_size == -1 && !compressed) { + return Err(Error::PartMissingOrCorrupt); + } + if object + .parts + .iter() + .any(|part| part.actual_size < -1 || (part.actual_size < 0 && !compressed)) + { + return Err(Error::PartMissingOrCorrupt); + } + let declared_actual_size = rustfs_utils::http::get_str(&object.user_defined, rustfs_utils::http::SUFFIX_ACTUAL_SIZE) + .filter(|value| !value.is_empty()); + if declared_actual_size + .as_deref() + .and_then(|value| value.parse::().ok()) + .is_some_and(|size| size < 0) + { + return Err(Error::PartMissingOrCorrupt); + } + let logical_size = match object.get_actual_size().map_err(Error::other)? { + size if size == -1 && compressed && declared_actual_size.is_none() => None, + size if size >= 0 => Some(u64::try_from(size).map_err(|_| Error::PartMissingOrCorrupt)?), + _ => return Err(Error::PartMissingOrCorrupt), + }; let persisted_part_size = if object.parts.is_empty() { u64::try_from(object.size).map_err(|_| Error::PartMissingOrCorrupt)? } else { @@ -1399,12 +1429,8 @@ pub fn quota_object_size(object: &ObjectInfo) -> Result { // Compressed streaming objects persist -1 when the transformed // part size is unknown. The physical part size remains a valid // quota floor; reject only non-negative values that overflow. - let actual_size = if part.actual_size < 0 { - if object.is_compressed() { - 0 - } else { - return Err(Error::PartMissingOrCorrupt); - } + let actual_size = if part.actual_size == -1 { + 0 } else { u64::try_from(part.actual_size).map_err(|_| Error::PartMissingOrCorrupt)? }; @@ -1412,7 +1438,7 @@ pub fn quota_object_size(object: &ObjectInfo) -> Result { total.checked_add(part_size).ok_or(Error::PartMissingOrCorrupt) })? }; - Ok(logical_size.max(persisted_part_size)) + Ok(logical_size.unwrap_or(0).max(persisted_part_size)) } type UsageVersionPage = StorageListObjectVersionsInfo; @@ -3320,6 +3346,80 @@ mod tests { ); } + #[test] + fn quota_object_size_accepts_compressed_unknown_actual_size_sentinel() { + let mut metadata = HashMap::new(); + rustfs_utils::http::insert_str( + &mut metadata, + rustfs_utils::http::SUFFIX_COMPRESSION, + "klauspost/compress/s2".to_string(), + ); + let object = ObjectInfo { + size: 400, + actual_size: -1, + user_defined: Arc::new(metadata), + ..Default::default() + }; + + assert_eq!(quota_object_size(&object).expect("compressed sentinel is valid"), 400); + } + + #[test] + fn quota_object_size_rejects_compressed_part_sum_overflow() { + let mut metadata = HashMap::new(); + rustfs_utils::http::insert_str( + &mut metadata, + rustfs_utils::http::SUFFIX_COMPRESSION, + "klauspost/compress/s2".to_string(), + ); + let object = ObjectInfo { + size: 1, + user_defined: Arc::new(metadata), + parts: Arc::new(vec![ + rustfs_filemeta::ObjectPartInfo { + actual_size: i64::MAX, + ..Default::default() + }, + rustfs_filemeta::ObjectPartInfo { + actual_size: 1, + ..Default::default() + }, + ]), + ..Default::default() + }; + + assert!(matches!(quota_object_size(&object), Err(Error::Io(_)))); + } + + #[test] + fn quota_object_size_rejects_negative_values_other_than_the_compressed_sentinel() { + let mut metadata = HashMap::new(); + rustfs_utils::http::insert_str( + &mut metadata, + rustfs_utils::http::SUFFIX_COMPRESSION, + "klauspost/compress/s2".to_string(), + ); + let corrupt_object = ObjectInfo { + size: 400, + actual_size: -2, + user_defined: Arc::new(metadata.clone()), + ..Default::default() + }; + assert!(matches!(quota_object_size(&corrupt_object), Err(Error::PartMissingOrCorrupt))); + + let corrupt_part = ObjectInfo { + size: 400, + user_defined: Arc::new(metadata), + parts: Arc::new(vec![rustfs_filemeta::ObjectPartInfo { + size: 400, + actual_size: -2, + ..Default::default() + }]), + ..Default::default() + }; + assert!(matches!(quota_object_size(&corrupt_part), Err(Error::PartMissingOrCorrupt))); + } + #[tokio::test] #[serial] async fn live_bucket_usage_refreshes_are_coalesced_only_while_in_flight() { diff --git a/crates/ecstore/src/object_api/types.rs b/crates/ecstore/src/object_api/types.rs index 1bbff7a7f..0194667e9 100644 --- a/crates/ecstore/src/object_api/types.rs +++ b/crates/ecstore/src/object_api/types.rs @@ -689,6 +689,9 @@ impl ObjectInfo { } pub fn get_actual_size(&self) -> std::io::Result { + if self.actual_size < -1 || (self.actual_size == -1 && !self.is_compressed()) { + return Err(std::io::Error::other("invalid negative actual size")); + } if self.actual_size > 0 { return Ok(self.actual_size); } @@ -700,10 +703,25 @@ impl ObjectInfo { let size = size_str.parse::().map_err(|e| std::io::Error::other(e.to_string()))?; return Ok(size); } - let mut actual_size = 0; - self.parts.iter().for_each(|part| { - actual_size += part.actual_size; - }); + if self.actual_size == -1 && self.parts.is_empty() { + return Ok(-1); + } + let mut actual_size = 0_i64; + let mut unknown = false; + for part in self.parts.iter() { + match part.actual_size { + -1 => unknown = true, + size if size >= 0 => { + actual_size = actual_size + .checked_add(size) + .ok_or_else(|| std::io::Error::other("compressed actual size overflow"))?; + } + _ => return Err(std::io::Error::other("invalid negative compressed part size")), + } + } + if unknown { + return Ok(-1); + } if actual_size == 0 && actual_size != self.size { return Err(std::io::Error::other(format!("invalid decompressed size {} {}", actual_size, self.size))); } @@ -718,6 +736,18 @@ impl ObjectInfo { Ok(self.size) } + /// Returns a non-negative size for client and replication boundaries. + /// + /// Compressed legacy metadata can retain the internal `-1` unknown-size + /// sentinel. Those boundaries cannot emit a negative length, so they use + /// the persisted physical size while quota accounting keeps the sentinel + /// distinction in [`crate::data_usage::quota_object_size`]. + pub fn get_actual_size_or_physical(&self) -> i64 { + self.get_actual_size() + .map(|size| if size >= 0 { size } else { self.size.max(0) }) + .unwrap_or_else(|_| self.size.max(0)) + } + pub fn from_file_info(fi: &FileInfo, bucket: &str, object: &str, versioned: bool) -> ObjectInfo { let mut version_id = fi.version_id; diff --git a/crates/ecstore/src/set_disk/mod.rs b/crates/ecstore/src/set_disk/mod.rs index cbf050a4b..e64e1704f 100644 --- a/crates/ecstore/src/set_disk/mod.rs +++ b/crates/ecstore/src/set_disk/mod.rs @@ -97,7 +97,7 @@ use crate::storage_api_contracts::{ CompletePart, ListMultipartsInfo, ListPartsInfo, MultipartInfo, MultipartOperations as _, MultipartUploadResult, PartInfo, }, namespace::NamespaceLocking as _, - object::{DeletedObject, HTTPPreconditions, ObjectIO as _, ObjectOperations as _, ObjectToDelete}, + object::{DeleteAccounting, DeletedObject, HTTPPreconditions, ObjectIO as _, ObjectOperations as _, ObjectToDelete}, range::HTTPRangeSpec, }; use crate::store::utils::is_reserved_or_invalid_bucket; diff --git a/crates/ecstore/src/set_disk/ops/object.rs b/crates/ecstore/src/set_disk/ops/object.rs index cb8f70766..faacca3a8 100644 --- a/crates/ecstore/src/set_disk/ops/object.rs +++ b/crates/ecstore/src/set_disk/ops/object.rs @@ -45,6 +45,7 @@ use crate::bucket::replication::{ DeleteReplicationConfigSnapshot, ReplicationLifecycleBridge, ReplicationStatusType, VersionPurgeStatusType, replication_state_to_filemeta, replication_status_from_filemeta, version_purge_status_to_filemeta, }; +use crate::data_usage::quota_object_size; use crate::diagnostics::get::GetObjectFailureReason; use crate::disk::{DataDirDeleteStatus, OldCurrentSize}; use crate::error::is_err_invalid_upload_id; @@ -5655,7 +5656,18 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { objects: Vec, opts: ObjectOptions, ) -> (Vec, Vec>) { + let (deleted, errors, _) = self.delete_objects_with_accounting(bucket, objects, opts).await; + (deleted, errors) + } + + async fn delete_objects_with_accounting( + &self, + bucket: &str, + objects: Vec, + opts: ObjectOptions, + ) -> (Vec, Vec>, Vec>) { let mut del_objects = vec![DeletedObject::default(); objects.len()]; + let mut accounting = vec![None; objects.len()]; let delete_config_snapshot = opts .delete_replication_config_snapshot .clone() @@ -5745,7 +5757,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { *item = Some(Error::other(message.clone())); } } - return (del_objects, del_errs); + return (del_objects, del_errs, accounting); } }, } @@ -5792,6 +5804,22 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { let source_missing = gerr .as_ref() .is_some_and(|err| is_err_object_not_found(err) || is_err_version_not_found(err)); + // Resolve accounting from the generation selected under this + // object's write lock. A request-layer pre-stat is only an + // optimization and cannot identify a concurrent overwrite. + let (accounting_size, accounting_version_id, removed_current_object) = if source_missing + || dobj.synthetic_version_id + || set_disk_delete_creates_delete_marker(&check_opts) + || goi.delete_marker + { + (None, None, false) + } else { + ( + quota_object_size(&goi).ok(), + goi.version_id.filter(|version_id| !version_id.is_nil()), + (dobj.version_id.is_none() || is_explicit_null_version(dobj.version_id)) && !dobj.synthetic_version_id, + ) + }; // Normalize both sides before comparing. `goi.version_id` is the // client-facing identity, where `from_file_info` synthesizes // `Some(Uuid::nil())` for a null version on a versioned or @@ -5920,7 +5948,12 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { }, replication_state: vr.replication_state_internal.clone(), ..Default::default() - } + }; + accounting[i] = Some(DeleteAccounting { + size: accounting_size, + version_id: accounting_version_id, + removed_current_object, + }); } // Only add to vers_map if we hold the lock @@ -5966,7 +5999,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { }); } } - return (del_objects, del_errs); + return (del_objects, del_errs, accounting); } let mut persisted_journal_entries = Vec::with_capacity(journal_entries.len()); @@ -6204,7 +6237,16 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { } } - (del_objects, del_errs) + // An accounting identity is actionable only when the delete result is + // successful. Never let a failed commit (including a partial quorum + // failure) reach the request-layer fast delta path. + for (index, err) in del_errs.iter().enumerate() { + if err.is_some() { + accounting[index] = None; + } + } + + (del_objects, del_errs, accounting) } #[tracing::instrument(skip(self))] @@ -6533,6 +6575,12 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { let mut obj_info = ObjectInfo::from_file_info(&dfi, bucket, object, opts.versioned || opts.version_suspended); obj_info.size = goi.size; + // Keep the committed source metadata on the internal delete result so + // the request layer can derive canonical accounting for this exact + // generation. Delete responses do not expose these fields. + obj_info.actual_size = goi.actual_size; + obj_info.user_defined = Arc::clone(&goi.user_defined); + obj_info.parts = Arc::clone(&goi.parts); obj_info.user_tags = Arc::clone(&goi.user_tags); self.invalidate_get_object_metadata_cache(bucket, object).await; Ok(obj_info) @@ -7824,6 +7872,113 @@ mod replication_quota_safety_tests { assert_eq!(stored.get_actual_size().expect("stored logical size should parse"), 1); } + #[tokio::test] + async fn delete_returns_canonical_compressed_accounting_size() { + let (_temp_dirs, disks, set_disks) = hermetic_set_disks(4).await; + let bucket = "compressed-delete-accounting"; + for disk in &disks { + disk.make_volume(bucket).await.expect("bucket volume should be created"); + } + + let mut user_defined = HashMap::new(); + insert_str( + &mut user_defined, + rustfs_utils::http::SUFFIX_COMPRESSION, + "klauspost/compress/s2".to_string(), + ); + insert_str(&mut user_defined, SUFFIX_ACTUAL_SIZE, "1000".to_string()); + let mut reader = PutObjReader::new( + HashReader::from_stream(Cursor::new(vec![0x5a; 400]), 400, 1000, None, None, false) + .expect("compressed fixture reader should be valid"), + ); + set_disks + .put_object( + bucket, + "object", + &mut reader, + &ObjectOptions { + user_defined, + ..Default::default() + }, + ) + .await + .expect("compressed object should be written"); + + let (deleted, errors, accounting) = set_disks + .delete_objects_with_accounting( + bucket, + vec![ObjectToDelete { + object_name: "object".to_string(), + ..Default::default() + }], + ObjectOptions { + object_lock_config_snapshot: Some(Arc::new(ObjectLockConfigSnapshot::new( + ObjectLockConfigState::ConfirmedAbsent, + ))), + ..Default::default() + }, + ) + .await; + + assert!(errors[0].is_none(), "compressed delete should succeed: {:?}", errors[0]); + assert!(deleted[0].found, "the committed object must be reported as found"); + assert_eq!(accounting[0].as_ref().and_then(|value| value.size), Some(1000)); + assert!(accounting[0].as_ref().is_some_and(|value| value.version_id.is_none())); + assert!(accounting[0].as_ref().is_some_and(|value| value.removed_current_object)); + } + + #[tokio::test] + async fn suspended_delete_marker_does_not_return_body_accounting() { + let (_temp_dirs, disks, set_disks) = hermetic_set_disks(4).await; + let bucket = "suspended-delete-accounting"; + for disk in &disks { + disk.make_volume(bucket).await.expect("bucket volume should be created"); + } + + let mut user_defined = HashMap::new(); + insert_str( + &mut user_defined, + rustfs_utils::http::SUFFIX_COMPRESSION, + "klauspost/compress/s2".to_string(), + ); + insert_str(&mut user_defined, SUFFIX_ACTUAL_SIZE, "1000".to_string()); + let mut reader = PutObjReader::new( + HashReader::from_stream(Cursor::new(vec![0x5a; 400]), 400, 1000, None, None, false) + .expect("compressed fixture reader should be valid"), + ); + let suspended_opts = ObjectOptions { + version_suspended: true, + delete_replication_config_snapshot: Some(Arc::new(DeleteReplicationConfigSnapshot::from_configs_for_test( + s3s::dto::VersioningConfiguration { + status: Some(s3s::dto::BucketVersioningStatus::from_static(s3s::dto::BucketVersioningStatus::SUSPENDED)), + ..Default::default() + }, + None, + ))), + user_defined, + object_lock_config_snapshot: Some(Arc::new(ObjectLockConfigSnapshot::new(ObjectLockConfigState::ConfirmedAbsent))), + ..Default::default() + }; + set_disks + .put_object(bucket, "object", &mut reader, &suspended_opts) + .await + .expect("compressed object should be written"); + + let (deleted, errors, accounting) = set_disks + .delete_objects_with_accounting( + bucket, + vec![ObjectToDelete { + object_name: "object".to_string(), + ..Default::default() + }], + suspended_opts, + ) + .await; + assert!(errors[0].is_none(), "suspended delete should create a marker: {:?}", errors[0]); + assert!(deleted[0].delete_marker); + assert!(accounting[0].is_none(), "a delete marker must not carry body accounting"); + } + #[tokio::test] async fn direct_put_cannot_persist_a_tiny_logical_size() { let (_temp_dirs, disks, set_disks) = hermetic_set_disks(4).await; diff --git a/crates/ecstore/src/storage_api_contracts/mod.rs b/crates/ecstore/src/storage_api_contracts/mod.rs index 11e7b800f..78a3f9c2d 100644 --- a/crates/ecstore/src/storage_api_contracts/mod.rs +++ b/crates/ecstore/src/storage_api_contracts/mod.rs @@ -62,8 +62,8 @@ pub(crate) mod object { use super::{Debug, Error, FileInfo, GetObjectReader, ObjectInfo, ObjectOptions, PutObjReader}; use crate::storage_api_contracts::range::HTTPRangeSpec; pub(crate) use rustfs_storage_api::{ - DeletedObject, HTTPPreconditions, ObjectIO, ObjectLockDeleteOptions, ObjectLockRetentionOptions, ObjectOperations, - ObjectPreconditionError, ObjectPreconditionPart, ObjectPreconditionState, ObjectToDelete, + DeleteAccounting, DeletedObject, HTTPPreconditions, ObjectIO, ObjectLockDeleteOptions, ObjectLockRetentionOptions, + ObjectOperations, ObjectPreconditionError, ObjectPreconditionPart, ObjectPreconditionState, ObjectToDelete, }; pub(crate) trait EcstoreObjectIO: diff --git a/crates/ecstore/src/store/object.rs b/crates/ecstore/src/store/object.rs index 312a25662..e5f2f0465 100644 --- a/crates/ecstore/src/store/object.rs +++ b/crates/ecstore/src/store/object.rs @@ -41,7 +41,7 @@ use crate::set_disk::{ }; use crate::storage_api_contracts::{ namespace::NamespaceLocking as _, - object::{ObjectIO as _, ObjectOperations as _}, + object::{DeleteAccounting, ObjectIO as _, ObjectOperations as _}, }; use parking_lot::Mutex as ParkingMutex; use rustfs_io_metrics::{ @@ -1216,6 +1216,14 @@ fn return_batch_delete_lock_error(objects: &[ObjectToDelete], err: Error) -> (Ve (del_objects, del_errs) } +fn return_batch_delete_lock_error_with_accounting( + objects: &[ObjectToDelete], + err: Error, +) -> (Vec, Vec>, Vec>) { + let (deleted, errors) = return_batch_delete_lock_error(objects, err); + (deleted, errors, vec![None; objects.len()]) +} + fn sorted_unique_delete_object_names(objects: &[ObjectToDelete]) -> Vec<&str> { let mut object_names: Vec<&str> = objects.iter().map(|object| object.object_name.as_str()).collect(); object_names.sort_unstable(); @@ -2312,6 +2320,22 @@ impl ECStore { result } + pub async fn delete_objects_with_tier_delete_journal_and_accounting( + self: &Arc, + bucket: &str, + objects: Vec, + opts: ObjectOptions, + ) -> (Vec, Vec>, Vec>) { + let result = self + .handle_delete_objects_with_journal_and_accounting(bucket, objects, opts, Some(Arc::clone(self))) + .await; + let success_count = result.1.iter().filter(|err| err.is_none()).count(); + if success_count > 0 { + list_objects::observe_list_objects_mutations(self, bucket, success_count).await; + } + result + } + #[instrument(skip(self))] pub(super) async fn handle_delete_object(&self, bucket: &str, object: &str, opts: ObjectOptions) -> Result { self.handle_delete_object_with_journal(bucket, object, opts, None).await @@ -2689,6 +2713,19 @@ impl ECStore { opts: ObjectOptions, tier_journal_api: Option>, ) -> (Vec, Vec>) { + let (deleted, errors, _) = self + .handle_delete_objects_with_journal_and_accounting(bucket, objects, opts, tier_journal_api) + .await; + (deleted, errors) + } + + pub(super) async fn handle_delete_objects_with_journal_and_accounting( + &self, + bucket: &str, + objects: Vec, + opts: ObjectOptions, + tier_journal_api: Option>, + ) -> (Vec, Vec>, Vec>) { // encode object name let objects: Vec = objects .iter() @@ -2701,6 +2738,7 @@ impl ECStore { // Default return value let mut del_objects = vec![DeletedObject::default(); objects.len()]; + let mut accounting = vec![None; objects.len()]; let mut del_errs = Vec::with_capacity(objects.len()); for _ in 0..objects.len() { @@ -2714,7 +2752,7 @@ impl ECStore { } else { match self.acquire_bucket_lifecycle_read_lock(bucket).await { Ok(guard) => Some(guard), - Err(err) => return return_batch_delete_lock_error(objects.as_slice(), err), + Err(err) => return return_batch_delete_lock_error_with_accounting(objects.as_slice(), err), } }; if let Some(guard) = _bucket_lifecycle_guard.as_ref() { @@ -2726,21 +2764,21 @@ impl ECStore { Err(err) => { let message = err.to_string(); let errors = (0..objects.len()).map(|_| Some(Error::other(message.clone()))).collect(); - return (del_objects, errors); + return (del_objects, errors, accounting); } } } if !is_meta_bucketname(bucket) && let Err(err) = get_cached_bucket_incarnation_id_in(&self.ctx, bucket).await { - return return_batch_delete_lock_error(objects.as_slice(), err); + return return_batch_delete_lock_error_with_accounting(objects.as_slice(), err); } let _object_lock_metadata_guard = if is_meta_bucketname(bucket) { None } else { Some(match acquire_bucket_metadata_transaction_read_lock_in(&self.ctx, bucket).await { Ok(guard) => guard, - Err(err) => return return_batch_delete_lock_error(objects.as_slice(), err), + Err(err) => return return_batch_delete_lock_error_with_accounting(objects.as_slice(), err), }) }; if let Some(guard) = _object_lock_metadata_guard.as_ref() { @@ -2750,7 +2788,7 @@ impl ECStore { let (state, incarnation_id, config_revision) = match get_object_lock_config_and_incarnation_from_disk_in(&self.ctx, bucket).await { Ok(snapshot) => snapshot, - Err(err) => return return_batch_delete_lock_error(objects.as_slice(), err), + Err(err) => return return_batch_delete_lock_error_with_accounting(objects.as_slice(), err), }; opts.object_lock_config_snapshot = Some(Arc::new(ObjectLockConfigSnapshot::for_store_bucket( self.id, @@ -2766,7 +2804,10 @@ impl ECStore { if let (Some(expected), Some(current)) = (opts.expected_bucket_incarnation_id, current_bucket_incarnation_id) && expected != current { - return return_batch_delete_lock_error(objects.as_slice(), StorageError::BucketNotFound(bucket.to_string())); + return return_batch_delete_lock_error_with_accounting( + objects.as_slice(), + StorageError::BucketNotFound(bucket.to_string()), + ); } #[cfg(test)] if current_bucket_incarnation_id.is_some() { @@ -2774,7 +2815,7 @@ impl ECStore { } let _object_lock_guards = match self.acquire_delete_objects_write_locks(bucket, &objects, &mut opts).await { Ok(guards) => guards, - Err(err) => return return_batch_delete_lock_error(objects.as_slice(), err), + Err(err) => return return_batch_delete_lock_error_with_accounting(objects.as_slice(), err), }; let mut futures = Vec::with_capacity(self.pools.len()); @@ -2783,22 +2824,24 @@ impl ECStore { if self.is_pool_rebalancing(pool.pool_idx).await { continue; } - futures.push(pool.delete_objects(bucket, objects.clone(), opts.clone())); + futures.push(pool.delete_objects_with_accounting(bucket, objects.clone(), opts.clone())); } let results = join_all(futures).await; for idx in 0..del_objects.len() { - for (dels, errs) in results.iter() { + for (dels, errs, pool_accounting) in results.iter() { if errs[idx].is_none() && dels[idx].found { del_errs[idx] = None; del_objects[idx] = dels[idx].clone(); + accounting[idx] = pool_accounting[idx].clone(); break; } if del_errs[idx].is_none() { del_errs[idx] = errs[idx].clone(); del_objects[idx] = dels[idx].clone(); + accounting[idx] = pool_accounting[idx].clone(); } } } @@ -2807,7 +2850,7 @@ impl ECStore { v.object_name = decode_dir_object(&v.object_name); }); - (del_objects, del_errs) + (del_objects, del_errs, accounting) // let mut futures = Vec::with_capacity(objects.len()); diff --git a/crates/storage-api/src/lib.rs b/crates/storage-api/src/lib.rs index e011327f9..1114349e1 100644 --- a/crates/storage-api/src/lib.rs +++ b/crates/storage-api/src/lib.rs @@ -76,6 +76,7 @@ pub use bucket::{BucketInfo, BucketOperations, BucketOptions, DeleteBucketOption pub use capability::{CapabilitySnapshotError, CapabilityState, CapabilityStatus}; pub use error::{StorageErrorCode, StorageResult}; pub use multipart::{CompletePart, ListMultipartsInfo, ListPartsInfo, MultipartInfo, MultipartUploadResult, PartInfo}; +pub use object::DeleteAccounting; pub use object::ObjectLockDeleteOptions; pub use object::{DeletedObject, ObjectToDelete}; pub use object::{ExpirationOptions, TransitionedObject}; diff --git a/crates/storage-api/src/object.rs b/crates/storage-api/src/object.rs index 9f0957bde..7959354dc 100644 --- a/crates/storage-api/src/object.rs +++ b/crates/storage-api/src/object.rs @@ -218,6 +218,17 @@ pub struct DeletedObject { pub force_delete_generation: Option, } +/// Accounting identity returned by the internal commit-time delete path. +/// +/// This is carried separately from [`DeletedObject`] so adding quota details +/// does not change the source shape of the public S3 delete result contract. +#[derive(Debug, Default, Clone, PartialEq, Eq)] +pub struct DeleteAccounting { + pub size: Option, + pub version_id: Option, + pub removed_current_object: bool, +} + impl DeletedObject { pub fn version_purge_status(&self) -> VersionPurgeStatusType { self.replication_state @@ -341,6 +352,19 @@ pub trait ObjectOperations: Send + Sync + fmt::Debug { objects: Vec, opts: Self::ObjectOptions, ) -> (Vec, Vec>); + /// Delete objects and optionally return commit-time accounting identities. + /// The default preserves the ordinary delete contract for implementations + /// that do not expose storage-level accounting details. + async fn delete_objects_with_accounting( + &self, + bucket: &str, + objects: Vec, + opts: Self::ObjectOptions, + ) -> (Vec, Vec>, Vec>) { + let object_count = objects.len(); + let (deleted, errors) = self.delete_objects(bucket, objects, opts).await; + (deleted, errors, vec![None; object_count]) + } async fn put_object_metadata( &self, bucket: &str, diff --git a/rustfs/src/app/bucket_usecase.rs b/rustfs/src/app/bucket_usecase.rs index 7d55035e9..b51254b23 100644 --- a/rustfs/src/app/bucket_usecase.rs +++ b/rustfs/src/app/bucket_usecase.rs @@ -1021,7 +1021,7 @@ fn build_list_objects_v2_metadata_output( object: Object { key: Some(encode_list_objects_v2_value(&object.name, encoding_type)), last_modified: object.mod_time.map(Timestamp::from), - size: Some(object.get_actual_size().unwrap_or_default()), + size: Some(object.get_actual_size_or_physical()), e_tag: object.etag.clone().map(|etag| to_s3s_etag(&etag)), storage_class: Some(ObjectStorageClass::from( object diff --git a/rustfs/src/app/object_usecase.rs b/rustfs/src/app/object_usecase.rs index a31ec1d19..d1978f2db 100644 --- a/rustfs/src/app/object_usecase.rs +++ b/rustfs/src/app/object_usecase.rs @@ -3969,6 +3969,55 @@ fn delete_creates_delete_marker(opts: &ObjectOptions) -> bool { opts.version_id.is_none() && opts.versioned && !opts.version_suspended } +fn delete_removes_current_object(opts: &ObjectOptions) -> bool { + delete_request_targets_current( + opts.version_id + .as_deref() + .and_then(|version_id| Uuid::parse_str(version_id).ok()), + ) +} + +fn delete_request_targets_current(version_id: Option) -> bool { + version_id.is_none() || version_id.is_some_and(|version_id| version_id.is_nil()) +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum DeleteMemoryUpdate { + DeleteMarker, + Object { size: u64, removed_current_object: bool }, +} + +fn delete_memory_update( + creates_delete_marker: bool, + committed_delete_marker: bool, + requested_current: bool, + accounting_size: Option, + removed_current_object: bool, +) -> Option { + if creates_delete_marker || (committed_delete_marker && requested_current) { + return Some(DeleteMemoryUpdate::DeleteMarker); + } + + (!committed_delete_marker) + .then_some(accounting_size) + .flatten() + .map(|size| DeleteMemoryUpdate::Object { + size, + removed_current_object, + }) +} + +async fn apply_delete_memory_update(bucket: &str, update: Option) { + match update { + Some(DeleteMemoryUpdate::DeleteMarker) => record_bucket_delete_marker_memory(bucket).await, + Some(DeleteMemoryUpdate::Object { + size, + removed_current_object, + }) => record_bucket_object_delete_memory(bucket, size, removed_current_object).await, + None => {} + } +} + /// `DeleteObjects` is idempotent. A raw filesystem `NotFound` can cross the /// distributed delete path instead of its usual typed missing-object error. fn is_delete_objects_not_found(error: &EcstoreError) -> bool { @@ -8409,8 +8458,6 @@ impl DefaultObjectUsecase { object: ObjectToDelete, versioned: bool, version_suspended: bool, - size: i64, - existing: Option, } // Phase 2 (bounded concurrency, backlog#929 / HP-8): collect the @@ -8428,32 +8475,23 @@ impl DefaultObjectUsecase { skip_stat, } = prepared; let synthetic_version_id = object.version_id.is_none() && is_dir_object(&object.object_name); - let (goi, source_missing) = if skip_stat { - (ObjectInfo::default(), false) - } else { + if !skip_stat { match store_ref.get_object_info(bucket_ref, &object.object_name, &opts).await { - Ok(res) => (res, false), - Err(err) if is_err_object_not_found(&err) || is_err_version_not_found(&err) => { - (ObjectInfo::default(), true) - } + Ok(_) => {} + Err(err) if is_err_object_not_found(&err) || is_err_version_not_found(&err) => {} Err(err) => return Err(ApiError::from(err)), } - }; - - let size = goi.size; + } if synthetic_version_id { object.version_id = Some(Uuid::nil()); } - let existing = (!skip_stat && !source_missing).then_some(goi); Ok::<_, ApiError>(AdmittedDelete { idx, object, versioned: opts.versioned, version_suspended: opts.version_suspended, - size, - existing, }) })) .buffered(DELETE_OBJECTS_PRE_STAT_CONCURRENCY) @@ -8464,15 +8502,11 @@ impl DefaultObjectUsecase { // per-key success/failure reporting is unchanged. let mut object_to_delete = Vec::new(); let mut object_to_delete_idx = Vec::new(); - let mut object_sizes = Vec::new(); - let mut existing_object_infos = Vec::new(); let mut object_versioning = Vec::new(); for admitted in admitted_deletes { - object_sizes.push(admitted.size); object_to_delete_idx.push(admitted.idx); object_versioning.push((admitted.versioned, admitted.version_suspended)); object_to_delete.push(admitted.object); - existing_object_infos.push(admitted.existing); } let cache_adapter = self.object_data_cache(); let cache_keys_before_delete = object_to_delete @@ -8489,8 +8523,8 @@ impl DefaultObjectUsecase { ..Default::default() }; apply_bucket_generation_guard(&req, &bucket, &mut storage_delete_opts)?; - let (dobjs, errs) = store - .delete_objects_with_tier_delete_journal(&bucket, object_to_delete.clone(), storage_delete_opts) + let (dobjs, errs, accounting) = store + .delete_objects_with_tier_delete_journal_and_accounting(&bucket, object_to_delete.clone(), storage_delete_opts) .await; let _manager = get_concurrency_manager(); @@ -8515,17 +8549,16 @@ impl DefaultObjectUsecase { delete_results[didx].delete_object = Some(deleted_object.clone()); let (versioned, version_suspended) = object_versioning[i]; let creates_delete_marker = object_to_delete[i].version_id.is_none() && versioned && !version_suspended; - if creates_delete_marker { - record_bucket_delete_marker_memory(&bucket).await; - } else { - let size = object_sizes[i].max(0) as u64; - record_bucket_object_delete_memory( - &bucket, - size, - existing_object_infos[i].is_some() && object_to_delete[i].version_id.is_none(), - ) - .await; - } + let committed_delete_marker = dobjs[i].delete_marker; + let delete_accounting = accounting.get(i).and_then(Option::as_ref); + let update = delete_memory_update( + creates_delete_marker, + committed_delete_marker, + delete_request_targets_current(object_to_delete[i].version_id), + delete_accounting.and_then(|value| value.size), + delete_accounting.is_some_and(|value| value.removed_current_object), + ); + apply_delete_memory_update(&bucket, update).await; } Err(error) => { delete_results[didx].error = Some(error); @@ -8803,12 +8836,24 @@ impl DefaultObjectUsecase { let _ = invalidate_object_data_cache_after_delete_success(&cache_adapter, &bucket, &key).await; } - // Fast in-memory update for immediate quota and admin usage consistency - if delete_creates_delete_marker(&opts) { - record_bucket_delete_marker_memory(&bucket).await; + // Fast in-memory update for immediate quota and admin usage consistency. + // Prefix/force deletes and synthetic directory entries do not carry one + // committed object identity; leave their cache delta to reconciliation. + let update = if force_delete || obj_info.name.is_empty() || synthetic_version_id { + None } else { - record_bucket_object_delete_memory(&bucket, obj_info.size.max(0) as u64, opts.version_id.is_none()).await; - } + // The storage commit returns this object's metadata while its + // generation lock is held. Never fall back to a pre-delete stat: + // an overwrite can commit between that stat and this delete. + delete_memory_update( + delete_creates_delete_marker(&opts), + obj_info.delete_marker, + opts.version_id.is_none(), + quota_object_size(&obj_info).ok(), + delete_removes_current_object(&opts), + ) + }; + apply_delete_memory_update(&bucket, update).await; if obj_info.name.is_empty() { if let Some((operation_id, target_arns, generation)) = force_delete_intent { @@ -17861,6 +17906,158 @@ mod tests { assert!(!can_skip_delete_objects_pre_stat(false, &delete_marker_creating_opts(), false)); } + #[test] + fn delete_accounting_recognizes_explicit_null_as_current_object() { + let opts = ObjectOptions { + version_id: Some(Uuid::nil().to_string()), + version_suspended: true, + ..Default::default() + }; + assert!(delete_removes_current_object(&opts)); + assert!(delete_request_targets_current(Some(Uuid::nil()))); + assert!(!delete_request_targets_current(Some(Uuid::new_v4()))); + assert!(!delete_removes_current_object(&ObjectOptions { + version_id: Some(Uuid::new_v4().to_string()), + ..Default::default() + })); + } + + #[test] + fn compressed_object_delete_restores_usage_baseline() { + let mut metadata = HashMap::new(); + insert_str(&mut metadata, SUFFIX_COMPRESSION, "klauspost/compress/s2".to_string()); + let object = ObjectInfo { + size: 400, + actual_size: 1000, + user_defined: Arc::new(metadata), + ..Default::default() + }; + let accounting_size = quota_object_size(&object).expect("logical compressed size should be canonical"); + + assert_eq!( + delete_memory_update(false, false, true, Some(accounting_size), true), + Some(DeleteMemoryUpdate::Object { + size: 1000, + removed_current_object: true, + }) + ); + } + + #[test] + fn invalid_accounting_metadata_is_reconciled_without_overflow() { + assert_eq!(delete_memory_update(false, false, true, None, true), None); + assert_eq!( + delete_memory_update(false, true, true, None, true), + Some(DeleteMemoryUpdate::DeleteMarker) + ); + } + + #[tokio::test] + #[serial_test::serial] + async fn compressed_delete_requests_restore_usage_baseline() { + use crate::app::storage_api::test::contract::bucket::{BucketOperations as _, DeleteBucketOptions, MakeBucketOptions}; + + let store = crate::app::gating_test_env::shared_gating_ecstore().await; + if current_app_context().is_none() { + crate::app::runtime_sources::install_test_app_context(Arc::clone(&store)).await; + } + let bucket = format!("compressed-delete-request-{}", Uuid::new_v4().simple()); + store + .make_bucket(&bucket, &MakeBucketOptions::default()) + .await + .expect("create compressed delete request bucket"); + + // Seed the process-local usage with the canonical logical bytes. The + // direct storage PUT below intentionally does not apply an app-layer + // usage delta; the two real DELETE requests must remove exactly this + // amount through their request-layer wiring. + crate::app::storage_api::test::data_usage::seed_bucket_usage_memory_for_test(&bucket, 2_000).await; + + for object in ["single", "batch"] { + let mut metadata = HashMap::new(); + insert_str(&mut metadata, SUFFIX_COMPRESSION, "klauspost/compress/s2".to_string()); + insert_str(&mut metadata, SUFFIX_ACTUAL_SIZE, "1000".to_string()); + let reader = HashReader::from_stream(std::io::Cursor::new(vec![0x5a; 400]), 400, 1000, None, None, false) + .expect("compressed fixture reader should be valid"); + let mut reader = PutObjReader::new(reader); + store + .put_object( + &bucket, + object, + &mut reader, + &ObjectOptions { + user_defined: metadata, + ..Default::default() + }, + ) + .await + .expect("compressed fixture object should be written"); + } + + let mut single_req = build_request( + DeleteObjectInput::builder() + .bucket(bucket.clone()) + .key("single".to_string()) + .build() + .expect("single delete input should build"), + Method::DELETE, + ); + single_req.extensions.insert(crate::storage::access::ReqInfo { + cred: Some(rustfs_credentials::Credentials::default()), + is_owner: true, + ..Default::default() + }); + DefaultObjectUsecase::from_global() + .execute_delete_object(single_req) + .await + .expect("single compressed delete should succeed"); + assert_eq!( + crate::app::storage_api::test::data_usage::get_bucket_usage_memory(&bucket).await, + Some(1_000), + "single delete must subtract the logical accounting size" + ); + + let mut batch_req = build_request( + DeleteObjectsInput::builder() + .bucket(bucket.clone()) + .delete(Delete { + objects: vec![ObjectIdentifier { + key: "batch".to_string(), + ..Default::default() + }], + quiet: None, + }) + .build() + .expect("batch delete input should build"), + Method::POST, + ); + batch_req.extensions.insert(crate::storage::access::ReqInfo { + cred: Some(rustfs_credentials::Credentials::default()), + is_owner: true, + ..Default::default() + }); + DefaultObjectUsecase::from_global() + .execute_delete_objects(batch_req) + .await + .expect("batch compressed delete should succeed"); + assert_eq!( + crate::app::storage_api::test::data_usage::get_bucket_usage_memory(&bucket).await, + Some(0), + "batch delete must subtract the committed logical accounting size" + ); + + store + .delete_bucket( + &bucket, + &DeleteBucketOptions { + force: true, + ..Default::default() + }, + ) + .await + .expect("clean up compressed delete request bucket"); + } + #[tokio::test] async fn execute_get_object_attributes_returns_internal_error_when_store_uninitialized() { let input = GetObjectAttributesInput::builder() diff --git a/rustfs/src/app/storage_api.rs b/rustfs/src/app/storage_api.rs index d90cfb6a8..a599cfde6 100644 --- a/rustfs/src/app/storage_api.rs +++ b/rustfs/src/app/storage_api.rs @@ -72,6 +72,11 @@ pub(crate) mod data_usage { compute_bucket_usage, live_bucket_usage_computations, seed_bucket_usage_memory_for_test, store_data_usage_in_backend, }; + #[cfg(test)] + pub(crate) async fn get_bucket_usage_memory(bucket: &str) -> Option { + crate::storage::storage_api::ecstore_data_usage::get_bucket_usage_memory(bucket).await + } + pub(crate) async fn record_bucket_object_delete_memory(bucket: &str, deleted_size: u64, removed_current_object: bool) { crate::storage::storage_api::ecstore_data_usage::record_bucket_object_delete_memory( bucket, @@ -1233,7 +1238,10 @@ pub(crate) mod test { pub(crate) use super::access::ReqInfo; pub(crate) use super::options::VERSIONING_CONFIG_LOOKUPS; - pub(crate) use super::{bucket, data_usage, ecfs, object_utils, runtime}; + pub(crate) use super::{bucket, ecfs, object_utils, runtime}; + pub(crate) mod data_usage { + pub(crate) use super::super::data_usage::*; + } pub(crate) use crate::storage::storage_api::test_consumer::{get_global_bucket_metadata_sys, set_bucket_metadata}; pub(crate) use crate::storage::storage_api::{ ECStore, Endpoint, Endpoints, PoolEndpoints, StorageObjectInfo, StorageObjectOptions, StoragePutObjReader, diff --git a/rustfs/src/storage/s3_api/bucket.rs b/rustfs/src/storage/s3_api/bucket.rs index e105adba4..6599ca107 100644 --- a/rustfs/src/storage/s3_api/bucket.rs +++ b/rustfs/src/storage/s3_api/bucket.rs @@ -296,7 +296,10 @@ pub(crate) fn build_list_objects_v2_output( let mut obj = Object { key: Some(key), last_modified: v.mod_time.map(Timestamp::from), - size: Some(v.get_actual_size().unwrap_or_default()), + // Compressed legacy objects may retain an unknown (-1) + // logical-size sentinel; never expose that internal value in + // an S3 response. + size: Some(v.get_actual_size_or_physical()), e_tag: v.etag.clone().map(|etag| to_s3s_etag(&etag)), storage_class: v.storage_class.clone().map(ObjectStorageClass::from), ..Default::default() @@ -656,6 +659,45 @@ mod tests { assert_eq!(output.common_prefixes.as_ref().map(std::vec::Vec::len), Some(2)); } + #[test] + fn list_objects_never_exposes_compressed_unknown_size_sentinel() { + let mut metadata = std::collections::HashMap::new(); + rustfs_utils::http::insert_str( + &mut metadata, + rustfs_utils::http::SUFFIX_COMPRESSION, + "klauspost/compress/s2".to_string(), + ); + let output = build_list_objects_v2_output( + ListObjectsV2Info { + objects: vec![ObjectInfo { + name: "legacy-compressed".to_string(), + size: 128, + actual_size: -1, + user_defined: std::sync::Arc::new(metadata), + ..Default::default() + }], + ..Default::default() + }, + false, + 1000, + "bucket".to_string(), + String::new(), + None, + None, + None, + None, + ); + + assert_eq!( + output + .contents + .as_ref() + .and_then(|objects| objects.first()) + .and_then(|object| object.size), + Some(128) + ); + } + #[test] fn list_responses_report_standard_for_legacy_label_only_file_metadata() { let version_id = Uuid::parse_str("11111111-2222-3333-4444-555555555555").expect("fixture version ID should be valid"); diff --git a/rustfs/src/storage/storage_api.rs b/rustfs/src/storage/storage_api.rs index 31668783a..e5277c0dd 100644 --- a/rustfs/src/storage/storage_api.rs +++ b/rustfs/src/storage/storage_api.rs @@ -429,6 +429,8 @@ pub(crate) mod ecstore_config { } pub(crate) mod ecstore_data_usage { + #[cfg(test)] + pub(crate) use rustfs_ecstore::api::data_usage::get_bucket_usage_memory; pub(crate) use rustfs_ecstore::api::data_usage::{ apply_bucket_usage_memory_overlay, init_compression_total_memory_from_backend, load_admin_data_usage_from_backend_cached, load_data_usage_from_backend, quota_object_size, record_bucket_delete_marker_memory, record_bucket_object_delete_memory,