From 9546baf1aba6a325ed5d070ce4d2bfdedb29c081 Mon Sep 17 00:00:00 2001 From: houseme Date: Thu, 13 Aug 2026 00:53:06 +0800 Subject: [PATCH] perf(get): share metadata cache hits (#6010) Keep fresh metadata fanout results owned while sharing cache-backed metadata through Arc, avoiding deep clones on eligible local cache hits without enabling unsafe distributed caching. Co-authored-by: heihutu --- crates/ecstore/src/set_disk/mod.rs | 51 +++++++++--- crates/ecstore/src/set_disk/ops/object.rs | 27 ++++--- crates/ecstore/src/set_disk/read.rs | 94 +++++++++++++++++----- crates/ecstore/src/set_disk/replication.rs | 10 ++- 4 files changed, 134 insertions(+), 48 deletions(-) diff --git a/crates/ecstore/src/set_disk/mod.rs b/crates/ecstore/src/set_disk/mod.rs index a360b9e0c..6f4577096 100644 --- a/crates/ecstore/src/set_disk/mod.rs +++ b/crates/ecstore/src/set_disk/mod.rs @@ -735,10 +735,41 @@ mod transition_matrix_tests; pub use ops::heal_walk::HealWalkVersion; +pub(in crate::set_disk) enum GetObjectMetadata { + Owned(T), + Shared(Arc), +} + +impl std::ops::Deref for GetObjectMetadata { + type Target = T; + + fn deref(&self) -> &Self::Target { + match self { + Self::Owned(value) => value, + Self::Shared(value) => value, + } + } +} + +impl GetObjectMetadata { + fn into_owned(self) -> T { + match self { + Self::Owned(value) => value, + Self::Shared(value) => Arc::try_unwrap(value).unwrap_or_else(|value| (*value).clone()), + } + } +} + +type GetObjectFileInfo = ( + GetObjectMetadata, + GetObjectMetadata>, + GetObjectMetadata>>, +); + pub(crate) struct PreparedGetObjectMetadata { - fi: FileInfo, - files: Vec, - disks: Vec>, + fi: GetObjectMetadata, + files: GetObjectMetadata>, + disks: GetObjectMetadata>>, object_info: Option, } @@ -807,9 +838,9 @@ mod prepared_get_object_metadata_tests { #[tokio::test] async fn prepared_metadata_is_consumed_exactly_once() { let metadata = PreparedGetObjectMetadata { - fi: FileInfo::default(), - files: Vec::new(), - disks: Vec::new(), + fi: GetObjectMetadata::Owned(FileInfo::default()), + files: GetObjectMetadata::Owned(Vec::new()), + disks: GetObjectMetadata::Owned(Vec::new()), object_info: None, }; @@ -2465,13 +2496,13 @@ impl Hash for GetObjectMetadataCacheKey { } } -#[derive(Clone, Debug)] +#[derive(Debug)] struct GetObjectMetadataCacheEntry { #[allow(dead_code)] // Kept for debugging; moka handles TTL internally created_at: Instant, - fi: FileInfo, - parts_metadata: Vec, - online_disks: Vec>, + fi: Arc, + parts_metadata: Arc>, + online_disks: Arc>>, read_quorum: usize, } diff --git a/crates/ecstore/src/set_disk/ops/object.rs b/crates/ecstore/src/set_disk/ops/object.rs index cfcfe0d7c..e5ebc2fad 100644 --- a/crates/ecstore/src/set_disk/ops/object.rs +++ b/crates/ecstore/src/set_disk/ops/object.rs @@ -750,8 +750,8 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks { 0, object_info.size, &mut output, - fi, - files, + fi.into_owned(), + files.into_owned(), &disks, self.set_index, self.pool_index, @@ -867,8 +867,8 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks { offset, length, &mut writer, - fi, - files, + fi.into_owned(), + files.into_owned(), &disks, set_index, pool_index, @@ -3198,7 +3198,8 @@ impl SetDisks { // Force the full quorum fanout (allow_early_stop=false): `disks` is the // write target below, and an early-stop subset would only carry read // quorum, failing write quorum on update_object_meta (backlog#872). - let (mut fi, _, disks) = self.get_object_fileinfo_gated(bucket, object, opts, false, false).await?; + let (fi, _, disks) = self.get_object_fileinfo_gated(bucket, object, opts, false, false).await?; + let mut fi = fi.into_owned(); fi.metadata.insert(AMZ_OBJECT_TAGGING.to_owned(), tags.to_owned()); if let Some(eval_metadata) = &opts.eval_metadata { @@ -3221,7 +3222,7 @@ impl SetDisks { }); } - self.update_object_meta(bucket, object, fi.clone(), disks.as_slice()).await?; + self.update_object_meta(bucket, object, fi.clone(), &disks).await?; Ok(ObjectInfo::from_file_info(&fi, bucket, object, opts.versioned || opts.version_suspended)) } @@ -4625,7 +4626,8 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { // _lock_guard = guard_opt; // } - let (mut fi, meta_arr, online_disks) = self.get_object_fileinfo(bucket, object, opts, true, false).await?; + let (fi, meta_arr, online_disks) = self.get_object_fileinfo(bucket, object, opts, true, false).await?; + let mut fi = fi.into_owned(); /*if err != nil { return Err(to_object_err(err, vec![bucket, object])); }*/ @@ -4739,7 +4741,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { cloned_fi.size, &mut writer, cloned_fi, - meta_arr, + meta_arr.into_owned(), &online_disks, set_index, pool_index, @@ -4863,7 +4865,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { }; self.invalidate_get_object_metadata_cache(bucket, object).await; let current = self.get_object_fileinfo(bucket, object, &commit_opts, true, false).await; - let (mut current_fi, _, _) = match current { + let (current_fi, _, _) = match current { Ok(current) => current, Err(err) => { drop(transition_lock_guard); @@ -4874,6 +4876,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { return Err(err); } }; + let mut current_fi = current_fi.into_owned(); let source_matches = current_fi.version_id == fi.version_id && current_fi.data_dir == fi.data_dir && current_fi.mod_time == fi.mod_time @@ -6232,9 +6235,9 @@ mod transition_commit_failure_tests { cache_key.clone(), Arc::new(GetObjectMetadataCacheEntry { created_at: Instant::now(), - fi: fi.clone(), - parts_metadata, - online_disks, + fi: Arc::new((*fi).clone()), + parts_metadata: Arc::new(parts_metadata.into_owned()), + online_disks: Arc::new(online_disks.into_owned()), read_quorum: 2, }), ) diff --git a/crates/ecstore/src/set_disk/read.rs b/crates/ecstore/src/set_disk/read.rs index 39bd74cd2..d4a98ccd4 100644 --- a/crates/ecstore/src/set_disk/read.rs +++ b/crates/ecstore/src/set_disk/read.rs @@ -116,9 +116,9 @@ impl SetDisks { .then_some(GET_METADATA_CACHE_REASON_DIST_ERASURE) } - async fn cached_get_object_fileinfo(&self, bucket: &str, object: &str) -> Option { + async fn cached_get_object_fileinfo(&self, bucket: &str, object: &str) -> Option> { match self.lookup_cached_get_object_fileinfo(bucket, object).await { - MetadataCacheLookup::Hit(entry) => Some((*entry).clone()), + MetadataCacheLookup::Hit(entry) => Some(entry), MetadataCacheLookup::Miss | MetadataCacheLookup::RejectedInsufficientQuorum => None, } } @@ -180,9 +180,9 @@ impl SetDisks { let key = GetObjectMetadataCacheKey::new(bucket, object, generation); let entry = Arc::new(GetObjectMetadataCacheEntry { created_at: Instant::now(), - fi: fi.clone(), - parts_metadata: parts_metadata.to_vec(), - online_disks: online_disks.to_vec(), + fi: Arc::new(fi.clone()), + parts_metadata: Arc::new(parts_metadata.to_vec()), + online_disks: Arc::new(online_disks.to_vec()), read_quorum, }); self.insert_get_object_metadata_cache_entry_after_insert(key, generation, entry, || {}) @@ -257,7 +257,7 @@ impl SetDisks { opts: &ObjectOptions, read_data: bool, caller_allows_early_stop: bool, - ) -> Result<(FileInfo, Vec, Vec>)> { + ) -> Result { self.get_object_fileinfo_gated(bucket, object, opts, read_data, caller_allows_early_stop) .await } @@ -274,7 +274,7 @@ impl SetDisks { opts: &ObjectOptions, read_data: bool, allow_early_stop: bool, - ) -> Result<(FileInfo, Vec, Vec>)> { + ) -> Result { let vid = opts.version_id.clone().unwrap_or_default(); let stage_metrics_enabled = rustfs_io_metrics::get_stage_metrics_enabled(); @@ -300,7 +300,11 @@ impl SetDisks { GET_STAGE_METADATA_CACHE_LOOKUP, metadata_cache_lookup_start, ); - return Ok((cached.fi.clone(), cached.parts_metadata.clone(), cached.online_disks.clone())); + return Ok(( + GetObjectMetadata::Shared(Arc::clone(&cached.fi)), + GetObjectMetadata::Shared(Arc::clone(&cached.parts_metadata)), + GetObjectMetadata::Shared(Arc::clone(&cached.online_disks)), + )); } MetadataCacheLookup::Miss => { rustfs_io_metrics::record_get_object_metadata_cache_decision( @@ -423,7 +427,11 @@ impl SetDisks { // let online_disks: Vec> = op_online_disks.iter().filter(|v| v.is_some()).cloned().collect(); - Ok((fi, parts_metadata, op_online_disks)) + Ok(( + GetObjectMetadata::Owned(fi), + GetObjectMetadata::Owned(parts_metadata), + GetObjectMetadata::Owned(op_online_disks), + )) } #[hotpath::measure(impl_type = "SetDisks")] @@ -2679,6 +2687,39 @@ mod metadata_cache_tests { assert_eq!(cached.read_quorum, 0); } + #[tokio::test] + async fn get_object_fileinfo_cache_hit_shares_cached_metadata() { + let set = new_metadata_cache_test_set().await; + let fi = valid_test_fileinfo("object"); + let parts_metadata = vec![fi.clone()]; + let online_disks = Vec::new(); + let generation = set.get_object_metadata_cache_generation("bucket", "object"); + set.cache_get_object_fileinfo(("bucket", "object"), generation, &fi, &parts_metadata, &online_disks, 0) + .await; + let cached = set + .cached_get_object_fileinfo("bucket", "object") + .await + .expect("fresh cache entry should be returned"); + + let (returned_fi, returned_parts_metadata, returned_online_disks) = set + .get_object_fileinfo("bucket", "object", &ObjectOptions::default(), true, false) + .await + .expect("cache-backed metadata lookup should succeed"); + + assert!( + matches!(returned_fi, GetObjectMetadata::Shared(ref value) if Arc::ptr_eq(value, &cached.fi)), + "cache hits must share FileInfo ownership" + ); + assert!( + matches!(returned_parts_metadata, GetObjectMetadata::Shared(ref value) if Arc::ptr_eq(value, &cached.parts_metadata)), + "cache hits must share the metadata vector" + ); + assert!( + matches!(returned_online_disks, GetObjectMetadata::Shared(ref value) if Arc::ptr_eq(value, &cached.online_disks)), + "cache hits must share the online-disk vector" + ); + } + #[tokio::test] async fn get_object_metadata_cache_rejects_deleted_and_invalid_fileinfo() { let set = new_metadata_cache_test_set().await; @@ -2718,9 +2759,9 @@ mod metadata_cache_tests { ), Arc::new(GetObjectMetadataCacheEntry { created_at: Instant::now(), - fi: fi.clone(), - parts_metadata: vec![fi], - online_disks: vec![None], + fi: Arc::new(fi.clone()), + parts_metadata: Arc::new(vec![fi]), + online_disks: Arc::new(vec![None]), read_quorum: 1, }), ) @@ -2734,9 +2775,6 @@ mod metadata_cache_tests { #[tokio::test] async fn get_object_metadata_cache_rejects_stale_entries() { - // moka handles TTL expiry automatically via time_to_live(250ms). - // This test verifies that entries inserted with the cache API are retrievable - // while fresh, and that the cache API works correctly. let set = new_metadata_cache_test_set().await; let fi = valid_test_fileinfo("object"); @@ -2748,6 +2786,14 @@ mod metadata_cache_tests { set.cached_get_object_fileinfo("bucket", "object").await.is_some(), "freshly inserted entry should be returned" ); + + tokio::time::timeout(GET_OBJECT_METADATA_CACHE_TTL + Duration::from_secs(1), async { + while set.cached_get_object_fileinfo("bucket", "object").await.is_some() { + tokio::time::sleep(Duration::from_millis(10)).await; + } + }) + .await + .expect("metadata cache entry should expire after its TTL"); } #[tokio::test] @@ -2809,9 +2855,13 @@ mod metadata_cache_tests { barrier.wait_until_paused().await; set.invalidate_get_object_metadata_cache(bucket, object).await; barrier.release(); - read.await + let (fi, parts_metadata, online_disks) = read + .await .expect("metadata read task should not panic") .expect("metadata fanout should still return its selected FileInfo"); + assert!(matches!(fi, GetObjectMetadata::Owned(_))); + assert!(matches!(parts_metadata, GetObjectMetadata::Owned(_))); + assert!(matches!(online_disks, GetObjectMetadata::Owned(_))); assert!( set.get_object_metadata_cache @@ -2858,9 +2908,9 @@ mod metadata_cache_tests { let key = GetObjectMetadataCacheKey::new("bucket", "object", generation); let entry = Arc::new(GetObjectMetadataCacheEntry { created_at: Instant::now(), - fi: fi.clone(), - parts_metadata: vec![fi], - online_disks: Vec::new(), + fi: Arc::new(fi.clone()), + parts_metadata: Arc::new(vec![fi]), + online_disks: Arc::new(Vec::new()), read_quorum: 0, }); @@ -2968,9 +3018,9 @@ mod metadata_cache_tests { let entry = |fi: FileInfo| { Arc::new(GetObjectMetadataCacheEntry { created_at: Instant::now(), - parts_metadata: vec![fi.clone()], - fi, - online_disks: Vec::new(), + parts_metadata: Arc::new(vec![fi.clone()]), + fi: Arc::new(fi), + online_disks: Arc::new(Vec::new()), read_quorum: 0, }) }; diff --git a/crates/ecstore/src/set_disk/replication.rs b/crates/ecstore/src/set_disk/replication.rs index 732666296..e1f1d512a 100644 --- a/crates/ecstore/src/set_disk/replication.rs +++ b/crates/ecstore/src/set_disk/replication.rs @@ -77,9 +77,10 @@ impl SetDisks { version_suspended: opts.version_suspended, ..Default::default() }; - let (mut fi, _, disks) = self + let (fi, _, disks) = self .get_object_fileinfo_gated(bucket, object, &read_opts, false, false) .await?; + let mut fi = fi.into_owned(); if let Some(expected_operation_id) = expected_operation_id { require_restore_operation_id(&fi.metadata, expected_operation_id)?; } @@ -101,7 +102,7 @@ impl SetDisks { bucket, object, fi.clone(), - disks.as_slice(), + &disks, &UpdateMetadataOpts { replace_user_metadata: true, ..Default::default() @@ -143,9 +144,10 @@ impl SetDisks { version_suspended: opts.version_suspended, ..Default::default() }; - let (mut fi, _, disks) = self + let (fi, _, disks) = self .get_object_fileinfo_gated(bucket, object, &read_opts, false, false) .await?; + let mut fi = fi.into_owned(); if let Some(expected_operation_id) = expected_operation_id { match restore_operation_id_from_metadata(&fi.metadata)? { Some(actual_operation_id) if actual_operation_id == expected_operation_id => {} @@ -170,7 +172,7 @@ impl SetDisks { bucket, object, fi, - disks.as_slice(), + &disks, &UpdateMetadataOpts { replace_user_metadata: true, ..Default::default()