diff --git a/crates/ecstore/src/rebalance.rs b/crates/ecstore/src/rebalance.rs index 66754cdfe..0404bcc2d 100644 --- a/crates/ecstore/src/rebalance.rs +++ b/crates/ecstore/src/rebalance.rs @@ -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)); diff --git a/crates/ecstore/src/rpc/peer_s3_client.rs b/crates/ecstore/src/rpc/peer_s3_client.rs index 7c906f361..2dbddc578 100644 --- a/crates/ecstore/src/rpc/peer_s3_client.rs +++ b/crates/ecstore/src/rpc/peer_s3_client.rs @@ -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 { + 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 { - 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> { - 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 = uniq_map + let buckets: Vec = 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 { - 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 { 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>) -> Result { 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 { as_clone.write().await[idx] = DriveState::Ok.to_string(); return None; @@ -1090,6 +1102,7 @@ mod tests { struct TestPeerS3Client { pools: Option>, make_bucket_result: Result<()>, + list_bucket_result: Result>, } #[async_trait] @@ -1103,7 +1116,7 @@ mod tests { } async fn list_bucket(&self, _opts: &BucketOptions) -> Result> { - 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>) -> 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>, + ) -> 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 { + 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 { + 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); + } } diff --git a/crates/ecstore/src/store/bucket.rs b/crates/ecstore/src/store/bucket.rs index ce3f60a4c..8b330575c 100644 --- a/crates/ecstore/src/store/bucket.rs +++ b/crates/ecstore/src/store/bucket.rs @@ -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