mirror of
https://github.com/rustfs/rustfs.git
synced 2026-07-26 16:28:15 +00:00
fix(ecstore): preserve bucket visibility during rebalance (#3418)
Co-authored-by: houseme <housemecn@gmail.com>
This commit is contained in:
+101
-18
@@ -1023,7 +1023,11 @@ impl ECStore {
|
||||
);
|
||||
return Ok(());
|
||||
}
|
||||
if complete_rebalance_pools_at_goal(meta, OffsetDateTime::now_utc()) {
|
||||
let now = OffsetDateTime::now_utc();
|
||||
if complete_rebalance_pools_at_goal(meta, now) {
|
||||
meta_to_save = Some(meta.clone());
|
||||
}
|
||||
if complete_rebalance_pools_with_empty_queue(meta, now) {
|
||||
meta_to_save = Some(meta.clone());
|
||||
}
|
||||
meta.cancel = Some(cancel_tx);
|
||||
@@ -1054,6 +1058,19 @@ impl ECStore {
|
||||
Vec::new()
|
||||
};
|
||||
|
||||
if !participants.iter().any(|participating| *participating) {
|
||||
debug!(
|
||||
event = EVENT_REBALANCE_STATE,
|
||||
component = LOG_COMPONENT_ECSTORE,
|
||||
subsystem = LOG_SUBSYSTEM_REBALANCE,
|
||||
state = "start_skipped",
|
||||
reason = "no_participants",
|
||||
"Skipped rebalance start because no pools are participating"
|
||||
);
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
let mut workers_started = 0usize;
|
||||
for (idx, participating) in participants.iter().enumerate() {
|
||||
if !*participating {
|
||||
debug!(
|
||||
@@ -1088,6 +1105,7 @@ impl ECStore {
|
||||
let pool_idx = idx;
|
||||
let store = self.clone();
|
||||
let rx_clone = rx.clone();
|
||||
workers_started += 1;
|
||||
tokio::spawn(async move {
|
||||
if let Err(err) = store.rebalance_buckets(rx_clone, pool_idx).await {
|
||||
error!("Rebalance failed for pool {}: {}", pool_idx, err);
|
||||
@@ -1104,6 +1122,18 @@ impl ECStore {
|
||||
});
|
||||
}
|
||||
|
||||
if workers_started == 0 {
|
||||
debug!(
|
||||
event = EVENT_REBALANCE_STATE,
|
||||
component = LOG_COMPONENT_ECSTORE,
|
||||
subsystem = LOG_SUBSYSTEM_REBALANCE,
|
||||
state = "start_skipped",
|
||||
reason = "no_local_participants",
|
||||
"Skipped rebalance start because no local pools are participating"
|
||||
);
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
info!(
|
||||
event = EVENT_REBALANCE_STATE,
|
||||
component = LOG_COMPONENT_ECSTORE,
|
||||
@@ -1610,6 +1640,23 @@ fn complete_rebalance_pools_at_goal(meta: &mut RebalanceMeta, now: OffsetDateTim
|
||||
changed
|
||||
}
|
||||
|
||||
fn complete_rebalance_pools_with_empty_queue(meta: &mut RebalanceMeta, now: OffsetDateTime) -> bool {
|
||||
let mut changed = false;
|
||||
|
||||
for pool_stat in meta.pool_stats.iter_mut() {
|
||||
if !is_rebalance_pool_started(pool_stat) || !pool_stat.buckets.is_empty() {
|
||||
continue;
|
||||
}
|
||||
|
||||
pool_stat.info.status = RebalStatus::Completed;
|
||||
pool_stat.info.end_time = Some(now);
|
||||
pool_stat.info.last_error = None;
|
||||
changed = true;
|
||||
}
|
||||
|
||||
changed
|
||||
}
|
||||
|
||||
fn has_deferred_rebalance_error(pool_stat: &RebalanceStats) -> bool {
|
||||
pool_stat
|
||||
.info
|
||||
@@ -2859,23 +2906,24 @@ mod rebalance_unit_tests {
|
||||
RebalSaveOpt, RebalStatus, RebalanceBucketOutcome, RebalanceEntryOutcome, RebalanceInfo, RebalanceMeta, RebalanceStats,
|
||||
RebalanceTerminalEvent, apply_rebalance_save_option, apply_rebalance_terminal_event, apply_stopped_at,
|
||||
classify_rebalance_terminal_event, clone_arc_by_index, clone_first_arc, clone_rebalance_pool_stats,
|
||||
complete_rebalance_pools_at_goal, defer_bucket_in_rebalance_queue, ensure_rebalance_listing_disks_available,
|
||||
ensure_rebalance_not_decommissioning, ensure_valid_rebalance_pool_index, has_deferred_rebalance_error,
|
||||
is_rebalance_stopped_terminal_event, is_transient_rebalance_error, load_rebalance_bucket_configs,
|
||||
mark_rebalance_bucket_done, merge_rebalance_meta, migrate_entry_version, migrate_entry_version_with_retry_wait,
|
||||
next_rebal_bucket_from_stat, rebalance_delete_marker_opts, rebalance_listing_retry_delay,
|
||||
rebalance_meta_load_no_data_error, rebalance_meta_load_unknown_format_error, rebalance_meta_load_unknown_version_error,
|
||||
rebalance_migration_retry_delay, resolve_load_rebalance_stats_update_result, resolve_next_rebalance_bucket,
|
||||
resolve_rebalance_bucket_error, resolve_rebalance_bucket_result, resolve_rebalance_entry_cleanup_delete_result,
|
||||
resolve_rebalance_file_info_versions_result, resolve_rebalance_meta_load_result, resolve_rebalance_meta_save_result,
|
||||
resolve_rebalance_migrate_result_error, resolve_rebalance_optional_bucket_config_result, resolve_rebalance_participants,
|
||||
resolve_rebalance_save_task_result, resolve_rebalance_stats_update_result, resolve_rebalance_terminal_error,
|
||||
resolve_rebalance_worker_result, send_rebalance_done_signal, should_accept_rebalance_stats_update,
|
||||
should_cleanup_rebalance_source_entry, should_count_rebalance_version_complete, should_defer_rebalance_entry_failure,
|
||||
should_ignore_rebalance_data_usage_cache, should_pool_participate, should_preserve_rebalance_stopped_state,
|
||||
should_retry_rebalance_listing, should_skip_rebalance_delete_marker, should_skip_start_rebalance,
|
||||
stop_rebalance_meta_snapshot, stop_rebalance_state, take_bucket_from_rebalance_queue, validate_start_rebalance_state,
|
||||
wait_rebalance_listing_retry, with_rebalance_entry_context,
|
||||
complete_rebalance_pools_at_goal, complete_rebalance_pools_with_empty_queue, defer_bucket_in_rebalance_queue,
|
||||
ensure_rebalance_listing_disks_available, ensure_rebalance_not_decommissioning, ensure_valid_rebalance_pool_index,
|
||||
has_deferred_rebalance_error, is_rebalance_stopped_terminal_event, is_transient_rebalance_error,
|
||||
load_rebalance_bucket_configs, mark_rebalance_bucket_done, merge_rebalance_meta, migrate_entry_version,
|
||||
migrate_entry_version_with_retry_wait, next_rebal_bucket_from_stat, rebalance_delete_marker_opts,
|
||||
rebalance_listing_retry_delay, rebalance_meta_load_no_data_error, rebalance_meta_load_unknown_format_error,
|
||||
rebalance_meta_load_unknown_version_error, rebalance_migration_retry_delay, resolve_load_rebalance_stats_update_result,
|
||||
resolve_next_rebalance_bucket, resolve_rebalance_bucket_error, resolve_rebalance_bucket_result,
|
||||
resolve_rebalance_entry_cleanup_delete_result, resolve_rebalance_file_info_versions_result,
|
||||
resolve_rebalance_meta_load_result, resolve_rebalance_meta_save_result, resolve_rebalance_migrate_result_error,
|
||||
resolve_rebalance_optional_bucket_config_result, resolve_rebalance_participants, resolve_rebalance_save_task_result,
|
||||
resolve_rebalance_stats_update_result, resolve_rebalance_terminal_error, resolve_rebalance_worker_result,
|
||||
send_rebalance_done_signal, should_accept_rebalance_stats_update, should_cleanup_rebalance_source_entry,
|
||||
should_count_rebalance_version_complete, should_defer_rebalance_entry_failure, should_ignore_rebalance_data_usage_cache,
|
||||
should_pool_participate, should_preserve_rebalance_stopped_state, should_retry_rebalance_listing,
|
||||
should_skip_rebalance_delete_marker, should_skip_start_rebalance, stop_rebalance_meta_snapshot, stop_rebalance_state,
|
||||
take_bucket_from_rebalance_queue, validate_start_rebalance_state, wait_rebalance_listing_retry,
|
||||
with_rebalance_entry_context,
|
||||
};
|
||||
use crate::data_movement;
|
||||
use crate::data_usage::DATA_USAGE_CACHE_NAME;
|
||||
@@ -5016,6 +5064,41 @@ mod rebalance_unit_tests {
|
||||
assert!(has_deferred_rebalance_error(&meta.pool_stats[0]));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_complete_rebalance_pools_with_empty_queue_marks_started_participants_completed() {
|
||||
let now = OffsetDateTime::from_unix_timestamp(1_000).unwrap();
|
||||
let mut meta = RebalanceMeta {
|
||||
pool_stats: vec![
|
||||
RebalanceStats {
|
||||
participating: true,
|
||||
buckets: Vec::new(),
|
||||
info: RebalanceInfo {
|
||||
status: RebalStatus::Started,
|
||||
last_error: Some("stale error".to_string()),
|
||||
..Default::default()
|
||||
},
|
||||
..Default::default()
|
||||
},
|
||||
RebalanceStats {
|
||||
participating: true,
|
||||
buckets: vec!["bucket-a".to_string()],
|
||||
info: RebalanceInfo {
|
||||
status: RebalStatus::Started,
|
||||
..Default::default()
|
||||
},
|
||||
..Default::default()
|
||||
},
|
||||
],
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
assert!(complete_rebalance_pools_with_empty_queue(&mut meta, now));
|
||||
assert_eq!(meta.pool_stats[0].info.status, RebalStatus::Completed);
|
||||
assert_eq!(meta.pool_stats[0].info.end_time, Some(now));
|
||||
assert!(meta.pool_stats[0].info.last_error.is_none());
|
||||
assert_eq!(meta.pool_stats[1].info.status, RebalStatus::Started);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_should_skip_start_rebalance_only_when_running_and_cancel_attached() {
|
||||
assert!(should_skip_start_rebalance(true, true));
|
||||
|
||||
@@ -130,16 +130,8 @@ impl S3PeerSys {
|
||||
|
||||
let mut pool_errs = Vec::new();
|
||||
for pool_idx in 0..self.pools_count {
|
||||
let mut per_pool_errs = vec![None; self.clients.len()];
|
||||
for (i, client) in self.clients.iter().enumerate() {
|
||||
if let Some(v) = client.get_pools()
|
||||
&& v.contains(&pool_idx)
|
||||
{
|
||||
per_pool_errs[i] = errs[i].clone();
|
||||
}
|
||||
}
|
||||
let qu = per_pool_errs.len() / 2;
|
||||
pool_errs.push(reduce_write_quorum_errs(&per_pool_errs, BUCKET_OP_IGNORED_ERRS, qu));
|
||||
let per_pool_errs = pool_participant_errors(&self.clients, &errs, pool_idx);
|
||||
pool_errs.push(reduce_pool_write_quorum_errs(&per_pool_errs));
|
||||
}
|
||||
|
||||
if !opts.recreate {
|
||||
@@ -165,28 +157,14 @@ impl S3PeerSys {
|
||||
let errs = join_all(futures).await;
|
||||
|
||||
for pool_idx in 0..self.pools_count {
|
||||
let mut per_pool_errs = vec![None; self.clients.len()];
|
||||
for (i, client) in self.clients.iter().enumerate() {
|
||||
if let Some(v) = client.get_pools()
|
||||
&& v.contains(&pool_idx)
|
||||
{
|
||||
per_pool_errs[i] = errs[i].clone();
|
||||
}
|
||||
}
|
||||
let qu = per_pool_errs.len() / 2;
|
||||
if let Some(pool_err) = reduce_write_quorum_errs(&per_pool_errs, BUCKET_OP_IGNORED_ERRS, qu) {
|
||||
let per_pool_errs = pool_participant_errors(&self.clients, &errs, pool_idx);
|
||||
if let Some(pool_err) = reduce_pool_write_quorum_errs(&per_pool_errs) {
|
||||
tracing::error!("heal_bucket per_pool_errs: {per_pool_errs:?}");
|
||||
tracing::error!("heal_bucket reduce_write_quorum_errs: {pool_err}");
|
||||
return Err(pool_err);
|
||||
}
|
||||
}
|
||||
|
||||
if let Some(err) = reduce_write_quorum_errs(&errs, BUCKET_OP_IGNORED_ERRS, (errs.len() / 2) + 1) {
|
||||
tracing::error!("heal_bucket errs: {errs:?}");
|
||||
tracing::error!("heal_bucket reduce_write_quorum_errs: {err}");
|
||||
return Err(err);
|
||||
}
|
||||
|
||||
for (i, err) in errs.iter().enumerate() {
|
||||
if err.is_none() {
|
||||
return Ok(heal_bucket_results.read().await[i].clone());
|
||||
@@ -251,18 +229,10 @@ impl S3PeerSys {
|
||||
|
||||
let mut result_map: HashMap<&String, BucketInfo> = HashMap::new();
|
||||
for i in 0..self.pools_count {
|
||||
let mut per_pool_errs = vec![None; self.clients.len()];
|
||||
for (j, cli) in self.clients.iter().enumerate() {
|
||||
let pools = cli.get_pools();
|
||||
let idx = i;
|
||||
if pools.unwrap_or_default().contains(&idx) {
|
||||
per_pool_errs[j] = errors[j].clone();
|
||||
}
|
||||
}
|
||||
let per_pool_errs = pool_participant_errors(&self.clients, &errors, i);
|
||||
let quorum = pool_write_quorum(per_pool_errs.len());
|
||||
|
||||
let quorum = per_pool_errs.len() / 2;
|
||||
|
||||
if let Some(pool_err) = reduce_write_quorum_errs(&per_pool_errs, BUCKET_OP_IGNORED_ERRS, quorum) {
|
||||
if let Some(pool_err) = reduce_pool_write_quorum_errs(&per_pool_errs) {
|
||||
tracing::error!("list_bucket per_pool_errs: {per_pool_errs:?}");
|
||||
tracing::error!("list_bucket reduce_write_quorum_errs: {pool_err}");
|
||||
return Err(pool_err);
|
||||
@@ -357,18 +327,8 @@ impl S3PeerSys {
|
||||
}
|
||||
|
||||
for i in 0..self.pools_count {
|
||||
let mut per_pool_errs = vec![None; self.clients.len()];
|
||||
for (j, cli) in self.clients.iter().enumerate() {
|
||||
let pools = cli.get_pools();
|
||||
let idx = i;
|
||||
if pools.unwrap_or_default().contains(&idx) {
|
||||
per_pool_errs[j] = errors[j].clone();
|
||||
}
|
||||
}
|
||||
|
||||
if let Some(pool_err) =
|
||||
reduce_write_quorum_errs(&per_pool_errs, BUCKET_OP_IGNORED_ERRS, (per_pool_errs.len() / 2) + 1)
|
||||
{
|
||||
let per_pool_errs = pool_participant_errors(&self.clients, &errors, i);
|
||||
if let Some(pool_err) = reduce_pool_write_quorum_errs(&per_pool_errs) {
|
||||
return Err(pool_err);
|
||||
}
|
||||
}
|
||||
@@ -399,6 +359,18 @@ impl LocalPeerS3Client {
|
||||
pools,
|
||||
}
|
||||
}
|
||||
|
||||
async fn local_disks_for_pools(&self) -> Vec<DiskStore> {
|
||||
let local_disks = all_local_disk().await;
|
||||
let Some(pools) = self.pools.as_ref() else {
|
||||
return local_disks;
|
||||
};
|
||||
|
||||
local_disks
|
||||
.into_iter()
|
||||
.filter(|disk| usize::try_from(disk.endpoint().pool_idx).is_ok_and(|pool_idx| pools.contains(&pool_idx)))
|
||||
.collect()
|
||||
}
|
||||
}
|
||||
|
||||
#[async_trait]
|
||||
@@ -408,11 +380,15 @@ impl PeerS3Client for LocalPeerS3Client {
|
||||
}
|
||||
|
||||
async fn heal_bucket(&self, bucket: &str, opts: &HealOpts) -> Result<HealResultItem> {
|
||||
heal_bucket_local(bucket, opts).await
|
||||
let disks = self.local_disks_for_pools().await.into_iter().map(Some).collect();
|
||||
heal_bucket_local_on_disks(bucket, opts, disks).await
|
||||
}
|
||||
|
||||
async fn list_bucket(&self, _opts: &BucketOptions) -> Result<Vec<BucketInfo>> {
|
||||
let local_disks = all_local_disk().await;
|
||||
let local_disks = self.local_disks_for_pools().await;
|
||||
if local_disks.is_empty() {
|
||||
return Err(Error::ErasureWriteQuorum);
|
||||
}
|
||||
|
||||
let mut futures = Vec::with_capacity(local_disks.len());
|
||||
for disk in local_disks.iter() {
|
||||
@@ -421,45 +397,62 @@ impl PeerS3Client for LocalPeerS3Client {
|
||||
|
||||
let results = join_all(futures).await;
|
||||
|
||||
let mut ress = Vec::new();
|
||||
let mut errs = Vec::new();
|
||||
let mut ress = Vec::with_capacity(local_disks.len());
|
||||
let mut errs = Vec::with_capacity(local_disks.len());
|
||||
|
||||
for result in results {
|
||||
match result {
|
||||
Ok(res) => {
|
||||
ress.push(res);
|
||||
ress.push(Some(res));
|
||||
errs.push(None);
|
||||
}
|
||||
Err(e) => errs.push(Some(e)),
|
||||
Err(e) => {
|
||||
ress.push(None);
|
||||
errs.push(Some(e));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
let mut uniq_map: HashMap<&String, &VolumeInfo> = HashMap::new();
|
||||
if let Some(err) = reduce_write_quorum_errs(&errs, BUCKET_OP_IGNORED_ERRS, (local_disks.len() / 2) + 1) {
|
||||
return Err(err);
|
||||
}
|
||||
|
||||
for info_list in ress.iter() {
|
||||
let quorum = (local_disks.len() / 2) + 1;
|
||||
let mut count_map: HashMap<&String, (usize, &VolumeInfo)> = HashMap::new();
|
||||
for info_list in ress.iter().flatten() {
|
||||
for info in info_list.iter() {
|
||||
if is_reserved_or_invalid_bucket(&info.name, false) {
|
||||
continue;
|
||||
}
|
||||
if !uniq_map.contains_key(&info.name) {
|
||||
uniq_map.insert(&info.name, info);
|
||||
}
|
||||
|
||||
let entry = count_map.entry(&info.name).or_insert((0, info));
|
||||
entry.0 += 1;
|
||||
}
|
||||
}
|
||||
|
||||
let buckets: Vec<BucketInfo> = uniq_map
|
||||
let buckets: Vec<BucketInfo> = count_map
|
||||
.values()
|
||||
.map(|&v| BucketInfo {
|
||||
name: v.name.clone(),
|
||||
created: v.created,
|
||||
..Default::default()
|
||||
.filter_map(|(count, info)| {
|
||||
if *count < quorum {
|
||||
return None;
|
||||
}
|
||||
|
||||
Some(BucketInfo {
|
||||
name: info.name.clone(),
|
||||
created: info.created,
|
||||
..Default::default()
|
||||
})
|
||||
})
|
||||
.collect();
|
||||
|
||||
Ok(buckets)
|
||||
}
|
||||
async fn make_bucket(&self, bucket: &str, opts: &MakeBucketOptions) -> Result<()> {
|
||||
let local_disks = all_local_disk().await;
|
||||
let local_disks = self.local_disks_for_pools().await;
|
||||
if local_disks.is_empty() {
|
||||
return Err(Error::ErasureWriteQuorum);
|
||||
}
|
||||
|
||||
let mut futures = Vec::with_capacity(local_disks.len());
|
||||
for disk in local_disks.iter() {
|
||||
futures.push(async move {
|
||||
@@ -494,7 +487,11 @@ impl PeerS3Client for LocalPeerS3Client {
|
||||
}
|
||||
|
||||
async fn get_bucket_info(&self, bucket: &str, _opts: &BucketOptions) -> Result<BucketInfo> {
|
||||
let local_disks = all_local_disk().await;
|
||||
let local_disks = self.local_disks_for_pools().await;
|
||||
if local_disks.is_empty() {
|
||||
return Err(Error::ErasureWriteQuorum);
|
||||
}
|
||||
|
||||
let mut futures = Vec::with_capacity(local_disks.len());
|
||||
for disk in local_disks.iter() {
|
||||
futures.push(disk.stat_volume(bucket));
|
||||
@@ -518,7 +515,10 @@ impl PeerS3Client for LocalPeerS3Client {
|
||||
}
|
||||
}
|
||||
|
||||
// TODO: reduceWriteQuorumErrs
|
||||
if let Some(err) = reduce_write_quorum_errs(&errs, BUCKET_OP_IGNORED_ERRS, (local_disks.len() / 2) + 1) {
|
||||
return Err(err);
|
||||
}
|
||||
|
||||
let mut versioned = false;
|
||||
if let Ok(sys) = metadata_sys::get(bucket).await {
|
||||
versioned = sys.versioning();
|
||||
@@ -537,7 +537,11 @@ impl PeerS3Client for LocalPeerS3Client {
|
||||
}
|
||||
|
||||
async fn delete_bucket(&self, bucket: &str, _opts: &DeleteBucketOptions) -> Result<()> {
|
||||
let local_disks = all_local_disk().await;
|
||||
let local_disks = self.local_disks_for_pools().await;
|
||||
if local_disks.is_empty() {
|
||||
return Err(Error::ErasureWriteQuorum);
|
||||
}
|
||||
|
||||
let mut futures = Vec::with_capacity(local_disks.len());
|
||||
|
||||
for disk in local_disks.iter() {
|
||||
@@ -912,6 +916,10 @@ impl PeerS3Client for RemotePeerS3Client {
|
||||
|
||||
pub async fn heal_bucket_local(bucket: &str, opts: &HealOpts) -> Result<HealResultItem> {
|
||||
let disks = clone_drives().await;
|
||||
heal_bucket_local_on_disks(bucket, opts, disks).await
|
||||
}
|
||||
|
||||
async fn heal_bucket_local_on_disks(bucket: &str, opts: &HealOpts, disks: Vec<Option<DiskStore>>) -> Result<HealResultItem> {
|
||||
let before_state = Arc::new(RwLock::new(vec![String::new(); disks.len()]));
|
||||
let after_state = Arc::new(RwLock::new(vec![String::new(); disks.len()]));
|
||||
|
||||
@@ -1011,8 +1019,12 @@ pub async fn heal_bucket_local(bucket: &str, opts: &HealOpts) -> Result<HealResu
|
||||
let errs_clone = errs.to_vec();
|
||||
futures.push(async move {
|
||||
if bs_clone.read().await[idx] == DriveState::Missing.to_string() {
|
||||
let Some(disk) = disk.as_ref() else {
|
||||
return Some(Error::DiskNotFound);
|
||||
};
|
||||
|
||||
info!("bucket not find, will recreate");
|
||||
match disk.as_ref().unwrap().make_volume(&bucket).await {
|
||||
match disk.make_volume(&bucket).await {
|
||||
Ok(_) => {
|
||||
as_clone.write().await[idx] = DriveState::Ok.to_string();
|
||||
return None;
|
||||
@@ -1090,6 +1102,7 @@ mod tests {
|
||||
struct TestPeerS3Client {
|
||||
pools: Option<Vec<usize>>,
|
||||
make_bucket_result: Result<()>,
|
||||
list_bucket_result: Result<Vec<BucketInfo>>,
|
||||
}
|
||||
|
||||
#[async_trait]
|
||||
@@ -1103,7 +1116,7 @@ mod tests {
|
||||
}
|
||||
|
||||
async fn list_bucket(&self, _opts: &BucketOptions) -> Result<Vec<BucketInfo>> {
|
||||
unreachable!("not used by quorum tests")
|
||||
self.list_bucket_result.clone()
|
||||
}
|
||||
|
||||
async fn delete_bucket(&self, _bucket: &str, _opts: &DeleteBucketOptions) -> Result<()> {
|
||||
@@ -1124,12 +1137,79 @@ mod tests {
|
||||
}
|
||||
|
||||
fn test_peer_with_make_bucket(pools: &[usize], make_bucket_result: Result<()>) -> Client {
|
||||
test_peer_with_results(pools, make_bucket_result, Ok(Vec::new()))
|
||||
}
|
||||
|
||||
fn test_peer_with_list_bucket(pools: &[usize], list_bucket_result: Result<Vec<BucketInfo>>) -> Client {
|
||||
test_peer_with_results(pools, Ok(()), list_bucket_result)
|
||||
}
|
||||
|
||||
fn test_peer_with_results(
|
||||
pools: &[usize],
|
||||
make_bucket_result: Result<()>,
|
||||
list_bucket_result: Result<Vec<BucketInfo>>,
|
||||
) -> Client {
|
||||
Arc::new(Box::new(TestPeerS3Client {
|
||||
pools: Some(pools.to_vec()),
|
||||
make_bucket_result,
|
||||
list_bucket_result,
|
||||
}))
|
||||
}
|
||||
|
||||
fn test_endpoint(path: &std::path::Path, pool_index: usize, set_index: usize, disk_index: usize) -> Endpoint {
|
||||
let mut endpoint = Endpoint::try_from(path.to_str().expect("disk path to str")).expect("endpoint");
|
||||
endpoint.set_pool_index(pool_index);
|
||||
endpoint.set_set_index(set_index);
|
||||
endpoint.set_disk_index(disk_index);
|
||||
endpoint
|
||||
}
|
||||
|
||||
async fn init_test_local_disks(temp_dir: &TempDir, disk_count: usize, cmd_line: &str) -> Vec<DiskStore> {
|
||||
init_test_local_disks_for_pools(temp_dir, &[(0, disk_count)], cmd_line).await
|
||||
}
|
||||
|
||||
async fn init_test_local_disks_for_pools(
|
||||
temp_dir: &TempDir,
|
||||
pool_disk_counts: &[(usize, usize)],
|
||||
cmd_line: &str,
|
||||
) -> Vec<DiskStore> {
|
||||
let total_disk_count = pool_disk_counts.iter().map(|(_, disk_count)| *disk_count).sum();
|
||||
let mut endpoints = Vec::with_capacity(total_disk_count);
|
||||
let mut pool_endpoints = Vec::with_capacity(pool_disk_counts.len());
|
||||
for (pool_idx, disk_count) in pool_disk_counts.iter().copied() {
|
||||
let mut endpoints_for_pool = Vec::with_capacity(disk_count);
|
||||
for disk_idx in 0..disk_count {
|
||||
let disk_path = temp_dir.path().join(format!("pool{pool_idx}-disk{disk_idx}"));
|
||||
std::fs::create_dir_all(&disk_path).expect("create disk path");
|
||||
let endpoint = test_endpoint(&disk_path, pool_idx, 0, disk_idx);
|
||||
endpoints.push(endpoint.clone());
|
||||
endpoints_for_pool.push(endpoint);
|
||||
}
|
||||
|
||||
pool_endpoints.push(PoolEndpoints {
|
||||
legacy: false,
|
||||
set_count: 1,
|
||||
drives_per_set: disk_count,
|
||||
endpoints: Endpoints::from(endpoints_for_pool),
|
||||
cmd_line: cmd_line.to_string(),
|
||||
platform: "test".to_string(),
|
||||
});
|
||||
}
|
||||
|
||||
let endpoint_pools = EndpointServerPools(pool_endpoints);
|
||||
|
||||
init_local_disks(endpoint_pools).await.expect("init local disks");
|
||||
|
||||
let mut disks = Vec::with_capacity(total_disk_count);
|
||||
for endpoint in endpoints.iter() {
|
||||
let disk = Arc::new(LocalDisk::new(endpoint, false).await.expect("local disk should be created"));
|
||||
let wrapper = crate::disk::Disk::Local(Box::new(LocalDiskWrapper::new(disk, false)));
|
||||
disks.push(Arc::new(wrapper) as DiskStore);
|
||||
}
|
||||
|
||||
disks
|
||||
}
|
||||
|
||||
fn test_remote_peer(addr: &str) -> RemotePeerS3Client {
|
||||
let node = Node {
|
||||
url: url::Url::parse(addr).expect("test peer URL should parse"),
|
||||
@@ -1197,28 +1277,8 @@ mod tests {
|
||||
reset_local_disk_globals().await;
|
||||
|
||||
let temp_dir = TempDir::new().expect("create temp dir for local peer listing regression");
|
||||
let disk_path = temp_dir.path().join("disk1");
|
||||
std::fs::create_dir_all(&disk_path).expect("create disk path");
|
||||
|
||||
let mut endpoint = Endpoint::try_from(disk_path.to_str().expect("disk path to str")).expect("endpoint");
|
||||
endpoint.set_pool_index(0);
|
||||
endpoint.set_set_index(0);
|
||||
endpoint.set_disk_index(0);
|
||||
|
||||
let endpoint_pools = EndpointServerPools(vec![PoolEndpoints {
|
||||
legacy: false,
|
||||
set_count: 1,
|
||||
drives_per_set: 1,
|
||||
endpoints: Endpoints::from(vec![endpoint.clone()]),
|
||||
cmd_line: "local-get-bucket-info-survives-prior-walk-timeout".to_string(),
|
||||
platform: "test".to_string(),
|
||||
}]);
|
||||
|
||||
init_local_disks(endpoint_pools).await.expect("init local disks");
|
||||
|
||||
let disk = Arc::new(LocalDisk::new(&endpoint, false).await.expect("local disk should be created"));
|
||||
let wrapper = crate::disk::Disk::Local(Box::new(LocalDiskWrapper::new(disk, false)));
|
||||
let disk_store: DiskStore = Arc::new(wrapper);
|
||||
let disks = init_test_local_disks(&temp_dir, 1, "local-get-bucket-info-survives-prior-walk-timeout").await;
|
||||
let disk_store = disks[0].clone();
|
||||
let bucket = "test-bucket";
|
||||
let object = "test-object";
|
||||
|
||||
@@ -1262,6 +1322,103 @@ mod tests {
|
||||
reset_local_disk_globals().await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn local_get_bucket_info_requires_local_write_quorum() {
|
||||
reset_local_disk_globals().await;
|
||||
|
||||
let temp_dir = TempDir::new().expect("create temp dir for partial bucket regression");
|
||||
let disks = init_test_local_disks(&temp_dir, 2, "local-get-bucket-info-requires-local-write-quorum").await;
|
||||
|
||||
disks[0]
|
||||
.make_volume("partial-bucket")
|
||||
.await
|
||||
.expect("bucket should be created on one disk");
|
||||
|
||||
let err = LocalPeerS3Client::new(None, Some(vec![0]))
|
||||
.get_bucket_info("partial-bucket", &BucketOptions::default())
|
||||
.await
|
||||
.expect_err("partial bucket should not satisfy local write quorum");
|
||||
|
||||
assert_eq!(err, Error::ErasureWriteQuorum);
|
||||
|
||||
reset_local_disk_globals().await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn local_peer_filters_disks_by_pool() {
|
||||
reset_local_disk_globals().await;
|
||||
|
||||
let temp_dir = TempDir::new().expect("create temp dir for pool filtered local peer regression");
|
||||
let disks = init_test_local_disks_for_pools(&temp_dir, &[(0, 2), (1, 2)], "local-peer-filters-disks-by-pool").await;
|
||||
let bucket = "pool0-bucket";
|
||||
|
||||
disks[0]
|
||||
.make_volume(bucket)
|
||||
.await
|
||||
.expect("bucket should be created on pool 0 disk 0");
|
||||
disks[1]
|
||||
.make_volume(bucket)
|
||||
.await
|
||||
.expect("bucket should be created on pool 0 disk 1");
|
||||
|
||||
let pool0_info = LocalPeerS3Client::new(None, Some(vec![0]))
|
||||
.get_bucket_info(bucket, &BucketOptions::default())
|
||||
.await
|
||||
.expect("pool 0 peer should see bucket on pool 0 disks");
|
||||
assert_eq!(pool0_info.name, bucket);
|
||||
|
||||
let pool1_err = LocalPeerS3Client::new(None, Some(vec![1]))
|
||||
.get_bucket_info(bucket, &BucketOptions::default())
|
||||
.await
|
||||
.expect_err("pool 1 peer should not count pool 0 disks");
|
||||
assert_eq!(pool1_err, Error::VolumeNotFound);
|
||||
|
||||
let pool1_buckets = LocalPeerS3Client::new(None, Some(vec![1]))
|
||||
.list_bucket(&BucketOptions::default())
|
||||
.await
|
||||
.expect("pool 1 local listing should succeed against its own disks");
|
||||
assert!(pool1_buckets.is_empty());
|
||||
|
||||
reset_local_disk_globals().await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn heal_bucket_local_recreates_missing_bucket_volumes() {
|
||||
reset_local_disk_globals().await;
|
||||
|
||||
let temp_dir = TempDir::new().expect("create temp dir for bucket heal regression");
|
||||
let disks = init_test_local_disks(&temp_dir, 2, "heal-bucket-local-recreates-missing-bucket-volumes").await;
|
||||
let bucket = "healed-bucket";
|
||||
|
||||
disks[0]
|
||||
.make_volume(bucket)
|
||||
.await
|
||||
.expect("bucket should be created on one disk");
|
||||
disks[1]
|
||||
.stat_volume(bucket)
|
||||
.await
|
||||
.expect_err("second disk should start missing the bucket");
|
||||
|
||||
heal_bucket_local(
|
||||
bucket,
|
||||
&HealOpts {
|
||||
recreate: true,
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await
|
||||
.expect("bucket heal should recreate missing volumes");
|
||||
|
||||
for disk in disks {
|
||||
disk.stat_volume(bucket).await.expect("bucket should exist after heal");
|
||||
}
|
||||
|
||||
reset_local_disk_globals().await;
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_reduce_pool_write_quorum_uses_only_pool_participants() {
|
||||
let clients = vec![
|
||||
@@ -1314,4 +1471,29 @@ mod tests {
|
||||
|
||||
assert_eq!(err, Error::VolumeExists);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_list_bucket_reduces_visibility_quorum_by_pool_participants() {
|
||||
let bucket = BucketInfo {
|
||||
name: "existing-bucket".to_string(),
|
||||
..Default::default()
|
||||
};
|
||||
let peer_sys = S3PeerSys {
|
||||
clients: vec![
|
||||
test_peer_with_list_bucket(&[0], Ok(vec![bucket.clone()])),
|
||||
test_peer_with_list_bucket(&[1], Ok(vec![bucket.clone()])),
|
||||
test_peer_with_list_bucket(&[2], Ok(vec![bucket.clone()])),
|
||||
test_peer_with_list_bucket(&[3], Ok(vec![bucket.clone()])),
|
||||
],
|
||||
pools_count: 4,
|
||||
};
|
||||
|
||||
let buckets = peer_sys
|
||||
.list_bucket(&BucketOptions::default())
|
||||
.await
|
||||
.expect("single-participant pools should still expose visible buckets");
|
||||
|
||||
assert_eq!(buckets.len(), 1);
|
||||
assert_eq!(buckets[0].name, bucket.name);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -100,6 +100,19 @@ impl ECStore {
|
||||
|
||||
if let Err(err) = self.peer_sys.make_bucket(bucket, opts).await {
|
||||
let err = to_object_err(err.into(), vec![bucket]);
|
||||
if is_err_bucket_exists(&err)
|
||||
&& let Err(heal_err) = self
|
||||
.handle_heal_bucket(
|
||||
bucket,
|
||||
&HealOpts {
|
||||
recreate: true,
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await
|
||||
{
|
||||
warn!("best-effort bucket heal after BucketExists failed: {heal_err}");
|
||||
}
|
||||
if !is_err_bucket_exists(&err) {
|
||||
error!("make bucket failed: {err}");
|
||||
let _ = self
|
||||
|
||||
Reference in New Issue
Block a user