diff --git a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs index 99896ed36..49708e34c 100644 --- a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs +++ b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs @@ -1868,11 +1868,11 @@ pub async fn put_restore_opts( ..Default::default() }); } - for (k, v) in &oi.user_defined { + for (k, v) in oi.user_defined.iter() { meta.insert(k.to_string(), v.clone()); } if !oi.user_tags.is_empty() { - meta.insert(AMZ_OBJECT_TAGGING.to_string(), oi.user_tags.clone()); + meta.insert(AMZ_OBJECT_TAGGING.to_string(), (*oi.user_tags).clone()); } let restore_expiry = lifecycle::expected_expiry_time(OffsetDateTime::now_utc(), rreq.days.unwrap_or(1)); meta.insert( @@ -1903,7 +1903,7 @@ impl LifecycleOps for ObjectInfo { fn to_lifecycle_opts(&self) -> lifecycle::ObjectOpts { lifecycle::ObjectOpts { name: self.name.clone(), - user_tags: self.user_tags.clone(), + user_tags: (*self.user_tags).clone(), version_id: self.version_id, mod_time: self.mod_time, size: self.size as usize, diff --git a/crates/ecstore/src/bucket/lifecycle/core.rs b/crates/ecstore/src/bucket/lifecycle/core.rs index 9ebda1f1b..cb4962b86 100644 --- a/crates/ecstore/src/bucket/lifecycle/core.rs +++ b/crates/ecstore/src/bucket/lifecycle/core.rs @@ -881,7 +881,7 @@ impl ObjectOpts { pub fn from_object_info(oi: &ObjectInfo) -> Self { Self { name: oi.name.clone(), - user_tags: oi.user_tags.clone(), + user_tags: (*oi.user_tags).clone(), mod_time: oi.mod_time, size: oi.size as usize, version_id: oi.version_id, @@ -894,7 +894,7 @@ impl ObjectOpts { restore_expires: oi.restore_expires, versioned: false, version_suspended: false, - user_defined: oi.user_defined.clone(), + user_defined: (*oi.user_defined).clone(), version_purge_status: oi.version_purge_status.clone(), replication_status: oi.replication_status.clone(), } diff --git a/crates/ecstore/src/bucket/replication/replication_pool.rs b/crates/ecstore/src/bucket/replication/replication_pool.rs index 6b0ab32a3..903c1c24f 100644 --- a/crates/ecstore/src/bucket/replication/replication_pool.rs +++ b/crates/ecstore/src/bucket/replication/replication_pool.rs @@ -1079,7 +1079,7 @@ pub async fn schedule_replication(oi: ObjectInfo, o: Arc, dsc: target_statuses: tgt_statuses, target_purge_statuses: purge_statuses, replication_timestamp: tm, - user_tags: oi.user_tags, + user_tags: (*oi.user_tags).clone(), checksum: None, retry_count: 0, event_type: "".to_string(), diff --git a/crates/ecstore/src/bucket/replication/replication_resyncer.rs b/crates/ecstore/src/bucket/replication/replication_resyncer.rs index 305dcacbc..6a8ca90f3 100644 --- a/crates/ecstore/src/bucket/replication/replication_resyncer.rs +++ b/crates/ecstore/src/bucket/replication/replication_resyncer.rs @@ -934,7 +934,7 @@ fn heal_should_use_check_replicate_delete(oi: &ObjectInfo) -> bool { pub async fn get_heal_replicate_object_info(oi: &ObjectInfo, rcfg: &ReplicationConfig) -> ReplicateObjectInfo { let mut oi = oi.clone(); - let mut user_defined = oi.user_defined.clone(); + let mut user_defined = (*oi.user_defined).clone(); if let Some(rc) = rcfg.config.as_ref() && !rc.role.is_empty() @@ -982,7 +982,7 @@ pub async fn get_heal_replicate_object_info(oi: &ObjectInfo, rcfg: &ReplicationC &oi.name, MustReplicateOptions::new( &user_defined, - oi.user_tags.clone(), + (*oi.user_tags).clone(), ReplicationStatusType::Empty, ReplicationType::Heal, ObjectOptions::default(), @@ -1020,7 +1020,7 @@ pub async fn get_heal_replicate_object_info(oi: &ObjectInfo, rcfg: &ReplicationC target_purge_statuses, replication_timestamp: None, ssec: false, // TODO: add ssec support - user_tags: oi.user_tags.clone(), + user_tags: (*oi.user_tags).clone(), checksum: oi.checksum.clone(), retry_count: 0, } @@ -1158,7 +1158,7 @@ impl ReplicationConfig { return self.resync_internal(oi, dsc, status); } - let mut user_defined = oi.user_defined.clone(); + let mut user_defined = (*oi.user_defined).clone(); user_defined.remove(AMZ_BUCKET_REPLICATION_STATUS); let dsc = must_replicate( @@ -1166,7 +1166,7 @@ impl ReplicationConfig { &oi.name, MustReplicateOptions::new( &user_defined, - oi.user_tags.clone(), + (*oi.user_tags).clone(), ReplicationStatusType::Empty, ReplicationType::ExistingObject, ObjectOptions::default(), @@ -1296,7 +1296,7 @@ impl MustReplicateOptions { } pub fn from_object_info(oi: &ObjectInfo, op_type: ReplicationType, opts: ObjectOptions) -> Self { - Self::new(&oi.user_defined, oi.user_tags.clone(), oi.replication_status.clone(), op_type, opts) + Self::new(&oi.user_defined, (*oi.user_tags).clone(), oi.replication_status.clone(), op_type, opts) } pub fn replication_status(&self) -> ReplicationStatusType { @@ -1411,7 +1411,7 @@ fn delete_replication_object_opts(dobj: &ObjectToDelete, oi: &ObjectInfo) -> Obj ObjectOpts { name: dobj.object_name.clone(), ssec: is_ssec_encrypted(&oi.user_defined), - user_tags: oi.user_tags.clone(), + user_tags: (*oi.user_tags).clone(), delete_marker: oi.delete_marker, version_id: dobj.version_id, op_type: ReplicationType::Delete, @@ -2995,7 +2995,7 @@ impl ReplicateObjectInfoExt for ReplicateObjectInfo { mod_time: self.mod_time, version_id: self.version_id, size: self.size, - user_tags: self.user_tags.clone(), + user_tags: Arc::new(self.user_tags.clone()), actual_size: self.actual_size, replication_status_internal: self.replication_status_internal.clone(), replication_status: self.replication_status.clone(), @@ -3221,7 +3221,7 @@ fn put_replication_opts(sc: &str, object_info: &ObjectInfo) -> Result<(PutObject } // Use case-insensitive lookup for headers - let lk_map = object_info.user_defined.clone(); + let lk_map = &*object_info.user_defined; if let Some(lang) = lk_map.lookup(CONTENT_LANGUAGE) { put_op.content_language = lang.to_string(); @@ -3500,7 +3500,7 @@ fn get_replication_action(oi1: &ObjectInfo, oi2: &HeadObjectOutput, op_type: Rep // compare metadata on both maps to see if meta is identical let mut compare_meta1 = HashMap::new(); - for (k, v) in &oi1.user_defined { + for (k, v) in oi1.user_defined.iter() { let mut found = false; for prefix in &compare_keys { if strings_has_prefix_fold(k, prefix) { diff --git a/crates/ecstore/src/config/com.rs b/crates/ecstore/src/config/com.rs index 9a9c0cb86..a30cf9394 100644 --- a/crates/ecstore/src/config/com.rs +++ b/crates/ecstore/src/config/com.rs @@ -1355,7 +1355,7 @@ mod tests { size: self.data.len() as i64, actual_size: self.data.len() as i64, is_dir: false, - user_defined: HashMap::new(), + user_defined: Arc::new(HashMap::new()), parity_blocks: 0, data_blocks: 0, version_id: None, @@ -1363,8 +1363,8 @@ mod tests { transitioned_object: Default::default(), restore_ongoing: false, restore_expires: None, - user_tags: String::new(), - parts: Vec::new(), + user_tags: Arc::new(String::new()), + parts: Arc::new(Vec::new()), is_latest: true, content_type: Some("application/json".to_string()), content_encoding: None, diff --git a/crates/ecstore/src/data_movement.rs b/crates/ecstore/src/data_movement.rs index f40840624..19342dbbb 100644 --- a/crates/ecstore/src/data_movement.rs +++ b/crates/ecstore/src/data_movement.rs @@ -94,7 +94,7 @@ fn data_movement_new_multipart_opts(object_info: &ObjectInfo, src_pool_idx: usiz ObjectOptions { versioned: object_info.version_id.is_some(), version_id: object_info.version_id.as_ref().map(|v| v.to_string()), - user_defined: object_info.user_defined.clone(), + user_defined: (*object_info.user_defined).clone(), preserve_etag: object_info.etag.clone(), src_pool_idx, data_movement: true, @@ -120,7 +120,7 @@ fn data_movement_put_object_opts(object_info: &ObjectInfo, src_pool_idx: usize) data_movement: true, version_id: object_info.version_id.as_ref().map(|v| v.to_string()), mod_time: object_info.mod_time, - user_defined: object_info.user_defined.clone(), + user_defined: (*object_info.user_defined).clone(), preserve_etag: object_info.etag.clone(), ..Default::default() } @@ -341,7 +341,7 @@ mod tests { let object_info = ObjectInfo { version_id: Some(version_id), etag: Some("etag-value".to_string()), - user_defined: std::collections::HashMap::from([("x-amz-meta-key".to_string(), "value".to_string())]), + user_defined: Arc::new(std::collections::HashMap::from([("x-amz-meta-key".to_string(), "value".to_string())])), ..Default::default() }; @@ -382,7 +382,7 @@ mod tests { version_id: Some(version_id), mod_time: Some(OffsetDateTime::UNIX_EPOCH), etag: Some("etag-value".to_string()), - user_defined: std::collections::HashMap::from([("x-amz-meta-key".to_string(), "value".to_string())]), + user_defined: Arc::new(std::collections::HashMap::from([("x-amz-meta-key".to_string(), "value".to_string())])), ..Default::default() }; diff --git a/crates/ecstore/src/rpc/peer_s3_client.rs b/crates/ecstore/src/rpc/peer_s3_client.rs index 3355de7de..3914b99cc 100644 --- a/crates/ecstore/src/rpc/peer_s3_client.rs +++ b/crates/ecstore/src/rpc/peer_s3_client.rs @@ -843,12 +843,13 @@ impl PeerS3Client for RemotePeerS3Client { }); let response = client.make_bucket(request).await?.into_inner(); - // TODO: deal with error if !response.success { return if let Some(err) = response.error { Err(err.into()) } else { - Err(Error::other("")) + Err(Error::other(format!( + "make_bucket({bucket}): peer returned failure without error details" + ))) }; } diff --git a/crates/ecstore/src/set_disk.rs b/crates/ecstore/src/set_disk.rs index 898aa8a92..5ef636323 100644 --- a/crates/ecstore/src/set_disk.rs +++ b/crates/ecstore/src/set_disk.rs @@ -1496,7 +1496,7 @@ impl ObjectOperations for SetDisks { } }; - fi.metadata = src_info.user_defined.clone(); + fi.metadata = (*src_info.user_defined).clone(); if let Some(etag) = &src_info.etag { fi.metadata.insert("etag".to_owned(), etag.clone()); @@ -1512,7 +1512,7 @@ impl ObjectOperations for SetDisks { for fi in metas.iter_mut() { if fi.is_valid() { - fi.metadata = src_info.user_defined.clone(); + fi.metadata = (*src_info.user_defined).clone(); if let Some(etag) = &src_info.etag { fi.metadata.insert("etag".to_owned(), etag.clone()); } @@ -2073,8 +2073,8 @@ impl ObjectOperations for SetDisks { check_object_lock_retention_update(bucket, object, &obj_info, opts)?; - for (k, v) in obj_info.user_defined { - fi.metadata.insert(k, v); + for (k, v) in obj_info.user_defined.iter() { + fi.metadata.insert(k.clone(), v.clone()); } if let Some(mt) = &opts.eval_metadata { @@ -2100,7 +2100,7 @@ impl ObjectOperations for SetDisks { #[tracing::instrument(skip(self))] async fn get_object_tags(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result { let oi = self.get_object_info(bucket, object, opts).await?; - Ok(oi.user_tags) + Ok((*oi.user_tags).clone()) } #[tracing::instrument(level = "debug", skip(self))] @@ -2172,7 +2172,7 @@ impl ObjectOperations for SetDisks { let dest_obj = dest_obj.unwrap(); let oi = ObjectInfo::from_file_info(&fi, bucket, object, opts.versioned || opts.version_suspended); - let mut transition_meta = oi.user_defined.clone(); + let mut transition_meta = (*oi.user_defined).clone(); transition_meta.insert("name".to_string(), object.to_string()); if let Some(content_type) = oi.content_type.as_ref().filter(|value| !value.is_empty()) { @@ -2340,9 +2340,9 @@ impl ObjectOperations for SetDisks { //} let mut uploaded_parts: Vec = vec![]; - let parts = oi.parts.clone(); + let parts = Arc::clone(&oi.parts); let mut part_offset: i64 = 0; - for part_info in &parts { + for part_info in parts.iter() { let mut part_opts = opts.clone(); part_opts.part_number = Some(part_info.number); if part_info.actual_size <= 0 { @@ -6026,7 +6026,7 @@ mod tests { ); let obj_info = ObjectInfo { - user_defined, + user_defined: Arc::new(user_defined), ..Default::default() }; let opts = ObjectOptions { @@ -6061,7 +6061,7 @@ mod tests { ); let obj_info = ObjectInfo { - user_defined, + user_defined: Arc::new(user_defined), ..Default::default() }; let opts = ObjectOptions { diff --git a/crates/ecstore/src/set_disk/replication.rs b/crates/ecstore/src/set_disk/replication.rs index f71c99885..a35e281da 100644 --- a/crates/ecstore/src/set_disk/replication.rs +++ b/crates/ecstore/src/set_disk/replication.rs @@ -25,7 +25,7 @@ impl SetDisks { let mut oi = obj_info.clone(); oi.metadata_only = true; - oi.user_defined.remove(X_AMZ_RESTORE.as_str()); + Arc::make_mut(&mut oi.user_defined).remove(X_AMZ_RESTORE.as_str()); let version_id = oi.version_id.map(|v| v.to_string()); let _obj = self diff --git a/crates/ecstore/src/set_disk/write.rs b/crates/ecstore/src/set_disk/write.rs index db34f6042..d8e7231de 100644 --- a/crates/ecstore/src/set_disk/write.rs +++ b/crates/ecstore/src/set_disk/write.rs @@ -383,30 +383,35 @@ impl SetDisks { } if let Some(err) = reduce_write_quorum_errs(&errs, OBJECT_OP_IGNORED_ERRS, write_quorum) { - // TODO: add concurrency + let mut revert_futures = Vec::with_capacity(disks.len()); for (i, err) in errs.iter().enumerate() { if err.is_some() { continue; } if let Some(disk) = disks[i].as_ref() { - let _ = disk - .delete( - bucket, - &path_join_buf(&[prefix, STORAGE_FORMAT_FILE]), - DeleteOptions { - recursive: true, - ..Default::default() - }, - ) - .await - .map_err(|e| { - warn!("write meta revert err {:?}", e); - e - }); + let disk = disk.clone(); + let bucket = bucket.to_string(); + let path = path_join_buf(&[prefix, STORAGE_FORMAT_FILE]); + revert_futures.push(async move { + if let Err(err) = disk + .delete( + &bucket, + &path, + DeleteOptions { + recursive: true, + ..Default::default() + }, + ) + .await + { + warn!("write meta revert err {:?}", err); + } + }); } } + join_all(revert_futures).await; return Err(err); } Ok(()) diff --git a/crates/ecstore/src/store/bucket.rs b/crates/ecstore/src/store/bucket.rs index 484e6b795..13bb0b7b6 100644 --- a/crates/ecstore/src/store/bucket.rs +++ b/crates/ecstore/src/store/bucket.rs @@ -14,6 +14,7 @@ use super::*; use crate::bucket::utils::is_meta_bucketname; +use crate::set_disk::get_lock_acquire_timeout; fn should_override_created_from_metadata(created: OffsetDateTime) -> bool { created != OffsetDateTime::UNIX_EPOCH @@ -28,7 +29,28 @@ impl ECStore { return Err(StorageError::BucketNameInvalid(err.to_string())); } - // TODO: nslock + let _ns_guard = if !opts.no_lock { + let ns_lock = self.new_ns_lock(bucket, bucket).await?; + Some( + ns_lock + .get_write_lock(get_lock_acquire_timeout()) + .await + .map_err(|e| match e { + rustfs_lock::error::LockError::QuorumNotReached { required, achieved } => { + StorageError::NamespaceLockQuorumUnavailable { + mode: "write", + bucket: bucket.to_string(), + object: bucket.to_string(), + required, + achieved, + } + } + other => StorageError::other(format!("make_bucket: failed to acquire write lock on {bucket}: {other}")), + })?, + ) + } else { + None + }; if let Err(err) = self.peer_sys.make_bucket(bucket, opts).await { let err = to_object_err(err.into(), vec![bucket]); @@ -111,7 +133,28 @@ impl ECStore { return Err(StorageError::BucketNameInvalid(err.to_string())); } - // TODO: nslock + let _ns_guard = if !opts.no_lock { + let ns_lock = self.new_ns_lock(bucket, bucket).await?; + Some( + ns_lock + .get_write_lock(get_lock_acquire_timeout()) + .await + .map_err(|e| match e { + rustfs_lock::error::LockError::QuorumNotReached { required, achieved } => { + StorageError::NamespaceLockQuorumUnavailable { + mode: "write", + bucket: bucket.to_string(), + object: bucket.to_string(), + required, + achieved, + } + } + other => StorageError::other(format!("delete_bucket: failed to acquire write lock on {bucket}: {other}")), + })?, + ) + } else { + None + }; // Check bucket exists before deletion (per S3 API spec) // If bucket doesn't exist, return NoSuchBucket error diff --git a/crates/ecstore/src/store/object.rs b/crates/ecstore/src/store/object.rs index ec5bdf0aa..0187120f4 100644 --- a/crates/ecstore/src/store/object.rs +++ b/crates/ecstore/src/store/object.rs @@ -423,7 +423,7 @@ impl ECStore { }; let put_opts = ObjectOptions { - user_defined: src_info.user_defined.clone(), + user_defined: (*src_info.user_defined).clone(), versioned: dst_opts.versioned, version_id: dst_opts.version_id.clone(), no_lock: dst_opts.no_lock, @@ -823,7 +823,7 @@ impl ECStore { } let (oi, _) = self.get_latest_accessible_object_info_with_idx(bucket, &object, opts).await?; - Ok(oi.user_tags) + Ok((*oi.user_tags).clone()) } #[instrument(level = "debug", skip(self))] diff --git a/crates/ecstore/src/store_api/readers.rs b/crates/ecstore/src/store_api/readers.rs index 0218bc3bb..c855288b9 100644 --- a/crates/ecstore/src/store_api/readers.rs +++ b/crates/ecstore/src/store_api/readers.rs @@ -961,7 +961,7 @@ mod tests { fn test_http_range_spec_from_object_info_valid_and_invalid_parts() { let object_info = ObjectInfo { size: 300, - parts: vec![ + parts: Arc::new(vec![ ObjectPartInfo { etag: String::new(), number: 1, @@ -983,7 +983,7 @@ mod tests { actual_size: 100, ..Default::default() }, - ], + ]), ..Default::default() }; @@ -999,7 +999,7 @@ mod tests { fn test_http_range_spec_from_object_info_uses_actual_size() { let object_info = ObjectInfo { size: 90, - parts: vec![ + parts: Arc::new(vec![ ObjectPartInfo { etag: String::new(), number: 1, @@ -1021,7 +1021,7 @@ mod tests { actual_size: 50, ..Default::default() }, - ], + ]), ..Default::default() }; @@ -1034,7 +1034,7 @@ mod tests { fn test_http_range_spec_from_object_info_falls_back_to_part_size_when_actual_size_missing() { let object_info = ObjectInfo { size: 90, - parts: vec![ + parts: Arc::new(vec![ ObjectPartInfo { etag: String::new(), number: 1, @@ -1056,7 +1056,7 @@ mod tests { actual_size: 0, ..Default::default() }, - ], + ]), ..Default::default() }; @@ -1155,7 +1155,7 @@ mod tests { ("X-Rustfs-Encryption-IV".to_string(), BASE64_STANDARD.encode(base_nonce)), ]); let object_info = ObjectInfo { - user_defined: metadata, + user_defined: Arc::new(metadata), ..Default::default() }; let material = resolve_encryption_material(&object_info, &HeaderMap::new()) @@ -1172,10 +1172,10 @@ mod tests { async fn test_get_object_reader_rejects_ssec_read_without_headers() { let object_info = ObjectInfo { size: 10, - user_defined: HashMap::from([ + user_defined: Arc::new(HashMap::from([ ("x-amz-server-side-encryption-customer-algorithm".to_string(), "AES256".to_string()), ("x-amz-server-side-encryption-customer-original-size".to_string(), "20".to_string()), - ]), + ])), ..Default::default() }; @@ -1201,10 +1201,10 @@ mod tests { async fn test_get_object_reader_restore_request_bypasses_encryption_range_rewrite() { let object_info = ObjectInfo { size: 10, - user_defined: HashMap::from([ + user_defined: Arc::new(HashMap::from([ ("x-rustfs-encryption-key".to_string(), "encrypted-key".to_string()), ("x-rustfs-encryption-original-size".to_string(), "20".to_string()), - ]), + ])), ..Default::default() }; @@ -1247,12 +1247,12 @@ mod tests { let object_info = ObjectInfo { size: encrypted.len() as i64, - user_defined: HashMap::from([ + user_defined: Arc::new(HashMap::from([ ("x-amz-server-side-encryption".to_string(), "AES256".to_string()), ("x-rustfs-encryption-key".to_string(), BASE64_STANDARD.encode(encrypted_dek.as_bytes())), ("x-rustfs-encryption-iv".to_string(), BASE64_STANDARD.encode(base_nonce)), ("x-rustfs-encryption-original-size".to_string(), plaintext.len().to_string()), - ]), + ])), ..Default::default() }; @@ -1293,12 +1293,12 @@ mod tests { let object_info = ObjectInfo { size: encrypted.len() as i64, - user_defined: HashMap::from([ + user_defined: Arc::new(HashMap::from([ ("x-amz-server-side-encryption".to_string(), "AES256".to_string()), ("x-rustfs-encryption-key".to_string(), BASE64_STANDARD.encode(encrypted_dek.as_bytes())), ("x-rustfs-encryption-iv".to_string(), BASE64_STANDARD.encode(base_nonce)), ("x-rustfs-encryption-original-size".to_string(), plaintext.len().to_string()), - ]), + ])), ..Default::default() }; let range = HTTPRangeSpec { @@ -1349,12 +1349,12 @@ mod tests { let object_info = ObjectInfo { size: encrypted.len() as i64, - user_defined: HashMap::from([ + user_defined: Arc::new(HashMap::from([ ("x-amz-server-side-encryption".to_string(), "AES256".to_string()), ("x-rustfs-encryption-key".to_string(), BASE64_STANDARD.encode(encrypted_dek.as_bytes())), ("x-rustfs-encryption-iv".to_string(), BASE64_STANDARD.encode(base_nonce)), ("x-rustfs-encryption-original-size".to_string(), plaintext.len().to_string()), - ]), + ])), ..Default::default() }; @@ -1386,18 +1386,18 @@ mod tests { let object_info = ObjectInfo { size: 3_000_000, - parts: vec![ObjectPartInfo { + parts: Arc::new(vec![ObjectPartInfo { etag: String::new(), number: 1, size: 3_000_000, actual_size: 4_194_304, index: Some(index.into_vec()), ..Default::default() - }], - user_defined: HashMap::from([ + }]), + user_defined: Arc::new(HashMap::from([ ("x-minio-internal-compression".to_string(), "gzip".to_string()), ("x-minio-internal-actual-size".to_string(), "4194304".to_string()), - ]), + ])), ..Default::default() }; @@ -1443,7 +1443,7 @@ mod tests { bucket: bucket.to_string(), name: object.to_string(), size: encrypted.len() as i64, - user_defined: HashMap::from([ + user_defined: Arc::new(HashMap::from([ ("x-amz-server-side-encryption-customer-algorithm".to_string(), "AES256".to_string()), ( "x-amz-server-side-encryption-customer-key-md5".to_string(), @@ -1453,7 +1453,7 @@ mod tests { "x-amz-server-side-encryption-customer-original-size".to_string(), plaintext.len().to_string(), ), - ]), + ])), ..Default::default() }; @@ -1496,7 +1496,7 @@ mod tests { bucket: bucket.to_string(), name: object.to_string(), size: encrypted.len() as i64, - user_defined: HashMap::from([ + user_defined: Arc::new(HashMap::from([ ("x-amz-server-side-encryption-customer-algorithm".to_string(), "AES256".to_string()), ( "x-amz-server-side-encryption-customer-key-md5".to_string(), @@ -1506,7 +1506,7 @@ mod tests { "x-amz-server-side-encryption-customer-original-size".to_string(), plaintext.len().to_string(), ), - ]), + ])), ..Default::default() }; let range = HTTPRangeSpec { @@ -1560,7 +1560,7 @@ mod tests { bucket: bucket.to_string(), name: object.to_string(), size: encrypted.len() as i64, - user_defined: HashMap::from([ + user_defined: Arc::new(HashMap::from([ ("x-amz-server-side-encryption-customer-algorithm".to_string(), "AES256".to_string()), ( "x-amz-server-side-encryption-customer-key-md5".to_string(), @@ -1572,7 +1572,7 @@ mod tests { ), ("x-minio-internal-compression".to_string(), CompressionAlgorithm::default().to_string()), ("x-minio-internal-actual-size".to_string(), plaintext.len().to_string()), - ]), + ])), ..Default::default() }; let range = HTTPRangeSpec { diff --git a/crates/ecstore/src/store_api/types.rs b/crates/ecstore/src/store_api/types.rs index 4f43dde66..5b69f1be9 100644 --- a/crates/ecstore/src/store_api/types.rs +++ b/crates/ecstore/src/store_api/types.rs @@ -293,7 +293,7 @@ pub struct ObjectInfo { // Actual size is the real size of the object uploaded by client. pub actual_size: i64, pub is_dir: bool, - pub user_defined: HashMap, + pub user_defined: Arc>, pub parity_blocks: usize, pub data_blocks: usize, pub version_id: Option, @@ -301,8 +301,8 @@ pub struct ObjectInfo { pub transitioned_object: TransitionedObject, pub restore_ongoing: bool, pub restore_expires: Option, - pub user_tags: String, - pub parts: Vec, + pub user_tags: Arc, + pub parts: Arc>, pub is_latest: bool, pub content_type: Option, pub content_encoding: Option, @@ -477,7 +477,11 @@ impl ObjectInfo { }; // tags - let user_tags = fi.metadata.get(AMZ_OBJECT_TAGGING).cloned().unwrap_or_default(); + let user_tags: Arc = fi + .metadata + .get(AMZ_OBJECT_TAGGING) + .map(|s| Arc::new(s.clone())) + .unwrap_or_default(); let inlined = fi.inline_data(); @@ -571,7 +575,7 @@ impl ObjectInfo { number: part.number, error: part.error.clone(), }) - .collect(); + .collect::>(); // TODO: part checksums @@ -585,7 +589,7 @@ impl ObjectInfo { delete_marker: fi.deleted, mod_time: fi.mod_time, size: fi.size, - parts, + parts: Arc::new(parts), is_latest: fi.is_latest, user_tags, content_type, @@ -595,7 +599,7 @@ impl ObjectInfo { successor_mod_time: fi.successor_mod_time, etag, inlined, - user_defined: metadata, + user_defined: Arc::new(metadata), transitioned_object, checksum: fi.checksum.clone(), storage_class, @@ -1207,7 +1211,7 @@ mod tests { let info = ObjectInfo { size: 100, actual_size: 0, - user_defined, + user_defined: Arc::new(user_defined), ..Default::default() }; @@ -1225,7 +1229,7 @@ mod tests { let info = ObjectInfo { size: 100, actual_size: 0, - user_defined, + user_defined: Arc::new(user_defined), ..Default::default() }; @@ -1277,8 +1281,8 @@ mod tests { let info = ObjectInfo { size: 12, actual_size: 0, - user_defined, - parts: vec![ + user_defined: Arc::new(user_defined), + parts: Arc::new(vec![ rustfs_filemeta::ObjectPartInfo { actual_size: 4, ..Default::default() @@ -1287,7 +1291,7 @@ mod tests { actual_size: 5, ..Default::default() }, - ], + ]), ..Default::default() }; @@ -1305,7 +1309,7 @@ mod tests { let info = ObjectInfo { size: 12, actual_size: 0, - user_defined, + user_defined: Arc::new(user_defined), ..Default::default() }; @@ -1329,7 +1333,7 @@ mod tests { }); let info = ObjectInfo { - user_defined, + user_defined: Arc::new(user_defined), ..Default::default() }; @@ -1359,7 +1363,7 @@ mod tests { }); let info = ObjectInfo { - user_defined, + user_defined: Arc::new(user_defined), ..Default::default() }; @@ -1372,10 +1376,72 @@ mod tests { user_defined.insert("X-Rustfs-Encryption-Key".to_string(), "encrypted-key".to_string()); let info = ObjectInfo { - user_defined, + user_defined: Arc::new(user_defined), ..Default::default() }; assert!(info.is_encrypted()); } + + #[test] + fn objectinfo_clone_shares_arc_data_and_is_correct() { + let mut ud = HashMap::new(); + ud.insert("content-type".to_string(), "application/octet-stream".to_string()); + ud.insert("x-custom-header".to_string(), "custom-value".to_string()); + + let original = ObjectInfo { + bucket: "test-bucket".to_string(), + name: "test-object".to_string(), + user_defined: Arc::new(ud), + user_tags: Arc::new("env=prod&team=storage".to_string()), + parts: Arc::new(vec![ + rustfs_filemeta::ObjectPartInfo { + number: 1, + size: 1024, + actual_size: 1024, + ..Default::default() + }, + rustfs_filemeta::ObjectPartInfo { + number: 2, + size: 512, + actual_size: 512, + ..Default::default() + }, + ]), + size: 1536, + etag: Some("abc123".to_string()), + ..Default::default() + }; + + let cloned = original.clone(); + + // Verify cloned values are correct + assert_eq!(cloned.bucket, "test-bucket"); + assert_eq!(cloned.name, "test-object"); + assert_eq!(cloned.size, 1536); + assert_eq!(cloned.etag, Some("abc123".to_string())); + + // Verify Arc fields share the same allocation + assert!(Arc::ptr_eq(&original.user_defined, &cloned.user_defined)); + assert!(Arc::ptr_eq(&original.user_tags, &cloned.user_tags)); + assert!(Arc::ptr_eq(&original.parts, &cloned.parts)); + + // Verify Arc-wrapped data is accessible through the clone + assert_eq!( + cloned.user_defined.get("content-type").map(String::as_str), + Some("application/octet-stream") + ); + assert_eq!(cloned.user_tags.as_str(), "env=prod&team=storage"); + assert_eq!(cloned.parts.len(), 2); + assert_eq!(cloned.parts[0].number, 1); + assert_eq!(cloned.parts[1].size, 512); + + // Verify default ObjectInfo clone also works + let default_obj = ObjectInfo::default(); + let default_cloned = default_obj.clone(); + assert!(default_obj.user_defined.is_empty()); + assert!(default_cloned.user_defined.is_empty()); + assert!(default_cloned.user_tags.is_empty()); + assert!(default_cloned.parts.is_empty()); + } } diff --git a/crates/ecstore/src/store_list_objects.rs b/crates/ecstore/src/store_list_objects.rs index 828f4b777..034b3ed27 100644 --- a/crates/ecstore/src/store_list_objects.rs +++ b/crates/ecstore/src/store_list_objects.rs @@ -1218,12 +1218,18 @@ async fn gather_results( } async fn select_from( + rx: &CancellationToken, in_channels: &mut [Receiver], idx: usize, top: &mut [Option], n_done: &mut usize, -) -> Result<()> { - match in_channels[idx].recv().await { +) -> Result { + let entry = tokio::select! { + entry = in_channels[idx].recv() => entry, + _ = rx.cancelled() => return Ok(false), + }; + + match entry { Some(entry) => { top[idx] = Some(entry); } @@ -1232,10 +1238,19 @@ async fn select_from( *n_done += 1; } } - Ok(()) + Ok(true) +} + +async fn send_or_cancel(rx: &CancellationToken, out_channel: &Sender, entry: MetaCacheEntry) -> Result { + tokio::select! { + result = out_channel.send(entry) => { + result.map_err(Error::other)?; + Ok(true) + } + _ = rx.cancelled() => Ok(false), + } } -// TODO: exit when cancel async fn merge_entry_channels( rx: CancellationToken, in_channels: Vec>, @@ -1253,7 +1268,9 @@ async fn merge_entry_channels( has_entry = in_channels[0].recv()=>{ if let Some(entry) = has_entry{ // warn!("merge_entry_channels entry {}", &entry.name); - out_channel.send(entry).await.map_err(Error::other)?; + if !send_or_cancel(&rx, &out_channel, entry).await? { + return Ok(()); + } } else { return Ok(()) } @@ -1272,7 +1289,9 @@ async fn merge_entry_channels( let in_channels_len = in_channels.len(); for idx in 0..in_channels_len { - select_from(&mut in_channels, idx, &mut top, &mut n_done).await?; + if !select_from(&rx, &mut in_channels, idx, &mut top, &mut n_done).await? { + return Ok(()); + } } let mut last = String::new(); @@ -1286,16 +1305,13 @@ async fn merge_entry_channels( let mut best_idx = 0; to_merge.clear(); - // FIXME: top move when select_from call - // let vtop = top.clone(); - - // let vtop = top.as_slice(); + // Note: `select_from` mutates `top[idx]` during the inner loop, but this is safe + // because each borrow from `top[other_idx]` is only used before any later + // `select_from` call that can mutate that slot. for other_idx in 1..top.len() { if let Some(other_entry) = &top[other_idx] { if let Some(best_entry) = &best { - // println!("get other_entry {:?}", other_entry.name); - if path::clean(&best_entry.name) == path::clean(&other_entry.name) { let dir_matches = best_entry.is_dir() && other_entry.is_dir(); let suffix_matches = @@ -1310,7 +1326,9 @@ async fn merge_entry_channels( // dir and object has the save name if other_entry.is_dir() { // TODO: read next entry to top - select_from(&mut in_channels, other_idx, &mut top, &mut n_done).await?; + if !select_from(&rx, &mut in_channels, other_idx, &mut top, &mut n_done).await? { + return Ok(()); + } continue; } @@ -1351,7 +1369,9 @@ async fn merge_entry_channels( let xl2 = match entry.clone().xl_meta() { Ok(res) => res, Err(_) => { - select_from(&mut in_channels, idx, &mut top, &mut n_done).await?; + if !select_from(&rx, &mut in_channels, idx, &mut top, &mut n_done).await? { + return Ok(()); + } continue; } @@ -1360,13 +1380,17 @@ async fn merge_entry_channels( versions.push(xl2.versions.clone()); if has_xl.is_none() { - select_from(&mut in_channels, best_idx, &mut top, &mut n_done).await?; + if !select_from(&rx, &mut in_channels, best_idx, &mut top, &mut n_done).await? { + return Ok(()); + } best_idx = idx; best = Some(entry.clone()); has_xl = Some(xl2); } else { - select_from(&mut in_channels, best_idx, &mut top, &mut n_done).await?; + if !select_from(&rx, &mut in_channels, best_idx, &mut top, &mut n_done).await? { + return Ok(()); + } } } } @@ -1391,11 +1415,15 @@ async fn merge_entry_channels( if let Some(best_entry) = &best && best_entry.name > last { - out_channel.send(best_entry.clone()).await.map_err(Error::other)?; + if !send_or_cancel(&rx, &out_channel, best_entry.clone()).await? { + return Ok(()); + } last = best_entry.name.clone(); } - select_from(&mut in_channels, best_idx, &mut top, &mut n_done).await?; + if !select_from(&rx, &mut in_channels, best_idx, &mut top, &mut n_done).await? { + return Ok(()); + } } } @@ -1557,8 +1585,8 @@ fn calc_common_counter(infos: &[DiskInfo], read_quorum: usize) -> u64 { mod test { use super::{ ENV_API_LIST_QUORUM, ListPathOptions, MAX_OBJECT_LIST, VersionMarker, gather_results, list_metadata_resolution_params, - list_quorum_from_env, max_keys_plus_one, normalize_list_quorum, parse_version_marker, version_marker_for_entries, - walk_result_from_set_errors, + list_quorum_from_env, max_keys_plus_one, merge_entry_channels, normalize_list_quorum, parse_version_marker, + version_marker_for_entries, walk_result_from_set_errors, }; use crate::error::StorageError; use rustfs_filemeta::{MetaCacheEntries, MetaCacheEntriesSorted, MetaCacheEntry}; @@ -2224,4 +2252,172 @@ mod test { // println!("get entry {:?}", entry) // } // } + + #[tokio::test] + async fn merge_entry_channels_produces_sorted_unique_output_from_two_channels() { + let (tx_a, rx_a) = mpsc::channel(4); + let (tx_b, rx_b) = mpsc::channel(4); + let (out_tx, mut out_rx) = mpsc::channel(8); + + // Send sorted entries from two channels with some overlap + tx_a.send(test_meta_entry("obj-a")).await.unwrap(); + tx_a.send(test_meta_entry("obj-c")).await.unwrap(); + tx_a.send(test_meta_entry("obj-e")).await.unwrap(); + drop(tx_a); + + tx_b.send(test_meta_entry("obj-b")).await.unwrap(); + tx_b.send(test_meta_entry("obj-d")).await.unwrap(); + drop(tx_b); + + let rx = CancellationToken::new(); + let handle = tokio::spawn(merge_entry_channels(rx, vec![rx_a, rx_b], out_tx, 1)); + + let mut results = Vec::new(); + while let Some(entry) = out_rx.recv().await { + results.push(entry.name.clone()); + } + + handle.await.unwrap().unwrap(); + + // Results should be sorted and deduplicated + assert_eq!(results, vec!["obj-a", "obj-b", "obj-c", "obj-d", "obj-e"]); + } + + #[tokio::test] + async fn merge_entry_channels_deduplicates_entries_across_channels() { + let (tx_a, rx_a) = mpsc::channel(4); + let (tx_b, rx_b) = mpsc::channel(4); + let (out_tx, mut out_rx) = mpsc::channel(8); + + // Both channels have the same entry + tx_a.send(test_meta_entry("obj-a")).await.unwrap(); + tx_a.send(test_meta_entry("obj-c")).await.unwrap(); + drop(tx_a); + + tx_b.send(test_meta_entry("obj-a")).await.unwrap(); + tx_b.send(test_meta_entry("obj-b")).await.unwrap(); + drop(tx_b); + + let rx = CancellationToken::new(); + let handle = tokio::spawn(merge_entry_channels(rx, vec![rx_a, rx_b], out_tx, 1)); + + let mut results = Vec::new(); + while let Some(entry) = out_rx.recv().await { + results.push(entry.name.clone()); + } + + handle.await.unwrap().unwrap(); + + // "obj-a" should appear only once despite being in both channels + assert_eq!(results, vec!["obj-a", "obj-b", "obj-c"]); + } + + #[tokio::test] + async fn merge_entry_channels_handles_single_channel() { + let (tx, rx) = mpsc::channel(4); + let (out_tx, mut out_rx) = mpsc::channel(8); + + tx.send(test_meta_entry("obj-a")).await.unwrap(); + tx.send(test_meta_entry("obj-b")).await.unwrap(); + drop(tx); + + let cancel = CancellationToken::new(); + let handle = tokio::spawn(merge_entry_channels(cancel, vec![rx], out_tx, 1)); + + let mut results = Vec::new(); + while let Some(entry) = out_rx.recv().await { + results.push(entry.name.clone()); + } + + handle.await.unwrap().unwrap(); + + assert_eq!(results, vec!["obj-a", "obj-b"]); + } + + #[tokio::test] + async fn merge_entry_channels_respects_cancellation() { + let (tx, rx) = mpsc::channel::(4); + let (out_tx, mut out_rx) = mpsc::channel(8); + + let cancel = CancellationToken::new(); + let cancel_clone = cancel.clone(); + + // Use a single channel so the cancellation check in tokio::select! is exercised. + let handle = tokio::spawn(merge_entry_channels(cancel_clone, vec![rx], out_tx, 1)); + + // Send an entry to prove the function is running and processing data. + tx.send(test_meta_entry("a")).await.unwrap(); + let received = timeout(Duration::from_millis(500), out_rx.recv()) + .await + .expect("should receive entry before cancellation") + .map(|e| e.name); + assert_eq!(received, Some("a".to_string())); + + // Cancel while the sender is still alive (channel is not closed). + cancel.cancel(); + + let result = timeout(Duration::from_secs(2), handle) + .await + .expect("merge should not hang after cancellation") + .expect("task should not panic"); + assert!(result.is_ok(), "merge should return Ok on cancellation"); + + // Keep tx alive until after the assertion so the channel doesn't close prematurely. + drop(tx); + } + + #[tokio::test] + async fn merge_entry_channels_respects_cancellation_with_multiple_live_channels() { + let (tx_a, rx_a) = mpsc::channel::(4); + let (tx_b, rx_b) = mpsc::channel::(4); + let (out_tx, _out_rx) = mpsc::channel(8); + + let cancel = CancellationToken::new(); + let cancel_clone = cancel.clone(); + + let handle = tokio::spawn(merge_entry_channels(cancel_clone, vec![rx_a, rx_b], out_tx, 1)); + cancel.cancel(); + + let result = timeout(Duration::from_secs(2), handle) + .await + .expect("multi-channel merge should not hang after cancellation") + .expect("task should not panic"); + assert!(result.is_ok(), "multi-channel merge should return Ok on cancellation"); + + drop(tx_a); + drop(tx_b); + } + + #[tokio::test] + async fn merge_entry_channels_respects_cancellation_when_output_is_full() { + let (tx_a, rx_a) = mpsc::channel::(4); + let (tx_b, rx_b) = mpsc::channel::(4); + let (out_tx, _out_rx) = mpsc::channel(1); + + out_tx + .send(test_meta_entry("already-buffered")) + .await + .expect("output channel should accept initial entry"); + tx_a.send(test_meta_entry("a")) + .await + .expect("input channel a should accept entry"); + tx_b.send(test_meta_entry("b")) + .await + .expect("input channel b should accept entry"); + + let cancel = CancellationToken::new(); + let cancel_clone = cancel.clone(); + + let handle = tokio::spawn(merge_entry_channels(cancel_clone, vec![rx_a, rx_b], out_tx, 1)); + cancel.cancel(); + + let result = timeout(Duration::from_secs(2), handle) + .await + .expect("merge should not hang when output is full and cancellation fires") + .expect("task should not panic"); + assert!(result.is_ok(), "merge should return Ok when cancelled during output send"); + + drop(tx_a); + drop(tx_b); + } } diff --git a/crates/protocols/src/swift/handler.rs b/crates/protocols/src/swift/handler.rs index f42a67461..2fc7e1681 100644 --- a/crates/protocols/src/swift/handler.rs +++ b/crates/protocols/src/swift/handler.rs @@ -721,10 +721,10 @@ async fn handle_authenticated_request( response = response.header("etag", etag); } - for (key, value) in info.user_defined { + for (key, value) in info.user_defined.iter() { if key != "content-type" { let header_name = format!("x-object-meta-{}", key); - response = response.header(header_name, value); + response = response.header(header_name, value.as_str()); } } @@ -790,10 +790,10 @@ async fn handle_authenticated_request( } // Add custom metadata headers (X-Object-Meta-*) - for (key, value) in info.user_defined { + for (key, value) in info.user_defined.iter() { if key != "content-type" { let header_name = format!("x-object-meta-{}", key); - response = response.header(header_name, value); + response = response.header(header_name, value.as_str()); } } @@ -829,10 +829,10 @@ async fn handle_authenticated_request( } // Add custom metadata headers (X-Object-Meta-*) - for (key, value) in info.user_defined { + for (key, value) in info.user_defined.iter() { if key != "content-type" { let header_name = format!("x-object-meta-{}", key); - response = response.header(header_name, value); + response = response.header(header_name, value.as_str()); } } @@ -1134,10 +1134,10 @@ async fn handle_object_get( response = response.header("etag", etag); } - for (key, value) in info.user_defined { + for (key, value) in info.user_defined.iter() { if key != "content-type" { let header_name = format!("x-object-meta-{}", key); - response = response.header(header_name, value); + response = response.header(header_name, value.as_str()); } } @@ -1200,13 +1200,13 @@ async fn handle_object_get( response = response.header("etag", etag); } - for (key, value) in info.user_defined { + for (key, value) in info.user_defined.iter() { if key == "x-delete-at" { // Add X-Delete-At header directly (not as X-Object-Meta-*) - response = response.header("x-delete-at", value); + response = response.header("x-delete-at", value.as_str()); } else if key != "content-type" { let header_name = format!("x-object-meta-{}", key); - response = response.header(header_name, value); + response = response.header(header_name, value.as_str()); } } @@ -1256,13 +1256,13 @@ async fn handle_object_head( response = response.header("etag", etag); } - for (key, value) in info.user_defined { + for (key, value) in info.user_defined.iter() { if key == "x-delete-at" { // Add X-Delete-At header directly (not as X-Object-Meta-*) - response = response.header("x-delete-at", value); + response = response.header("x-delete-at", value.as_str()); } else if key != "content-type" { let header_name = format!("x-object-meta-{}", key); - response = response.header(header_name, value); + response = response.header(header_name, value.as_str()); } } diff --git a/crates/protocols/src/swift/object.rs b/crates/protocols/src/swift/object.rs index 1c59da07f..3e277bfba 100644 --- a/crates/protocols/src/swift/object.rs +++ b/crates/protocols/src/swift/object.rs @@ -867,7 +867,7 @@ pub async fn copy_object( // 10. Prepare metadata for destination object // Start with source metadata - let mut new_metadata = src_info.user_defined.clone(); + let mut new_metadata = (*src_info.user_defined).clone(); // 11. If custom metadata headers provided, use those instead (Swift behavior) let mut has_custom_meta = false; diff --git a/rustfs/src/app/bucket_usecase.rs b/rustfs/src/app/bucket_usecase.rs index a2edd9c9a..f16220392 100644 --- a/rustfs/src/app/bucket_usecase.rs +++ b/rustfs/src/app/bucket_usecase.rs @@ -472,7 +472,7 @@ fn build_list_object_versions_m_output( None }; let user_tags = if permission.tags_allowed && !object.user_tags.is_empty() { - Some(object.user_tags.clone()) + Some((*object.user_tags).clone()) } else { None }; @@ -569,7 +569,7 @@ fn build_list_objects_v2m_output( None }; let user_tags = if permission.tags_allowed && !object.user_tags.is_empty() { - Some(object.user_tags.clone()) + Some((*object.user_tags).clone()) } else { None }; @@ -2101,6 +2101,7 @@ impl DefaultBucketUsecase { mod tests { use super::*; use http::{Extensions, HeaderMap, Method, Uri}; + use std::sync::Arc; fn build_request(input: T, method: Method) -> S3Request { S3Request { @@ -2717,11 +2718,11 @@ mod tests { name: "obj-a".to_string(), mod_time: Some(datetime!(2025-01-01 00:00 UTC)), size: 11, - user_defined: HashMap::from([("project".to_string(), "alpha".to_string())]), + user_defined: Arc::new(HashMap::from([("project".to_string(), "alpha".to_string())])), parity_blocks: 2, data_blocks: 4, version_id: Some(Uuid::nil()), - user_tags: "env=prod".to_string(), + user_tags: Arc::new("env=prod".to_string()), is_latest: true, etag: Some("0123456789abcdef0123456789abcdef".to_string()), ..Default::default() @@ -2731,7 +2732,7 @@ mod tests { name: "obj-b".to_string(), mod_time: Some(datetime!(2025-01-02 00:00 UTC)), delete_marker: true, - user_defined: HashMap::from([("marker".to_string(), "true".to_string())]), + user_defined: Arc::new(HashMap::from([("marker".to_string(), "true".to_string())])), version_id: None, ..Default::default() }, @@ -2825,8 +2826,8 @@ mod tests { name: "logs and more/object one.txt".to_string(), mod_time: Some(datetime!(2025-01-04 00:00 UTC)), size: 7, - user_defined: HashMap::from([("secret".to_string(), "value".to_string())]), - user_tags: "env=prod".to_string(), + user_defined: Arc::new(HashMap::from([("secret".to_string(), "value".to_string())])), + user_tags: Arc::new("env=prod".to_string()), parity_blocks: 1, data_blocks: 2, ..Default::default() @@ -2905,10 +2906,10 @@ mod tests { name: "logs/obj a.txt".to_string(), mod_time: Some(datetime!(2025-01-03 00:00 UTC)), size: 11, - user_defined: HashMap::from([("project".to_string(), "alpha".to_string())]), + user_defined: Arc::new(HashMap::from([("project".to_string(), "alpha".to_string())])), parity_blocks: 2, data_blocks: 4, - user_tags: "env=prod".to_string(), + user_tags: Arc::new("env=prod".to_string()), etag: Some("0123456789abcdef0123456789abcdef".to_string()), ..Default::default() }], @@ -2981,8 +2982,8 @@ mod tests { name: "logs and more/object one.txt".to_string(), mod_time: Some(datetime!(2025-01-05 00:00 UTC)), size: 13, - user_defined: HashMap::from([("secret".to_string(), "value".to_string())]), - user_tags: "env=prod".to_string(), + user_defined: Arc::new(HashMap::from([("secret".to_string(), "value".to_string())])), + user_tags: Arc::new("env=prod".to_string()), parity_blocks: 1, data_blocks: 2, ..Default::default() diff --git a/rustfs/src/app/object_usecase.rs b/rustfs/src/app/object_usecase.rs index d30160590..4d4b9727d 100644 --- a/rustfs/src/app/object_usecase.rs +++ b/rustfs/src/app/object_usecase.rs @@ -559,7 +559,7 @@ fn delete_replication_state_from_config( ) -> Option { let opts = ReplicationObjectOpts { name: obj_info.name.clone(), - user_tags: obj_info.user_tags.clone(), + user_tags: (*obj_info.user_tags).clone(), version_id, delete_marker: obj_info.delete_marker, op_type: ReplicationType::Delete, @@ -1347,7 +1347,7 @@ impl DefaultObjectUsecase { check_preconditions(&req.headers, &info)?; debug!(object_size = info.size, part_count = info.parts.len(), "GET object metadata snapshot"); - for part in &info.parts { + for part in info.parts.iter() { debug!( part_number = part.number, part_size = part.size, @@ -2810,7 +2810,10 @@ impl DefaultObjectUsecase { src_info.metadata_only = true; } - strip_managed_encryption_metadata(&mut src_info.user_defined); + // Extract user_defined from Arc for mutation; it will be re-wrapped after all edits. + let mut user_defined = (*src_info.user_defined).clone(); + + strip_managed_encryption_metadata(&mut user_defined); let actual_size = src_info.get_actual_size().map_err(ApiError::from)?; @@ -2824,27 +2827,27 @@ impl DefaultObjectUsecase { insert_str(&mut compress_metadata, SUFFIX_COMPRESSION, CompressionAlgorithm::default().to_string()); insert_str(&mut compress_metadata, SUFFIX_ACTUAL_SIZE, actual_size.to_string()); } else { - remove_str(&mut src_info.user_defined, SUFFIX_COMPRESSION); - remove_str(&mut src_info.user_defined, SUFFIX_ACTUAL_SIZE); - remove_str(&mut src_info.user_defined, SUFFIX_COMPRESSION_SIZE); + remove_str(&mut user_defined, SUFFIX_COMPRESSION); + remove_str(&mut user_defined, SUFFIX_ACTUAL_SIZE); + remove_str(&mut user_defined, SUFFIX_COMPRESSION_SIZE); } // Handle MetadataDirective REPLACE: replace user metadata while preserving system metadata. // System metadata (compression, encryption) is added after this block to ensure // it's not cleared by the REPLACE operation. if metadata_directive.as_ref().map(|d| d.as_str()) == Some(MetadataDirective::REPLACE) { - src_info.user_defined.clear(); + user_defined.clear(); if let Some(metadata) = metadata { - src_info.user_defined.extend(metadata); + user_defined.extend(metadata); } if let Some(ct) = content_type { src_info.content_type = Some(ct.clone()); - src_info.user_defined.insert("content-type".to_string(), ct); + user_defined.insert("content-type".to_string(), ct); } } let has_explicit_object_lock_retention = object_lock_mode.is_some() || object_lock_retain_until_date.is_some(); - remove_object_lock_metadata_for_copy(&mut src_info.user_defined); + remove_object_lock_metadata_for_copy(&mut user_defined); if let Some(object_lock_metadata) = build_put_like_object_lock_metadata( &bucket, object_lock_legal_hold_status, @@ -2853,9 +2856,9 @@ impl DefaultObjectUsecase { ) .await? { - src_info.user_defined.extend(object_lock_metadata); + user_defined.extend(object_lock_metadata); } - apply_bucket_default_lock_retention(&bucket, &mut src_info.user_defined, has_explicit_object_lock_retention).await?; + apply_bucket_default_lock_retention(&bucket, &mut user_defined, has_explicit_object_lock_retention).await?; let mut reader = if should_compress { let hrd = HashReader::from_stream(gr.stream, length, actual_size, None, None, false).map_err(ApiError::from)?; @@ -2892,7 +2895,7 @@ impl DefaultObjectUsecase { reader = HashReader::from_reader(encrypted_reader, HashReader::SIZE_PRESERVE_LAYER, actual_size, None, None, false) .map_err(ApiError::from)?; - src_info.user_defined.extend(encryption_material_to_metadata(&material)); + user_defined.extend(encryption_material_to_metadata(&material)); } src_info.put_object_reader = Some(PutObjReader::new(reader)); @@ -2900,9 +2903,11 @@ impl DefaultObjectUsecase { // check quota for (k, v) in compress_metadata { - src_info.user_defined.insert(k, v); + user_defined.insert(k, v); } + src_info.user_defined = Arc::new(user_defined); + self.check_bucket_quota(&bucket, QuotaOperation::CopyObject, src_info.size as u64) .await?; let has_bucket_metadata = self.bucket_metadata_sys().is_some(); @@ -3899,7 +3904,7 @@ impl DefaultObjectUsecase { } let restore_expiry = lifecycle::expected_expiry_time(OffsetDateTime::now_utc(), *rreq.days.as_ref().unwrap_or(&1)); - let mut metadata = obj_info.user_defined.clone(); + let mut metadata = (*obj_info.user_defined).clone(); let mut header = HeaderMap::new(); @@ -3931,7 +3936,7 @@ impl DefaultObjectUsecase { .to_string(), ); } - obj_info.user_defined = metadata; + obj_info.user_defined = Arc::new(metadata); store .clone() @@ -5032,7 +5037,7 @@ mod tests { let metadata = HashMap::new(); let standard_info = ObjectInfo { storage_class: Some(storageclass::STANDARD.to_string()), - user_defined: metadata.clone(), + user_defined: Arc::new(metadata.clone()), ..Default::default() }; assert!(response_storage_class(&standard_info, &metadata).is_none()); @@ -5041,7 +5046,7 @@ mod tests { metadata.insert(AMZ_STORAGE_CLASS.to_string(), storageclass::STANDARD_IA.to_string()); let infrequent_access_info = ObjectInfo { storage_class: Some(storageclass::STANDARD_IA.to_string()), - user_defined: metadata.clone(), + user_defined: Arc::new(metadata.clone()), ..Default::default() }; assert_eq!(