diff --git a/Cargo.lock b/Cargo.lock index fafbc34ba..ba722d3d5 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -9565,6 +9565,7 @@ dependencies = [ name = "rustfs-ecstore" version = "1.0.0-rc.3" dependencies = [ + "ahash", "arc-swap", "async-channel", "async-recursion", diff --git a/crates/ecstore/Cargo.toml b/crates/ecstore/Cargo.toml index fb223489c..e56d6cffb 100644 --- a/crates/ecstore/Cargo.toml +++ b/crates/ecstore/Cargo.toml @@ -235,6 +235,9 @@ faster-hex = { workspace = true } ratelimit = { workspace = true } aws-smithy-http-client = { workspace = true, default-features = false, features = ["rustls-aws-lc"] } +# High-performance hashing +ahash = { workspace = true, features = ["serde"] } + # Observability and Metrics metrics = { workspace = true } diff --git a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs index 5e13ba918..0a6ab0746 100644 --- a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs +++ b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs @@ -12,6 +12,7 @@ // See the License for the specific language governing permissions and // limitations under the License. +use ahash::AHashMap; use super::{metadata_boundary, object_lock_boundary, runtime_boundary as runtime_sources}; use crate::bucket::lifecycle::bucket_lifecycle_audit::{ LcAuditEvent, LcEventSrc, emit_non_transitioned_expiration_event, emit_transition_complete_event, @@ -2870,7 +2871,7 @@ fn spawn_transition_transaction_recovery_once(api: Arc) { struct StaleMultipartUploadCandidate { path: String, initiated: OffsetDateTime, - metadata: Option>, + metadata: Option>, } fn parse_stale_uploads_duration(env_key: &str, default: StdDuration) -> StdDuration { @@ -2915,9 +2916,9 @@ async fn stale_upload_current_size(set: &Arc, metadata: &HashMap( set: &Arc, - metadata: &HashMap, + metadata: &HashMap, upload_dir: &str, no_lock: bool, ) -> Option { @@ -2950,9 +2951,9 @@ async fn stale_upload_current_size_with_opts( ) } -async fn stale_upload_lifecycle_due( +async fn stale_upload_lifecycle_due( set: &Arc, - metadata: &HashMap, + metadata: &HashMap, initiated: OffsetDateTime, upload_dir: &str, no_lock: bool, @@ -2978,7 +2979,7 @@ async fn stale_upload_lifecycle_due( .unwrap_or_default(), is_latest: true, delete_marker: false, - user_defined: metadata.clone(), + user_defined: metadata.iter().map(|(k, v)| (k.clone(), v.clone())).collect(), ..Default::default() }; diff --git a/crates/ecstore/src/core/pools.rs b/crates/ecstore/src/core/pools.rs index 0f188db18..2645b4912 100644 --- a/crates/ecstore/src/core/pools.rs +++ b/crates/ecstore/src/core/pools.rs @@ -5796,7 +5796,7 @@ fn decommission_remote_tiered_opts( versioned: version_id.is_some(), version_id, mod_time: version.mod_time, - user_defined: version.metadata.clone(), + user_defined: version.metadata.iter().map(|(k, v)| (k.clone(), v.clone())).collect(), src_pool_idx, data_movement: true, incl_free_versions: version.tier_free_version(), diff --git a/crates/ecstore/src/object_api/types.rs b/crates/ecstore/src/object_api/types.rs index d91d96a0d..ef7874ca6 100644 --- a/crates/ecstore/src/object_api/types.rs +++ b/crates/ecstore/src/object_api/types.rs @@ -1008,7 +1008,7 @@ impl ObjectInfo { successor_mod_time: fi.successor_mod_time, etag, inlined, - user_defined: Arc::new(metadata), + user_defined: Arc::new(metadata.iter().map(|(k, v)| (k.clone(), v.clone())).collect()), transitioned_object, transition_version_state: fi.transition_version_state, checksum: fi.checksum.clone(), @@ -1316,7 +1316,7 @@ impl ObjectInfo { if part > 0 && let Some(checksums) = self.parts.iter().find(|p| p.number == part).and_then(|p| p.checksums.clone()) { - return Ok((checksums, true)); + return Ok((checksums.iter().map(|(k, v)| (k.clone(), v.clone())).collect(), true)); } if let Some(data) = &self.checksum { diff --git a/crates/ecstore/src/services/rebalance/migration.rs b/crates/ecstore/src/services/rebalance/migration.rs index cf2dabd8f..ee4d635f4 100644 --- a/crates/ecstore/src/services/rebalance/migration.rs +++ b/crates/ecstore/src/services/rebalance/migration.rs @@ -57,7 +57,7 @@ fn rebalance_remote_tiered_opts( versioned: version_id.is_some(), version_id, mod_time: version.mod_time, - user_defined: version.metadata.clone(), + user_defined: version.metadata.iter().map(|(k, v)| (k.clone(), v.clone())).collect(), src_pool_idx, data_movement: true, include_part_checksums: true, diff --git a/crates/ecstore/src/set_disk/ops/multipart.rs b/crates/ecstore/src/set_disk/ops/multipart.rs index 2c25ebeb8..163f99e8b 100644 --- a/crates/ecstore/src/set_disk/ops/multipart.rs +++ b/crates/ecstore/src/set_disk/ops/multipart.rs @@ -1280,7 +1280,7 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks { mod_time: Some(OffsetDateTime::now_utc()), actual_size, index: index_op, - checksums: if checksums.is_empty() { None } else { Some(checksums) }, + checksums: if checksums.is_empty() { None } else { Some(checksums.iter().map(|(k, v)| (k.clone(), v.clone())).collect()) }, ..Default::default() }; @@ -1468,7 +1468,7 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks { max_parts, part_number_marker, user_defined: { - let mut metadata = fi.metadata.clone(); + let mut metadata: HashMap = fi.metadata.iter().map(|(k, v)| (k.clone(), v.clone())).collect(); strip_internal_multipart_metadata(&mut metadata); metadata }, @@ -1727,7 +1727,7 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks { let mod_time = opts.mod_time.unwrap_or_else(OffsetDateTime::now_utc); for f in parts_metadatas.iter_mut() { - f.metadata = user_defined.clone(); + f.metadata = user_defined.iter().map(|(k, v)| (k.clone(), v.clone())).collect(); f.mod_time = Some(mod_time); f.fresh = true; } @@ -1816,7 +1816,7 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks { upload_id: upload_id.to_owned(), user_defined: { strip_internal_multipart_metadata(&mut fi.metadata); - fi.metadata.clone() + fi.metadata.iter().map(|(k, v)| (k.clone(), v.clone())).collect() }, ..Default::default() }) diff --git a/crates/ecstore/src/set_disk/ops/object.rs b/crates/ecstore/src/set_disk/ops/object.rs index 594814257..30be81498 100644 --- a/crates/ecstore/src/set_disk/ops/object.rs +++ b/crates/ecstore/src/set_disk/ops/object.rs @@ -19,6 +19,7 @@ //! bounds are unchanged, and the impls reach shared primitives through the //! SetDisks core (io_primitives) via inherent calls. +use ahash::AHashMap; use super::super::*; use super::bitrot_self_verify::{BitrotSelfVerifyTarget, drop_failed_writer_disks, verify_written_bitrot_shards}; use crate::bucket::utils::is_meta_bucketname; @@ -2697,7 +2698,7 @@ impl SetDisks { ))); } - fi.metadata = user_defined.into(); + fi.metadata = user_defined.iter().map(|(k, v)| (k.clone(), v.clone())).collect(); fi.mod_time = mod_time; fi.size = w_size as i64; fi.versioned = opts.versioned || opts.version_suspended; @@ -5882,7 +5883,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { } else { None }; - let mut replacement_metadata: AHashMap = (*src_info.user_defined).clone().into(); + let mut replacement_metadata: AHashMap = (*src_info.user_defined).iter().map(|(k, v)| (k.clone(), v.clone())).collect(); if let Some(part_checksums) = preserved_part_checksums { rustfs_utils::http::insert_str(&mut replacement_metadata, rustfs_utils::http::SUFFIX_PART_CHECKSUMS, part_checksums); }