diff --git a/crates/ecstore/src/pools.rs b/crates/ecstore/src/pools.rs index 21ca3372b..f48077583 100644 --- a/crates/ecstore/src/pools.rs +++ b/crates/ecstore/src/pools.rs @@ -142,6 +142,13 @@ fn ensure_decommission_not_rebalancing(rebalance_running: bool) -> Result<()> { Ok(()) } +fn decommission_meta_bucket_options() -> MakeBucketOptions { + MakeBucketOptions { + force_create: true, + ..Default::default() + } +} + fn is_decommission_active(complete: bool, failed: bool, canceled: bool) -> bool { !complete && !failed && !canceled } @@ -2036,9 +2043,10 @@ impl ECStore { path_join(&[PathBuf::from(RUSTFS_META_BUCKET), PathBuf::from(BUCKET_META_PREFIX)]), ]; + let meta_bucket_opts = decommission_meta_bucket_options(); for bk in meta_buckets.iter() { if let Err(err) = self - .make_bucket(bk.to_string_lossy().to_string().as_str(), &MakeBucketOptions::default()) + .make_bucket(bk.to_string_lossy().to_string().as_str(), &meta_bucket_opts) .await && !is_err_bucket_exists(&err) { @@ -2737,12 +2745,13 @@ mod pools_tests { use super::{ DecomBucketInfo, DecommissionTerminalState, PoolDecommissionInfo, PoolMeta, PoolStatus, bind_decommission_cancelers, cancel_decommission_canceler, classify_decommission_terminal_state, count_decommission_item, - decommission_cancel_signal_result, decommission_item_size, decommission_start_guard_state, dedup_indices, - ensure_decommission_cancel_allowed, ensure_decommission_listing_disks_available, ensure_decommission_not_rebalancing, - ensure_decommission_start_allowed, ensure_decommission_terminal_operation_supported, - ensure_valid_decommission_pool_index, get_by_index, has_active_decommission_canceler, is_decommission_active, - is_decommission_cancel_terminal, load_decommission_entry_versions, mark_decommission_bucket_done, - require_decommission_store, resolve_decommission_bucket_done_save_result, resolve_decommission_bucket_state, + decommission_cancel_signal_result, decommission_item_size, decommission_meta_bucket_options, + decommission_start_guard_state, dedup_indices, ensure_decommission_cancel_allowed, + ensure_decommission_listing_disks_available, ensure_decommission_not_rebalancing, ensure_decommission_start_allowed, + ensure_decommission_terminal_operation_supported, ensure_valid_decommission_pool_index, get_by_index, + has_active_decommission_canceler, is_decommission_active, is_decommission_cancel_terminal, + load_decommission_entry_versions, mark_decommission_bucket_done, require_decommission_store, + resolve_decommission_bucket_done_save_result, resolve_decommission_bucket_state, resolve_decommission_check_after_list_result, resolve_decommission_entry_cleanup_delete_result, resolve_decommission_entry_reload_result, resolve_decommission_listing_worker_result, resolve_decommission_optional_bucket_config_result, resolve_decommission_pool_meta_reload_result, @@ -3376,6 +3385,13 @@ mod pools_tests { assert!(ensure_decommission_not_rebalancing(false).is_ok()); } + #[test] + fn test_decommission_meta_bucket_options_are_idempotent() { + let opts = decommission_meta_bucket_options(); + + assert!(opts.force_create); + } + #[test] fn test_is_decommission_active_true_only_when_not_terminal() { assert!(is_decommission_active(false, false, false)); diff --git a/crates/ecstore/src/rebalance.rs b/crates/ecstore/src/rebalance.rs index 567958650..aba7aaa27 100644 --- a/crates/ecstore/src/rebalance.rs +++ b/crates/ecstore/src/rebalance.rs @@ -780,6 +780,7 @@ impl ECStore { let cancel_tx = CancellationToken::new(); let rx = cancel_tx.clone(); + let mut meta_to_save = None; { let mut rebalance_meta = self.rebalance_meta.write().await; @@ -792,11 +793,19 @@ impl ECStore { info!("start_rebalance: already in progress, skip duplicate start"); return Ok(()); } + if complete_rebalance_pools_at_goal(meta, OffsetDateTime::now_utc()) { + meta_to_save = Some(meta.clone()); + } meta.cancel = Some(cancel_tx); drop(rebalance_meta); } + if let Some(meta) = meta_to_save { + let pool = clone_first_arc(self.pools.as_slice(), "start_rebalance: no pools available")?; + resolve_rebalance_meta_save_result(meta.save(pool).await, "start_rebalance complete pools at goal")?; + } + let participants = if let Some(ref meta) = *self.rebalance_meta.read().await { resolve_rebalance_participants(meta.pool_stats.as_slice(), self.pools.len()) } else { @@ -1126,6 +1135,30 @@ fn should_pool_participate(init_free_space: u64, init_capacity: u64, percent_fre init_capacity > 0 && percent_free_ratio(init_free_space, init_capacity) < percent_free_goal } +fn complete_rebalance_pools_at_goal(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) { + continue; + } + + if rebalance_goal_reached( + pool_stat.init_free_space, + pool_stat.init_capacity, + pool_stat.bytes, + meta.percent_free_goal, + ) { + pool_stat.info.status = RebalStatus::Completed; + pool_stat.info.end_time = Some(now); + pool_stat.info.last_error = None; + changed = true; + } + } + + changed +} + fn resolve_rebalance_worker_result( set_idx: usize, worker_result: std::result::Result, tokio::task::JoinError>, @@ -1838,12 +1871,12 @@ mod rebalance_unit_tests { GetObjectReader, HTTPRangeSpec, MigrationBackend, MigrationVersionResult, ObjectInfo, ObjectOptions, RebalSaveOpt, RebalStatus, 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, ensure_rebalance_listing_disks_available, ensure_rebalance_not_decommissioning, - ensure_valid_rebalance_pool_index, is_rebalance_stopped_terminal_event, load_rebalance_bucket_configs, - mark_rebalance_bucket_done, migrate_entry_version, next_rebal_bucket_from_stat, rebalance_delete_marker_opts, - rebalance_meta_load_no_data_error, rebalance_meta_load_unknown_format_error, rebalance_meta_load_unknown_version_error, - resolve_load_rebalance_stats_update_result, resolve_next_rebalance_bucket, resolve_rebalance_bucket_error, - resolve_rebalance_bucket_result, resolve_rebalance_entry_cleanup_delete_result, + clone_rebalance_pool_stats, complete_rebalance_pools_at_goal, ensure_rebalance_listing_disks_available, + ensure_rebalance_not_decommissioning, ensure_valid_rebalance_pool_index, is_rebalance_stopped_terminal_event, + load_rebalance_bucket_configs, mark_rebalance_bucket_done, migrate_entry_version, next_rebal_bucket_from_stat, + rebalance_delete_marker_opts, rebalance_meta_load_no_data_error, rebalance_meta_load_unknown_format_error, + rebalance_meta_load_unknown_version_error, 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, @@ -3248,6 +3281,44 @@ mod rebalance_unit_tests { assert!(!should_pool_participate(300, 1_000, 0.3)); } + #[test] + fn test_complete_rebalance_pools_at_goal_marks_started_participants_completed() { + let now = OffsetDateTime::from_unix_timestamp(1_000).unwrap(); + let mut meta = RebalanceMeta { + percent_free_goal: 0.5, + pool_stats: vec![ + RebalanceStats { + participating: true, + init_free_space: 400, + init_capacity: 1_000, + bytes: 50, + info: RebalanceInfo { + status: RebalStatus::Started, + ..Default::default() + }, + ..Default::default() + }, + RebalanceStats { + participating: true, + init_free_space: 100, + init_capacity: 1_000, + bytes: 0, + info: RebalanceInfo { + status: RebalStatus::Started, + ..Default::default() + }, + ..Default::default() + }, + ], + ..Default::default() + }; + + assert!(complete_rebalance_pools_at_goal(&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_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)); diff --git a/crates/ecstore/src/rpc/peer_s3_client.rs b/crates/ecstore/src/rpc/peer_s3_client.rs index d4dd4e8ff..3355de7de 100644 --- a/crates/ecstore/src/rpc/peer_s3_client.rs +++ b/crates/ecstore/src/rpc/peer_s3_client.rs @@ -49,6 +49,32 @@ use tracing::{debug, info, warn}; type Client = Arc>; +fn pool_participant_errors(clients: &[Client], errors: &[Option], pool_idx: usize) -> Vec> { + clients + .iter() + .zip(errors.iter()) + .filter_map(|(client, err)| { + if client.get_pools().unwrap_or_default().contains(&pool_idx) { + Some(err.clone()) + } else { + None + } + }) + .collect() +} + +fn pool_write_quorum(participant_count: usize) -> usize { + (participant_count / 2) + 1 +} + +fn reduce_pool_write_quorum_errs(per_pool_errs: &[Option]) -> Option { + if per_pool_errs.is_empty() { + return Some(Error::ErasureWriteQuorum); + } + + reduce_write_quorum_errs(per_pool_errs, BUCKET_OP_IGNORED_ERRS, pool_write_quorum(per_pool_errs.len())) +} + #[async_trait] pub trait PeerS3Client: Debug + Sync + Send + 'static { async fn heal_bucket(&self, bucket: &str, opts: &HealOpts) -> Result; @@ -190,18 +216,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) { tracing::error!("make_bucket per_pool_errs: {per_pool_errs:?}"); tracing::error!("make_bucket reduce_write_quorum_errs: {pool_err}"); return Err(pool_err); @@ -1031,6 +1047,50 @@ async fn clone_drives() -> Vec> { mod tests { use super::*; + #[derive(Debug)] + struct TestPeerS3Client { + pools: Option>, + make_bucket_result: Result<()>, + } + + #[async_trait] + impl PeerS3Client for TestPeerS3Client { + async fn heal_bucket(&self, _bucket: &str, _opts: &HealOpts) -> Result { + unreachable!("not used by quorum tests") + } + + async fn make_bucket(&self, _bucket: &str, _opts: &MakeBucketOptions) -> Result<()> { + self.make_bucket_result.clone() + } + + async fn list_bucket(&self, _opts: &BucketOptions) -> Result> { + unreachable!("not used by quorum tests") + } + + async fn delete_bucket(&self, _bucket: &str, _opts: &DeleteBucketOptions) -> Result<()> { + unreachable!("not used by quorum tests") + } + + async fn get_bucket_info(&self, _bucket: &str, _opts: &BucketOptions) -> Result { + unreachable!("not used by quorum tests") + } + + fn get_pools(&self) -> Option> { + self.pools.clone() + } + } + + fn test_peer(pools: &[usize]) -> Client { + test_peer_with_make_bucket(pools, Ok(())) + } + + fn test_peer_with_make_bucket(pools: &[usize], make_bucket_result: Result<()>) -> Client { + Arc::new(Box::new(TestPeerS3Client { + pools: Some(pools.to_vec()), + make_bucket_result, + })) + } + fn test_remote_peer(addr: &str) -> RemotePeerS3Client { let node = Node { url: url::Url::parse(addr).expect("test peer URL should parse"), @@ -1091,4 +1151,57 @@ mod tests { client.cancel_token.cancel(); } + + #[test] + fn test_reduce_pool_write_quorum_uses_only_pool_participants() { + let clients = vec![ + test_peer(&[0]), + test_peer(&[0]), + test_peer(&[0]), + test_peer(&[0]), + test_peer(&[1]), + test_peer(&[1]), + test_peer(&[1]), + test_peer(&[1]), + ]; + let errors = vec![ + Some(Error::VolumeExists), + Some(Error::VolumeExists), + Some(Error::VolumeExists), + Some(Error::VolumeExists), + None, + None, + None, + None, + ]; + + let per_pool_errs = pool_participant_errors(&clients, &errors, 0); + let err = reduce_pool_write_quorum_errs(&per_pool_errs).expect("all pool participants returned VolumeExists"); + + assert_eq!(err, Error::VolumeExists); + } + + #[tokio::test] + async fn test_make_bucket_reduces_quorum_by_pool_participants() { + let peer_sys = S3PeerSys { + clients: vec![ + test_peer_with_make_bucket(&[0], Err(Error::VolumeExists)), + test_peer_with_make_bucket(&[0], Err(Error::VolumeExists)), + test_peer_with_make_bucket(&[0], Err(Error::VolumeExists)), + test_peer_with_make_bucket(&[0], Err(Error::VolumeExists)), + test_peer(&[1]), + test_peer(&[1]), + test_peer(&[1]), + test_peer(&[1]), + ], + pools_count: 2, + }; + + let err = peer_sys + .make_bucket("existing-bucket", &MakeBucketOptions::default()) + .await + .expect_err("existing bucket should surface as VolumeExists, not quorum failure"); + + assert_eq!(err, Error::VolumeExists); + } }