Compare commits

...

5 Commits

Author SHA1 Message Date
overtrue 2cf18b6913 fix(ecstore): type decommission completion result 2026-08-23 04:50:38 +08:00
overtrue 087be055b5 Merge remote-tracking branch 'origin/main' into overtrue/fix-1904-unresolved-ledger 2026-08-23 04:33:37 +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
overtrue c18f9f6b74 fix(ecstore): persist unresolved decommission entries 2026-08-23 04:20:50 +08:00
15 changed files with 976 additions and 125 deletions
Generated
+1
View File
@@ -12688,6 +12688,7 @@ dependencies = [
"js-sys",
"rand 0.10.2",
"serde_core",
"sha1_smol",
"wasm-bindgen",
]
+2 -2
View File
@@ -241,8 +241,8 @@ pub mod cache {
pub mod capacity {
pub use crate::core::pools::{
PoolDecommissionInfo, PoolStatus, get_total_usable_capacity, get_total_usable_capacity_free, path2_bucket_object,
path2_bucket_object_with_base_path,
DecommissionUnresolvedEntry, PoolDecommissionInfo, PoolStatus, get_total_usable_capacity, get_total_usable_capacity_free,
path2_bucket_object, path2_bucket_object_with_base_path,
};
pub use crate::store::utils::is_reserved_or_invalid_bucket;
}
+397 -67
View File
@@ -804,6 +804,88 @@ fn ensure_decommission_generation(meta: &PoolMeta, idx: usize, generation: Offse
}
}
fn record_decommission_unresolved_entry(
meta: &mut PoolMeta,
idx: usize,
generation: OffsetDateTime,
entry: DecommissionUnresolvedEntry,
) -> Result<bool> {
ensure_decommission_generation(meta, idx, generation)?;
if entry.pool_index != idx || entry.source_generation != generation {
return Err(Error::other("decommission unresolved entry does not match the active pool generation"));
}
let pool_count = meta.pools.len();
let Some(pool) = meta.pools.get_mut(idx) else {
return Err(invalid_decommission_pool_index_error(pool_count, idx));
};
let Some(info) = pool.decommission.as_mut() else {
return Err(decommission_metadata_not_initialized_error("record decommission unresolved entry"));
};
let existing = info.unresolved_entries.iter_mut().find(|existing| {
existing.bucket == entry.bucket
&& existing.object == entry.object
&& existing.pool_index == entry.pool_index
&& existing.set_index == entry.set_index
&& existing.source_generation == entry.source_generation
});
let changed = match existing {
Some(existing) if existing == &entry => false,
Some(existing) => {
*existing = entry;
true
}
None => {
info.unresolved_entries.push(entry);
true
}
};
if changed {
pool.last_update = OffsetDateTime::now_utc();
}
Ok(changed)
}
fn reconcile_decommission_unresolved_entries_for_completion(
meta: &mut PoolMeta,
idx: usize,
verified_generation: Option<OffsetDateTime>,
) -> Result<()> {
let pool_count = meta.pools.len();
let Some(pool) = meta.pools.get(idx) else {
return Err(invalid_decommission_pool_index_error(pool_count, idx));
};
let Some(info) = pool.decommission.as_ref() else {
return Err(decommission_metadata_not_initialized_error("reconcile decommission unresolved entries"));
};
if info.unresolved_entries.is_empty() {
return Ok(());
}
let Some(generation) = verified_generation else {
return Err(Error::other(format!(
"failed to complete decommission for pool {idx}: {} unresolved listing entries remain",
info.unresolved_entries.len()
)));
};
ensure_decommission_generation(meta, idx, generation)?;
if info
.unresolved_entries
.iter()
.any(|entry| entry.source_generation != generation)
{
return Err(Error::other(format!(
"failed to complete decommission for pool {idx}: unresolved listing ledger contains a different generation"
)));
}
let Some(info) = meta.pools.get_mut(idx).and_then(|pool| pool.decommission.as_mut()) else {
return Err(decommission_metadata_not_initialized_error("reconcile decommission unresolved entries"));
};
info.unresolved_entries.clear();
Ok(())
}
async fn run_decommission_side_effect<T, F, Fut>(
rx: &CancellationToken,
operation_gate: &Arc<tokio::sync::RwLock<()>>,
@@ -981,18 +1063,15 @@ fn resolve_decommission_listing_error(listing_error: Option<Error>, entry_error:
}
}
fn decommission_unresolved_listing_error(
bucket: &str,
prefix: &str,
candidate: Option<&str>,
candidate_count: usize,
disk_error_count: usize,
pool_index: usize,
set_index: usize,
) -> Error {
let location = candidate.unwrap_or(prefix);
fn decommission_unresolved_listing_error(entry: &DecommissionUnresolvedEntry) -> Error {
Error::other(format!(
"decommission listing could not resolve metadata for {bucket}/{location} on pool {pool_index} set {set_index} ({candidate_count} candidate(s), {disk_error_count} disk error(s))"
"decommission listing could not resolve metadata for {bucket}/{object} on pool {pool_index} set {set_index} ({candidate_count} candidate(s), {disk_error_count} disk error(s))",
bucket = entry.bucket,
object = entry.object,
pool_index = entry.pool_index,
set_index = entry.set_index,
candidate_count = entry.candidate_count,
disk_error_count = entry.disk_error_count,
))
}
@@ -1004,22 +1083,32 @@ fn resolve_decommission_partial_listing_entry(
disk_error_count: usize,
pool_index: usize,
set_index: usize,
) -> Result<MetaCacheEntry> {
source_generation: OffsetDateTime,
) -> std::result::Result<MetaCacheEntry, DecommissionUnresolvedEntry> {
let candidate_count = entries.as_ref().iter().flatten().count();
if let Some(entry) = entries.resolve(resolver) {
return Ok(entry);
}
let candidate = entries.as_ref().iter().flatten().map(|entry| entry.name.as_str()).next();
Err(decommission_unresolved_listing_error(
bucket,
prefix,
candidate,
let object = entries
.as_ref()
.iter()
.flatten()
.map(|entry| entry.name.as_str())
.next()
.unwrap_or(prefix)
.to_string();
Err(DecommissionUnresolvedEntry {
bucket: bucket.to_string(),
object,
candidate_count,
disk_error_count,
pool_index,
set_index,
))
source_generation,
observed_at: OffsetDateTime::now_utc(),
reason: "metadata_resolution_failed".to_string(),
})
}
fn resolve_decommission_pool_meta_reload_result(result: Result<()>, stage: &str) -> Result<()> {
@@ -1624,6 +1713,8 @@ struct PersistedPoolDecommissionInfo {
pub terminal_reload_attempt_at: Option<OffsetDateTime>,
#[serde(rename = "terminalReloadFailures", default)]
pub terminal_reload_failures: Vec<String>,
#[serde(rename = "unresolvedEntries", default)]
pub unresolved_entries: Vec<DecommissionUnresolvedEntry>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
@@ -1755,6 +1846,7 @@ impl TryFrom<PersistedPoolDecommissionInfo> for PoolDecommissionInfo {
bytes_failed: value.bytes_failed,
terminal_reload_attempt_at: value.terminal_reload_attempt_at,
terminal_reload_failures: value.terminal_reload_failures,
unresolved_entries: value.unresolved_entries,
progress_save_item_baseline: value.items_decommissioned.saturating_add(value.items_decommission_failed),
progress_save_retry_after: None,
})
@@ -1787,6 +1879,7 @@ impl TryFrom<LegacyPoolDecommissionInfo> for PoolDecommissionInfo {
bytes_failed: value.bytes_failed,
terminal_reload_attempt_at: None,
terminal_reload_failures: Vec::new(),
unresolved_entries: Vec::new(),
progress_save_item_baseline: value.items_decommissioned.saturating_add(value.items_decommission_failed),
progress_save_retry_after: None,
})
@@ -1835,6 +1928,7 @@ impl From<&PoolDecommissionInfo> for PersistedPoolDecommissionInfo {
bytes_failed: value.bytes_failed,
terminal_reload_attempt_at: value.terminal_reload_attempt_at,
terminal_reload_failures: value.terminal_reload_failures.clone(),
unresolved_entries: value.unresolved_entries.clone(),
}
}
}
@@ -2026,7 +2120,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,
{
@@ -2435,6 +2529,26 @@ pub fn path2_bucket_object_with_base_path(base_path: &str, path: &str) -> (Strin
path_to_bucket_object_with_base_path(base_path, path)
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct DecommissionUnresolvedEntry {
pub bucket: String,
pub object: String,
#[serde(rename = "poolIndex")]
pub pool_index: usize,
#[serde(rename = "setIndex")]
pub set_index: usize,
#[serde(rename = "sourceGeneration", with = "time::serde::rfc3339")]
pub source_generation: OffsetDateTime,
#[serde(rename = "candidateCount")]
pub candidate_count: usize,
#[serde(rename = "diskErrorCount")]
pub disk_error_count: usize,
#[serde(rename = "observedAt", with = "time::serde::rfc3339")]
pub observed_at: OffsetDateTime,
pub reason: String,
}
#[derive(Debug, Clone, Serialize, Deserialize, Default)]
pub struct PoolDecommissionInfo {
#[serde(rename = "startTime", with = "time::serde::rfc3339::option")]
@@ -2480,6 +2594,8 @@ pub struct PoolDecommissionInfo {
#[serde(rename = "terminalReloadFailures", default)]
pub terminal_reload_failures: Vec<String>,
#[serde(skip)]
pub unresolved_entries: Vec<DecommissionUnresolvedEntry>,
#[serde(skip)]
pub progress_save_item_baseline: usize,
#[serde(skip)]
pub progress_save_retry_after: Option<OffsetDateTime>,
@@ -2513,6 +2629,7 @@ impl PoolDecommissionInfo {
|| self.bytes_failed > 0
|| self.terminal_reload_attempt_at.is_some()
|| !self.terminal_reload_failures.is_empty()
|| !self.unresolved_entries.is_empty()
}
fn counted_items(&self) -> usize {
@@ -2830,6 +2947,21 @@ impl ECStore {
snapshot.save(self.pools.clone()).await
}
async fn persist_decommission_unresolved_entry(
&self,
idx: usize,
generation: OffsetDateTime,
entry: DecommissionUnresolvedEntry,
) -> Result<()> {
{
let mut pool_meta = self.pool_meta.write().await;
record_decommission_unresolved_entry(&mut pool_meta, idx, generation, entry)?;
}
self.save_current_pool_meta()
.await
.map_err(|err| Error::other(format!("decommission unresolved entry ledger save failed: {err}")))
}
async fn save_decommission_progress_checkpoint(&self, idx: usize, generation: OffsetDateTime) -> Result<bool> {
// Lock order: save gate, then the short pool metadata read/write sections. Peer
// reloads are intentionally performed by the caller after both locks are released.
@@ -3678,6 +3810,7 @@ impl ECStore {
let list_bi = bi.clone();
let list_outstanding = outstanding.clone();
let list_entry_error = entry_error.clone();
let list_store = self.clone();
let mut listing = tokio::spawn(async move {
run_decommission_listing_with_retry_and_drain(
list_rx.clone(),
@@ -3691,8 +3824,9 @@ impl ECStore {
let rx = list_rx_for_list.clone();
let bucket = list_bi.clone();
let entry_error = list_entry_error.clone();
let store = list_store.clone();
async move {
set.list_objects_to_decommission(rx, bucket, callback, entry_error, idx, set_idx)
set.list_objects_to_decommission(store, rx, bucket, callback, entry_error, idx, set_idx, generation)
.await
}
},
@@ -4593,7 +4727,7 @@ impl ECStore {
state = "verifying_completion",
"Decommission completion verification started"
);
if let Err(err) = self.check_after_decommission(idx).await {
if let Err(err) = self.check_after_decommission(idx, generation).await {
resolve_decommission_terminal_mark_result(
self.decommission_failed_for_operation(idx, canceler).await,
"failed",
@@ -4620,7 +4754,7 @@ impl ECStore {
"Decommission marking completed state"
);
resolve_decommission_terminal_mark_result(
self.complete_decommission_for_operation(idx, canceler).await,
self.complete_decommission_for_operation(idx, canceler, generation).await,
"completed",
&cmd_line,
)?;
@@ -4756,14 +4890,25 @@ impl ECStore {
#[tracing::instrument(skip(self))]
pub async fn complete_decommission(&self, idx: usize) -> Result<()> {
self.complete_decommission_with_owner(idx, None).await
self.complete_decommission_with_owner(idx, None, None).await
}
async fn complete_decommission_for_operation(&self, idx: usize, owner: &DecommissionCanceler) -> Result<()> {
self.complete_decommission_with_owner(idx, Some(owner)).await
async fn complete_decommission_for_operation(
&self,
idx: usize,
owner: &DecommissionCanceler,
verified_generation: OffsetDateTime,
) -> Result<()> {
self.complete_decommission_with_owner(idx, Some(owner), Some(verified_generation))
.await
}
async fn complete_decommission_with_owner(&self, idx: usize, owner: Option<&DecommissionCanceler>) -> Result<()> {
async fn complete_decommission_with_owner(
&self,
idx: usize,
owner: Option<&DecommissionCanceler>,
verified_generation: Option<OffsetDateTime>,
) -> Result<()> {
ensure_decommission_terminal_operation_supported(self.single_pool(), "complete decommission")?;
let _start_guard = self.start_gate.lock().await;
@@ -4775,11 +4920,13 @@ impl ECStore {
let previous_pool_meta = pool_meta.clone();
let Some(changed) =
update_decommission_for_operation(cancelers.as_slice(), &mut pool_meta, idx, owner, |pool_meta| {
pool_meta.decommission_complete(idx)
reconcile_decommission_unresolved_entries_for_completion(pool_meta, idx, verified_generation)?;
Ok::<bool, Error>(pool_meta.decommission_complete(idx))
})
else {
return Ok(());
};
let changed = changed?;
let terminal_canceler = if let Some(owner) = owner {
Some(owner.clone())
} else {
@@ -5139,7 +5286,7 @@ impl ECStore {
Ok(ret)
}
async fn check_after_decommission(self: &Arc<Self>, idx: usize) -> Result<()> {
async fn check_after_decommission(self: &Arc<Self>, idx: usize, generation: OffsetDateTime) -> Result<()> {
let buckets = self.get_buckets_to_decommission().await?;
let pool = self.pools[idx].clone();
@@ -5238,7 +5385,16 @@ impl ECStore {
});
let list_result = set
.list_objects_to_decommission(callback_rx, bucket_info.clone(), callback, entry_error.clone(), idx, set_index)
.list_objects_to_decommission(
self.clone(),
callback_rx,
bucket_info.clone(),
callback,
entry_error.clone(),
idx,
set_index,
generation,
)
.await;
let entry_error = entry_error.lock().await.clone();
resolve_decommission_check_after_list_result(list_result, entry_error)?;
@@ -5607,6 +5763,17 @@ mod tests {
#[test]
fn pool_meta_persists_decommission_resume_queues() {
let start_time = OffsetDateTime::now_utc();
let unresolved_entry = DecommissionUnresolvedEntry {
bucket: "bucket-b".to_string(),
object: "prefix/unresolved.txt".to_string(),
pool_index: 0,
set_index: 1,
source_generation: start_time,
candidate_count: 2,
disk_error_count: 1,
observed_at: start_time,
reason: "metadata_resolution_failed".to_string(),
};
let pool_meta = PoolMeta {
version: POOL_META_VERSION,
pools: vec![PoolStatus {
@@ -5627,6 +5794,7 @@ mod tests {
bytes_failed: 128,
terminal_reload_attempt_at: Some(start_time),
terminal_reload_failures: vec!["complete_decommission: peer node-a failed".to_string()],
unresolved_entries: vec![unresolved_entry.clone()],
..Default::default()
}),
}],
@@ -5666,6 +5834,7 @@ mod tests {
restored_decommission.terminal_reload_failures,
vec!["complete_decommission: peer node-a failed".to_string()]
);
assert_eq!(restored_decommission.unresolved_entries, vec![unresolved_entry]);
assert!(restored_decommission.queued);
assert_eq!(restored_decommission.items_since_last_progress_save(), 0);
}
@@ -5757,6 +5926,7 @@ mod tests {
assert!(decommission.bucket.is_empty());
assert!(decommission.prefix.is_empty());
assert!(decommission.object.is_empty());
assert!(decommission.unresolved_entries.is_empty());
}
#[test]
@@ -6046,28 +6216,32 @@ async fn record_decommission_entry_error(
entry_error: &Arc<tokio::sync::Mutex<Option<Error>>>,
rx: &CancellationToken,
err: Error,
) {
) -> bool {
if rx.is_cancelled() {
return;
return false;
}
let mut first_err = entry_error.lock().await;
if first_err.is_none() && !rx.is_cancelled() {
*first_err = Some(err);
rx.cancel();
return true;
}
false
}
impl SetDisks {
#[tracing::instrument(skip(self, rx, cb_func, entry_error))]
#[tracing::instrument(skip(self, store, rx, cb_func, entry_error))]
async fn list_objects_to_decommission(
self: &Arc<Self>,
store: Arc<ECStore>,
rx: CancellationToken,
bucket_info: DecomBucketInfo,
cb_func: ListCallback,
entry_error: Arc<tokio::sync::Mutex<Option<Error>>>,
pool_index: usize,
set_index: usize,
source_generation: OffsetDateTime,
) -> Result<()> {
let (disks, _) = self.get_online_disks_with_healing(false).await;
ensure_decommission_listing_disks_available(!disks.is_empty(), &bucket_info.name)?;
@@ -6088,6 +6262,8 @@ impl SetDisks {
let unresolved_prefix = bucket_info.prefix.clone();
let unresolved_pool_index = pool_index;
let unresolved_set_index = set_index;
let unresolved_generation = source_generation;
let unresolved_store = store;
list_path_raw(
rx,
@@ -6107,8 +6283,10 @@ impl SetDisks {
let prefix = unresolved_prefix.clone();
let unresolved_error = unresolved_error.clone();
let unresolved_rx = unresolved_rx.clone();
let unresolved_store = unresolved_store.clone();
let pool_index = unresolved_pool_index;
let set_index = unresolved_set_index;
let source_generation = unresolved_generation;
let disk_error_count = errs.iter().flatten().count();
if unresolved_rx.is_cancelled() {
return Box::pin(async {});
@@ -6122,6 +6300,7 @@ impl SetDisks {
disk_error_count,
pool_index,
set_index,
source_generation,
) {
Ok(entry) => {
warn!("decommission_pool: list_objects_to_decommission get {}", &entry.name);
@@ -6129,10 +6308,11 @@ impl SetDisks {
cb_func(entry).await;
})
}
Err(err) => Box::pin(async move {
Err(unresolved_entry) => Box::pin(async move {
if unresolved_rx.is_cancelled() {
return;
}
let err = decommission_unresolved_listing_error(&unresolved_entry);
warn!(
event = EVENT_DECOMMISSION_BUCKET,
component = LOG_COMPONENT_ECSTORE,
@@ -6143,6 +6323,13 @@ impl SetDisks {
error = %err,
"Decommission listing failed closed on unresolved metadata"
);
let err = match unresolved_store
.persist_decommission_unresolved_entry(pool_index, source_generation, unresolved_entry)
.await
{
Ok(()) => err,
Err(ledger_err) => Error::other(format!("{err}; {ledger_err}")),
};
record_decommission_entry_error(&unresolved_error, &unresolved_rx, err).await;
}),
}
@@ -6351,47 +6538,52 @@ mod pools_tests {
use super::{
DECOMMISSION_ENTRY_CONCURRENCY_DEFAULT_CAP, DECOMMISSION_ENTRY_CONCURRENCY_HARD_CAP, DECOMMISSION_ENTRY_QUEUE_HARD_CAP,
DECOMMISSION_PROGRESS_SAVE_INTERVAL, DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD, DecomBucketInfo, DecommissionCanceler,
DecommissionEntryEnqueueResult, DecommissionStartPoolState, DecommissionTerminalState, ListCallback,
PoolDecommissionInfo, PoolMeta, PoolSpaceInfo, PoolStatus, QueuedDecommissionEntry, apply_decommission_status_space_info,
await_decommission_worker, bind_decommission_cancelers, bind_missing_decommission_cancelers,
cancel_decommission_canceler, clamp_decommission_entry_concurrency, classify_decommission_terminal_state,
count_decommission_item, decommission_cancel_signal_result, decommission_entry_queue_capacity, decommission_item_size,
decommission_meta_bucket_options, decommission_start_pool_state, dedup_indices, default_decommission_bucket_concurrency,
default_decommission_entry_concurrency, drain_decommission_entry_queue, enqueue_decommission_entry,
ensure_decommission_cancel_allowed, ensure_decommission_clear_allowed, ensure_decommission_generation,
ensure_decommission_listing_disks_available, ensure_decommission_not_rebalancing, ensure_decommission_start_allowed,
ensure_decommission_start_keeps_active_pool, ensure_decommission_start_local_leader,
DecommissionEntryEnqueueResult, DecommissionStartPoolState, DecommissionTerminalState, DecommissionUnresolvedEntry,
ListCallback, POOL_META_VERSION, PoolDecommissionInfo, PoolMeta, PoolSpaceInfo, PoolStatus, QueuedDecommissionEntry,
apply_decommission_status_space_info, await_decommission_worker, bind_decommission_cancelers,
bind_missing_decommission_cancelers, cancel_decommission_canceler, clamp_decommission_entry_concurrency,
classify_decommission_terminal_state, count_decommission_item, decommission_cancel_signal_result,
decommission_entry_queue_capacity, decommission_item_size, decommission_meta_bucket_options,
decommission_start_pool_state, decommission_unresolved_listing_error, dedup_indices,
default_decommission_bucket_concurrency, default_decommission_entry_concurrency, drain_decommission_entry_queue,
enqueue_decommission_entry, ensure_decommission_cancel_allowed, ensure_decommission_clear_allowed,
ensure_decommission_generation, ensure_decommission_listing_disks_available, ensure_decommission_not_rebalancing,
ensure_decommission_start_allowed, ensure_decommission_start_keeps_active_pool, ensure_decommission_start_local_leader,
ensure_decommission_start_pool_states, ensure_decommission_start_rebalance_meta_allowed,
ensure_decommission_start_target_capacity, ensure_decommission_terminal_operation_supported,
ensure_local_decommission_pool_leaders, ensure_valid_decommission_pool_index, first_resumable_decommission_queue_indices,
get_by_index, guard_decommission_cancelers, has_active_decommission_canceler, is_decommission_active,
is_decommission_cancel_requested, load_decommission_entry_versions, local_decommission_queue_prefix,
mark_decommission_bucket_done, merge_pool_status_refresh, missing_decommission_worker_prefix,
observe_decommission_terminal_reload_result, pool_meta_has_active_decommission, require_decommission_store,
reserve_decommission_start_cancelers, resolve_decommission_bucket_done_save_result, resolve_decommission_bucket_state,
resolve_decommission_check_after_list_result, resolve_decommission_entry_cleanup_delete_result,
resolve_decommission_entry_exact_versions, resolve_decommission_entry_reload_result,
resolve_decommission_listing_worker_result, resolve_decommission_optional_bucket_config_result,
resolve_decommission_partial_listing_entry, resolve_decommission_pool_meta_reload_result,
resolve_decommission_preflight_heal_result, resolve_decommission_progress_save_result,
resolve_decommission_terminal_mark_after_error_result, resolve_decommission_terminal_mark_result,
resolve_decommission_update_after_result, resolve_start_decommission_pool_meta_reload_result,
rollback_start_decommission_pool_meta, run_decommission_buckets_bounded, run_decommission_listing_with_retry,
run_decommission_listing_with_retry_and_drain, run_decommission_side_effect, should_cleanup_decommission_source_entry,
should_continue_decommission_queue, should_count_decommission_version_complete,
should_preserve_decommission_canceled_state, should_reject_decommission_cancel_as_terminal,
should_retry_decommission_cancel_reload, should_retry_decommission_listing, should_skip_canceled_decommission_routine,
spawn_decommission_index_cancelers, split_decommission_buckets, take_and_cancel_decommission_canceler,
take_decommission_canceler, track_decommission_current_object, track_decommission_current_object_stage,
update_decommission_for_operation, validate_start_decommission_request, wait_decommission_listing_retry,
wait_decommission_worker_drain, with_decommission_entry_context,
observe_decommission_terminal_reload_result, pool_meta_has_active_decommission,
reconcile_decommission_unresolved_entries_for_completion, record_decommission_unresolved_entry,
require_decommission_store, reserve_decommission_start_cancelers, resolve_decommission_bucket_done_save_result,
resolve_decommission_bucket_state, resolve_decommission_check_after_list_result,
resolve_decommission_entry_cleanup_delete_result, resolve_decommission_entry_exact_versions,
resolve_decommission_entry_reload_result, resolve_decommission_listing_worker_result,
resolve_decommission_optional_bucket_config_result, resolve_decommission_partial_listing_entry,
resolve_decommission_pool_meta_reload_result, resolve_decommission_preflight_heal_result,
resolve_decommission_progress_save_result, resolve_decommission_terminal_mark_after_error_result,
resolve_decommission_terminal_mark_result, resolve_decommission_update_after_result,
resolve_start_decommission_pool_meta_reload_result, rollback_start_decommission_pool_meta,
run_decommission_buckets_bounded, run_decommission_listing_with_retry, run_decommission_listing_with_retry_and_drain,
run_decommission_side_effect, should_cleanup_decommission_source_entry, should_continue_decommission_queue,
should_count_decommission_version_complete, should_preserve_decommission_canceled_state,
should_reject_decommission_cancel_as_terminal, should_retry_decommission_cancel_reload,
should_retry_decommission_listing, should_skip_canceled_decommission_routine, spawn_decommission_index_cancelers,
split_decommission_buckets, take_and_cancel_decommission_canceler, take_decommission_canceler,
track_decommission_current_object, track_decommission_current_object_stage, update_decommission_for_operation,
validate_start_decommission_request, wait_decommission_listing_retry, wait_decommission_worker_drain,
with_decommission_entry_context,
};
use crate::bucket::metadata_sys;
use crate::data_movement;
use crate::disk::endpoint::Endpoint;
use crate::disk::{STORAGE_FORMAT_FILE, endpoint::Endpoint};
use crate::error::{Error, StorageError};
use crate::layout::endpoints::{EndpointServerPools, Endpoints, PoolEndpoints};
use crate::runtime::instance::InstanceContext;
use crate::services::rebalance::{RebalStatus, RebalanceInfo, RebalanceMeta, RebalanceStats};
use crate::storage_api_contracts::bucket::MakeBucketOptions;
use crate::store::ECStore;
use rustfs_filemeta::{FileInfo, FileInfoVersions, MetaCacheEntry, ObjectPartInfo};
use rustfs_filemeta::{MetaCacheEntries, MetadataResolutionParams};
@@ -7622,7 +7814,8 @@ mod pools_tests {
#[test]
fn test_resolve_decommission_partial_listing_entry_rejects_unresolved_metadata() {
let err = resolve_decommission_partial_listing_entry(
let generation = OffsetDateTime::now_utc();
let unresolved_entry = resolve_decommission_partial_listing_entry(
MetaCacheEntries(vec![None]),
MetadataResolutionParams {
dir_quorum: 2,
@@ -7635,23 +7828,160 @@ mod pools_tests {
1,
2,
3,
generation,
)
.expect_err("unresolved partial listing must fail closed");
let message = err.to_string();
assert_eq!(unresolved_entry.bucket, "bucket-a");
assert_eq!(unresolved_entry.object, "prefix/");
assert_eq!(unresolved_entry.pool_index, 2);
assert_eq!(unresolved_entry.set_index, 3);
assert_eq!(unresolved_entry.source_generation, generation);
assert_eq!(unresolved_entry.candidate_count, 0);
assert_eq!(unresolved_entry.disk_error_count, 1);
assert_eq!(unresolved_entry.reason, "metadata_resolution_failed");
let message = decommission_unresolved_listing_error(&unresolved_entry).to_string();
assert!(message.contains("decommission listing could not resolve metadata"));
assert!(message.contains("bucket-a/prefix/"));
assert!(message.contains("pool 2 set 3"));
assert!(message.contains("1 disk error(s)"));
}
#[test]
fn unresolved_entry_ledger_blocks_unverified_completion_and_reconciles_after_verified_sweep() {
let generation = OffsetDateTime::now_utc();
let mut pool_meta = PoolMeta {
version: POOL_META_VERSION,
pools: vec![PoolStatus {
id: 0,
cmd_line: "/data/pool".to_string(),
last_update: generation,
decommission: Some(PoolDecommissionInfo {
start_time: Some(generation),
..Default::default()
}),
}],
dont_save: true,
};
let unresolved_entry = DecommissionUnresolvedEntry {
bucket: "bucket-a".to_string(),
object: "object-a".to_string(),
pool_index: 0,
set_index: 0,
source_generation: generation,
candidate_count: 1,
disk_error_count: 0,
observed_at: generation,
reason: "metadata_resolution_failed".to_string(),
};
assert!(
record_decommission_unresolved_entry(&mut pool_meta, 0, generation, unresolved_entry)
.expect("active generation should accept unresolved entry")
);
let err = reconcile_decommission_unresolved_entries_for_completion(&mut pool_meta, 0, None)
.expect_err("unverified completion must retain unresolved entries");
assert!(err.to_string().contains("1 unresolved listing entries remain"));
assert_eq!(
pool_meta.pools[0]
.decommission
.as_ref()
.expect("decommission should exist")
.unresolved_entries
.len(),
1
);
reconcile_decommission_unresolved_entries_for_completion(&mut pool_meta, 0, Some(generation))
.expect("a successful same-generation final sweep should reconcile stale entries");
assert!(
pool_meta.pools[0]
.decommission
.as_ref()
.expect("decommission should exist")
.unresolved_entries
.is_empty()
);
}
#[tokio::test]
async fn final_sweep_persists_unresolved_entry_from_real_listing() {
let (dirs, store) = metadata_sys::test_support::isolated_store_over_temp_disks().await;
let bucket = "decommission-final-sweep-unresolved";
let object = "corrupt-object";
store
.peer_sys
.make_bucket(bucket, &MakeBucketOptions::default())
.await
.expect("test bucket should be created");
metadata_sys::init_bucket_metadata_sys(store.clone(), vec![bucket.to_string()]).await;
let generation = OffsetDateTime::now_utc();
{
let mut pool_meta = store.pool_meta.write().await;
pool_meta.dont_save = false;
pool_meta.pools[0].decommission = Some(PoolDecommissionInfo {
start_time: Some(generation),
..Default::default()
});
}
for (disk_index, dir) in dirs.iter().enumerate() {
let object_dir = dir.path().join(bucket).join(object);
tokio::fs::create_dir_all(&object_dir)
.await
.expect("corrupt object directory should be created");
tokio::fs::write(object_dir.join(STORAGE_FORMAT_FILE), format!("corrupt-xl-meta-{disk_index}").into_bytes())
.await
.expect("divergent corrupt metadata should be written");
}
let err = store
.check_after_decommission(0, generation)
.await
.expect_err("real final sweep must fail closed on unresolved listing metadata");
assert!(err.to_string().contains("decommission listing could not resolve metadata"));
let in_memory_entries = store.pool_meta.read().await.pools[0]
.decommission
.as_ref()
.expect("decommission should remain active")
.unresolved_entries
.clone();
assert_eq!(in_memory_entries.len(), 1);
let unresolved_entry = &in_memory_entries[0];
assert_eq!(unresolved_entry.bucket, bucket);
assert_eq!(unresolved_entry.object, object);
assert_eq!(unresolved_entry.pool_index, 0);
assert_eq!(unresolved_entry.set_index, 0);
assert_eq!(unresolved_entry.source_generation, generation);
assert_eq!(unresolved_entry.candidate_count, 4);
assert_eq!(unresolved_entry.disk_error_count, 0);
assert_eq!(unresolved_entry.reason, "metadata_resolution_failed");
let mut restored = PoolMeta::default();
restored
.load_for_startup(store.pools[0].clone())
.await
.expect("persisted pool metadata should reload after the final-sweep failure");
assert_eq!(
restored.pools[0]
.decommission
.as_ref()
.expect("reloaded decommission should exist")
.unresolved_entries,
in_memory_entries
);
}
#[tokio::test]
async fn test_record_decommission_entry_error_cancels_listing_and_preserves_first_error() {
let entry_error = Arc::new(tokio::sync::Mutex::new(None));
let rx = CancellationToken::new();
record_decommission_entry_error(&entry_error, &rx, Error::SlowDown).await;
record_decommission_entry_error(&entry_error, &rx, Error::OperationCanceled).await;
assert!(record_decommission_entry_error(&entry_error, &rx, Error::SlowDown).await);
assert!(!record_decommission_entry_error(&entry_error, &rx, Error::OperationCanceled).await);
assert!(rx.is_cancelled());
assert!(matches!(*entry_error.lock().await, Some(Error::SlowDown)));
@@ -7663,7 +7993,7 @@ mod pools_tests {
let rx = CancellationToken::new();
rx.cancel();
record_decommission_entry_error(&entry_error, &rx, Error::SlowDown).await;
assert!(!record_decommission_entry_error(&entry_error, &rx, Error::SlowDown).await);
assert!(entry_error.lock().await.is_none());
}
+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};
}
+24 -2
View File
@@ -16,7 +16,8 @@
use super::storage_api::admin_usecase::admin::get_server_info;
use super::storage_api::admin_usecase::capacity::{
PoolDecommissionInfo, PoolStatus, RebalStatus, get_total_usable_capacity, get_total_usable_capacity_free,
DecommissionUnresolvedEntry, PoolDecommissionInfo, PoolStatus, RebalStatus, get_total_usable_capacity,
get_total_usable_capacity_free,
};
use super::storage_api::admin_usecase::contract::StorageAdminApi;
use super::storage_api::admin_usecase::contract::bucket::{BucketOperations as _, BucketOptions};
@@ -107,6 +108,8 @@ pub struct AdminPoolDecommissionInfo {
pub bytes_failed: usize,
#[serde(rename = "waitingReason")]
pub waiting_reason: Option<String>,
#[serde(rename = "unresolvedEntries", skip_serializing_if = "Vec::is_empty")]
pub unresolved_entries: Vec<DecommissionUnresolvedEntry>,
}
#[derive(Debug, Clone, serde::Serialize)]
@@ -619,6 +622,7 @@ impl DefaultAdminUsecase {
bytes_done: info.bytes_done,
bytes_failed: info.bytes_failed,
waiting_reason,
unresolved_entries: info.unresolved_entries,
}
}
@@ -676,7 +680,7 @@ impl DefaultAdminUsecase {
#[cfg(test)]
mod tests {
use super::super::storage_api::admin_usecase::capacity::{PoolDecommissionInfo, PoolStatus};
use super::super::storage_api::admin_usecase::capacity::{DecommissionUnresolvedEntry, PoolDecommissionInfo, PoolStatus};
use super::*;
use time::OffsetDateTime;
use tracing_subscriber::{Layer, Registry, layer::Context, prelude::*};
@@ -987,6 +991,17 @@ mod tests {
items_decommission_failed: 1,
bytes_done: 1024,
bytes_failed: 64,
unresolved_entries: vec![DecommissionUnresolvedEntry {
bucket: "bucket-a".to_string(),
object: "prefix/unresolved.txt".to_string(),
pool_index: 3,
set_index: 1,
source_generation: OffsetDateTime::UNIX_EPOCH,
candidate_count: 2,
disk_error_count: 1,
observed_at: OffsetDateTime::UNIX_EPOCH,
reason: "metadata_resolution_failed".to_string(),
}],
..Default::default()
}),
},
@@ -1010,6 +1025,13 @@ mod tests {
assert_eq!(value["decommissionInfo"]["objectsDecommissionedFailed"], 1);
assert_eq!(value["decommissionInfo"]["bytesDecommissioned"], 1024);
assert_eq!(value["decommissionInfo"]["bytesDecommissionedFailed"], 64);
assert_eq!(value["decommissionInfo"]["unresolvedEntries"][0]["bucket"], "bucket-a");
assert_eq!(value["decommissionInfo"]["unresolvedEntries"][0]["object"], "prefix/unresolved.txt");
assert_eq!(
value["decommissionInfo"]["unresolvedEntries"][0]["sourceGeneration"],
"1970-01-01T00:00:00Z"
);
assert_eq!(value["decommissionInfo"]["unresolvedEntries"][0]["reason"], "metadata_resolution_failed");
assert_eq!(value["decommissionInfo"]["waitingReason"], "queued");
}
+1
View File
@@ -31,6 +31,7 @@ pub(crate) mod admin {
}
pub(crate) mod capacity {
pub(crate) type DecommissionUnresolvedEntry = crate::storage::storage_api::ecstore_capacity::DecommissionUnresolvedEntry;
pub(crate) type PoolDecommissionInfo = crate::storage::storage_api::ecstore_capacity::PoolDecommissionInfo;
pub(crate) type PoolStatus = crate::storage::storage_api::ecstore_capacity::PoolStatus;
pub(crate) type RebalStatus = crate::storage::storage_api::ecstore_rebalance::RebalStatus;
+1 -1
View File
@@ -398,7 +398,7 @@ pub(crate) mod ecstore_bucket {
pub(crate) mod ecstore_capacity {
pub(crate) use rustfs_ecstore::api::capacity::{
PoolDecommissionInfo, PoolStatus, get_total_usable_capacity, get_total_usable_capacity_free,
DecommissionUnresolvedEntry, PoolDecommissionInfo, PoolStatus, get_total_usable_capacity, get_total_usable_capacity_free,
is_reserved_or_invalid_bucket,
};
}