mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-29 00:17:11 +00:00
fix(rebalance): fence writer pool lookups (#5845)
This commit is contained in:
@@ -602,7 +602,7 @@ mod tests {
|
|||||||
error::{Error, Result, StorageError},
|
error::{Error, Result, StorageError},
|
||||||
layout::endpoints::{EndpointServerPools, Endpoints, PoolEndpoints},
|
layout::endpoints::{EndpointServerPools, Endpoints, PoolEndpoints},
|
||||||
object_api::{GetObjectReader, ObjectInfo, ObjectOptions, PutObjReader},
|
object_api::{GetObjectReader, ObjectInfo, ObjectOptions, PutObjReader},
|
||||||
services::rebalance::RebalanceMeta,
|
services::rebalance::{RebalStatus, RebalanceInfo, RebalanceMeta, RebalanceStats},
|
||||||
storage_api_contracts::{
|
storage_api_contracts::{
|
||||||
bucket::{BucketOperations as _, MakeBucketOptions},
|
bucket::{BucketOperations as _, MakeBucketOptions},
|
||||||
multipart::MultipartOperations as _,
|
multipart::MultipartOperations as _,
|
||||||
@@ -1155,6 +1155,183 @@ mod tests {
|
|||||||
(instance_ctx, store, shutdown)
|
(instance_ctx, store, shutdown)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn active_rebalance_meta_for_pool(pool_count: usize, active_pool_idx: usize) -> RebalanceMeta {
|
||||||
|
let now = OffsetDateTime::now_utc();
|
||||||
|
let mut pool_stats = vec![RebalanceStats::default(); pool_count];
|
||||||
|
pool_stats[active_pool_idx] = RebalanceStats {
|
||||||
|
participating: true,
|
||||||
|
info: RebalanceInfo {
|
||||||
|
start_time: Some(now),
|
||||||
|
status: RebalStatus::Started,
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
|
||||||
|
RebalanceMeta {
|
||||||
|
id: uuid::Uuid::new_v4().to_string(),
|
||||||
|
pool_stats,
|
||||||
|
..Default::default()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
#[serial_test::serial(storage_class_env)]
|
||||||
|
async fn tag_updates_skip_active_rebalance_source_pool() {
|
||||||
|
let temp_dir = tempfile::tempdir().expect("create writer-fencing store dir");
|
||||||
|
let (_ctx, store, shutdown) =
|
||||||
|
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "writer-fencing-tags", &[4, 4])).await;
|
||||||
|
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
|
||||||
|
|
||||||
|
let bucket = format!("writer-fencing-tags-{}", uuid::Uuid::new_v4());
|
||||||
|
let object = "tagged-object.bin";
|
||||||
|
store
|
||||||
|
.make_bucket(&bucket, &MakeBucketOptions::default())
|
||||||
|
.await
|
||||||
|
.expect("create writer fencing bucket");
|
||||||
|
|
||||||
|
let old_time = OffsetDateTime::from_unix_timestamp(1_700_000_000).expect("fixed timestamp should be valid");
|
||||||
|
let newer_time = old_time + time::Duration::seconds(10);
|
||||||
|
let mut source_reader = PutObjReader::from_vec(b"source-body".to_vec());
|
||||||
|
store.pools[0]
|
||||||
|
.put_object(
|
||||||
|
&bucket,
|
||||||
|
object,
|
||||||
|
&mut source_reader,
|
||||||
|
&ObjectOptions {
|
||||||
|
mod_time: Some(newer_time),
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.expect("write newer source object");
|
||||||
|
let mut target_reader = PutObjReader::from_vec(b"target-body".to_vec());
|
||||||
|
store.pools[1]
|
||||||
|
.put_object(
|
||||||
|
&bucket,
|
||||||
|
object,
|
||||||
|
&mut target_reader,
|
||||||
|
&ObjectOptions {
|
||||||
|
mod_time: Some(old_time),
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.expect("write older target object");
|
||||||
|
|
||||||
|
*store.rebalance_meta.write().await = Some(active_rebalance_meta_for_pool(store.pools.len(), 0));
|
||||||
|
assert!(store.is_pool_rebalancing(0).await, "pool 0 must be marked as an active rebalance source");
|
||||||
|
|
||||||
|
let tags = "rebalance=target";
|
||||||
|
assert_ne!(
|
||||||
|
store.pools[0]
|
||||||
|
.get_object_tags(&bucket, object, &ObjectOptions::default())
|
||||||
|
.await
|
||||||
|
.expect("source object tags should be readable before update"),
|
||||||
|
tags,
|
||||||
|
"source object must start without the target tag"
|
||||||
|
);
|
||||||
|
assert_ne!(
|
||||||
|
store.pools[1]
|
||||||
|
.get_object_tags(&bucket, object, &ObjectOptions::default())
|
||||||
|
.await
|
||||||
|
.expect("target object tags should be readable before update"),
|
||||||
|
tags,
|
||||||
|
"target object must start without the target tag"
|
||||||
|
);
|
||||||
|
let selected_pool = store
|
||||||
|
.get_pool_idx_existing_with_opts(
|
||||||
|
&bucket,
|
||||||
|
object,
|
||||||
|
&ObjectOptions {
|
||||||
|
no_lock: true,
|
||||||
|
metadata_chg: true,
|
||||||
|
skip_decommissioned: true,
|
||||||
|
skip_rebalancing: true,
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.expect("writer lookup should select an existing non-rebalancing pool");
|
||||||
|
assert_eq!(selected_pool, 1, "writer lookup must skip active rebalance pool 0");
|
||||||
|
|
||||||
|
let updated = store
|
||||||
|
.put_object_tags(&bucket, object, tags, &ObjectOptions::default())
|
||||||
|
.await
|
||||||
|
.expect("tag update should use the non-rebalancing target pool");
|
||||||
|
assert_eq!(
|
||||||
|
updated.mod_time,
|
||||||
|
Some(old_time),
|
||||||
|
"tag update must return the non-rebalancing pool object rather than the newer active source"
|
||||||
|
);
|
||||||
|
|
||||||
|
let target_tags = store.pools[1]
|
||||||
|
.get_object_tags(&bucket, object, &ObjectOptions::default())
|
||||||
|
.await
|
||||||
|
.expect("target object tags should be readable");
|
||||||
|
assert_eq!(target_tags, tags, "non-rebalancing pool must receive writer tag updates");
|
||||||
|
|
||||||
|
shutdown.cancel();
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
#[serial_test::serial(storage_class_env)]
|
||||||
|
async fn multipart_listing_skips_active_rebalance_source_pool() {
|
||||||
|
let temp_dir = tempfile::tempdir().expect("create multipart writer-fencing store dir");
|
||||||
|
let (_ctx, store, shutdown) =
|
||||||
|
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "writer-fencing-multipart", &[4, 4])).await;
|
||||||
|
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
|
||||||
|
|
||||||
|
let bucket = format!("writer-fencing-multipart-{}", uuid::Uuid::new_v4());
|
||||||
|
store
|
||||||
|
.make_bucket(&bucket, &MakeBucketOptions::default())
|
||||||
|
.await
|
||||||
|
.expect("create multipart writer fencing bucket");
|
||||||
|
|
||||||
|
let incarnation = store.bucket_incarnation_id(&bucket).await.expect("read bucket incarnation");
|
||||||
|
let lifecycle_guard = store
|
||||||
|
.acquire_bucket_lifecycle_read_lock(&bucket)
|
||||||
|
.await
|
||||||
|
.expect("acquire multipart test lifecycle fence");
|
||||||
|
let mut upload_opts = ObjectOptions {
|
||||||
|
expected_bucket_incarnation_id: Some(incarnation),
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
upload_opts.add_bucket_lifecycle_lock_guard(&lifecycle_guard);
|
||||||
|
let source_upload = store.pools[0]
|
||||||
|
.new_multipart_upload(&bucket, "source-only.bin", &upload_opts)
|
||||||
|
.await
|
||||||
|
.expect("create source upload");
|
||||||
|
let target_upload = store.pools[1]
|
||||||
|
.new_multipart_upload(&bucket, "target-visible.bin", &upload_opts)
|
||||||
|
.await
|
||||||
|
.expect("create target upload");
|
||||||
|
|
||||||
|
*store.rebalance_meta.write().await = Some(active_rebalance_meta_for_pool(store.pools.len(), 0));
|
||||||
|
assert!(store.is_pool_rebalancing(0).await, "pool 0 must be marked as an active rebalance source");
|
||||||
|
|
||||||
|
let listed = store
|
||||||
|
.list_multipart_uploads(&bucket, "", None, None, None, 100)
|
||||||
|
.await
|
||||||
|
.expect("list multipart uploads");
|
||||||
|
let listed_uploads: Vec<(&str, &str)> = listed
|
||||||
|
.uploads
|
||||||
|
.iter()
|
||||||
|
.map(|upload| (upload.object.as_str(), upload.upload_id.as_str()))
|
||||||
|
.collect();
|
||||||
|
|
||||||
|
assert!(
|
||||||
|
!listed_uploads.contains(&("source-only.bin", source_upload.upload_id.as_str())),
|
||||||
|
"active source pool upload must be hidden from multipart listing"
|
||||||
|
);
|
||||||
|
assert!(
|
||||||
|
listed_uploads.contains(&("target-visible.bin", target_upload.upload_id.as_str())),
|
||||||
|
"non-rebalancing pool upload must remain visible"
|
||||||
|
);
|
||||||
|
|
||||||
|
shutdown.cancel();
|
||||||
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
#[serial_test::serial(storage_class_env)]
|
#[serial_test::serial(storage_class_env)]
|
||||||
async fn data_movement_conflicts_preserve_newer_target_and_abort_staging() {
|
async fn data_movement_conflicts_preserve_newer_target_and_abort_staging() {
|
||||||
|
|||||||
@@ -221,7 +221,7 @@ impl ECStore {
|
|||||||
}
|
}
|
||||||
|
|
||||||
for pool in self.pools.iter() {
|
for pool in self.pools.iter() {
|
||||||
if self.is_suspended(pool.pool_idx).await {
|
if self.is_suspended(pool.pool_idx).await || self.is_pool_rebalancing(pool.pool_idx).await {
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
return match pool
|
return match pool
|
||||||
@@ -284,7 +284,7 @@ impl ECStore {
|
|||||||
let mut source_truncated = false;
|
let mut source_truncated = false;
|
||||||
|
|
||||||
for pool in self.pools.iter() {
|
for pool in self.pools.iter() {
|
||||||
if self.is_suspended(pool.pool_idx).await {
|
if self.is_suspended(pool.pool_idx).await || self.is_pool_rebalancing(pool.pool_idx).await {
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
let res = list_pool_multipart_uploads_for_incarnation(
|
let res = list_pool_multipart_uploads_for_incarnation(
|
||||||
@@ -433,7 +433,7 @@ impl ECStore {
|
|||||||
}
|
}
|
||||||
|
|
||||||
for pool in self.pools.iter() {
|
for pool in self.pools.iter() {
|
||||||
if self.is_suspended(pool.pool_idx).await {
|
if self.is_suspended(pool.pool_idx).await || self.is_pool_rebalancing(pool.pool_idx).await {
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
let err = match pool.put_object_part(bucket, object, upload_id, part_id, data, opts).await {
|
let err = match pool.put_object_part(bucket, object, upload_id, part_id, data, opts).await {
|
||||||
@@ -472,7 +472,7 @@ impl ECStore {
|
|||||||
}
|
}
|
||||||
|
|
||||||
for pool in self.pools.iter() {
|
for pool in self.pools.iter() {
|
||||||
if self.is_suspended(pool.pool_idx).await {
|
if self.is_suspended(pool.pool_idx).await || self.is_pool_rebalancing(pool.pool_idx).await {
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -510,7 +510,7 @@ impl ECStore {
|
|||||||
}
|
}
|
||||||
|
|
||||||
for pool in self.pools.iter() {
|
for pool in self.pools.iter() {
|
||||||
if self.is_suspended(pool.pool_idx).await {
|
if self.is_suspended(pool.pool_idx).await || self.is_pool_rebalancing(pool.pool_idx).await {
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -551,7 +551,7 @@ impl ECStore {
|
|||||||
}
|
}
|
||||||
|
|
||||||
for pool in self.pools.iter() {
|
for pool in self.pools.iter() {
|
||||||
if self.is_suspended(pool.pool_idx).await {
|
if self.is_suspended(pool.pool_idx).await || self.is_pool_rebalancing(pool.pool_idx).await {
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -597,6 +597,10 @@ fn version_aware_lookup_opts(opts: &ObjectOptions, no_lock: bool) -> ObjectOptio
|
|||||||
}
|
}
|
||||||
|
|
||||||
fn data_movement_pool_lookup_opts(opts: &ObjectOptions, no_lock: bool) -> ObjectOptions {
|
fn data_movement_pool_lookup_opts(opts: &ObjectOptions, no_lock: bool) -> ObjectOptions {
|
||||||
|
writer_pool_lookup_opts(opts, no_lock)
|
||||||
|
}
|
||||||
|
|
||||||
|
fn writer_pool_lookup_opts(opts: &ObjectOptions, no_lock: bool) -> ObjectOptions {
|
||||||
let mut lookup_opts = version_aware_lookup_opts(opts, no_lock);
|
let mut lookup_opts = version_aware_lookup_opts(opts, no_lock);
|
||||||
lookup_opts.skip_decommissioned = true;
|
lookup_opts.skip_decommissioned = true;
|
||||||
lookup_opts.skip_rebalancing = true;
|
lookup_opts.skip_rebalancing = true;
|
||||||
@@ -607,6 +611,7 @@ fn data_movement_pool_lookup_opts(opts: &ObjectOptions, no_lock: bool) -> Object
|
|||||||
fn transition_restore_pool_opts(opts: &ObjectOptions) -> ObjectOptions {
|
fn transition_restore_pool_opts(opts: &ObjectOptions) -> ObjectOptions {
|
||||||
let mut lookup_opts = opts.clone();
|
let mut lookup_opts = opts.clone();
|
||||||
lookup_opts.skip_decommissioned = true;
|
lookup_opts.skip_decommissioned = true;
|
||||||
|
lookup_opts.skip_rebalancing = true;
|
||||||
lookup_opts
|
lookup_opts
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -1358,7 +1363,7 @@ impl ECStore {
|
|||||||
|
|
||||||
if cp_src_dst_same {
|
if cp_src_dst_same {
|
||||||
let pool_idx = self
|
let pool_idx = self
|
||||||
.get_pool_info_existing_with_opts(src_bucket, &src_object, &version_aware_lookup_opts(src_opts, true))
|
.get_pool_info_existing_with_opts(src_bucket, &src_object, &writer_pool_lookup_opts(src_opts, true))
|
||||||
.await?
|
.await?
|
||||||
.0
|
.0
|
||||||
.index;
|
.index;
|
||||||
@@ -1671,7 +1676,7 @@ impl ECStore {
|
|||||||
return Ok(ObjectInfo::default());
|
return Ok(ObjectInfo::default());
|
||||||
}
|
}
|
||||||
|
|
||||||
let gopts = version_aware_lookup_opts(&opts, true);
|
let gopts = writer_pool_lookup_opts(&opts, true);
|
||||||
|
|
||||||
if opts.data_movement {
|
if opts.data_movement {
|
||||||
let existing_pool_info = self.get_pool_info_existing_with_opts(bucket, object, &gopts).await;
|
let existing_pool_info = self.get_pool_info_existing_with_opts(bucket, object, &gopts).await;
|
||||||
@@ -1787,6 +1792,10 @@ impl ECStore {
|
|||||||
}
|
}
|
||||||
|
|
||||||
for pool in self.pools.iter() {
|
for pool in self.pools.iter() {
|
||||||
|
if self.is_suspended(pool.pool_idx).await || self.is_pool_rebalancing(pool.pool_idx).await {
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
|
||||||
match pool.delete_object(bucket, object, opts.clone()).await {
|
match pool.delete_object(bucket, object, opts.clone()).await {
|
||||||
Ok(res) => {
|
Ok(res) => {
|
||||||
if let (Some(api), Some(je)) = (tier_journal_api.as_ref(), journal_entry.as_ref()) {
|
if let (Some(api), Some(je)) = (tier_journal_api.as_ref(), journal_entry.as_ref()) {
|
||||||
@@ -2090,7 +2099,7 @@ impl ECStore {
|
|||||||
..Default::default()
|
..Default::default()
|
||||||
};
|
};
|
||||||
let (_, idx) = self
|
let (_, idx) = self
|
||||||
.get_latest_accessible_object_info_with_idx(bucket, object.as_str(), &opts)
|
.get_latest_accessible_object_info_with_idx(bucket, object.as_str(), &writer_pool_lookup_opts(&opts, opts.no_lock))
|
||||||
.await?;
|
.await?;
|
||||||
|
|
||||||
let _ = self.pools[idx].add_partial(bucket, object.as_str(), version_id).await;
|
let _ = self.pools[idx].add_partial(bucket, object.as_str(), version_id).await;
|
||||||
@@ -2181,7 +2190,7 @@ impl ECStore {
|
|||||||
}
|
}
|
||||||
|
|
||||||
let (_, idx) = self
|
let (_, idx) = self
|
||||||
.get_latest_accessible_object_info_with_idx(bucket, object.as_str(), &opts)
|
.get_latest_accessible_object_info_with_idx(bucket, object.as_str(), &writer_pool_lookup_opts(&opts, opts.no_lock))
|
||||||
.await?;
|
.await?;
|
||||||
|
|
||||||
self.pools[idx]
|
self.pools[idx]
|
||||||
@@ -2224,7 +2233,7 @@ impl ECStore {
|
|||||||
}
|
}
|
||||||
|
|
||||||
let (_, idx) = self
|
let (_, idx) = self
|
||||||
.get_latest_accessible_object_info_with_idx(bucket, object.as_str(), &opts)
|
.get_latest_accessible_object_info_with_idx(bucket, object.as_str(), &writer_pool_lookup_opts(&opts, opts.no_lock))
|
||||||
.await?;
|
.await?;
|
||||||
|
|
||||||
let result = self.pools[idx].put_object_metadata(bucket, object.as_str(), &opts).await;
|
let result = self.pools[idx].put_object_metadata(bucket, object.as_str(), &opts).await;
|
||||||
@@ -2259,7 +2268,7 @@ impl ECStore {
|
|||||||
}
|
}
|
||||||
|
|
||||||
let (_, idx) = self
|
let (_, idx) = self
|
||||||
.get_latest_accessible_object_info_with_idx(bucket, object.as_str(), opts)
|
.get_latest_accessible_object_info_with_idx(bucket, object.as_str(), &writer_pool_lookup_opts(opts, opts.no_lock))
|
||||||
.await?;
|
.await?;
|
||||||
|
|
||||||
self.pools[idx].put_object_tags(bucket, object.as_str(), tags, opts).await
|
self.pools[idx].put_object_tags(bucket, object.as_str(), tags, opts).await
|
||||||
@@ -2294,7 +2303,7 @@ impl ECStore {
|
|||||||
}
|
}
|
||||||
|
|
||||||
let (_, idx) = self
|
let (_, idx) = self
|
||||||
.get_latest_accessible_object_info_with_idx(bucket, object.as_str(), opts)
|
.get_latest_accessible_object_info_with_idx(bucket, object.as_str(), &writer_pool_lookup_opts(opts, opts.no_lock))
|
||||||
.await?;
|
.await?;
|
||||||
|
|
||||||
self.pools[idx].delete_object_tags(bucket, object.as_str(), opts).await
|
self.pools[idx].delete_object_tags(bucket, object.as_str(), opts).await
|
||||||
@@ -2923,6 +2932,23 @@ mod tests {
|
|||||||
assert_eq!(lookup_opts.version_id.as_deref(), Some("vid-1"));
|
assert_eq!(lookup_opts.version_id.as_deref(), Some("vid-1"));
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn writer_pool_lookup_opts_skips_rebalance_sources() {
|
||||||
|
let lookup_opts = writer_pool_lookup_opts(
|
||||||
|
&ObjectOptions {
|
||||||
|
version_id: Some("vid-1".to_string()),
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
true,
|
||||||
|
);
|
||||||
|
|
||||||
|
assert!(lookup_opts.no_lock);
|
||||||
|
assert!(lookup_opts.metadata_chg);
|
||||||
|
assert!(lookup_opts.skip_decommissioned);
|
||||||
|
assert!(lookup_opts.skip_rebalancing);
|
||||||
|
assert_eq!(lookup_opts.version_id.as_deref(), Some("vid-1"));
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn data_movement_pool_lookup_opts_keeps_no_lock_for_tiered_moves() {
|
fn data_movement_pool_lookup_opts_keeps_no_lock_for_tiered_moves() {
|
||||||
let lookup_opts = data_movement_pool_lookup_opts(
|
let lookup_opts = data_movement_pool_lookup_opts(
|
||||||
@@ -2948,6 +2974,7 @@ mod tests {
|
|||||||
});
|
});
|
||||||
|
|
||||||
assert!(lookup_opts.skip_decommissioned);
|
assert!(lookup_opts.skip_decommissioned);
|
||||||
|
assert!(lookup_opts.skip_rebalancing);
|
||||||
assert!(!lookup_opts.no_lock);
|
assert!(!lookup_opts.no_lock);
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -2959,6 +2986,7 @@ mod tests {
|
|||||||
});
|
});
|
||||||
|
|
||||||
assert!(lookup_opts.skip_decommissioned);
|
assert!(lookup_opts.skip_decommissioned);
|
||||||
|
assert!(lookup_opts.skip_rebalancing);
|
||||||
assert!(lookup_opts.no_lock);
|
assert!(lookup_opts.no_lock);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -541,14 +541,22 @@ impl ECStore {
|
|||||||
opts: &ObjectOptions,
|
opts: &ObjectOptions,
|
||||||
) -> Result<(ObjectInfo, usize)> {
|
) -> Result<(ObjectInfo, usize)> {
|
||||||
let mut futures = Vec::with_capacity(self.pools.len());
|
let mut futures = Vec::with_capacity(self.pools.len());
|
||||||
for pool in self.pools.iter() {
|
for (idx, pool) in self.pools.iter().enumerate() {
|
||||||
futures.push(pool.get_object_info(bucket, object, opts));
|
if opts.skip_decommissioned && self.is_suspended(idx).await {
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
|
||||||
|
if opts.skip_rebalancing && self.is_pool_rebalancing(idx).await {
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
|
||||||
|
futures.push(async move { (idx, pool.get_object_info(bucket, object, opts).await) });
|
||||||
}
|
}
|
||||||
|
|
||||||
let results = join_all(futures).await;
|
let results = join_all(futures).await;
|
||||||
let mut candidates = Vec::with_capacity(self.pools.len());
|
let mut candidates = Vec::with_capacity(self.pools.len());
|
||||||
|
|
||||||
for (idx, result) in results.into_iter().enumerate() {
|
for (idx, result) in results {
|
||||||
match result {
|
match result {
|
||||||
Ok(res) => {
|
Ok(res) => {
|
||||||
candidates.push(LatestObjectInfoCandidate {
|
candidates.push(LatestObjectInfoCandidate {
|
||||||
|
|||||||
Reference in New Issue
Block a user