diff --git a/Cargo.lock b/Cargo.lock index 637202fd5..b7b1ec220 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -12688,6 +12688,7 @@ dependencies = [ "js-sys", "rand 0.10.2", "serde_core", + "sha1_smol", "wasm-bindgen", ] diff --git a/crates/ecstore/src/core/pools.rs b/crates/ecstore/src/core/pools.rs index 9d8e0d01a..4ca460c20 100644 --- a/crates/ecstore/src/core/pools.rs +++ b/crates/ecstore/src/core/pools.rs @@ -2026,7 +2026,7 @@ impl PoolMeta { self.load_no_lock(pool).await } - async fn load_no_lock(&mut self, pool: Arc) -> Result<()> + pub(crate) async fn load_no_lock(&mut self, pool: Arc) -> Result<()> where S: EcstoreObjectIO, { diff --git a/crates/ecstore/src/core/sets.rs b/crates/ecstore/src/core/sets.rs index d9b354a08..0c79c080f 100644 --- a/crates/ecstore/src/core/sets.rs +++ b/crates/ecstore/src/core/sets.rs @@ -988,14 +988,11 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for Sets { } } -#[async_trait::async_trait] -impl crate::storage_api_contracts::heal::HealOperations for Sets { - type Error = Error; - type HealResultItem = HealResultItem; - type HealOptions = HealOpts; - - #[tracing::instrument(skip(self))] - async fn heal_format(&self, dry_run: bool) -> Result<(HealResultItem, Option)> { +impl Sets { + pub(crate) async fn heal_format_with_fence(&self, dry_run: bool, fence_lost: F) -> Result<(HealResultItem, Option)> + where + F: Fn() -> bool + Send + Sync, + { let (disks, init_errs) = init_storage_disks_with_errors( &self.endpoints.endpoints, &DiskOption { @@ -1068,6 +1065,9 @@ impl crate::storage_api_contracts::heal::HealOperations for Sets { // Save new formats `format.json` on unformatted disks. for (index, (fm, disk)) in tmp_new_formats.iter_mut().zip(disks.iter()).enumerate() { if fm.is_some() && disk.is_some() { + if fence_lost() { + return Ok((res, Some(StorageError::SlowDown))); + } if let Err(err) = save_format_file(disk, fm).await { if let Some(disk) = disk.as_ref() { let _ = disk.close().await; @@ -1101,6 +1101,18 @@ impl crate::storage_api_contracts::heal::HealOperations for Sets { } Ok((res, None)) } +} + +#[async_trait::async_trait] +impl crate::storage_api_contracts::heal::HealOperations for Sets { + type Error = Error; + type HealResultItem = HealResultItem; + type HealOptions = HealOpts; + + #[tracing::instrument(skip(self))] + async fn heal_format(&self, dry_run: bool) -> Result<(HealResultItem, Option)> { + self.heal_format_with_fence(dry_run, || false).await + } #[tracing::instrument(skip(self))] async fn heal_bucket(&self, bucket: &str, opts: &HealOpts) -> Result { let mut result = HealResultItem { diff --git a/crates/ecstore/src/store/heal.rs b/crates/ecstore/src/store/heal.rs index ffac77751..a7aed758d 100644 --- a/crates/ecstore/src/store/heal.rs +++ b/crates/ecstore/src/store/heal.rs @@ -13,7 +13,12 @@ // limitations under the License. use super::*; +use crate::core::pools::POOL_META_NAME; +use crate::services::rebalance::{REBAL_META_NAME, RebalStatus}; +use crate::set_disk::get_lock_acquire_timeout; use crate::storage_api_contracts::heal::HealOperations as _; +use crate::storage_api_contracts::namespace::NamespaceLocking as _; +use rustfs_lock::NamespaceLockGuard; use tracing::trace; const LOG_COMPONENT_ECSTORE: &str = "ecstore"; @@ -30,7 +35,119 @@ fn invalid_heal_pool_index(pool_idx: usize, pool_count: usize) -> Error { ) } +#[derive(Debug, Clone, Copy)] +enum HealFormatPoolSkip { + Completed, + Retryable, +} + +fn classify_heal_format_pool( + pool_idx: usize, + pool_cmd_line: &str, + pool_meta: &PoolMeta, + rebalance_meta: Option<&RebalanceMeta>, +) -> Option { + let Some(pool) = pool_meta.pools.get(pool_idx) else { + return Some(HealFormatPoolSkip::Retryable); + }; + + if pool.id != pool_idx || pool_cmd_line.is_empty() || pool.cmd_line.is_empty() || pool.cmd_line != pool_cmd_line { + return Some(HealFormatPoolSkip::Retryable); + } + + if let Some(decommission) = pool.decommission.as_ref() { + if decommission.complete { + return Some(HealFormatPoolSkip::Completed); + } + if decommission.failed || decommission.canceled || decommission.queued || pool_meta.is_suspended(pool_idx) { + return Some(HealFormatPoolSkip::Retryable); + } + } + + if let Some(meta) = rebalance_meta { + let Some(pool_stats) = meta.pool_stats.get(pool_idx) else { + return Some(HealFormatPoolSkip::Retryable); + }; + if pool_stats.info.stopping || (pool_stats.participating && pool_stats.info.status == RebalStatus::Started) { + return Some(HealFormatPoolSkip::Retryable); + } + } + + None +} + +fn heal_format_pool_skip_error(skip: HealFormatPoolSkip) -> Error { + match skip { + HealFormatPoolSkip::Completed => StorageError::NoHealRequired, + HealFormatPoolSkip::Retryable => StorageError::SlowDown, + } +} + +fn heal_format_fence_lost_error() -> Error { + StorageError::SlowDown +} + impl ECStore { + async fn acquire_heal_format_fence( + &self, + ) -> Result<(NamespaceLockGuard, NamespaceLockGuard, PoolMeta, Option)> { + let metadata_pool = self + .pools + .first() + .cloned() + .ok_or_else(|| Error::other("heal format requires at least one storage pool"))?; + + // Metadata fence order is part of the decommission/rebalance protocol: + // pool.bin must always be acquired before rebalance.bin. + let pool_lock = metadata_pool.new_ns_lock(RUSTFS_META_BUCKET, POOL_META_NAME).await?; + let pool_guard = pool_lock.get_write_lock(get_lock_acquire_timeout()).await?; + let rebalance_lock = metadata_pool.new_ns_lock(RUSTFS_META_BUCKET, REBAL_META_NAME).await?; + let rebalance_guard = rebalance_lock.get_write_lock(get_lock_acquire_timeout()).await?; + + if pool_guard.is_lock_lost() || rebalance_guard.is_lock_lost() { + return Err(heal_format_fence_lost_error()); + } + + let mut pool_meta = PoolMeta::default(); + pool_meta.load_no_lock(metadata_pool.clone()).await?; + if pool_meta.pools.len() != self.pools.len() + || pool_meta.pools.iter().enumerate().any(|(pool_idx, pool)| { + pool.id != pool_idx || pool.cmd_line.is_empty() || pool.cmd_line != self.pools[pool_idx].endpoints.cmd_line + }) + { + return Err(heal_format_fence_lost_error()); + } + + let mut rebalance_meta = RebalanceMeta::new(); + let rebalance_meta = match rebalance_meta + .load_with_opts( + metadata_pool, + ObjectOptions { + no_lock: true, + ..Default::default() + }, + ) + .await + { + Ok(()) => Some(rebalance_meta), + Err(Error::ConfigNotFound) => None, + Err(err) => return Err(err), + }; + + if rebalance_meta + .as_ref() + .is_some_and(|meta| meta.pool_stats.len() != self.pools.len()) + { + return Err(heal_format_fence_lost_error()); + } + + if pool_guard.is_lock_lost() || rebalance_guard.is_lock_lost() { + return Err(heal_format_fence_lost_error()); + } + + Ok((pool_guard, rebalance_guard, pool_meta, rebalance_meta)) + } + fn get_pools_for_heal_object(&self, opts: &HealOpts) -> Result>> { match opts.pool { Some(pool_idx) => Ok(vec![ @@ -52,9 +169,26 @@ impl ECStore { }; let mut count_no_heal = 0; + let mut count_completed = 0; let mut first_error = None; - for pool in self.pools.iter() { - let (mut result, err) = pool.heal_format(dry_run).await?; + for (pool_idx, pool) in self.pools.iter().enumerate() { + let (pool_guard, rebalance_guard, pool_meta, rebalance_meta) = self.acquire_heal_format_fence().await?; + if pool_guard.is_lock_lost() || rebalance_guard.is_lock_lost() { + first_error.get_or_insert(heal_format_fence_lost_error()); + break; + } + if let Some(skip) = classify_heal_format_pool(pool_idx, &pool.endpoints.cmd_line, &pool_meta, rebalance_meta.as_ref()) + { + if matches!(skip, HealFormatPoolSkip::Completed) { + count_completed += 1; + } else { + first_error.get_or_insert(heal_format_pool_skip_error(skip)); + } + continue; + } + + let fence_lost = || pool_guard.is_lock_lost() || rebalance_guard.is_lock_lost(); + let (mut result, err) = pool.heal_format_with_fence(dry_run, fence_lost).await?; if let Some(err) = err { match err { StorageError::NoHealRequired => { @@ -69,11 +203,18 @@ impl ECStore { r.set_count += result.set_count; r.before.drives.append(&mut result.before.drives); r.after.drives.append(&mut result.after.drives); + + // A lease can be lost after the final write; fail closed before + // reporting the pool as successfully healed. + if pool_guard.is_lock_lost() || rebalance_guard.is_lock_lost() { + first_error.get_or_insert(heal_format_fence_lost_error()); + break; + } } if let Some(err) = first_error { return Ok((r, Some(err))); } - if count_no_heal == self.pools.len() { + if count_no_heal + count_completed == self.pools.len() { info!( event = EVENT_HEAL_FORMAT_COMPLETED, component = LOG_COMPONENT_ECSTORE, @@ -302,6 +443,7 @@ mod tests { use crate::disk::{DeleteOptions, DiskOption, format::FormatV3, new_disk}; use crate::layout::endpoints::{EndpointServerPools, Endpoints, PoolEndpoints}; use crate::runtime::instance::InstanceContext; + use crate::services::rebalance::{RebalanceInfo, RebalanceStats}; use crate::storage_api_contracts::bucket::{BucketOperations, MakeBucketOptions}; use crate::storage_api_contracts::object::{ObjectIO as _, ObjectOperations}; use crate::store::init_format::{load_format_erasure, save_format_file}; @@ -353,6 +495,164 @@ mod tests { } } + fn pool_meta_with_decommission(info: PoolDecommissionInfo) -> PoolMeta { + PoolMeta { + pools: vec![PoolStatus { + id: 0, + cmd_line: "pool-0".to_string(), + last_update: OffsetDateTime::UNIX_EPOCH, + decommission: Some(info), + }], + ..Default::default() + } + } + + #[test] + fn heal_format_pool_state_barriers_are_classified() { + let active = pool_meta_with_decommission(PoolDecommissionInfo { + start_time: Some(OffsetDateTime::UNIX_EPOCH), + ..Default::default() + }); + assert!(matches!( + classify_heal_format_pool(0, "pool-0", &active, None), + Some(HealFormatPoolSkip::Retryable) + )); + + for info in [ + PoolDecommissionInfo { + failed: true, + ..Default::default() + }, + PoolDecommissionInfo { + canceled: true, + ..Default::default() + }, + ] { + assert!(matches!( + classify_heal_format_pool(0, "pool-0", &pool_meta_with_decommission(info), None), + Some(HealFormatPoolSkip::Retryable) + )); + } + + let completed = pool_meta_with_decommission(PoolDecommissionInfo { + complete: true, + ..Default::default() + }); + assert!(matches!( + classify_heal_format_pool(0, "pool-0", &completed, None), + Some(HealFormatPoolSkip::Completed) + )); + } + + #[test] + fn heal_format_pool_rebalance_barriers_and_identity_are_fail_closed() { + let identity_meta = pool_meta_with_decommission(PoolDecommissionInfo::default()); + let rebalance = RebalanceMeta { + pool_stats: vec![RebalanceStats { + participating: true, + info: RebalanceInfo { + status: RebalStatus::Started, + ..Default::default() + }, + ..Default::default() + }], + ..Default::default() + }; + assert!(matches!( + classify_heal_format_pool(0, "pool-0", &identity_meta, Some(&rebalance)), + Some(HealFormatPoolSkip::Retryable) + )); + + let stopping = RebalanceMeta { + pool_stats: vec![RebalanceStats { + info: RebalanceInfo { + stopping: true, + ..Default::default() + }, + ..Default::default() + }], + ..Default::default() + }; + assert!(matches!( + classify_heal_format_pool(0, "pool-0", &identity_meta, Some(&stopping)), + Some(HealFormatPoolSkip::Retryable) + )); + + let identity = pool_meta_with_decommission(PoolDecommissionInfo::default()); + assert!(matches!( + classify_heal_format_pool(0, "pool-new", &identity, None), + Some(HealFormatPoolSkip::Retryable) + )); + + let identity_without_decommission = PoolMeta { + pools: vec![PoolStatus { + id: 0, + cmd_line: "pool-0".to_string(), + last_update: OffsetDateTime::UNIX_EPOCH, + decommission: None, + }], + ..Default::default() + }; + assert!(matches!( + classify_heal_format_pool(0, "pool-new", &identity_without_decommission, None), + Some(HealFormatPoolSkip::Retryable) + )); + + assert!(matches!( + classify_heal_format_pool(0, "", &identity_meta, None), + Some(HealFormatPoolSkip::Retryable) + )); + + assert!(matches!( + classify_heal_format_pool(0, "pool-0", &PoolMeta::default(), None), + Some(HealFormatPoolSkip::Retryable) + )); + + let stopped = RebalanceMeta { + stopped_at: Some(OffsetDateTime::UNIX_EPOCH), + pool_stats: vec![RebalanceStats { + participating: true, + info: RebalanceInfo { + status: RebalStatus::Stopped, + ..Default::default() + }, + ..Default::default() + }], + ..Default::default() + }; + assert!(classify_heal_format_pool(0, "pool-0", &identity_meta, Some(&stopped)).is_none()); + + let stopping_after_stop = RebalanceMeta { + stopped_at: Some(OffsetDateTime::UNIX_EPOCH), + pool_stats: vec![RebalanceStats { + participating: true, + info: RebalanceInfo { + status: RebalStatus::Started, + stopping: true, + ..Default::default() + }, + ..Default::default() + }], + ..Default::default() + }; + assert!(matches!( + classify_heal_format_pool(0, "pool-0", &identity_meta, Some(&stopping_after_stop)), + Some(HealFormatPoolSkip::Retryable) + )); + } + + #[test] + fn skipped_heal_format_pool_is_never_reported_as_success() { + assert!(matches!( + heal_format_pool_skip_error(HealFormatPoolSkip::Retryable), + StorageError::SlowDown + )); + assert!(matches!( + heal_format_pool_skip_error(HealFormatPoolSkip::Completed), + StorageError::NoHealRequired + )); + } + async fn multi_pool_heal_store() -> (tempfile::TempDir, Arc, CancellationToken) { let temp_dir = tempfile::tempdir().expect("multi-pool heal test directory should be created"); let mut pool_endpoints = Vec::new(); @@ -889,6 +1189,18 @@ mod tests { bucket_fence_registry: std::sync::Arc::default(), }; + let err = store + .handle_heal_format(false) + .await + .expect_err("missing pool metadata must fail closed before format writes"); + assert!(matches!(err, StorageError::SlowDown)); + + let pool_meta = PoolMeta::new(&store.pools, &PoolMeta::default()); + pool_meta + .save(store.pools.clone()) + .await + .expect("pool metadata should be persisted before format heal"); + let (result, err) = store .handle_heal_format(false) .await @@ -902,5 +1214,22 @@ mod tests { .await .expect("the later pool should be healed despite the first pool error"); assert_eq!(healed.erasure.this, recoverable_format.erasure.sets[0][2]); + + let mut completed_meta = PoolMeta::new(&store.pools, &PoolMeta::default()); + for status in &mut completed_meta.pools { + status.decommission = Some(PoolDecommissionInfo { + complete: true, + ..Default::default() + }); + } + completed_meta + .save(store.pools.clone()) + .await + .expect("completed pool metadata should be persisted"); + let (_, err) = store + .handle_heal_format(false) + .await + .expect("completed pools should be reported as a no-op"); + assert!(matches!(err, Some(StorageError::NoHealRequired))); } } diff --git a/crates/heal/src/heal/task/heal_erasure_set.rs b/crates/heal/src/heal/task/heal_erasure_set.rs index e0cfe5a90..7b25bcba2 100644 --- a/crates/heal/src/heal/task/heal_erasure_set.rs +++ b/crates/heal/src/heal/task/heal_erasure_set.rs @@ -231,6 +231,10 @@ impl HealTask { "Heal erasure set format repair skipped because no format heal was required" ); } else { + let error = e; + if error.is_recoverable_heal() { + return Err(error); + } error!( target: "rustfs::heal::task", event = EVENT_HEAL_ERASURE_SET_RESULT, @@ -239,7 +243,7 @@ impl HealTask { task_id = %self.id, set_disk_id, result = "format_failed", - error = %e, + error = %error, "Heal erasure set failed" ); { @@ -247,7 +251,7 @@ impl HealTask { progress.update_progress(4, 4, 0, 0); } return Err(Error::TaskExecutionFailed { - message: format!("Failed to heal disk format for {set_disk_id}: {e}"), + message: format!("Failed to heal disk format for {set_disk_id}: {error}"), }); } } else { @@ -284,6 +288,9 @@ impl HealTask { Err(Error::TaskCancelled) => return Err(Error::TaskCancelled), Err(Error::TaskTimeout) => return Err(Error::TaskTimeout), Err(e) => { + if e.is_recoverable_heal() { + return Err(e); + } error!( target: "rustfs::heal::task", event = EVENT_HEAL_ERASURE_SET_RESULT, diff --git a/crates/heal/src/heal/task/tests.rs b/crates/heal/src/heal/task/tests.rs index 464ff3606..f2b442205 100644 --- a/crates/heal/src/heal/task/tests.rs +++ b/crates/heal/src/heal/task/tests.rs @@ -547,6 +547,7 @@ struct MockStorage { heal_object_outcome: Mutex>, heal_object_outcomes: Mutex>>, format_no_heal_required: Mutex, + format_error: Mutex>, global_format_calls: Mutex, replacement_format_calls: Mutex)>>, replacement_targets_ready: Mutex, @@ -867,6 +868,9 @@ impl HealStorageAPI for MockStorage { async fn heal_format(&self, _dry_run: bool) -> Result<(HealResultItem, Option)> { *self.global_format_calls.lock().unwrap() += 1; + if let Some(error) = self.format_error.lock().unwrap().take() { + return Err(error); + } let no_heal_required = *self.format_no_heal_required.lock().unwrap(); if no_heal_required { Ok((HealResultItem::default(), Some(Error::Storage(EcstoreError::NoHealRequired)))) @@ -2052,6 +2056,30 @@ async fn test_erasure_set_heal_continues_after_format_no_heal_required() { ); } +#[tokio::test] +async fn erasure_set_format_slowdown_is_propagated() { + let storage = Arc::new(MockStorage { + format_error: Mutex::new(Some(Error::Storage(EcstoreError::SlowDown))), + ..Default::default() + }); + let request = HealRequest::new( + HealType::ErasureSet { + buckets: Vec::new(), + set_disk_id: "pool_0_set_0".to_string(), + }, + HealOptions::default(), + HealPriority::Normal, + ); + let task = HealTask::from_request(request, storage); + + let error = task + .execute() + .await + .expect_err("format SlowDown must remain recoverable for the task manager"); + + assert!(matches!(error, Error::Storage(EcstoreError::SlowDown))); +} + #[tokio::test] async fn erasure_set_bucket_prepass_failure_stops_before_object_heal() { let temp = TempDir::new().expect("temporary directory should be created"); diff --git a/crates/test-utils/src/lib.rs b/crates/test-utils/src/lib.rs index 59b4f3e0c..0188c7064 100644 --- a/crates/test-utils/src/lib.rs +++ b/crates/test-utils/src/lib.rs @@ -245,6 +245,18 @@ impl TestECStoreEnvBuilder { .await .expect("build test ECStore"); + // The production bootstrap only persists pool.bin from the elected + // first cluster node. Test stores intentionally have no cluster + // election, but heal-format still requires that durable fence before + // it can write any disk format. Materialize the validated topology + // here so the shared fixture models a ready single-node store. + let mut pool_meta = ecstore.pool_meta.read().await.clone(); + pool_meta.dont_save = false; + pool_meta + .save(ecstore.pools.clone()) + .await + .expect("persist test pool metadata"); + if self.init_bucket_metadata { let buckets_list = ecstore .list_bucket(&BucketOptions { diff --git a/rustfs/Cargo.toml b/rustfs/Cargo.toml index 28f4ad259..af5514870 100644 --- a/rustfs/Cargo.toml +++ b/rustfs/Cargo.toml @@ -322,7 +322,7 @@ thiserror = { workspace = true } tracing.workspace = true url = { workspace = true } urlencoding = { workspace = true } -uuid = { workspace = true, features = ["v4", "fast-rng", "macro-diagnostics"] } +uuid = { workspace = true, features = ["v4", "v5", "fast-rng", "macro-diagnostics"] } zip = { workspace = true } libc = { workspace = true } rand = { workspace = true, features = ["serde"] } @@ -345,7 +345,7 @@ libsystemd.workspace = true libmimalloc-sys.workspace = true [dev-dependencies] -uuid = { workspace = true, features = ["v4", "fast-rng", "macro-diagnostics"] } +uuid = { workspace = true, features = ["v4", "v5", "fast-rng", "macro-diagnostics"] } serial_test = { workspace = true } tempfile = { workspace = true } aws-config = { workspace = true } diff --git a/rustfs/src/admin/handlers/site_replication.rs b/rustfs/src/admin/handlers/site_replication.rs index 1bb0379df..23c8b680b 100644 --- a/rustfs/src/admin/handlers/site_replication.rs +++ b/rustfs/src/admin/handlers/site_replication.rs @@ -41,7 +41,7 @@ use crate::admin::storage_api::config::save_admin_config; use crate::admin::storage_api::contract::bucket::{ BucketOperations, BucketOptions, DeleteBucketOptions, MakeBucketOptions, SRBucketDeleteOp, }; -use crate::admin::storage_api::error::Error as StorageError; +use crate::admin::storage_api::error::{Error as StorageError, is_err_bucket_not_found}; use crate::admin::storage_api::runtime::ECStore; use crate::admin::utils::{encode_compatible_admin_payload, read_compatible_admin_body}; use crate::auth::constant_time_eq; @@ -55,6 +55,7 @@ use crate::storage::storage_api::{ use base64::Engine; use base64::engine::general_purpose::STANDARD as BASE64_STANDARD; use base64::engine::general_purpose::URL_SAFE_NO_PAD; +use futures::StreamExt; use hmac::{Hmac, Mac}; use http::header::{CONTENT_TYPE, HOST}; use http::{HeaderMap, HeaderValue, Uri}; @@ -2096,6 +2097,18 @@ async fn remote_add_preflight_info(site: &PeerSite) -> S3Result Option { query_pairs(uri).get("bootstrapToken").cloned() } -fn bootstrap_bucket_make_op_path(bucket: &SRBucketInfo) -> String { +/// Query for a peer `make-with-versioning` bucket op. `versioningEnabled` +/// always travels so the outbound query matches MinIO's site-replication +/// make-bucket wire contract: MinIO's own create-bucket hook sends +/// `versioningEnabled=true` on this op. RustFS's inbound handler +/// force-enables versioning either way. +fn make_with_versioning_bucket_op_path(bucket: &str, created_at: Option<&str>, lock_enabled: bool) -> String { let mut query = form_urlencoded::Serializer::new(String::new()); - query.append_pair("bucket", &bucket.bucket); - query.append_pair("operation", "make-with-versioning"); - if let Some(created_at) = bucket - .created_at - .and_then(|value| value.format(&time::format_description::well_known::Rfc3339).ok()) - { - query.append_pair("createdAt", &created_at); + query.append_pair("bucket", bucket); + query.append_pair("operation", SITE_REPLICATION_BUCKET_OP_MAKE_WITH_VERSIONING); + query.append_pair("versioningEnabled", "true"); + if let Some(created_at) = created_at { + query.append_pair("createdAt", created_at); } - if bucket.object_lock_config.is_some() { + if lock_enabled { query.append_pair("lockEnabled", "true"); } - format!("/rustfs/admin/v3/site-replication/peer/bucket-ops?{}", query.finish()) + format!("{SITE_REPLICATION_PEER_BUCKET_OPS_PATH}?{}", query.finish()) +} + +fn bootstrap_bucket_make_op_path(bucket: &SRBucketInfo) -> String { + let created_at = bucket + .created_at + .and_then(|value| value.format(&time::format_description::well_known::Rfc3339).ok()); + make_with_versioning_bucket_op_path(&bucket.bucket, created_at.as_deref(), bucket.object_lock_config.is_some()) } fn bootstrap_bucket_meta_item(bucket: &SRBucketInfo, item_type: &str, updated_at: Option) -> SRBucketMeta { @@ -4246,16 +4269,7 @@ async fn broadcast_site_replication_make_bucket( .format(&time::format_description::well_known::Rfc3339) .unwrap_or_default(); - let path = { - let mut query = form_urlencoded::Serializer::new(String::new()); - query.append_pair("bucket", bucket); - query.append_pair("operation", "make-with-versioning"); - query.append_pair("createdAt", &created_at); - if lock_enabled { - query.append_pair("lockEnabled", "true"); - } - format!("/rustfs/admin/v3/site-replication/peer/bucket-ops?{}", query.finish()) - }; + let path = make_with_versioning_bucket_op_path(bucket, Some(&created_at), lock_enabled); let path = if let Some(token) = bootstrap_token { with_site_replication_bootstrap_token(&path, token) } else { @@ -10206,13 +10220,25 @@ impl Operation for SiteReplicationStatusHandler { } } +/// `POST /v3/site-replication/devnull` — peer link-check upload drain. +/// MinIO streams multi-megabyte probe bodies here during site netperf link +/// checks and expects an unbounded discard (its handler copies to io.Discard); +/// buffering through the 1MB admin body cap turned any larger probe into a +/// 400 and a false link failure. Stream and discard instead — no size cap. +async fn drain_site_replication_devnull(mut input: Body) -> S3Result<()> { + while let Some(chunk) = input.next().await { + chunk.map_err(|e| s3_error!(InvalidRequest, "failed to read devnull stream: {}", e))?; + } + Ok(()) +} + pub struct SiteReplicationDevNullHandler {} #[async_trait::async_trait] impl Operation for SiteReplicationDevNullHandler { async fn call(&self, req: S3Request, _params: Params<'_, '_>) -> S3Result> { validate_site_replication_admin_request(&req, AdminAction::SiteReplicationOperationAction).await?; - let _ = read_plain_admin_body(req.input).await?; + drain_site_replication_devnull(req.input).await?; Ok(empty_response(StatusCode::NO_CONTENT)) } } @@ -10471,6 +10497,19 @@ impl Operation for SRPeerJoinHandler { } } +/// Outcome of a peer-driven `purge-deleted-bucket` replay. A bucket that is +/// already gone means the purge raced an earlier replay or a local delete — +/// that is success — but any other failure must reach the sender like the +/// sibling delete branches do: swallowing it answered 200 while the bucket +/// survived on this site. +fn purge_deleted_bucket_result(result: Result<(), StorageError>) -> S3Result<()> { + match result { + Ok(()) => Ok(()), + Err(err) if is_err_bucket_not_found(&err) => Ok(()), + Err(err) => Err(ApiError::from(err).into()), + } +} + pub struct SRPeerBucketOpsHandler {} #[async_trait::async_trait] @@ -10570,16 +10609,18 @@ impl Operation for SRPeerBucketOpsHandler { .map_err(ApiError::from)?; } "purge-deleted-bucket" => { - let _ = store - .delete_bucket( - &bucket, - &DeleteBucketOptions { - force: true, - srdelete_op: SRBucketDeleteOp::Purge, - ..Default::default() - }, - ) - .await; + purge_deleted_bucket_result( + store + .delete_bucket( + &bucket, + &DeleteBucketOptions { + force: true, + srdelete_op: SRBucketDeleteOp::Purge, + ..Default::default() + }, + ) + .await, + )?; } _ => return Err(s3_error!(InvalidRequest, "unsupported site replication bucket operation")), } @@ -13925,6 +13966,54 @@ mod tests { assert!(!query_flag(&uri, "missing")); } + /// A5 red-light: a `purge-deleted-bucket` replay must report success when + /// the bucket is already gone, and must propagate every other failure — + /// the swallowed error answered 200 while the bucket survived. + #[test] + fn test_purge_deleted_bucket_result_tolerates_only_missing_bucket() { + assert!(purge_deleted_bucket_result(Ok(())).is_ok()); + assert!(purge_deleted_bucket_result(Err(StorageError::BucketNotFound("photos".to_string()))).is_ok()); + assert!(purge_deleted_bucket_result(Err(StorageError::VolumeNotFound)).is_ok()); + let err = purge_deleted_bucket_result(Err(StorageError::StorageFull)) + .expect_err("non-not-found delete failures must propagate"); + assert_ne!(*err.code(), S3ErrorCode::NoSuchBucket); + } + + /// C5 red-light: the site-replication devnull drain must accept bodies + /// beyond the 1MB admin body cap — MinIO's link check streams large + /// probe bodies and treats a 400 as a broken link. + #[tokio::test] + async fn test_site_replication_devnull_drains_body_beyond_admin_cap() { + let body = Body::from(vec![0u8; MAX_ADMIN_REQUEST_BODY_SIZE + 1]); + drain_site_replication_devnull(body) + .await + .expect("devnull must drain bodies larger than the admin body cap"); + } + + /// A3 red-light: `versioningEnabled` must travel on every outbound + /// make-with-versioning bucket op so the query matches MinIO's + /// site-replication make-bucket wire contract (MinIO's own hook sends + /// `versioningEnabled=true` on this op). + #[test] + fn test_make_with_versioning_op_paths_send_versioning_enabled() { + let bucket = SRBucketInfo { + bucket: "photos".to_string(), + created_at: Some(OffsetDateTime::UNIX_EPOCH), + object_lock_config: Some(BASE64_STANDARD.encode("")), + ..Default::default() + }; + let bootstrap = bootstrap_bucket_make_op_path(&bucket); + assert!(bootstrap.contains("operation=make-with-versioning"), "{bootstrap}"); + assert!(bootstrap.contains("versioningEnabled=true"), "{bootstrap}"); + assert!(bootstrap.contains("createdAt="), "{bootstrap}"); + assert!(bootstrap.contains("lockEnabled=true"), "{bootstrap}"); + + // The broadcast path (create-bucket hook) shares the same builder. + let broadcast = make_with_versioning_bucket_op_path("photos", Some("1970-01-01T00:00:00Z"), false); + assert!(broadcast.contains("versioningEnabled=true"), "{broadcast}"); + assert!(!broadcast.contains("lockEnabled"), "{broadcast}"); + } + #[tokio::test] #[serial] async fn test_add_bootstrap_scope_only_allows_expected_bucket_setup_until_guard_drops() { diff --git a/rustfs/src/admin/site_replication_identity.rs b/rustfs/src/admin/site_replication_identity.rs index dc6440b4d..24784160b 100644 --- a/rustfs/src/admin/site_replication_identity.rs +++ b/rustfs/src/admin/site_replication_identity.rs @@ -13,9 +13,9 @@ // limitations under the License. use rustfs_madmin::{PeerInfo, SyncStatus}; -use std::collections::{BTreeMap, hash_map::DefaultHasher}; -use std::hash::{Hash, Hasher}; +use std::collections::BTreeMap; use url::Url; +use uuid::Uuid; fn has_http_scheme(endpoint: &str) -> bool { endpoint.get(..7).is_some_and(|prefix| prefix.eq_ignore_ascii_case("http://")) @@ -66,10 +66,12 @@ pub fn site_identity_key(endpoint: &str) -> String { .unwrap_or_else(|| trimmed.to_ascii_lowercase()) } +/// Fallback deployment ID for a peer that reported none. UUIDv5 over the +/// canonical endpoint: the ID is persisted in site-replication state and +/// broadcast to peers, so it must be identical across Rust toolchains +/// (`DefaultHasher` is not) and across spellings of the same endpoint. pub fn deployment_id_for_endpoint(endpoint: &str) -> String { - let mut hasher = DefaultHasher::new(); - endpoint.hash(&mut hasher); - format!("{:016x}", hasher.finish()) + Uuid::new_v5(&Uuid::NAMESPACE_URL, canonical_endpoint(endpoint).as_bytes()).to_string() } pub fn same_identity_endpoint(left: &str, right: &str) -> bool { @@ -174,6 +176,23 @@ mod tests { } } + /// B8 red-light: the fallback deployment ID must be a toolchain-stable + /// UUIDv5 over the canonical endpoint — `DefaultHasher` output is not + /// guaranteed stable across Rust releases, yet the ID is persisted in + /// site-replication state and broadcast to peers. + #[test] + fn deployment_id_for_endpoint_is_stable_uuid_v5_over_canonical_endpoint() { + let endpoint = "https://node-a.example.com:9000"; + let id = deployment_id_for_endpoint(endpoint); + let parsed = uuid::Uuid::parse_str(&id).expect("fallback deployment ID must be a UUID"); + assert_eq!(parsed.get_version_num(), 5, "fallback deployment ID must be UUIDv5"); + // Deterministic for the same endpoint and for spelling variants that + // share a canonical form; distinct endpoints stay distinct. + assert_eq!(id, deployment_id_for_endpoint(endpoint)); + assert_eq!(id, deployment_id_for_endpoint(" HTTPS://Node-A.Example.Com:9000/ ")); + assert_ne!(id, deployment_id_for_endpoint("https://node-b.example.com:9000")); + } + #[test] fn canonical_endpoint_accepts_case_insensitive_scheme() { assert_eq!( diff --git a/rustfs/src/admin/storage_api.rs b/rustfs/src/admin/storage_api.rs index 318718b49..bdf28f15e 100644 --- a/rustfs/src/admin/storage_api.rs +++ b/rustfs/src/admin/storage_api.rs @@ -51,7 +51,7 @@ mod ecstore_disk { } mod ecstore_error { - pub(crate) use crate::storage::storage_api::ecstore_error::StorageError; + pub(crate) use crate::storage::storage_api::ecstore_error::{StorageError, is_err_bucket_not_found}; } #[allow(unused_imports)] @@ -919,6 +919,7 @@ pub(crate) mod contract { } pub(crate) mod error { + pub(crate) use super::ecstore_error::is_err_bucket_not_found; pub(crate) use super::{Error, StorageError}; }