Compare commits

..

3 Commits

8 changed files with 495 additions and 383 deletions
+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;
}
+396 -112
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(),
}
}
}
@@ -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,13 +5286,10 @@ 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();
self.ensure_decommission_multipart_uploads_drained(idx, &pool, &buckets)
.await?;
for (set_index, set) in pool.disk_set.iter().enumerate() {
for bucket_info in &buckets {
let mut lifecycle_config = None;
@@ -5241,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)?;
@@ -5259,49 +5412,6 @@ impl ECStore {
Ok(())
}
async fn ensure_decommission_multipart_uploads_drained(
&self,
idx: usize,
pool: &Sets,
buckets: &[DecomBucketInfo],
) -> Result<()> {
let mut bucket_names = buckets
.iter()
.filter(|bucket| bucket.name != RUSTFS_META_BUCKET)
.map(|bucket| bucket.name.as_str())
.collect::<Vec<_>>();
bucket_names.sort_unstable();
bucket_names.dedup();
// Take one bucket fence at a time so cross-bucket COPY cannot form an
// ABBA cycle. Suspension prevents new source uploads after each fence.
for bucket in bucket_names {
let lifecycle_guard = self.acquire_bucket_lifecycle_write_lock(bucket).await?;
if lifecycle_guard.is_lock_lost() {
return Err(Error::other(format!(
"decommission multipart drain lost the bucket lifecycle fence for `{bucket}`"
)));
}
for set in &pool.disk_set {
if let Some(upload_path) = set.first_multipart_upload_path_for_decommission(bucket).await? {
return Err(Error::other(format!(
"pool {idx} still contains multipart upload `{upload_path}` for bucket `{bucket}`; resolve it before retrying decommission"
)));
}
}
}
Ok(())
}
#[cfg(test)]
pub(crate) async fn ensure_decommission_multipart_uploads_drained_for_test(self: &Arc<Self>, idx: usize) -> Result<()> {
let buckets = self.get_buckets_to_decommission().await?;
let pool = self.pools[idx].clone();
self.ensure_decommission_multipart_uploads_drained(idx, pool.as_ref(), &buckets)
.await
}
#[tracing::instrument(skip(self, rd))]
async fn decommission_object(
self: Arc<Self>,
@@ -5653,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 {
@@ -5673,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()
}),
}],
@@ -5712,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);
}
@@ -5803,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]
@@ -6092,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)?;
@@ -6134,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,
@@ -6153,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 {});
@@ -6168,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);
@@ -6175,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,
@@ -6189,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;
}),
}
@@ -6397,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};
@@ -7668,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,
@@ -7681,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)));
@@ -7709,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());
}
+47 -64
View File
@@ -556,69 +556,6 @@ async fn multipart_upload_paths_on_disk(disk: DiskStore, bucket: &str) -> disk::
}
impl SetDisks {
async fn discover_multipart_upload_paths(
&self,
orig_bucket: &str,
error_path: &str,
) -> Result<(Vec<Option<DiskStore>>, Vec<String>, usize)> {
let disks = self.disks.read().await.clone();
if disks.is_empty() {
return Err(Error::ErasureReadQuorum);
}
let discovery_quorum = if self.default_parity_count == 0 {
disks.len()
} else {
(disks.len() / 2).max(1)
};
let mut discovery_errors = (0..disks.len()).map(|_| Some(DiskError::DiskNotFound)).collect::<Vec<_>>();
let mut candidate_counts = HashMap::<String, usize>::new();
let mut discovery_tasks = JoinSet::new();
for (index, disk) in disks.iter().enumerate() {
let disk = disk.clone();
let orig_bucket = orig_bucket.to_string();
discovery_tasks.spawn(async move {
let result = match disk {
Some(disk) => multipart_upload_paths_on_disk(disk, &orig_bucket).await,
None => Err(DiskError::DiskNotFound),
};
(index, result)
});
}
while let Some(task_result) = discovery_tasks.join_next().await {
let Ok((index, result)) = task_result else {
continue;
};
match result {
Ok(paths) => {
discovery_errors[index] = None;
for path in paths {
*candidate_counts.entry(path).or_insert(0) += 1;
}
}
Err(err) => discovery_errors[index] = Some(err),
}
}
if let Some(err) = reduce_read_quorum_errs(&discovery_errors, OBJECT_OP_IGNORED_ERRS, discovery_quorum) {
return Err(to_object_err(err.into(), vec![orig_bucket, error_path]));
}
let mut candidate_paths = candidate_counts
.into_iter()
.filter_map(|(path, count)| (count >= discovery_quorum).then_some(path))
.collect::<Vec<_>>();
candidate_paths.sort_unstable();
Ok((disks, candidate_paths, discovery_quorum))
}
pub(crate) async fn first_multipart_upload_path_for_decommission(&self, bucket: &str) -> Result<Option<String>> {
let (_, paths, _) = self
.discover_multipart_upload_paths(bucket, RUSTFS_META_MULTIPART_BUCKET)
.await?;
Ok(paths.into_iter().next())
}
async fn acquire_multipart_upload_read_lock(
&self,
op: &'static str,
@@ -810,7 +747,53 @@ impl SetDisks {
max_uploads: usize,
expected_incarnation_id: Option<Uuid>,
) -> Result<ListMultipartsInfo> {
let (disks, candidate_paths, discovery_quorum) = self.discover_multipart_upload_paths(bucket, prefix).await?;
let disks = self.disks.read().await.clone();
if disks.is_empty() {
return Err(Error::ErasureReadQuorum);
}
let discovery_quorum = if self.default_parity_count == 0 {
disks.len()
} else {
(disks.len() / 2).max(1)
};
let mut discovery_errors = (0..disks.len()).map(|_| Some(DiskError::DiskNotFound)).collect::<Vec<_>>();
let mut candidate_counts = HashMap::<String, usize>::new();
let mut discovery_tasks = JoinSet::new();
for (index, disk) in disks.iter().enumerate() {
let disk = disk.clone();
let bucket = bucket.to_string();
discovery_tasks.spawn(async move {
let result = match disk {
Some(disk) => multipart_upload_paths_on_disk(disk, &bucket).await,
None => Err(DiskError::DiskNotFound),
};
(index, result)
});
}
while let Some(task_result) = discovery_tasks.join_next().await {
let Ok((index, result)) = task_result else {
continue;
};
match result {
Ok(paths) => {
discovery_errors[index] = None;
for path in paths {
*candidate_counts.entry(path).or_insert(0) += 1;
}
}
Err(err) => discovery_errors[index] = Some(err),
}
}
if let Some(err) = reduce_read_quorum_errs(&discovery_errors, OBJECT_OP_IGNORED_ERRS, discovery_quorum) {
return Err(to_object_err(err.into(), vec![bucket, prefix]));
}
let candidate_paths = candidate_counts
.into_iter()
.filter_map(|(path, count)| (count >= discovery_quorum).then_some(path))
.collect::<Vec<_>>();
let listed_uploads = stream::iter(candidate_paths)
.map(|upload_path| {
let disks = &disks;
-171
View File
@@ -1589,177 +1589,6 @@ mod tests {
shutdown.cancel();
}
#[tokio::test]
#[serial_test::serial(storage_class_env)]
async fn suspended_decommission_source_multipart_remains_operable_until_drained() {
let temp_dir = tempfile::tempdir().expect("create decommission multipart drain store dir");
let (_ctx, store, shutdown) =
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "decommission-multipart-drain", &[4, 4])).await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
let bucket = format!("decommission-multipart-drain-{}", uuid::Uuid::new_v4());
let complete_object = "complete.bin";
let abort_object = "abort.bin";
store
.make_bucket(&bucket, &MakeBucketOptions::default())
.await
.expect("create decommission multipart drain bucket");
let incarnation = store.bucket_incarnation_id(&bucket).await.expect("read bucket incarnation");
let lifecycle_guard = store
.acquire_bucket_lifecycle_read_lock(&bucket)
.await
.expect("acquire multipart creation lifecycle fence");
let mut upload_opts = ObjectOptions {
expected_bucket_incarnation_id: Some(incarnation),
..Default::default()
};
upload_opts.add_bucket_lifecycle_lock_guard(&lifecycle_guard);
let complete_upload = store.pools[0]
.new_multipart_upload(&bucket, complete_object, &upload_opts)
.await
.expect("create source upload to complete");
let abort_upload = store.pools[0]
.new_multipart_upload(&bucket, abort_object, &upload_opts)
.await
.expect("create source upload to abort");
drop(lifecycle_guard);
mark_test_pool_decommissioning(&store, 0).await;
let err = store
.ensure_decommission_multipart_uploads_drained_for_test(0)
.await
.expect_err("an unresolved source multipart upload must block final decommission");
let drain_error = err.to_string();
assert!(
drain_error.contains("still contains multipart upload") && drain_error.contains(&bucket),
"the drain error must identify both the upload path and user bucket: {drain_error}"
);
let listed = store
.list_multipart_uploads(&bucket, "", None, None, None, 100)
.await
.expect("list uploads from suspended decommission source");
assert!(
listed
.uploads
.iter()
.any(|upload| upload.upload_id.as_str() == complete_upload.upload_id.as_str()),
"the upload selected before suspension must remain visible"
);
store
.get_multipart_info(&bucket, complete_object, &complete_upload.upload_id, &ObjectOptions::default())
.await
.expect("read upload metadata from suspended decommission source");
let mut part_reader = PutObjReader::from_vec(b"multipart body".to_vec());
let part = store
.put_object_part(
&bucket,
complete_object,
&complete_upload.upload_id,
1,
&mut part_reader,
&ObjectOptions::default(),
)
.await
.expect("write part to suspended decommission source");
let parts = store
.list_object_parts(&bucket, complete_object, &complete_upload.upload_id, None, 100, &ObjectOptions::default())
.await
.expect("list parts from suspended decommission source");
assert_eq!(parts.parts.len(), 1);
assert_eq!(parts.parts[0].etag.as_deref(), part.etag.as_deref());
store
.clone()
.complete_multipart_upload(
&bucket,
complete_object,
&complete_upload.upload_id,
vec![crate::storage_api_contracts::multipart::CompletePart {
part_num: part.part_num,
etag: part.etag,
..Default::default()
}],
&ObjectOptions::default(),
)
.await
.expect("complete upload on suspended decommission source");
store
.abort_multipart_upload(&bucket, abort_object, &abort_upload.upload_id, &ObjectOptions::default())
.await
.expect("abort upload on suspended decommission source");
store
.ensure_decommission_multipart_uploads_drained_for_test(0)
.await
.expect("final decommission gate should open after all source uploads are resolved");
assert_pool_object_present(&store.pools[0], &bucket, complete_object).await;
shutdown.cancel();
}
#[tokio::test]
#[serial_test::serial(storage_class_env)]
async fn active_multipart_upload_routes_before_faulted_suspended_source() {
let temp_dir = tempfile::tempdir().expect("create active-first multipart routing store dir");
let (_ctx, store, shutdown) =
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "active-first-multipart-routing", &[4, 4]))
.await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
let bucket = format!("active-first-multipart-routing-{}", uuid::Uuid::new_v4());
let object = "target-upload.bin";
store
.make_bucket(&bucket, &MakeBucketOptions::default())
.await
.expect("create active-first multipart routing bucket");
let incarnation = store.bucket_incarnation_id(&bucket).await.expect("read bucket incarnation");
let lifecycle_guard = store
.acquire_bucket_lifecycle_read_lock(&bucket)
.await
.expect("acquire multipart creation lifecycle fence");
let mut upload_opts = ObjectOptions {
expected_bucket_incarnation_id: Some(incarnation),
..Default::default()
};
upload_opts.add_bucket_lifecycle_lock_guard(&lifecycle_guard);
let upload = store.pools[1]
.new_multipart_upload(&bucket, object, &upload_opts)
.await
.expect("create upload in active target pool");
drop(lifecycle_guard);
mark_test_pool_decommissioning(&store, 0).await;
let source_set = store.pools[0].get_disks_by_key(object);
let original_source_disks = {
let mut disks = source_set.disks.write().await;
let original = disks.clone();
disks.fill(None);
original
};
let source_result = store.pools[0]
.get_multipart_info(&bucket, object, &upload.upload_id, &ObjectOptions::default())
.await;
let routed_result = store
.get_multipart_info(&bucket, object, &upload.upload_id, &ObjectOptions::default())
.await;
*source_set.disks.write().await = original_source_disks;
assert!(
matches!(&source_result, Err(StorageError::ErasureReadQuorum)),
"the suspended source must expose the injected hard read failure: {source_result:?}"
);
let routed = routed_result.expect("the active target UploadID must be resolved before the faulted suspended source");
assert_eq!(routed.upload_id, upload.upload_id);
shutdown.cancel();
}
#[tokio::test]
#[serial_test::serial(storage_class_env)]
async fn delete_objects_skips_active_rebalance_source_pool() {
+24 -31
View File
@@ -196,25 +196,6 @@ async fn list_pool_multipart_uploads_for_incarnation(
}
impl ECStore {
async fn existing_multipart_pool_order(&self) -> Vec<usize> {
// A draining source must not hide a valid UploadID in an active target,
// while physical order within each phase preserves fail-closed errors.
let mut active = Vec::with_capacity(self.pools.len());
let mut draining = Vec::new();
for (idx, pool) in self.pools.iter().enumerate() {
if self.is_pool_rebalancing(pool.pool_idx).await {
continue;
}
if self.is_suspended(pool.pool_idx).await {
draining.push(idx);
} else {
active.push(idx);
}
}
active.extend(draining);
active
}
#[allow(clippy::too_many_arguments)]
pub async fn list_multipart_uploads_for_bucket_incarnation(
&self,
@@ -309,8 +290,10 @@ impl ECStore {
.await;
}
for pool_idx in self.existing_multipart_pool_order().await {
let pool = &self.pools[pool_idx];
for pool in self.pools.iter() {
if self.is_suspended(pool.pool_idx).await || self.is_pool_rebalancing(pool.pool_idx).await {
continue;
}
return match pool
.list_object_parts(bucket, object, upload_id, part_number_marker, max_parts, opts)
.await
@@ -370,8 +353,10 @@ impl ECStore {
let mut common_prefixes = HashSet::new();
let mut source_truncated = false;
for pool_idx in self.existing_multipart_pool_order().await {
let pool = &self.pools[pool_idx];
for pool in self.pools.iter() {
if self.is_suspended(pool.pool_idx).await || self.is_pool_rebalancing(pool.pool_idx).await {
continue;
}
let res = list_pool_multipart_uploads_for_incarnation(
pool,
bucket,
@@ -538,8 +523,10 @@ impl ECStore {
.await;
}
for pool_idx in self.existing_multipart_pool_order().await {
let pool = &self.pools[pool_idx];
for pool in self.pools.iter() {
if self.is_suspended(pool.pool_idx).await || self.is_pool_rebalancing(pool.pool_idx).await {
continue;
}
let err = match pool.put_object_part(bucket, object, upload_id, part_id, data, opts).await {
Ok(res) => return Ok(res),
Err(err) => {
@@ -599,8 +586,10 @@ impl ECStore {
return self.pools[0].get_multipart_info(bucket, object, upload_id, opts).await;
}
for pool_idx in self.existing_multipart_pool_order().await {
let pool = &self.pools[pool_idx];
for pool in self.pools.iter() {
if self.is_suspended(pool.pool_idx).await || self.is_pool_rebalancing(pool.pool_idx).await {
continue;
}
return match pool.get_multipart_info(bucket, object, upload_id, opts).await {
Ok(res) => Ok(res),
@@ -635,8 +624,10 @@ impl ECStore {
return self.pools[0].abort_multipart_upload(bucket, object, upload_id, opts).await;
}
for pool_idx in self.existing_multipart_pool_order().await {
let pool = &self.pools[pool_idx];
for pool in self.pools.iter() {
if self.is_suspended(pool.pool_idx).await || self.is_pool_rebalancing(pool.pool_idx).await {
continue;
}
let err = match pool.abort_multipart_upload(bucket, object, upload_id, opts).await {
Ok(_) => return Ok(()),
@@ -694,8 +685,10 @@ impl ECStore {
.await;
}
for pool_idx in self.existing_multipart_pool_order().await {
let pool = &self.pools[pool_idx];
for pool in self.pools.iter() {
if self.is_suspended(pool.pool_idx).await || self.is_pool_rebalancing(pool.pool_idx).await {
continue;
}
let pool = pool.clone();
let err = match pool
+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,
};
}