mirror of
https://github.com/rustfs/rustfs.git
synced 2026-09-01 17:58:22 +00:00
perf(ecstore): reuse prepared metadata across pools (#6889)
Co-authored-by: heihutu <heihutu@gmail.com>
This commit is contained in:
@@ -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<dyn GetObjectBodyCacheHook>);
|
||||
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);
|
||||
}
|
||||
|
||||
|
||||
@@ -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<PreparedPoolReadFallbackBarrierState>,
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
static PREPARED_POOL_READ_FALLBACK_BARRIER: std::sync::OnceLock<
|
||||
std::sync::Mutex<Option<Arc<PreparedPoolReadFallbackBarrierState>>>,
|
||||
> = 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::<Vec<_>>(),
|
||||
)
|
||||
} 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::<Vec<_>>();
|
||||
// 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::<Vec<_>>(),
|
||||
);
|
||||
}
|
||||
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::<Vec<_>>()
|
||||
}
|
||||
};
|
||||
|
||||
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,
|
||||
|
||||
@@ -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<LatestObjectInfoCandidate>,
|
||||
@@ -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")
|
||||
|
||||
Reference in New Issue
Block a user