Compare commits

..

3 Commits

Author SHA1 Message Date
overtrue bf00d8fd27 test(e2e): add platform-safe selection updates 2026-08-23 09:04:20 +08:00
cxymds 20d1266496 fix(heal): fence format repair during pool transitions (#6342)
* fix(heal): fence format repair during pool transitions

* fix(heal): fence format writes during transitions
2026-08-23 04:24:48 +08:00
唐小鸭 f7003dfddd fix(admin): four site-replication interop correctness fixes (B5-rc T2) (#6399)
* fix(admin): send versioningEnabled on site replication make-bucket ops

The outbound make-with-versioning bucket-op query only carried
operation/createdAt/lockEnabled. MinIO's own create-bucket hook sends
versioningEnabled=true on this op, so align the outbound query with
MinIO's site-replication make-bucket wire contract. Route both outbound
builders (bootstrap plan and create-bucket hook) through one shared
builder that always appends versioningEnabled=true. RustFS's own inbound
handler force-enables versioning either way, so RustFS-to-RustFS
behavior is unchanged; the MinIO release verified against
(RELEASE.2025-09-07) also force-enables versioning regardless of the
flag, so this aligns the wire contract rather than changing observable
behavior there.

* fix(admin): propagate purge-deleted-bucket errors in site replication

The purge-deleted-bucket branch of the peer bucket-ops handler dropped
the delete_bucket error and answered 200, so a peer-driven purge that
failed (disk full, quorum loss) was reported as success while the
bucket survived on this site. Tolerate only bucket-not-found (the purge
raced an earlier replay or a local delete) and propagate every other
error through ApiError like the sibling delete branches do.

* fix(admin): derive fallback site deployment ID with UUIDv5

deployment_id_for_endpoint used DefaultHasher, whose algorithm is not
guaranteed stable across Rust releases. The fallback fires when a peer
response carries an empty deploymentID; the result is persisted in
site-replication state, used for collision disambiguation, and
broadcast to peers, so a toolchain bump could re-derive a different ID
for the same endpoint. Note that the add preflight currently rejects
that case upstream of this fallback. Derive UUIDv5 (NAMESPACE_URL) over
the canonical endpoint instead, and log a structured warn when a peer
metainfo response arrives without a deploymentID. Already persisted
fallback IDs are non-empty and therefore never re-derived, so existing
state is unaffected.

* fix(admin): stream site replication devnull body without 1MB cap

The site-replication devnull endpoint buffered the request body through
read_plain_admin_body, which enforces the 1MB admin body cap. MinIO
peers stream multi-megabyte probe bodies to this endpoint during site
netperf link checks and expect an unbounded discard, so any larger
probe got a 400 and was misreported as a broken link. Stream and
discard the body chunk by chunk with no size cap instead, mirroring
MinIO's io.Discard drain. The response stays 204 with an empty body.
2026-08-23 04:24:04 +08:00
14 changed files with 640 additions and 73 deletions
+1 -6
View File
@@ -178,14 +178,9 @@ jobs:
- name: Install Python tools
run: |
python3 -m pip install --user --upgrade pip "awscurl==0.44" "tox==4.60.0"
python3 -m pip install --user --upgrade pip awscurl tox
echo "$HOME/.local/bin" >> "$GITHUB_PATH"
- name: Verify Python tools
run: |
test "$(python3 -c 'import importlib.metadata as m; print(m.version("awscurl"))')" = "0.44"
test "$(python3 -c 'import importlib.metadata as m; print(m.version("tox"))')" = "4.60.0"
- name: Enable buildx
uses: docker/setup-buildx-action@8d2750c68a42422c14e847fe6c8ac0403b4cbd6f # v3
Generated
+1
View File
@@ -12688,6 +12688,7 @@ dependencies = [
"js-sys",
"rand 0.10.2",
"serde_core",
"sha1_smol",
"wasm-bindgen",
]
+6 -1
View File
@@ -278,4 +278,9 @@ listed by `cargo nextest list -p e2e_test`. Regenerate it when adding or
moving e2e tests so acceptance numbers in the test-strategy issues
(backlog#1147#1155) stay auditable. When a profile membership change is
intentional, review its JSON listing before updating the matching
`.config/e2e-*-selection.txt` test-ID digest.
`.config/e2e-*-selection.txt` test-ID digest. Update only the platform that
produced the listing:
```bash
python3 scripts/check_test_wiring.py --update-profile e2e-full /path/to/listing.json linux
```
+1 -1
View File
@@ -2026,7 +2026,7 @@ impl PoolMeta {
self.load_no_lock(pool).await
}
async fn load_no_lock<S>(&mut self, pool: Arc<S>) -> Result<()>
pub(crate) async fn load_no_lock<S>(&mut self, pool: Arc<S>) -> Result<()>
where
S: EcstoreObjectIO,
{
+20 -8
View File
@@ -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<Error>)> {
impl Sets {
pub(crate) async fn heal_format_with_fence<F>(&self, dry_run: bool, fence_lost: F) -> Result<(HealResultItem, Option<Error>)>
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<Error>)> {
self.heal_format_with_fence(dry_run, || false).await
}
#[tracing::instrument(skip(self))]
async fn heal_bucket(&self, bucket: &str, opts: &HealOpts) -> Result<HealResultItem> {
let mut result = HealResultItem {
+332 -3
View File
@@ -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<HealFormatPoolSkip> {
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<RebalanceMeta>)> {
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<Vec<Arc<Sets>>> {
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<ECStore>, 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)));
}
}
@@ -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,
+28
View File
@@ -547,6 +547,7 @@ struct MockStorage {
heal_object_outcome: Mutex<Option<MockHealObjectOutcome>>,
heal_object_outcomes: Mutex<HashMap<String, VecDeque<MockHealObjectOutcome>>>,
format_no_heal_required: Mutex<bool>,
format_error: Mutex<Option<Error>>,
global_format_calls: Mutex<u32>,
replacement_format_calls: Mutex<Vec<(usize, usize, Vec<String>)>>,
replacement_targets_ready: Mutex<bool>,
@@ -867,6 +868,9 @@ impl HealStorageAPI for MockStorage {
async fn heal_format(&self, _dry_run: bool) -> Result<(HealResultItem, Option<Error>)> {
*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");
+12
View File
@@ -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 {
+2 -2
View File
@@ -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 }
+121 -32
View File
@@ -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<SiteReplicationA
format!("invalid site replication metainfo from `{}`: {e}", site.endpoint),
)
})?;
if info.deployment_id.is_empty() {
// The peer will be tracked under a locally derived fallback ID
// (deployment_id_for_endpoint) instead of its real deployment ID.
warn!(
event = EVENT_ADMIN_SITE_REPLICATION_STATE,
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_SITE_REPLICATION,
result = "peer_deployment_id_missing",
peer_endpoint = %site.endpoint,
"admin site replication state"
);
}
let idp_body = send_peer_admin_get_request_with_client(
&client,
@@ -2206,20 +2219,30 @@ fn site_replication_bootstrap_token(uri: &Uri) -> Option<String> {
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<OffsetDateTime>) -> 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<Body>, _params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
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("<ObjectLockConfiguration/>")),
..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() {
+24 -5
View File
@@ -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!(
+2 -1
View File
@@ -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};
}
+81 -12
View File
@@ -219,7 +219,7 @@ def check_s3_tests_runner(root: Path) -> list[str]:
return []
def profile_selection(root: Path, profile: str) -> str:
def profile_selection_entries(root: Path, profile: str) -> tuple[Path, list[str], dict[str, str]]:
if not re.fullmatch(r"e2e-[a-z0-9-]+", profile):
raise ValueError(f"invalid e2e profile name: {profile}")
path = root / f".config/{profile}-selection.txt"
@@ -227,6 +227,11 @@ def profile_selection(root: Path, profile: str) -> str:
values = dict(line.split("=", 1) for line in lines if "=" in line)
if len(values) != len(lines) or any(not re.fullmatch(r"sha256(?:-[a-z0-9]+)?", key) for key in values):
raise ValueError(f"{path.relative_to(root).as_posix()}: invalid sha256 entry")
return path, lines, values
def profile_selection(root: Path, profile: str) -> str:
path, _, values = profile_selection_entries(root, profile)
key = f"sha256-{sys.platform}"
digest = values.get(key, values.get("sha256", ""))
if not re.fullmatch(r"[0-9a-f]{64}", digest):
@@ -234,6 +239,29 @@ def profile_selection(root: Path, profile: str) -> str:
return digest
def profile_listing_digest(listing: Path) -> tuple[int, str]:
data = json.loads(listing.read_text())
selected = sorted(
f"{suite_id}::{test_name}"
for suite_id, suite in data["rust-suites"].items()
for test_name, testcase in suite["testcases"].items()
if testcase.get("filter-match", {}).get("status") == "matches"
)
return len(selected), hashlib.sha256(("\n".join(selected) + "\n").encode()).hexdigest()
def update_profile_selection(root: Path, profile: str, listing: Path, platform: str) -> tuple[int, str, str]:
if not re.fullmatch(r"[a-z0-9]+", platform):
raise ValueError(f"invalid platform name: {platform}")
path, lines, values = profile_selection_entries(root, profile)
key = "sha256" if "sha256" in values else f"sha256-{platform}"
if key not in values:
raise ValueError(f"{path.relative_to(root).as_posix()}: missing {key} entry")
count, digest = profile_listing_digest(listing)
path.write_text("\n".join(f"{key}={digest}" if line.startswith(f"{key}=") else line for line in lines) + "\n")
return count, digest, key
def check_profile_definitions(root: Path) -> list[str]:
config = tomllib.loads((root / ".config/nextest.toml").read_text())
profiles = {
@@ -547,22 +575,15 @@ def check_scheduled_alerts(root: Path) -> list[str]:
def check_profile_listing(root: Path, profile: str, listing: Path) -> list[str]:
try:
expected_digest = profile_selection(root, profile)
data = json.loads(listing.read_text())
selected = sorted(
f"{suite_id}::{test_name}"
for suite_id, suite in data["rust-suites"].items()
for test_name, testcase in suite["testcases"].items()
if testcase.get("filter-match", {}).get("status") == "matches"
)
digest = hashlib.sha256(("\n".join(selected) + "\n").encode()).hexdigest()
count, digest = profile_listing_digest(listing)
except (FileNotFoundError, KeyError, TypeError, ValueError, json.JSONDecodeError) as error:
return [f"cannot read {profile} nextest listing: {error}"]
if digest != expected_digest:
return [
f"{profile} selection changed: count={len(selected)} sha256={digest}; "
f"{profile} selection changed: count={count} sha256={digest}; "
f"expected sha256={expected_digest}"
]
print(f"{profile} selection OK: {len(selected)} tests, sha256={digest}")
print(f"{profile} selection OK: {count} tests, sha256={digest}")
return []
@@ -707,6 +728,45 @@ class SelfTests(unittest.TestCase):
with mock.patch.object(sys, "platform", "linux"):
self.assertEqual(len(check_profile_listing(root, "e2e-full", listing)), 1)
def test_update_profile_selection_changes_only_requested_platform(self) -> None:
with tempfile.TemporaryDirectory() as tmp:
root = Path(tmp)
(root / ".config").mkdir()
darwin_digest = "a" * 64
(root / ".config/e2e-full-selection.txt").write_text(
f"sha256-darwin={darwin_digest}\nsha256-linux={'b' * 64}\n"
)
listing = root / "listing.json"
listing.write_text(
json.dumps(
{
"rust-suites": {
"suite": {
"testcases": {"linux": {"filter-match": {"status": "matches"}}}
}
}
}
)
)
count, digest, key = update_profile_selection(root, "e2e-full", listing, "linux")
self.assertEqual(count, 1)
self.assertEqual(key, "sha256-linux")
self.assertEqual(
(root / ".config/e2e-full-selection.txt").read_text(),
f"sha256-darwin={darwin_digest}\nsha256-linux={digest}\n",
)
with mock.patch.object(sys, "platform", "linux"):
self.assertEqual(check_profile_listing(root, "e2e-full", listing), [])
selection = root / ".config/e2e-smoke-selection.txt"
selection.write_text(f"sha256={'a' * 64}\n")
count, digest, key = update_profile_selection(root, "e2e-smoke", listing, "linux")
self.assertEqual((count, key), (1, "sha256"))
self.assertEqual(selection.read_text(), f"sha256={digest}\n")
def test_scheduled_alerts_require_completion_watchdog(self) -> None:
with tempfile.TemporaryDirectory() as tmp:
root = Path(tmp)
@@ -984,9 +1044,18 @@ def main() -> int:
print(f"ERROR: {error}", file=sys.stderr)
return 1
return 0
if len(sys.argv) == 5 and sys.argv[1] == "--update-profile":
try:
count, digest, key = update_profile_selection(ROOT, sys.argv[2], Path(sys.argv[3]), sys.argv[4])
except (FileNotFoundError, KeyError, TypeError, ValueError, json.JSONDecodeError) as error:
print(f"ERROR: cannot update {sys.argv[2]} selection: {error}", file=sys.stderr)
return 1
print(f"Updated .config/{sys.argv[2]}-selection.txt: count={count} {key}={digest}")
return 0
if sys.argv[1:]:
print(
"usage: check_test_wiring.py [--self-test | --check-profile PROFILE LISTING]",
"usage: check_test_wiring.py [--self-test | --check-profile PROFILE LISTING | "
"--update-profile PROFILE LISTING PLATFORM]",
file=sys.stderr,
)
return 2