From 2f9c75d04fce2f52a4e411b2227c86ec6d1795b1 Mon Sep 17 00:00:00 2001 From: houseme Date: Mon, 31 Aug 2026 00:15:10 +0800 Subject: [PATCH] perf(ecstore): reuse prepared metadata across pools (#6889) Co-authored-by: heihutu --- crates/ecstore/src/store/object.rs | 380 ++++++++++++++++-- crates/ecstore/src/store/rebalance.rs | 229 ++++++++++- crates/ecstore/src/store/rebalance/support.rs | 33 +- 3 files changed, 614 insertions(+), 28 deletions(-) diff --git a/crates/ecstore/src/store/object.rs b/crates/ecstore/src/store/object.rs index 6e3ce4214..51211590d 100644 --- a/crates/ecstore/src/store/object.rs +++ b/crates/ecstore/src/store/object.rs @@ -1817,16 +1817,21 @@ impl ECStore { let metadata = pool.prepare_get_object_reader_metadata(bucket, &object, &opts).await?; (metadata, pool) } else { - let (_, pool_idx) = self - .get_latest_accessible_object_info_with_idx(bucket, &object, &opts) - .await?; - let pool = self - .pools - .get(pool_idx) - .cloned() - .ok_or_else(|| Error::other(format!("resolved GET pool index {pool_idx} is out of bounds")))?; - let metadata = pool.prepare_get_object_reader_metadata(bucket, &object, &opts).await?; - (metadata, pool) + // Keep the large multi-pool selection future off the caller stack. + // Debug builds otherwise exceed the common 2 MiB worker stack. + Box::pin(async { + let (metadata, pool_idx) = self.prepare_latest_object_metadata_with_idx(bucket, &object, &opts).await?; + if let Some(error) = latest_object_access_delete_marker_error(bucket, &object, metadata.object_info(), &opts) { + return Err(error); + } + let pool = self + .pools + .get(pool_idx) + .cloned() + .ok_or_else(|| Error::other(format!("resolved GET pool index {pool_idx} is out of bounds")))?; + Ok((metadata, pool)) + }) + .await? }; Ok(PreparedGetObjectReader { @@ -2518,12 +2523,18 @@ impl ECStore { .get_object_reader(bucket, object.as_ref(), range, h, &opts) .await? } else { - let (_, idx) = self - .get_latest_accessible_object_info_with_idx(bucket, &object, &opts) - .await?; - self.pools[idx] - .get_object_reader(bucket, object.as_ref(), range, h, &opts) - .await? + // Keep selection plus prepared-open state off the caller stack. + // Debug builds otherwise exceed the common 2 MiB worker stack. + Box::pin(async { + let (metadata, idx) = self.prepare_latest_object_metadata_with_idx(bucket, &object, &opts).await?; + if let Some(error) = latest_object_access_delete_marker_error(bucket, &object, metadata.object_info(), &opts) { + return Err(error); + } + self.pools[idx] + .get_object_reader_with_prepared_metadata(bucket, object.as_ref(), range, h, &opts, metadata) + .await + }) + .await? }; Ok(Self::attach_read_lock_guard(reader, read_lock_guard)) @@ -3914,6 +3925,7 @@ mod tests { ReplicationState, ReplicationStatusType, VersionPurgeStatusType, replication_state_to_filemeta, replication_statuses_map, version_purge_statuses_map, }; + use crate::core::pools::{PoolDecommissionInfo, PoolStatus}; use crate::core::sets::make_local_two_set_sets_with_ctx; use crate::ecstore_validation_blackbox::{make_local_set_disks, make_local_set_disks_with_ctx}; use crate::layout::{ @@ -6085,13 +6097,62 @@ mod tests { #[tokio::test] #[serial_test::serial(body_cache_hook)] - async fn prepared_reader_resolves_object_from_second_pool() { + async fn prepared_reader_reuses_metadata_across_three_pools() { + let (_first_dirs, first_set) = make_local_set_disks(4, 2).await; + let (_second_dirs, second_set) = make_local_set_disks(4, 2).await; + let (_third_dirs, third_set) = make_local_set_disks(4, 2).await; + let store = new_prepared_reader_test_store(&[first_set, second_set, third_set]).await; + let bucket = "prepared-reader-three-pools"; + let object = "object.bin"; + let payload = b"prepared-reader-three-pool-payload-".repeat(40_000); + let opts = ObjectOptions { + no_lock: true, + ..Default::default() + }; + + for pool in &store.pools { + pool.make_bucket(bucket, &MakeBucketOptions::default()) + .await + .expect("bucket should be created in each pool"); + } + let mut put_reader = PutObjReader::from_vec(payload.clone()); + store.pools[2] + .put_object(bucket, object, &mut put_reader, &opts) + .await + .expect("object should be written only to the third pool"); + + let calls = disk_call_counters::observe(object); + let prepared = store + .prepare_get_object_reader(bucket, object, None, HeaderMap::new(), &opts) + .await + .expect("prepared reader should resolve the third-pool object"); + assert_eq!(prepared.object_info().size, payload.len() as i64); + let metadata_calls = calls.total(disk_call_counters::KIND_READ_VERSION); + assert_eq!(metadata_calls, 12, "three 4-disk pools must fan out metadata exactly once each"); + let mut reader = prepared.into_reader().await.expect("prepared body reader should open"); + assert_eq!( + calls.total(disk_call_counters::KIND_READ_VERSION), + metadata_calls, + "the selected pool must reuse its prepared metadata" + ); + let mut restored = Vec::new(); + reader + .stream + .read_to_end(&mut restored) + .await + .expect("prepared body should stream"); + assert_eq!(restored, payload); + } + + #[tokio::test] + #[serial_test::serial(body_cache_hook)] + async fn legacy_reader_reuses_selected_pool_metadata() { let (_first_dirs, first_set) = make_local_set_disks(4, 2).await; let (_second_dirs, second_set) = make_local_set_disks(4, 2).await; let store = new_prepared_reader_test_store(&[first_set, second_set]).await; - let bucket = "prepared-reader-second-pool"; + let bucket = "legacy-reader-second-pool"; let object = "object.bin"; - let payload = b"prepared-reader-second-pool-payload-".repeat(40_000); + let payload = b"legacy-reader-second-pool-payload-".repeat(40_000); let opts = ObjectOptions { no_lock: true, ..Default::default() @@ -6108,18 +6169,287 @@ mod tests { .await .expect("object should be written only to the second pool"); - let prepared = store - .prepare_get_object_reader(bucket, object, None, HeaderMap::new(), &opts) + clear_get_object_body_cache_hook(); + let hook = Arc::new(CountingMissHook { + calls: AtomicUsize::new(0), + }); + register_get_object_body_cache_hook(Arc::clone(&hook) as Arc); + let _hook_guard = BodyCacheHookGuard; + + let calls = disk_call_counters::observe(object); + let mut reader = store + .handle_get_object_reader(bucket, object, None, HeaderMap::new(), &opts) .await - .expect("prepared reader should resolve the second-pool object"); - assert_eq!(prepared.object_info().size, payload.len() as i64); - let mut reader = prepared.into_reader().await.expect("prepared body reader should open"); + .expect("legacy reader should resolve the second-pool object"); + assert_eq!( + hook.calls.load(Ordering::Relaxed), + 1, + "legacy reader must probe the body cache exactly once" + ); + assert_eq!(reader.body_source, GetObjectBodySource::HookMissed); + assert!( + calls.total(disk_call_counters::KIND_READ_VERSION) <= 8, + "legacy reader must fan out each 4-disk pool at most once" + ); let mut restored = Vec::new(); reader .stream .read_to_end(&mut restored) .await - .expect("prepared body should stream"); + .expect("legacy reader body should stream"); + assert_eq!(restored, payload); + } + + fn prepared_pool_test_status(id: usize, suspended: bool) -> PoolStatus { + PoolStatus { + id, + cmd_line: format!("prepared-pool-{id}"), + last_update: OffsetDateTime::now_utc(), + decommission: suspended.then(|| PoolDecommissionInfo { + start_time: Some(OffsetDateTime::now_utc()), + ..Default::default() + }), + } + } + + #[tokio::test] + #[serial_test::serial(body_cache_hook)] + async fn prepared_reader_refetches_when_final_pool_state_changes_winner() { + let (_dirs, set_disks) = make_local_set_disks(4, 2).await; + let store = Arc::new(new_prepared_reader_test_store(&[Arc::clone(&set_disks), Arc::clone(&set_disks)]).await); + let bucket = "prepared-reader-pool-state-fallback"; + let object = "object.bin"; + let payload = b"pool-state fallback payload".repeat(8_000); + let opts = ObjectOptions { + no_lock: true, + ..Default::default() + }; + + set_disks + .make_bucket(bucket, &MakeBucketOptions::default()) + .await + .expect("bucket should be created"); + let mut put_reader = PutObjReader::from_vec(payload.clone()); + set_disks + .put_object(bucket, object, &mut put_reader, &opts) + .await + .expect("shared object should be written"); + + let calls = disk_call_counters::observe(object); + let barrier = crate::store::rebalance::PreparedPoolReadFallbackBarrier::install(object, false); + let read_store = Arc::clone(&store); + let read_opts = opts.clone(); + let read = tokio::spawn(async move { + read_store + .prepare_get_object_reader(bucket, object, None, HeaderMap::new(), &read_opts) + .await + }); + barrier.wait_after_fanout().await; + *store.pool_meta.write().await = PoolMeta { + pools: vec![prepared_pool_test_status(0, false), prepared_pool_test_status(1, true)], + ..Default::default() + }; + barrier.release_after_fanout(); + + let prepared = read + .await + .expect("prepared read task should not panic") + .expect("final active pool should be refetched"); + assert!(Arc::ptr_eq(&prepared.pool, &store.pools[0])); + assert_eq!( + calls.total(disk_call_counters::KIND_READ_VERSION), + 12, + "two initial 4-disk fanouts plus one fallback refetch are required" + ); + let mut reader = prepared.into_reader().await.expect("fallback body reader should open"); + let mut restored = Vec::new(); + reader + .stream + .read_to_end(&mut restored) + .await + .expect("fallback body should stream"); + assert_eq!(restored, payload); + } + + #[tokio::test] + #[serial_test::serial(body_cache_hook)] + async fn prepared_reader_fallback_rejects_generation_change_before_refetch() { + let (_dirs, set_disks) = make_local_set_disks(4, 2).await; + let store = Arc::new(new_prepared_reader_test_store(&[Arc::clone(&set_disks), Arc::clone(&set_disks)]).await); + let bucket = "prepared-reader-pool-state-generation-change"; + let object = "object.bin"; + let opts = ObjectOptions { + no_lock: true, + ..Default::default() + }; + + set_disks + .make_bucket(bucket, &MakeBucketOptions::default()) + .await + .expect("bucket should be created"); + let mut initial_reader = PutObjReader::from_vec(b"initial generation".to_vec()); + set_disks + .put_object(bucket, object, &mut initial_reader, &opts) + .await + .expect("initial object should be written"); + + let barrier = crate::store::rebalance::PreparedPoolReadFallbackBarrier::install(object, true); + let read_store = Arc::clone(&store); + let read_opts = opts.clone(); + let read = tokio::spawn(async move { + read_store + .prepare_get_object_reader(bucket, object, None, HeaderMap::new(), &read_opts) + .await + }); + barrier.wait_after_fanout().await; + *store.pool_meta.write().await = PoolMeta { + pools: vec![prepared_pool_test_status(0, false), prepared_pool_test_status(1, true)], + ..Default::default() + }; + barrier.release_after_fanout(); + barrier.wait_before_refetch().await; + + let mut replacement_reader = PutObjReader::from_vec(b"replacement generation".to_vec()); + set_disks + .put_object(bucket, object, &mut replacement_reader, &opts) + .await + .expect("replacement generation should be written before fallback refetch"); + barrier.release_before_refetch(); + + let error = match read.await.expect("prepared read task should not panic") { + Ok(_) => panic!("changed fallback generation must not be accepted"), + Err(error) => error, + }; + assert_eq!(error, Error::ErasureReadQuorum); + } + + #[tokio::test] + #[serial_test::serial(body_cache_hook)] + async fn prepared_reader_rejects_latest_delete_marker_without_refetching_metadata() { + let ctx = Arc::new(crate::runtime::instance::InstanceContext::new()); + let (_first_dirs, first_set) = make_local_set_disks_with_ctx(4, 2, Arc::clone(&ctx)).await; + let (_second_dirs, second_set) = make_local_set_disks_with_ctx(4, 2, Arc::clone(&ctx)).await; + let store = new_prepared_reader_test_store_with_ctx(&[Arc::clone(&first_set), Arc::clone(&second_set)], ctx).await; + let bucket = "prepared-reader-latest-delete-marker"; + let object = "versioned-object.bin"; + let versioned_opts = ObjectOptions { + no_lock: true, + versioned: true, + object_lock_config_snapshot: Some(Arc::new(ObjectLockConfigSnapshot::new(ObjectLockConfigState::ConfirmedAbsent))), + ..Default::default() + }; + + for set_disks in [&first_set, &second_set] { + set_disks + .make_bucket(bucket, &MakeBucketOptions::default()) + .await + .expect("bucket should be created"); + } + let mut older_reader = PutObjReader::from_vec(b"older visible generation".to_vec()); + first_set + .put_object(bucket, object, &mut older_reader, &versioned_opts) + .await + .expect("older object should be written"); + let mut hidden_reader = PutObjReader::from_vec(b"hidden generation".to_vec()); + second_set + .put_object(bucket, object, &mut hidden_reader, &versioned_opts) + .await + .expect("newer object should be written"); + let marker = second_set + .delete_object(bucket, object, versioned_opts.clone()) + .await + .expect("delete marker should be committed"); + assert!(marker.delete_marker); + + let calls = disk_call_counters::observe(object); + let error = match store + .prepare_get_object_reader( + bucket, + object, + None, + HeaderMap::new(), + &ObjectOptions { + no_lock: true, + versioned: true, + ..Default::default() + }, + ) + .await + { + Ok(_) => panic!("latest delete marker should hide the older live object"), + Err(error) => error, + }; + + assert!(is_err_object_not_found(&error)); + assert!( + calls.total(disk_call_counters::KIND_READ_VERSION) <= 8, + "delete-marker resolution must fan out each pool at most once" + ); + } + + #[tokio::test] + #[serial_test::serial(body_cache_hook)] + async fn prepared_reader_explicit_version_reuses_the_matching_pool_metadata() { + let (_first_dirs, first_set) = make_local_set_disks(4, 2).await; + let (_second_dirs, second_set) = make_local_set_disks(4, 2).await; + let store = new_prepared_reader_test_store(&[Arc::clone(&first_set), Arc::clone(&second_set)]).await; + let bucket = "prepared-reader-explicit-version"; + let object = "versioned-object.bin"; + let payload = b"explicit version from first pool".repeat(8_000); + let versioned_opts = ObjectOptions { + no_lock: true, + versioned: true, + object_lock_config_snapshot: Some(Arc::new(ObjectLockConfigSnapshot::new(ObjectLockConfigState::ConfirmedAbsent))), + ..Default::default() + }; + + for set_disks in [&first_set, &second_set] { + set_disks + .make_bucket(bucket, &MakeBucketOptions::default()) + .await + .expect("bucket should be created"); + } + let mut first_reader = PutObjReader::from_vec(payload.clone()); + let first = first_set + .put_object(bucket, object, &mut first_reader, &versioned_opts) + .await + .expect("requested version should be written to the first pool"); + let mut second_reader = PutObjReader::from_vec(b"different pool version".to_vec()); + second_set + .put_object(bucket, object, &mut second_reader, &versioned_opts) + .await + .expect("a different version should be written to the second pool"); + + let requested_version = first + .version_id + .expect("versioned PUT should return a version id") + .to_string(); + let read_opts = ObjectOptions { + no_lock: true, + versioned: true, + version_id: Some(requested_version), + ..Default::default() + }; + let calls = disk_call_counters::observe(object); + let prepared = store + .prepare_get_object_reader(bucket, object, None, HeaderMap::new(), &read_opts) + .await + .expect("explicit version should resolve from the matching pool"); + assert_eq!(prepared.object_info().version_id, first.version_id); + let metadata_calls = calls.total(disk_call_counters::KIND_READ_VERSION); + assert!(metadata_calls <= 8, "explicit-version lookup must fan out each pool at most once"); + + let mut reader = prepared + .into_reader() + .await + .expect("prepared explicit-version body should open"); + assert_eq!(calls.total(disk_call_counters::KIND_READ_VERSION), metadata_calls); + let mut restored = Vec::new(); + reader + .stream + .read_to_end(&mut restored) + .await + .expect("explicit-version body should stream"); assert_eq!(restored, payload); } diff --git a/crates/ecstore/src/store/rebalance.rs b/crates/ecstore/src/store/rebalance.rs index ddcd918ac..6c2fe8642 100644 --- a/crates/ecstore/src/store/rebalance.rs +++ b/crates/ecstore/src/store/rebalance.rs @@ -18,18 +18,117 @@ use crate::core::pools::merge_pool_status_refresh; use crate::layout::pool_space::{ServerPoolsAvailableSpace, build_server_pools_available_space}; use crate::runtime::sources as runtime_sources; use crate::storage_api_contracts::{admin::StorageAdminApi, namespace::NamespaceLocking as _, object::ObjectOperations as _}; +use futures::stream::{FuturesUnordered, StreamExt}; pub(in crate::store) mod support; const LOG_COMPONENT_ECSTORE: &str = "ecstore"; const LOG_SUBSYSTEM_POOLS: &str = "pools"; const EVENT_POOL_META_RELOAD: &str = "pool_meta_reload"; + +#[cfg(test)] +struct PreparedPoolReadFallbackBarrierState { + object: String, + pause_before_refetch: bool, + fanout_arrived: tokio::sync::Notify, + fanout_release: tokio::sync::Notify, + refetch_arrived: tokio::sync::Notify, + refetch_release: tokio::sync::Notify, +} + +#[cfg(test)] +pub(in crate::store) struct PreparedPoolReadFallbackBarrier { + state: Arc, +} + +#[cfg(test)] +static PREPARED_POOL_READ_FALLBACK_BARRIER: std::sync::OnceLock< + std::sync::Mutex>>, +> = std::sync::OnceLock::new(); + +#[cfg(test)] +impl PreparedPoolReadFallbackBarrier { + pub(in crate::store) fn install(object: &str, pause_before_refetch: bool) -> Self { + let state = Arc::new(PreparedPoolReadFallbackBarrierState { + object: object.to_string(), + pause_before_refetch, + fanout_arrived: tokio::sync::Notify::new(), + fanout_release: tokio::sync::Notify::new(), + refetch_arrived: tokio::sync::Notify::new(), + refetch_release: tokio::sync::Notify::new(), + }); + *PREPARED_POOL_READ_FALLBACK_BARRIER + .get_or_init(|| std::sync::Mutex::new(None)) + .lock() + .expect("prepared pool read fallback barrier must not be poisoned") = Some(Arc::clone(&state)); + Self { state } + } + + pub(in crate::store) async fn wait_after_fanout(&self) { + self.state.fanout_arrived.notified().await; + } + + pub(in crate::store) fn release_after_fanout(&self) { + self.state.fanout_release.notify_one(); + } + + pub(in crate::store) async fn wait_before_refetch(&self) { + self.state.refetch_arrived.notified().await; + } + + pub(in crate::store) fn release_before_refetch(&self) { + self.state.refetch_release.notify_one(); + } +} + +#[cfg(test)] +impl Drop for PreparedPoolReadFallbackBarrier { + fn drop(&mut self) { + self.state.fanout_release.notify_waiters(); + self.state.refetch_release.notify_waiters(); + if let Some(barrier) = PREPARED_POOL_READ_FALLBACK_BARRIER.get() { + *barrier + .lock() + .expect("prepared pool read fallback barrier must not be poisoned") = None; + } + } +} + +#[cfg(test)] +async fn pause_prepared_pool_read_after_fanout(object: &str) { + let state = PREPARED_POOL_READ_FALLBACK_BARRIER + .get_or_init(|| std::sync::Mutex::new(None)) + .lock() + .expect("prepared pool read fallback barrier must not be poisoned") + .as_ref() + .filter(|state| state.object == object) + .cloned(); + if let Some(state) = state { + state.fanout_arrived.notify_one(); + state.fanout_release.notified().await; + } +} + +#[cfg(test)] +async fn pause_prepared_pool_read_before_refetch(object: &str) { + let state = PREPARED_POOL_READ_FALLBACK_BARRIER + .get_or_init(|| std::sync::Mutex::new(None)) + .lock() + .expect("prepared pool read fallback barrier must not be poisoned") + .as_ref() + .filter(|state| state.object == object && state.pause_before_refetch) + .cloned(); + if let Some(state) = state { + state.refetch_arrived.notify_one(); + state.refetch_release.notified().await; + } +} #[cfg(test)] use support::resolve_latest_object_info_candidates; use support::{ LatestObjectInfoCandidate, PoolErr, PoolObjInfo, RebalanceDeletePoolResult, pool_lookup_not_found_error, rebalance_disk_set_lookup_error, resolve_latest_object_info_candidates_with_pool_state, resolve_rebalance_delete_from_all_pools_result, resolve_rebalance_delete_from_all_pools_results, - resolve_store_rebalance_pool_meta_reload_result, + resolve_store_rebalance_pool_meta_reload_result, validate_prepared_pool_refetch_identity, }; #[derive(Debug, Default, Eq, PartialEq)] @@ -675,6 +774,134 @@ impl ECStore { resolve_latest_object_info_candidates_with_pool_state(candidates, &suspended_pools, bucket, object, opts) } + pub(super) async fn prepare_latest_object_metadata_with_idx( + &self, + bucket: &str, + object: &str, + opts: &ObjectOptions, + ) -> Result<(crate::set_disk::PreparedGetObjectMetadata, usize)> { + let suspended_pools = if opts.skip_decommissioned { + let pool_meta = self.pool_meta.read().await; + Some( + (0..self.pools.len()) + .map(|idx| pool_meta.is_suspended(idx)) + .collect::>(), + ) + } else { + None + }; + let mut futures = FuturesUnordered::new(); + for (idx, pool) in self.pools.iter().enumerate() { + if suspended_pools.as_ref().is_some_and(|pools| pools[idx]) { + continue; + } + + if opts.skip_rebalancing && self.is_pool_rebalancing(idx).await { + continue; + } + + futures.push(async move { + let result = pool + .prepare_get_object_reader_metadata(bucket, object, opts) + .await + .map_err(|err| to_object_err(err, vec![bucket, object])); + (idx, result) + }); + } + + let mut candidates = (0..self.pools.len()).map(|_| None).collect::>(); + // Retain one provisional winner. Other pools only need their lightweight + // identity for final conflict checks; if pool state changes while the + // fanout runs, the final winner is refetched and revalidated below. + let mut latest_prepared = None; + let mut latest_mod_time = None; + let mut provisional_dynamic_pool_state = None; + while let Some((idx, result)) = futures.next().await { + match result { + Ok(metadata) => { + let mod_time = metadata.object_info().mod_time.unwrap_or(OffsetDateTime::UNIX_EPOCH); + let info = metadata.object_info().clone(); + let retain = match (latest_mod_time, latest_prepared.as_ref()) { + (None, _) => true, + (Some(current), _) if mod_time > current => true, + (Some(current), _) if mod_time < current => false, + (Some(_), Some((current_idx, _))) => { + if suspended_pools.is_none() && provisional_dynamic_pool_state.is_none() { + let pool_meta = self.pool_meta.read().await; + provisional_dynamic_pool_state = Some( + (0..self.pools.len()) + .map(|pool_idx| pool_meta.is_suspended(pool_idx)) + .collect::>(), + ); + } + let provisional_pool_state = suspended_pools + .as_ref() + .or(provisional_dynamic_pool_state.as_ref()) + .ok_or_else(|| Error::other("GET pool state snapshot is unavailable"))?; + let new_key = (provisional_pool_state.get(idx).copied().unwrap_or(false), std::cmp::Reverse(idx)); + let current_key = ( + provisional_pool_state.get(*current_idx).copied().unwrap_or(false), + std::cmp::Reverse(*current_idx), + ); + new_key < current_key + } + (Some(_), None) => true, + }; + if retain { + if latest_mod_time.is_none_or(|current| mod_time > current) { + latest_mod_time = Some(mod_time); + } + latest_prepared = Some((idx, metadata)); + } + candidates[idx] = Some(LatestObjectInfoCandidate { + info: Some(info), + idx, + err: None, + }); + } + Err(err) => { + candidates[idx] = Some(LatestObjectInfoCandidate { + info: None, + idx, + err: Some(err), + }); + } + } + } + + #[cfg(test)] + pause_prepared_pool_read_after_fanout(object).await; + + let suspended_pools = match suspended_pools { + Some(pools) => pools, + None => { + let pool_meta = self.pool_meta.read().await; + (0..self.pools.len()) + .map(|idx| pool_meta.is_suspended(idx)) + .collect::>() + } + }; + + let candidates = candidates.into_iter().flatten().collect(); + let (winner_info, winner_idx) = + resolve_latest_object_info_candidates_with_pool_state(candidates, &suspended_pools, bucket, object, opts)?; + if let Some((prepared_idx, metadata)) = latest_prepared + && prepared_idx == winner_idx + { + return Ok((metadata, winner_idx)); + } + + let pool = self.pools.get(winner_idx).ok_or(Error::ErasureReadQuorum)?; + #[cfg(test)] + pause_prepared_pool_read_before_refetch(object).await; + let metadata = pool + .prepare_get_object_reader_metadata(bucket, object, opts) + .await + .map_err(|err| to_object_err(err, vec![bucket, object]))?; + validate_prepared_pool_refetch_identity(&winner_info, metadata.object_info())?; + Ok((metadata, winner_idx)) + } + pub(super) async fn delete_object_from_all_pools( &self, bucket: &str, diff --git a/crates/ecstore/src/store/rebalance/support.rs b/crates/ecstore/src/store/rebalance/support.rs index e1a773a54..376ac05b9 100644 --- a/crates/ecstore/src/store/rebalance/support.rs +++ b/crates/ecstore/src/store/rebalance/support.rs @@ -218,7 +218,7 @@ fn same_user_defined_identity(left: &ObjectInfo, right: &ObjectInfo) -> bool { /// excluded. The selected winner still carries the chosen pool's layout, while /// the remaining read-visible fields must agree before the pool index can /// provide a deterministic tie-break. -fn same_latest_object_info_identity(left: &ObjectInfo, right: &ObjectInfo) -> bool { +pub(super) fn same_latest_object_info_identity(left: &ObjectInfo, right: &ObjectInfo) -> bool { let same_read_surface = left.bucket == right.bucket && left.name == right.name && left.is_dir == right.is_dir @@ -277,6 +277,14 @@ fn same_latest_object_info_identity(left: &ObjectInfo, right: &ObjectInfo) -> bo } } +pub(super) fn validate_prepared_pool_refetch_identity(expected: &ObjectInfo, refetched: &ObjectInfo) -> Result<()> { + if same_latest_object_info_identity(expected, refetched) { + Ok(()) + } else { + Err(Error::ErasureReadQuorum) + } +} + #[cfg(test)] pub(super) fn resolve_latest_object_info_candidates( candidates: Vec, @@ -328,7 +336,11 @@ pub(super) fn resolve_latest_object_info_candidates_with_pool_state( return Err(Error::ErasureReadQuorum); } - return Ok((winner_info.clone(), winner.idx)); + let winner = latest_candidates.swap_remove(0); + let Some(winner_info) = winner.info else { + return Err(Error::ErasureReadQuorum); + }; + return Ok((winner_info, winner.idx)); } for candidate in candidates { @@ -347,6 +359,23 @@ pub(super) fn resolve_latest_object_info_candidates_with_pool_state( mod tests { use super::*; + #[test] + fn prepared_pool_refetch_identity_fails_closed_on_generation_change() { + let expected = ObjectInfo { + mod_time: Some(OffsetDateTime::from_unix_timestamp(10).expect("test timestamp should be valid")), + version_id: Some(uuid::Uuid::from_u128(1)), + etag: Some("etag-a".to_string()), + ..Default::default() + }; + let mut refetched = expected.clone(); + refetched.etag = Some("etag-b".to_string()); + + let error = validate_prepared_pool_refetch_identity(&expected, &refetched) + .expect_err("refetched metadata from a changed generation must fail closed"); + + assert_eq!(error, Error::ErasureReadQuorum); + } + #[test] fn rebalance_delete_result_preserves_precondition_failed() { let err = resolve_rebalance_delete_from_all_pools_result(Err(Error::PreconditionFailed), "bucket", "object")