mirror of
https://github.com/rustfs/rustfs.git
synced 2026-07-26 16:28:15 +00:00
fix(ecstore): repair decommission pool quorum (#2847)
Co-authored-by: houseme <housemecn@gmail.com>
This commit is contained in:
@@ -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));
|
||||
|
||||
@@ -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<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));
|
||||
|
||||
@@ -49,6 +49,32 @@ use tracing::{debug, info, warn};
|
||||
|
||||
type Client = Arc<Box<dyn PeerS3Client>>;
|
||||
|
||||
fn pool_participant_errors(clients: &[Client], errors: &[Option<Error>], pool_idx: usize) -> Vec<Option<Error>> {
|
||||
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<Error>]) -> Option<Error> {
|
||||
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<HealResultItem>;
|
||||
@@ -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<Option<DiskStore>> {
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[derive(Debug)]
|
||||
struct TestPeerS3Client {
|
||||
pools: Option<Vec<usize>>,
|
||||
make_bucket_result: Result<()>,
|
||||
}
|
||||
|
||||
#[async_trait]
|
||||
impl PeerS3Client for TestPeerS3Client {
|
||||
async fn heal_bucket(&self, _bucket: &str, _opts: &HealOpts) -> Result<HealResultItem> {
|
||||
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<Vec<BucketInfo>> {
|
||||
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<BucketInfo> {
|
||||
unreachable!("not used by quorum tests")
|
||||
}
|
||||
|
||||
fn get_pools(&self) -> Option<Vec<usize>> {
|
||||
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);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user