mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-23 04:39:04 +00:00
Compare commits
4 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 82fb0a8843 | |||
| 307510749e | |||
| 5496e14960 | |||
| ec3b7a7dc6 |
@@ -236,6 +236,15 @@ pub struct DataUsageInfo {
|
||||
/// without relying on synchronized clocks.
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub usage_snapshot_authoritative_baseline: Option<DataUsageSnapshotIdentity>,
|
||||
/// Per-set freshness for an observational aggregate. A set entry is
|
||||
/// never sufficient to make the aggregate authoritative; it only records
|
||||
/// which last-known-good generation contributed to the view.
|
||||
#[serde(default, skip_serializing_if = "Vec::is_empty")]
|
||||
pub usage_snapshot_set_states: Vec<DataUsageSnapshotSetState>,
|
||||
/// An observational view may contain only the sets that completed this
|
||||
/// cycle (or retained a compatible last-known-good cache).
|
||||
#[serde(default)]
|
||||
pub usage_snapshot_partial: bool,
|
||||
/// Deprecated kept here for backward compatibility reasons
|
||||
pub bucket_sizes: HashMap<String, u64>,
|
||||
/// Per-disk snapshot information when available
|
||||
@@ -252,6 +261,22 @@ pub struct DataUsageSnapshotIdentity {
|
||||
pub scanner_epoch: Option<u64>,
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug, Default, Serialize, Deserialize, PartialEq, Eq)]
|
||||
pub struct DataUsageSnapshotSetState {
|
||||
pub pool_index: u64,
|
||||
pub set_index: u64,
|
||||
#[serde(default)]
|
||||
pub scanner_cycle: Option<u64>,
|
||||
#[serde(default)]
|
||||
pub scanner_epoch: Option<u64>,
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub scan_plan_digest: Option<[u8; 32]>,
|
||||
#[serde(default)]
|
||||
pub complete: bool,
|
||||
#[serde(default)]
|
||||
pub tombstone: bool,
|
||||
}
|
||||
|
||||
impl DataUsageInfo {
|
||||
pub fn snapshot_identity(&self) -> DataUsageSnapshotIdentity {
|
||||
DataUsageSnapshotIdentity {
|
||||
@@ -291,7 +316,7 @@ pub fn data_usage_snapshot_is_newer(candidate: &DataUsageInfo, baseline: &DataUs
|
||||
/// rollback delete/recreate fences the previous bucket incarnation too.
|
||||
pub fn observed_data_usage_is_newer(observed: &DataUsageInfo, authoritative: &DataUsageInfo) -> bool {
|
||||
observed.usage_snapshot_converged == Some(false)
|
||||
&& observed.is_complete_bucket_usage_snapshot()
|
||||
&& (observed.is_complete_bucket_usage_snapshot() || observed.is_valid_partial_snapshot())
|
||||
&& observed.usage_snapshot_authoritative_baseline.as_ref() == Some(&authoritative.snapshot_identity())
|
||||
&& data_usage_snapshot_is_newer(observed, authoritative)
|
||||
}
|
||||
@@ -1436,6 +1461,39 @@ impl DataUsageInfo {
|
||||
&& u64::try_from(self.buckets_usage.len()).ok() == Some(self.buckets_count)
|
||||
}
|
||||
|
||||
/// Validate provenance before an observational view can be selected for
|
||||
/// admin display. Partial data is accepted only with unique set states,
|
||||
/// a plan digest for every state, and at least one usable generation.
|
||||
pub fn is_valid_partial_snapshot(&self) -> bool {
|
||||
if !self.usage_snapshot_partial
|
||||
|| self.usage_snapshot_converged != Some(false)
|
||||
|| self.last_update.is_none()
|
||||
|| self.scanner_cycle.is_none()
|
||||
|| self.scanner_epoch.is_none()
|
||||
|| self.usage_snapshot_set_states.is_empty()
|
||||
|| u64::try_from(self.buckets_usage.len()).ok() != Some(self.buckets_count)
|
||||
{
|
||||
return false;
|
||||
}
|
||||
|
||||
let mut previous = None;
|
||||
let mut plan_digest = None;
|
||||
let mut has_source = false;
|
||||
for state in &self.usage_snapshot_set_states {
|
||||
if state.scan_plan_digest.is_none()
|
||||
|| plan_digest.is_some_and(|digest| Some(digest) != state.scan_plan_digest)
|
||||
|| state.scanner_cycle.is_some() != state.scanner_epoch.is_some()
|
||||
|| previous.is_some_and(|(pool, set)| (pool, set) >= (state.pool_index, state.set_index))
|
||||
{
|
||||
return false;
|
||||
}
|
||||
previous = Some((state.pool_index, state.set_index));
|
||||
plan_digest = state.scan_plan_digest;
|
||||
has_source |= state.scanner_cycle.is_some() && !state.tombstone;
|
||||
}
|
||||
has_source
|
||||
}
|
||||
|
||||
/// Add object metadata to data usage statistics
|
||||
pub fn add_object(&mut self, object_path: &str, meta_object: &rustfs_filemeta::MetaObject) {
|
||||
// This method is kept for backward compatibility
|
||||
@@ -2263,6 +2321,55 @@ mod tests {
|
||||
assert!(!observed_data_usage_is_newer(&candidate(2, 9, Some(false), true), &authoritative));
|
||||
assert!(!observed_data_usage_is_newer(&candidate(2, 11, Some(true), true), &authoritative));
|
||||
assert!(!observed_data_usage_is_newer(&candidate(2, 11, Some(false), false), &authoritative));
|
||||
|
||||
let mut partial = candidate(2, 11, Some(false), false);
|
||||
partial.usage_snapshot_partial = true;
|
||||
partial.usage_snapshot_set_states = vec![DataUsageSnapshotSetState {
|
||||
pool_index: 0,
|
||||
set_index: 0,
|
||||
scanner_cycle: Some(10),
|
||||
scanner_epoch: Some(2),
|
||||
scan_plan_digest: Some([1; 32]),
|
||||
complete: false,
|
||||
tombstone: false,
|
||||
}];
|
||||
assert!(observed_data_usage_is_newer(&partial, &authoritative));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn mixed_topology_snapshot_is_rejected() {
|
||||
let mut partial = DataUsageInfo {
|
||||
last_update: Some(SystemTime::UNIX_EPOCH + Duration::from_secs(2)),
|
||||
scanner_cycle: Some(11),
|
||||
scanner_epoch: Some(2),
|
||||
buckets_count: 0,
|
||||
usage_snapshot_converged: Some(false),
|
||||
usage_snapshot_partial: true,
|
||||
usage_snapshot_set_states: vec![
|
||||
DataUsageSnapshotSetState {
|
||||
pool_index: 0,
|
||||
set_index: 0,
|
||||
scanner_cycle: Some(11),
|
||||
scanner_epoch: Some(2),
|
||||
scan_plan_digest: Some([1; 32]),
|
||||
complete: true,
|
||||
tombstone: false,
|
||||
},
|
||||
DataUsageSnapshotSetState {
|
||||
pool_index: 1,
|
||||
set_index: 0,
|
||||
scanner_cycle: Some(10),
|
||||
scanner_epoch: Some(2),
|
||||
scan_plan_digest: Some([2; 32]),
|
||||
complete: false,
|
||||
tombstone: false,
|
||||
},
|
||||
],
|
||||
..Default::default()
|
||||
};
|
||||
assert!(!partial.is_valid_partial_snapshot());
|
||||
partial.usage_snapshot_set_states[1].scan_plan_digest = Some([1; 32]);
|
||||
assert!(partial.is_valid_partial_snapshot());
|
||||
}
|
||||
|
||||
#[test]
|
||||
|
||||
@@ -241,8 +241,8 @@ pub mod cache {
|
||||
|
||||
pub mod capacity {
|
||||
pub use crate::core::pools::{
|
||||
DecommissionUnresolvedEntry, PoolDecommissionInfo, PoolStatus, get_total_usable_capacity, get_total_usable_capacity_free,
|
||||
path2_bucket_object, path2_bucket_object_with_base_path,
|
||||
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;
|
||||
}
|
||||
|
||||
@@ -804,88 +804,6 @@ 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<()>>,
|
||||
@@ -1063,15 +981,18 @@ fn resolve_decommission_listing_error(listing_error: Option<Error>, entry_error:
|
||||
}
|
||||
}
|
||||
|
||||
fn decommission_unresolved_listing_error(entry: &DecommissionUnresolvedEntry) -> 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);
|
||||
Error::other(format!(
|
||||
"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,
|
||||
"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))"
|
||||
))
|
||||
}
|
||||
|
||||
@@ -1083,32 +1004,22 @@ fn resolve_decommission_partial_listing_entry(
|
||||
disk_error_count: usize,
|
||||
pool_index: usize,
|
||||
set_index: usize,
|
||||
source_generation: OffsetDateTime,
|
||||
) -> std::result::Result<MetaCacheEntry, DecommissionUnresolvedEntry> {
|
||||
) -> Result<MetaCacheEntry> {
|
||||
let candidate_count = entries.as_ref().iter().flatten().count();
|
||||
if let Some(entry) = entries.resolve(resolver) {
|
||||
return Ok(entry);
|
||||
}
|
||||
|
||||
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,
|
||||
let candidate = entries.as_ref().iter().flatten().map(|entry| entry.name.as_str()).next();
|
||||
Err(decommission_unresolved_listing_error(
|
||||
bucket,
|
||||
prefix,
|
||||
candidate,
|
||||
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<()> {
|
||||
@@ -1713,8 +1624,6 @@ 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)]
|
||||
@@ -1846,7 +1755,6 @@ 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,
|
||||
})
|
||||
@@ -1879,7 +1787,6 @@ 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,
|
||||
})
|
||||
@@ -1928,7 +1835,6 @@ 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(),
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -2529,26 +2435,6 @@ 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")]
|
||||
@@ -2594,8 +2480,6 @@ 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>,
|
||||
@@ -2629,7 +2513,6 @@ 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 {
|
||||
@@ -2947,21 +2830,6 @@ 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.
|
||||
@@ -3810,7 +3678,6 @@ 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(),
|
||||
@@ -3824,9 +3691,8 @@ 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(store, rx, bucket, callback, entry_error, idx, set_idx, generation)
|
||||
set.list_objects_to_decommission(rx, bucket, callback, entry_error, idx, set_idx)
|
||||
.await
|
||||
}
|
||||
},
|
||||
@@ -4727,7 +4593,7 @@ impl ECStore {
|
||||
state = "verifying_completion",
|
||||
"Decommission completion verification started"
|
||||
);
|
||||
if let Err(err) = self.check_after_decommission(idx, generation).await {
|
||||
if let Err(err) = self.check_after_decommission(idx).await {
|
||||
resolve_decommission_terminal_mark_result(
|
||||
self.decommission_failed_for_operation(idx, canceler).await,
|
||||
"failed",
|
||||
@@ -4754,7 +4620,7 @@ impl ECStore {
|
||||
"Decommission marking completed state"
|
||||
);
|
||||
resolve_decommission_terminal_mark_result(
|
||||
self.complete_decommission_for_operation(idx, canceler, generation).await,
|
||||
self.complete_decommission_for_operation(idx, canceler).await,
|
||||
"completed",
|
||||
&cmd_line,
|
||||
)?;
|
||||
@@ -4890,25 +4756,14 @@ impl ECStore {
|
||||
|
||||
#[tracing::instrument(skip(self))]
|
||||
pub async fn complete_decommission(&self, idx: usize) -> Result<()> {
|
||||
self.complete_decommission_with_owner(idx, None, None).await
|
||||
self.complete_decommission_with_owner(idx, None).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_for_operation(&self, idx: usize, owner: &DecommissionCanceler) -> Result<()> {
|
||||
self.complete_decommission_with_owner(idx, Some(owner)).await
|
||||
}
|
||||
|
||||
async fn complete_decommission_with_owner(
|
||||
&self,
|
||||
idx: usize,
|
||||
owner: Option<&DecommissionCanceler>,
|
||||
verified_generation: Option<OffsetDateTime>,
|
||||
) -> Result<()> {
|
||||
async fn complete_decommission_with_owner(&self, idx: usize, owner: Option<&DecommissionCanceler>) -> Result<()> {
|
||||
ensure_decommission_terminal_operation_supported(self.single_pool(), "complete decommission")?;
|
||||
let _start_guard = self.start_gate.lock().await;
|
||||
|
||||
@@ -4920,13 +4775,11 @@ 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| {
|
||||
reconcile_decommission_unresolved_entries_for_completion(pool_meta, idx, verified_generation)?;
|
||||
Ok::<bool, Error>(pool_meta.decommission_complete(idx))
|
||||
pool_meta.decommission_complete(idx)
|
||||
})
|
||||
else {
|
||||
return Ok(());
|
||||
};
|
||||
let changed = changed?;
|
||||
let terminal_canceler = if let Some(owner) = owner {
|
||||
Some(owner.clone())
|
||||
} else {
|
||||
@@ -5286,7 +5139,7 @@ impl ECStore {
|
||||
Ok(ret)
|
||||
}
|
||||
|
||||
async fn check_after_decommission(self: &Arc<Self>, idx: usize, generation: OffsetDateTime) -> Result<()> {
|
||||
async fn check_after_decommission(self: &Arc<Self>, idx: usize) -> Result<()> {
|
||||
let buckets = self.get_buckets_to_decommission().await?;
|
||||
let pool = self.pools[idx].clone();
|
||||
|
||||
@@ -5385,16 +5238,7 @@ impl ECStore {
|
||||
});
|
||||
|
||||
let list_result = set
|
||||
.list_objects_to_decommission(
|
||||
self.clone(),
|
||||
callback_rx,
|
||||
bucket_info.clone(),
|
||||
callback,
|
||||
entry_error.clone(),
|
||||
idx,
|
||||
set_index,
|
||||
generation,
|
||||
)
|
||||
.list_objects_to_decommission(callback_rx, bucket_info.clone(), callback, entry_error.clone(), idx, set_index)
|
||||
.await;
|
||||
let entry_error = entry_error.lock().await.clone();
|
||||
resolve_decommission_check_after_list_result(list_result, entry_error)?;
|
||||
@@ -5763,17 +5607,6 @@ 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 {
|
||||
@@ -5794,7 +5627,6 @@ 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()
|
||||
}),
|
||||
}],
|
||||
@@ -5834,7 +5666,6 @@ 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);
|
||||
}
|
||||
@@ -5926,7 +5757,6 @@ mod tests {
|
||||
assert!(decommission.bucket.is_empty());
|
||||
assert!(decommission.prefix.is_empty());
|
||||
assert!(decommission.object.is_empty());
|
||||
assert!(decommission.unresolved_entries.is_empty());
|
||||
}
|
||||
|
||||
#[test]
|
||||
@@ -6216,32 +6046,28 @@ async fn record_decommission_entry_error(
|
||||
entry_error: &Arc<tokio::sync::Mutex<Option<Error>>>,
|
||||
rx: &CancellationToken,
|
||||
err: Error,
|
||||
) -> bool {
|
||||
) {
|
||||
if rx.is_cancelled() {
|
||||
return false;
|
||||
return;
|
||||
}
|
||||
|
||||
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, store, rx, cb_func, entry_error))]
|
||||
#[tracing::instrument(skip(self, 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)?;
|
||||
@@ -6262,8 +6088,6 @@ 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,
|
||||
@@ -6283,10 +6107,8 @@ 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 {});
|
||||
@@ -6300,7 +6122,6 @@ impl SetDisks {
|
||||
disk_error_count,
|
||||
pool_index,
|
||||
set_index,
|
||||
source_generation,
|
||||
) {
|
||||
Ok(entry) => {
|
||||
warn!("decommission_pool: list_objects_to_decommission get {}", &entry.name);
|
||||
@@ -6308,11 +6129,10 @@ impl SetDisks {
|
||||
cb_func(entry).await;
|
||||
})
|
||||
}
|
||||
Err(unresolved_entry) => Box::pin(async move {
|
||||
Err(err) => 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,
|
||||
@@ -6323,13 +6143,6 @@ 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;
|
||||
}),
|
||||
}
|
||||
@@ -6538,52 +6351,47 @@ 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, 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,
|
||||
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,
|
||||
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,
|
||||
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,
|
||||
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,
|
||||
};
|
||||
use crate::bucket::metadata_sys;
|
||||
use crate::data_movement;
|
||||
use crate::disk::{STORAGE_FORMAT_FILE, endpoint::Endpoint};
|
||||
use crate::disk::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};
|
||||
@@ -7814,8 +7622,7 @@ mod pools_tests {
|
||||
|
||||
#[test]
|
||||
fn test_resolve_decommission_partial_listing_entry_rejects_unresolved_metadata() {
|
||||
let generation = OffsetDateTime::now_utc();
|
||||
let unresolved_entry = resolve_decommission_partial_listing_entry(
|
||||
let err = resolve_decommission_partial_listing_entry(
|
||||
MetaCacheEntries(vec![None]),
|
||||
MetadataResolutionParams {
|
||||
dir_quorum: 2,
|
||||
@@ -7828,160 +7635,23 @@ mod pools_tests {
|
||||
1,
|
||||
2,
|
||||
3,
|
||||
generation,
|
||||
)
|
||||
.expect_err("unresolved partial listing must fail closed");
|
||||
|
||||
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();
|
||||
let message = err.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();
|
||||
|
||||
assert!(record_decommission_entry_error(&entry_error, &rx, Error::SlowDown).await);
|
||||
assert!(!record_decommission_entry_error(&entry_error, &rx, Error::OperationCanceled).await);
|
||||
record_decommission_entry_error(&entry_error, &rx, Error::SlowDown).await;
|
||||
record_decommission_entry_error(&entry_error, &rx, Error::OperationCanceled).await;
|
||||
|
||||
assert!(rx.is_cancelled());
|
||||
assert!(matches!(*entry_error.lock().await, Some(Error::SlowDown)));
|
||||
@@ -7993,7 +7663,7 @@ mod pools_tests {
|
||||
let rx = CancellationToken::new();
|
||||
rx.cancel();
|
||||
|
||||
assert!(!record_decommission_entry_error(&entry_error, &rx, Error::SlowDown).await);
|
||||
record_decommission_entry_error(&entry_error, &rx, Error::SlowDown).await;
|
||||
|
||||
assert!(entry_error.lock().await.is_none());
|
||||
}
|
||||
|
||||
@@ -73,6 +73,16 @@ struct CachedBucketUsage {
|
||||
// mutation. A strictly later generation is required before the mutation
|
||||
// evidence can be discarded.
|
||||
pending_scanner_position: Option<(u64, u64)>,
|
||||
// Deletes are visible to admin immediately, but quota admission keeps
|
||||
// them pending until a complete scanner generation reconciles the set.
|
||||
// This marker intentionally remains process-local: the delete request
|
||||
// updates this overlay before the scanner writes a durable snapshot. If
|
||||
// the process restarts first, loading the persisted complete snapshot
|
||||
// restores the pre-reconciliation (larger) baseline, which is
|
||||
// conservative for quota admission. A persisted post-delete snapshot is
|
||||
// necessarily a complete scanner reconciliation and therefore creates a
|
||||
// fresh cache entry with no pending hold.
|
||||
pending_negative_delta: u64,
|
||||
}
|
||||
|
||||
type UsageMemoryCache = Arc<RwLock<HashMap<String, CachedBucketUsage>>>;
|
||||
@@ -948,7 +958,12 @@ async fn load_observed_data_usage_snapshot(store: Arc<ECStore>) -> Option<DataUs
|
||||
};
|
||||
|
||||
match parse_usage_snapshot(&data) {
|
||||
Ok(info) if info.usage_snapshot_converged == Some(false) && info.is_complete_bucket_usage_snapshot() => Some(info),
|
||||
Ok(info)
|
||||
if info.usage_snapshot_converged == Some(false)
|
||||
&& (info.is_complete_bucket_usage_snapshot() || info.is_valid_partial_snapshot()) =>
|
||||
{
|
||||
Some(info)
|
||||
}
|
||||
Ok(_) => {
|
||||
error!(
|
||||
event = "data_usage_snapshot_load_failed",
|
||||
@@ -993,7 +1008,7 @@ async fn load_admin_data_usage_from_backend(store: Arc<ECStore>) -> Result<DataU
|
||||
}
|
||||
|
||||
fn discard_incomplete_bucket_usage(data_usage_info: &mut DataUsageInfo) {
|
||||
if !data_usage_info.is_complete_bucket_usage_snapshot() {
|
||||
if !data_usage_info.is_complete_bucket_usage_snapshot() && !data_usage_info.usage_snapshot_partial {
|
||||
data_usage_info.usage_snapshot_complete = false;
|
||||
data_usage_info.buckets_usage.clear();
|
||||
data_usage_info.bucket_sizes.clear();
|
||||
@@ -1643,6 +1658,7 @@ fn cached_bucket_usage_from_backend(usage: BucketUsageInfo, updated_at: SystemTi
|
||||
dirty: false,
|
||||
stale_snapshot_pending: false,
|
||||
pending_scanner_position: None,
|
||||
pending_negative_delta: 0,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1656,6 +1672,7 @@ fn cached_bucket_usage_now(usage: BucketUsageInfo) -> CachedBucketUsage {
|
||||
dirty: false,
|
||||
stale_snapshot_pending: false,
|
||||
pending_scanner_position: None,
|
||||
pending_negative_delta: 0,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1808,6 +1825,7 @@ pub async fn record_bucket_object_delete_memory(bucket: &str, deleted_size: u64,
|
||||
.or_insert_with(|| cached_bucket_usage_now(BucketUsageInfo::default()));
|
||||
|
||||
entry.usage.size = entry.usage.size.saturating_sub(deleted_size);
|
||||
entry.pending_negative_delta = entry.pending_negative_delta.saturating_add(deleted_size);
|
||||
if removed_current_object {
|
||||
entry.usage.objects_count = entry.usage.objects_count.saturating_sub(1);
|
||||
entry.usage.versions_count = entry.usage.versions_count.saturating_sub(1);
|
||||
@@ -1863,7 +1881,7 @@ pub async fn get_bucket_usage_memory(bucket: &str) -> Option<u64> {
|
||||
cache
|
||||
.get(bucket)
|
||||
.filter(|cached| cached.authoritative)
|
||||
.map(|cached| cached.usage.size)
|
||||
.map(|cached| cached.usage.size.saturating_add(cached.pending_negative_delta))
|
||||
}
|
||||
|
||||
async fn update_usage_cache_if_needed() {
|
||||
@@ -2943,6 +2961,45 @@ mod tests {
|
||||
assert_eq!(selected.usage_snapshot_converged, Some(true));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn persisted_authoritative_stalls_but_memory_overlay_remains_visible() {
|
||||
let authoritative = DataUsageInfo {
|
||||
last_update: Some(SystemTime::UNIX_EPOCH),
|
||||
scanner_epoch: Some(4),
|
||||
scanner_cycle: Some(10),
|
||||
usage_snapshot_complete: true,
|
||||
..Default::default()
|
||||
};
|
||||
let mut partial = authoritative.clone();
|
||||
partial.last_update = Some(SystemTime::UNIX_EPOCH + Duration::from_secs(1));
|
||||
partial.scanner_cycle = Some(11);
|
||||
partial.usage_snapshot_complete = false;
|
||||
partial.usage_snapshot_partial = true;
|
||||
partial.usage_snapshot_converged = Some(false);
|
||||
partial.usage_snapshot_authoritative_baseline = Some(authoritative.snapshot_identity());
|
||||
partial.usage_snapshot_set_states = vec![rustfs_data_usage::DataUsageSnapshotSetState {
|
||||
pool_index: 0,
|
||||
set_index: 0,
|
||||
scanner_cycle: Some(10),
|
||||
scanner_epoch: Some(4),
|
||||
scan_plan_digest: Some([1; 32]),
|
||||
complete: false,
|
||||
tombstone: false,
|
||||
}];
|
||||
partial.buckets_usage.insert(
|
||||
"bucket".to_string(),
|
||||
BucketUsageInfo {
|
||||
size: 100,
|
||||
..Default::default()
|
||||
},
|
||||
);
|
||||
partial.buckets_count = 1;
|
||||
|
||||
let (selected, _) = select_admin_data_usage_snapshot(authoritative, true, Some(partial));
|
||||
assert!(selected.usage_snapshot_partial);
|
||||
assert_eq!(selected.buckets_usage.get("bucket").map(|usage| usage.size), Some(100));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn authoritative_save_cleanup_removes_observed_snapshot_best_effort() {
|
||||
let store = UsageCasStore::default();
|
||||
@@ -4665,6 +4722,55 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn partial_usage_is_observational_not_authoritative_for_quota() {
|
||||
clear_usage_memory_cache_for_test().await;
|
||||
|
||||
let mut partial = data_usage_info_for_test("bucket-a", 10, 100, SystemTime::now());
|
||||
partial.usage_snapshot_complete = false;
|
||||
partial.usage_snapshot_partial = true;
|
||||
replace_bucket_usage_memory_from_info(&partial).await;
|
||||
|
||||
assert_eq!(get_bucket_usage_memory("bucket-a").await, None);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn stale_quota_uses_complete_baseline_plus_positive_deltas() {
|
||||
clear_usage_memory_cache_for_test().await;
|
||||
|
||||
let baseline = data_usage_info_for_test("bucket-a", 1, 100, SystemTime::now());
|
||||
replace_bucket_usage_memory_from_info(&baseline).await;
|
||||
record_bucket_object_write_memory("bucket-a", None, 25).await;
|
||||
|
||||
assert_eq!(get_bucket_usage_memory("bucket-a").await, Some(125));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn negative_delta_waits_for_set_reconciliation() {
|
||||
clear_usage_memory_cache_for_test().await;
|
||||
|
||||
let baseline = data_usage_info_for_test("bucket-a", 1, 100, SystemTime::UNIX_EPOCH + Duration::from_secs(100));
|
||||
replace_bucket_usage_memory_from_info(&baseline).await;
|
||||
record_bucket_object_delete_memory("bucket-a", 25, true).await;
|
||||
|
||||
assert_eq!(get_bucket_usage_memory("bucket-a").await, Some(100));
|
||||
|
||||
// Simulate a process restart: the request-path overlay is gone, but
|
||||
// the persisted authoritative snapshot is still the pre-reconciliation
|
||||
// baseline. Quota must remain conservative until a complete scanner
|
||||
// result proves the delete.
|
||||
clear_usage_memory_cache_for_test().await;
|
||||
replace_bucket_usage_memory_from_info(&baseline).await;
|
||||
assert_eq!(get_bucket_usage_memory("bucket-a").await, Some(100));
|
||||
|
||||
let reconciled = data_usage_info_for_test("bucket-a", 0, 75, SystemTime::UNIX_EPOCH + Duration::from_secs(101));
|
||||
replace_bucket_usage_memory_from_info(&reconciled).await;
|
||||
assert_eq!(get_bucket_usage_memory("bucket-a").await, Some(75));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn memory_overlay_counts_versioned_overwrite_as_new_version() {
|
||||
|
||||
@@ -28,8 +28,9 @@ use rustfs_common::heal_channel::HealScanMode;
|
||||
use rustfs_config::ENV_SCANNER_CACHE_SAVE_TIMEOUT_SECS;
|
||||
pub use rustfs_data_usage::{
|
||||
AllTierStats, BucketTargetUsageInfo, BucketUsageInfo, DATA_USAGE_OBJECT_NAME, DATA_USAGE_OBSERVED_OBJECT_NAME,
|
||||
DataUsageEntry, DataUsageHash, DataUsageHashMap, DataUsageInfo, LEGACY_DATA_USAGE_OBJECT_NAME, PrefixUsageEntry,
|
||||
PrefixUsageQuery, PrefixUsageSummary, ReplTargetSizeSummary, SizeSummary, TierStats, hash_path, prefix_usage_in_cache,
|
||||
DataUsageEntry, DataUsageHash, DataUsageHashMap, DataUsageInfo, DataUsageSnapshotSetState, LEGACY_DATA_USAGE_OBJECT_NAME,
|
||||
PrefixUsageEntry, PrefixUsageQuery, PrefixUsageSummary, ReplTargetSizeSummary, SizeSummary, TierStats, hash_path,
|
||||
prefix_usage_in_cache,
|
||||
};
|
||||
use rustfs_utils::path::{SLASH_SEPARATOR, path_join_buf};
|
||||
use tokio::time::{Duration, Instant, sleep, timeout};
|
||||
@@ -344,6 +345,18 @@ pub struct DataUsageCacheInfo {
|
||||
pub scan_plan_digest: Option<DataUsageScanPlanDigest>,
|
||||
#[serde(default)]
|
||||
pub cache_key_format: u16,
|
||||
/// Whether the entries retained while a set scan was incomplete come
|
||||
/// from a prior complete set snapshot. This is observational input only.
|
||||
#[serde(default)]
|
||||
pub lkg_snapshot_complete: bool,
|
||||
#[serde(default)]
|
||||
pub lkg_next_cycle: Option<u64>,
|
||||
#[serde(default)]
|
||||
pub lkg_last_update: Option<SystemTime>,
|
||||
#[serde(default)]
|
||||
pub lkg_leader_epoch: Option<u64>,
|
||||
#[serde(default)]
|
||||
pub lkg_scan_plan_digest: Option<DataUsageScanPlanDigest>,
|
||||
}
|
||||
|
||||
impl Serialize for DataUsageCacheInfo {
|
||||
@@ -353,7 +366,7 @@ impl Serialize for DataUsageCacheInfo {
|
||||
{
|
||||
// Keep this metadata map-encoded so older readers can ignore fields
|
||||
// appended by newer scanner versions during rolling upgrades.
|
||||
let mut state = serializer.serialize_map(Some(16))?;
|
||||
let mut state = serializer.serialize_map(Some(21))?;
|
||||
state.serialize_entry("name", &self.name)?;
|
||||
state.serialize_entry("next_cycle", &self.next_cycle)?;
|
||||
state.serialize_entry("leader_epoch", &self.leader_epoch)?;
|
||||
@@ -370,6 +383,11 @@ impl Serialize for DataUsageCacheInfo {
|
||||
state.serialize_entry("snapshot_complete", &self.snapshot_complete)?;
|
||||
state.serialize_entry("scan_plan_digest", &self.scan_plan_digest)?;
|
||||
state.serialize_entry("cache_key_format", &self.cache_key_format)?;
|
||||
state.serialize_entry("lkg_snapshot_complete", &self.lkg_snapshot_complete)?;
|
||||
state.serialize_entry("lkg_next_cycle", &self.lkg_next_cycle)?;
|
||||
state.serialize_entry("lkg_last_update", &self.lkg_last_update)?;
|
||||
state.serialize_entry("lkg_leader_epoch", &self.lkg_leader_epoch)?;
|
||||
state.serialize_entry("lkg_scan_plan_digest", &self.lkg_scan_plan_digest)?;
|
||||
state.end()
|
||||
}
|
||||
}
|
||||
|
||||
@@ -2274,9 +2274,8 @@ async fn final_data_usage_publication_defer_reason(
|
||||
}
|
||||
}
|
||||
ScannerCycleStatus::Deferred(reason) => Some(reason),
|
||||
// Incomplete cycles do not publish a usage snapshot. Keep the
|
||||
// decision permissive so existing partial-cycle handling remains
|
||||
// unchanged if a future scanner path emits a bookkeeping update.
|
||||
// Incomplete cycles may publish a non-authoritative observational
|
||||
// snapshot when at least one set has a usable current/LKG view.
|
||||
ScannerCycleStatus::Incomplete => None,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -198,7 +198,7 @@ where
|
||||
data_usage_info.usage_snapshot_authoritative_baseline = Some(authoritative.snapshot_identity());
|
||||
}
|
||||
|
||||
if !data_usage_info.is_complete_bucket_usage_snapshot() {
|
||||
if !data_usage_info.is_complete_bucket_usage_snapshot() && !data_usage_info.usage_snapshot_partial {
|
||||
error!(
|
||||
target: "rustfs::scanner",
|
||||
event = EVENT_SCANNER_PERSIST_STATE,
|
||||
|
||||
@@ -18,8 +18,8 @@ use crate::scanner_folder::{ScannerItem, scan_data_folder};
|
||||
use crate::sleeper::SCANNER_SLEEPER;
|
||||
use crate::{
|
||||
DATA_USAGE_CACHE_NAME, DATA_USAGE_ROOT, DataUsageCache, DataUsageCacheInfo, DataUsageCachePrepareOutcome,
|
||||
DataUsageCacheSource, DataUsageEntry, DataUsageEntryInfo, DataUsageInfo, DataUsageScanPlanDigest, ScannerError, SizeSummary,
|
||||
TierStats,
|
||||
DataUsageCacheSource, DataUsageEntry, DataUsageEntryInfo, DataUsageInfo, DataUsageScanPlanDigest, DataUsageSnapshotSetState,
|
||||
ScannerError, SizeSummary, TierStats,
|
||||
};
|
||||
use futures::future::join_all;
|
||||
use metrics::counter;
|
||||
@@ -278,6 +278,17 @@ async fn publish_usage_snapshot(
|
||||
Ok(true)
|
||||
}
|
||||
|
||||
async fn publish_observational_snapshot(
|
||||
updates: &mpsc::Sender<DataUsageInfo>,
|
||||
mut data_usage_info: DataUsageInfo,
|
||||
) -> Result<bool> {
|
||||
data_usage_info.usage_snapshot_complete = false;
|
||||
data_usage_info.usage_snapshot_partial = true;
|
||||
data_usage_info.usage_snapshot_converged = Some(false);
|
||||
send_data_usage_update(updates, data_usage_info).await?;
|
||||
Ok(true)
|
||||
}
|
||||
|
||||
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
|
||||
enum ScannerCycleActivityStatus {
|
||||
Unchanged,
|
||||
|
||||
@@ -188,7 +188,7 @@ pub(super) fn completed_data_usage_info(
|
||||
}
|
||||
|
||||
let mut total = DataUsageEntry::default();
|
||||
let mut buckets_usage = HashMap::with_capacity(all_buckets.len());
|
||||
let mut bucket_entries = HashMap::with_capacity(all_buckets.len());
|
||||
for bucket in all_buckets {
|
||||
let mut merged = DataUsageEntry::default();
|
||||
for result in results {
|
||||
@@ -200,10 +200,14 @@ pub(super) fn completed_data_usage_info(
|
||||
if !total.checked_merge(&merged) {
|
||||
return None;
|
||||
}
|
||||
buckets_usage.insert(bucket.clone(), checked_bucket_usage_info(&merged)?);
|
||||
bucket_entries.insert(bucket.clone(), merged);
|
||||
}
|
||||
|
||||
let merged_last_update = results.iter().filter_map(|result| result.info.last_update).max()?;
|
||||
let buckets_usage = bucket_entries
|
||||
.iter()
|
||||
.map(|(bucket, entry)| Some((bucket.clone(), checked_bucket_usage_info(entry)?)))
|
||||
.collect::<Option<HashMap<_, _>>>()?;
|
||||
let bucket_sizes = buckets_usage
|
||||
.iter()
|
||||
.map(|(bucket, usage)| (bucket.clone(), usage.size))
|
||||
@@ -225,6 +229,145 @@ pub(super) fn completed_data_usage_info(
|
||||
Some((data_usage_info, merged_last_update))
|
||||
}
|
||||
|
||||
/// Build a non-authoritative view from the set snapshots that completed this
|
||||
/// cycle plus compatible per-set last-known-good caches. The caller must
|
||||
/// persist this result only on the observational object; a missing set is
|
||||
/// intentionally represented by an incomplete state and is never treated as
|
||||
/// an empty set.
|
||||
pub(super) fn observational_data_usage_info(
|
||||
results: &[DataUsageCache],
|
||||
expected_sources: &HashSet<DataUsageCacheSource>,
|
||||
all_buckets: &[String],
|
||||
expected_plan_digest: DataUsageScanPlanDigest,
|
||||
scanner_cycle: u64,
|
||||
leader_epoch: u64,
|
||||
) -> Option<(DataUsageInfo, SystemTime)> {
|
||||
let mut by_source = HashMap::with_capacity(results.len());
|
||||
for result in results {
|
||||
let source = result.info.source?;
|
||||
if !expected_sources.contains(&source) || by_source.insert(source, result).is_some() {
|
||||
return None;
|
||||
}
|
||||
}
|
||||
|
||||
let mut usable = Vec::new();
|
||||
let mut set_states = Vec::with_capacity(expected_sources.len());
|
||||
let mut sources = expected_sources.iter().copied().collect::<Vec<_>>();
|
||||
sources.sort_by_key(|source| (source.pool_index, source.set_index));
|
||||
for source in sources {
|
||||
let result = by_source.get(&source).copied();
|
||||
let current = result.filter(|result| {
|
||||
result.info.snapshot_complete
|
||||
&& result.info.next_cycle == scanner_cycle
|
||||
&& result.info.leader_epoch == leader_epoch
|
||||
&& result.info.scan_plan_digest == Some(expected_plan_digest)
|
||||
});
|
||||
let lkg = result.filter(|result| {
|
||||
!result.info.snapshot_complete
|
||||
&& result.info.lkg_snapshot_complete
|
||||
&& result.info.lkg_scan_plan_digest == Some(expected_plan_digest)
|
||||
&& result.info.lkg_leader_epoch.is_some_and(|epoch| {
|
||||
epoch < leader_epoch
|
||||
|| (epoch == leader_epoch && result.info.lkg_next_cycle.is_some_and(|cycle| cycle <= scanner_cycle))
|
||||
})
|
||||
});
|
||||
let current_snapshot = current.is_some();
|
||||
let selected = current.or(lkg);
|
||||
if let Some(selected) = selected {
|
||||
let (cycle, epoch, digest, last_update, complete) = if current_snapshot {
|
||||
(
|
||||
Some(selected.info.next_cycle),
|
||||
Some(selected.info.leader_epoch),
|
||||
selected.info.scan_plan_digest.map(|digest| digest.0),
|
||||
selected.info.last_update,
|
||||
true,
|
||||
)
|
||||
} else {
|
||||
(
|
||||
selected.info.lkg_next_cycle,
|
||||
selected.info.lkg_leader_epoch,
|
||||
selected.info.lkg_scan_plan_digest.map(|digest| digest.0),
|
||||
selected.info.lkg_last_update,
|
||||
false,
|
||||
)
|
||||
};
|
||||
set_states.push(DataUsageSnapshotSetState {
|
||||
pool_index: u64::try_from(source.pool_index).ok()?,
|
||||
set_index: u64::try_from(source.set_index).ok()?,
|
||||
scanner_cycle: cycle,
|
||||
scanner_epoch: epoch,
|
||||
scan_plan_digest: digest,
|
||||
complete,
|
||||
tombstone: false,
|
||||
});
|
||||
usable.push((selected, last_update));
|
||||
} else {
|
||||
set_states.push(DataUsageSnapshotSetState {
|
||||
pool_index: u64::try_from(source.pool_index).ok()?,
|
||||
set_index: u64::try_from(source.set_index).ok()?,
|
||||
scanner_cycle: None,
|
||||
scanner_epoch: None,
|
||||
scan_plan_digest: Some(expected_plan_digest.0),
|
||||
complete: false,
|
||||
tombstone: false,
|
||||
});
|
||||
}
|
||||
}
|
||||
if usable.is_empty() {
|
||||
return None;
|
||||
}
|
||||
|
||||
let mut total = DataUsageEntry::default();
|
||||
let mut bucket_entries = HashMap::with_capacity(all_buckets.len());
|
||||
let mut merged_last_update = None;
|
||||
for (result, last_update) in usable {
|
||||
if let Some(update) = last_update {
|
||||
merged_last_update = Some(merged_last_update.map_or(update, |current: SystemTime| current.max(update)));
|
||||
}
|
||||
for bucket in all_buckets {
|
||||
let Some(entry) = result.checked_flatten(bucket) else {
|
||||
continue;
|
||||
};
|
||||
let bucket_entry = bucket_entries.entry(bucket.clone()).or_insert_with(DataUsageEntry::default);
|
||||
if !bucket_entry.checked_merge(&entry) {
|
||||
return None;
|
||||
}
|
||||
if !total.checked_merge(&entry) {
|
||||
return None;
|
||||
}
|
||||
}
|
||||
}
|
||||
let merged_last_update = merged_last_update?;
|
||||
let buckets_usage = bucket_entries
|
||||
.iter()
|
||||
.map(|(bucket, entry)| Some((bucket.clone(), checked_bucket_usage_info(entry)?)))
|
||||
.collect::<Option<HashMap<_, _>>>()?;
|
||||
Some((
|
||||
DataUsageInfo {
|
||||
last_update: Some(merged_last_update),
|
||||
scanner_cycle: Some(scanner_cycle),
|
||||
scanner_epoch: Some(leader_epoch),
|
||||
objects_total_count: u64::try_from(total.objects).ok()?,
|
||||
versions_total_count: u64::try_from(total.versions).ok()?,
|
||||
delete_markers_total_count: u64::try_from(total.delete_markers).ok()?,
|
||||
objects_total_size: u64::try_from(total.size).ok()?,
|
||||
tier_stats: total.all_tier_stats.filter(|tiers| !tiers.is_empty()),
|
||||
buckets_count: u64::try_from(buckets_usage.len()).ok()?,
|
||||
bucket_sizes: buckets_usage
|
||||
.iter()
|
||||
.map(|(bucket, usage)| (bucket.clone(), usage.size))
|
||||
.collect(),
|
||||
buckets_usage,
|
||||
usage_snapshot_complete: false,
|
||||
usage_snapshot_partial: true,
|
||||
usage_snapshot_converged: Some(false),
|
||||
usage_snapshot_set_states: set_states,
|
||||
..Default::default()
|
||||
},
|
||||
merged_last_update,
|
||||
))
|
||||
}
|
||||
|
||||
pub(super) async fn send_cache_root_entry_info(
|
||||
bucket_result_tx: &mpsc::Sender<DataUsageEntryInfo>,
|
||||
cache: &DataUsageCache,
|
||||
|
||||
@@ -40,6 +40,21 @@ impl ScannerIOCache for SetDisks {
|
||||
let set_label = self.set_index.to_string();
|
||||
|
||||
let source = DataUsageCacheSource::new(self.pool_index, self.set_index);
|
||||
let mut old_cache = DataUsageCache::default();
|
||||
if let Err(e) = old_cache.load(self.clone(), DATA_USAGE_CACHE_NAME).await {
|
||||
warn!(
|
||||
target: "rustfs::scanner::io",
|
||||
event = EVENT_SCANNER_CACHE_PERSIST_STATE,
|
||||
component = LOG_COMPONENT_SCANNER,
|
||||
subsystem = LOG_SUBSYSTEM_IO,
|
||||
pool = self.pool_index,
|
||||
set = self.set_index,
|
||||
cache_name = DATA_USAGE_CACHE_NAME,
|
||||
state = "old_cache_load_failed",
|
||||
error = %e,
|
||||
"Scanner old data usage cache load failed; rebuilding from bucket caches"
|
||||
);
|
||||
}
|
||||
if buckets.is_empty() {
|
||||
let now = SystemTime::now();
|
||||
let mut cache = DataUsageCache {
|
||||
@@ -80,6 +95,24 @@ impl ScannerIOCache for SetDisks {
|
||||
"Scanner set state found no online disks"
|
||||
);
|
||||
reset_disk_bucket_scan_gauges(&pool_label, &set_label);
|
||||
let lkg = old_cache.info.snapshot_complete.then(|| old_cache.clone());
|
||||
let mut incomplete_scope = lkg.clone().unwrap_or_default();
|
||||
incomplete_scope.info.name = DATA_USAGE_ROOT.to_string();
|
||||
incomplete_scope.info.next_cycle = want_cycle;
|
||||
incomplete_scope.info.last_update = None;
|
||||
incomplete_scope.info.leader_epoch = leader_epoch;
|
||||
incomplete_scope.info.source = Some(source);
|
||||
incomplete_scope.info.snapshot_complete = false;
|
||||
incomplete_scope.info.scan_plan_digest = Some(scan_plan_digest);
|
||||
incomplete_scope.info.cache_key_format = DATA_USAGE_CACHE_KEY_FORMAT;
|
||||
if let Some(lkg) = lkg {
|
||||
incomplete_scope.info.lkg_snapshot_complete = true;
|
||||
incomplete_scope.info.lkg_next_cycle = Some(lkg.info.next_cycle);
|
||||
incomplete_scope.info.lkg_last_update = lkg.info.last_update;
|
||||
incomplete_scope.info.lkg_leader_epoch = Some(lkg.info.leader_epoch);
|
||||
incomplete_scope.info.lkg_scan_plan_digest = lkg.info.scan_plan_digest;
|
||||
}
|
||||
let _ = updates.send(incomplete_scope).await;
|
||||
return Ok(());
|
||||
}
|
||||
// Preserve the original set topology across capability filtering. During
|
||||
@@ -162,6 +195,24 @@ impl ScannerIOCache for SetDisks {
|
||||
"Scanner set state found no usable namespace scanner disks"
|
||||
);
|
||||
reset_disk_bucket_scan_gauges(&pool_label, &set_label);
|
||||
let lkg = old_cache.info.snapshot_complete.then(|| old_cache.clone());
|
||||
let mut incomplete_scope = lkg.clone().unwrap_or_default();
|
||||
incomplete_scope.info.name = DATA_USAGE_ROOT.to_string();
|
||||
incomplete_scope.info.next_cycle = want_cycle;
|
||||
incomplete_scope.info.last_update = None;
|
||||
incomplete_scope.info.leader_epoch = leader_epoch;
|
||||
incomplete_scope.info.source = Some(source);
|
||||
incomplete_scope.info.snapshot_complete = false;
|
||||
incomplete_scope.info.scan_plan_digest = Some(scan_plan_digest);
|
||||
incomplete_scope.info.cache_key_format = DATA_USAGE_CACHE_KEY_FORMAT;
|
||||
if let Some(lkg) = lkg {
|
||||
incomplete_scope.info.lkg_snapshot_complete = true;
|
||||
incomplete_scope.info.lkg_next_cycle = Some(lkg.info.next_cycle);
|
||||
incomplete_scope.info.lkg_last_update = lkg.info.last_update;
|
||||
incomplete_scope.info.lkg_leader_epoch = Some(lkg.info.leader_epoch);
|
||||
incomplete_scope.info.lkg_scan_plan_digest = lkg.info.scan_plan_digest;
|
||||
}
|
||||
let _ = updates.send(incomplete_scope).await;
|
||||
return Ok(());
|
||||
}
|
||||
let set_disk_inventory = Arc::new(scanner_set_disk_inventory(self.as_ref()).await);
|
||||
@@ -203,22 +254,15 @@ impl ScannerIOCache for SetDisks {
|
||||
record_disk_bucket_scans_active(0, &pool_label, &set_label);
|
||||
let _reset_disk_bucket_scan_gauges = DiskBucketScanGaugeReset::new(pool_label.clone(), set_label.clone());
|
||||
|
||||
let mut old_cache = DataUsageCache::default();
|
||||
if let Err(e) = old_cache.load(self.clone(), DATA_USAGE_CACHE_NAME).await {
|
||||
warn!(
|
||||
target: "rustfs::scanner::io",
|
||||
event = EVENT_SCANNER_CACHE_PERSIST_STATE,
|
||||
component = LOG_COMPONENT_SCANNER,
|
||||
subsystem = LOG_SUBSYSTEM_IO,
|
||||
pool = self.pool_index,
|
||||
set = self.set_index,
|
||||
cache_name = DATA_USAGE_CACHE_NAME,
|
||||
state = "old_cache_load_failed",
|
||||
error = %e,
|
||||
"Scanner old data usage cache load failed; rebuilding from bucket caches"
|
||||
);
|
||||
}
|
||||
match old_cache.prepare_for_scan(
|
||||
let old_lkg = old_cache.info.snapshot_complete.then(|| {
|
||||
(
|
||||
old_cache.info.next_cycle,
|
||||
old_cache.info.last_update,
|
||||
old_cache.info.leader_epoch,
|
||||
old_cache.info.scan_plan_digest,
|
||||
)
|
||||
});
|
||||
let prepare_outcome = match old_cache.prepare_for_scan(
|
||||
DATA_USAGE_ROOT,
|
||||
want_cycle,
|
||||
leader_epoch,
|
||||
@@ -259,7 +303,16 @@ impl ScannerIOCache for SetDisks {
|
||||
);
|
||||
return Ok(());
|
||||
}
|
||||
DataUsageCachePrepareOutcome::Reused | DataUsageCachePrepareOutcome::Reset => {}
|
||||
outcome => outcome,
|
||||
};
|
||||
if matches!(prepare_outcome, DataUsageCachePrepareOutcome::Reused)
|
||||
&& let Some((cycle, last_update, epoch, digest)) = old_lkg
|
||||
{
|
||||
old_cache.info.lkg_snapshot_complete = true;
|
||||
old_cache.info.lkg_next_cycle = Some(cycle);
|
||||
old_cache.info.lkg_last_update = last_update;
|
||||
old_cache.info.lkg_leader_epoch = Some(epoch);
|
||||
old_cache.info.lkg_scan_plan_digest = digest;
|
||||
}
|
||||
|
||||
let mut cache = DataUsageCache {
|
||||
@@ -1099,23 +1152,29 @@ impl ScannerIOCache for SetDisks {
|
||||
cache.info.next_cycle = want_cycle;
|
||||
cache.info.last_update.get_or_insert_with(SystemTime::now);
|
||||
cache.info.snapshot_complete = true;
|
||||
cache.info.lkg_snapshot_complete = false;
|
||||
cache.info.lkg_next_cycle = None;
|
||||
cache.info.lkg_last_update = None;
|
||||
cache.info.lkg_leader_epoch = None;
|
||||
cache.info.lkg_scan_plan_digest = None;
|
||||
cache.clone()
|
||||
};
|
||||
let _ = persist_and_publish_cache_snapshot(self.clone(), &updates, cache_snapshot, cache_cycle_floor.as_ref()).await;
|
||||
} else {
|
||||
let incomplete_scope = DataUsageCache {
|
||||
info: DataUsageCacheInfo {
|
||||
name: DATA_USAGE_ROOT.to_string(),
|
||||
next_cycle: want_cycle,
|
||||
leader_epoch,
|
||||
source: Some(source),
|
||||
snapshot_complete: false,
|
||||
scan_plan_digest: Some(scan_plan_digest),
|
||||
cache_key_format: DATA_USAGE_CACHE_KEY_FORMAT,
|
||||
..Default::default()
|
||||
},
|
||||
cache: HashMap::new(),
|
||||
};
|
||||
let mut incomplete_scope = cache_mutex.lock().await.clone();
|
||||
incomplete_scope.info.name = DATA_USAGE_ROOT.to_string();
|
||||
incomplete_scope.info.next_cycle = want_cycle;
|
||||
incomplete_scope.info.last_update = None;
|
||||
incomplete_scope.info.leader_epoch = leader_epoch;
|
||||
incomplete_scope.info.source = Some(source);
|
||||
incomplete_scope.info.snapshot_complete = false;
|
||||
incomplete_scope.info.scan_plan_digest = Some(scan_plan_digest);
|
||||
incomplete_scope.info.cache_key_format = DATA_USAGE_CACHE_KEY_FORMAT;
|
||||
incomplete_scope.info.lkg_snapshot_complete = old_cache.info.lkg_snapshot_complete;
|
||||
incomplete_scope.info.lkg_next_cycle = old_cache.info.lkg_next_cycle;
|
||||
incomplete_scope.info.lkg_last_update = old_cache.info.lkg_last_update;
|
||||
incomplete_scope.info.lkg_leader_epoch = old_cache.info.lkg_leader_epoch;
|
||||
incomplete_scope.info.lkg_scan_plan_digest = old_cache.info.lkg_scan_plan_digest;
|
||||
if let Err(e) = updates.send(incomplete_scope).await {
|
||||
error!(
|
||||
target: "rustfs::scanner::io",
|
||||
|
||||
@@ -234,6 +234,7 @@ impl ScannerIOCycle for ECStore {
|
||||
let active_set_scans_clone = active_set_scans.clone();
|
||||
|
||||
let (tx, mut rx) = mpsc::channel::<DataUsageCache>(1);
|
||||
let failed_scope_tx = tx.clone();
|
||||
|
||||
// Spawn task to receive and store results
|
||||
let receiver_fut = tokio::spawn(async move {
|
||||
@@ -314,6 +315,21 @@ impl ScannerIOCycle for ECStore {
|
||||
state = "set_scan_failed",
|
||||
"Scanner set scan failed; continuing cycle"
|
||||
);
|
||||
let _ = failed_scope_tx
|
||||
.send(DataUsageCache {
|
||||
info: DataUsageCacheInfo {
|
||||
name: DATA_USAGE_ROOT.to_string(),
|
||||
next_cycle: want_cycle_clone,
|
||||
leader_epoch,
|
||||
source: Some(source),
|
||||
snapshot_complete: false,
|
||||
scan_plan_digest: Some(scan_plan_digest),
|
||||
cache_key_format: DATA_USAGE_CACHE_KEY_FORMAT,
|
||||
..Default::default()
|
||||
},
|
||||
cache: HashMap::new(),
|
||||
})
|
||||
.await;
|
||||
let mut first_err = first_err_mutex_clone.lock().await;
|
||||
record_set_scan_failure(&mut first_err, e);
|
||||
}
|
||||
@@ -370,6 +386,19 @@ impl ScannerIOCycle for ECStore {
|
||||
budget_elapsed,
|
||||
ctx.is_cancelled(),
|
||||
);
|
||||
let observational_usage = completed_usage
|
||||
.is_none()
|
||||
.then(|| {
|
||||
observational_data_usage_info(
|
||||
&results,
|
||||
&expected_sources,
|
||||
&all_bucket_names,
|
||||
scan_plan_digest,
|
||||
want_cycle,
|
||||
leader_epoch,
|
||||
)
|
||||
})
|
||||
.flatten();
|
||||
let structurally_complete_snapshot = result.is_ok() && completed_all_sets && completed_usage.is_some();
|
||||
let cycle_status = classify_nsscanner_cycle(
|
||||
structurally_complete_snapshot,
|
||||
@@ -381,6 +410,10 @@ impl ScannerIOCycle for ECStore {
|
||||
);
|
||||
if let Some((data_usage_info, _)) = completed_usage {
|
||||
publish_usage_snapshot(&updates, cycle_status, data_usage_info).await?;
|
||||
} else if !ctx.is_cancelled()
|
||||
&& let Some((data_usage_info, _)) = observational_usage
|
||||
{
|
||||
publish_observational_snapshot(&updates, data_usage_info).await?;
|
||||
}
|
||||
let dirty_usage_clear = should_clear_dirty_usage_snapshot(
|
||||
result.is_ok(),
|
||||
|
||||
@@ -105,6 +105,160 @@ fn completed_data_usage_info_for_test(
|
||||
completed_data_usage_info(results, &expected_sources, all_buckets, true, budget_elapsed, cancelled)
|
||||
}
|
||||
|
||||
fn lkg_root_cache(bucket: &str, objects: usize, source: DataUsageCacheSource) -> DataUsageCache {
|
||||
let mut cache = completed_root_cache(bucket, objects, 10, source);
|
||||
cache.info.snapshot_complete = false;
|
||||
cache.info.next_cycle = 8;
|
||||
cache.info.leader_epoch = 3;
|
||||
cache.info.lkg_snapshot_complete = true;
|
||||
cache.info.lkg_next_cycle = Some(7);
|
||||
cache.info.lkg_last_update = cache.info.last_update;
|
||||
cache.info.lkg_leader_epoch = Some(3);
|
||||
cache.info.lkg_scan_plan_digest = Some(TEST_PLAN_DIGEST);
|
||||
cache
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn partial_usage_is_observational_not_authoritative_for_quota() {
|
||||
let all_buckets = vec!["bucket".to_string()];
|
||||
let current_source = DataUsageCacheSource::new(0, 0);
|
||||
let stalled_source = DataUsageCacheSource::new(1, 0);
|
||||
let mut current = completed_root_cache("bucket", 2, 20, current_source);
|
||||
current.info.next_cycle = 8;
|
||||
current.info.leader_epoch = 3;
|
||||
let stalled = lkg_root_cache("bucket", 1, stalled_source);
|
||||
let expected = HashSet::from([current_source, stalled_source]);
|
||||
|
||||
assert!(
|
||||
completed_data_usage_info(&[current.clone(), stalled.clone()], &expected, &all_buckets, true, false, false).is_none()
|
||||
);
|
||||
let (observed, _) = observational_data_usage_info(&[current, stalled], &expected, &all_buckets, TEST_PLAN_DIGEST, 8, 3)
|
||||
.expect("a completed set should produce an observational view");
|
||||
assert!(observed.usage_snapshot_partial);
|
||||
assert!(!observed.usage_snapshot_complete);
|
||||
assert_eq!(observed.usage_snapshot_converged, Some(false));
|
||||
assert_eq!(observed.usage_snapshot_set_states.len(), 2);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn lkg_scope_does_not_count_as_current_cycle_completion() {
|
||||
let source = DataUsageCacheSource::new(0, 0);
|
||||
let mut lkg = lkg_root_cache("bucket", 1, source);
|
||||
lkg.info.last_update = None;
|
||||
let expected = HashSet::from([source]);
|
||||
assert!(!scanner_results_form_complete_snapshot(&[lkg], &expected));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn stale_quota_uses_complete_baseline_plus_positive_deltas() {
|
||||
let all_buckets = vec!["bucket".to_string()];
|
||||
let source = DataUsageCacheSource::new(0, 0);
|
||||
let mut current = completed_root_cache("bucket", 3, 20, source);
|
||||
current.info.next_cycle = 8;
|
||||
current.info.leader_epoch = 3;
|
||||
let expected = HashSet::from([source]);
|
||||
let (observed, _) = observational_data_usage_info(&[current], &expected, &all_buckets, TEST_PLAN_DIGEST, 8, 3)
|
||||
.expect("complete set data is a valid observational baseline");
|
||||
assert_eq!(observed.objects_total_size, 30);
|
||||
assert_eq!(observed.usage_snapshot_set_states[0].complete, true);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn negative_delta_waits_for_set_reconciliation() {
|
||||
let all_buckets = vec!["bucket".to_string()];
|
||||
let source = DataUsageCacheSource::new(0, 0);
|
||||
let mut stalled = lkg_root_cache("bucket", 4, source);
|
||||
stalled.info.lkg_scan_plan_digest = Some(DataUsageScanPlanDigest([9; 32]));
|
||||
let expected = HashSet::from([source]);
|
||||
assert!(observational_data_usage_info(&[stalled], &expected, &all_buckets, TEST_PLAN_DIGEST, 8, 3).is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn set_membership_add_remove_uses_generation_and_tombstone() {
|
||||
let state = DataUsageSnapshotSetState {
|
||||
pool_index: 1,
|
||||
set_index: 2,
|
||||
scanner_cycle: Some(9),
|
||||
scanner_epoch: Some(4),
|
||||
scan_plan_digest: Some(TEST_PLAN_DIGEST.0),
|
||||
complete: false,
|
||||
tombstone: true,
|
||||
};
|
||||
let encoded = serde_json::to_vec(&state).expect("set state should serialize");
|
||||
let decoded: DataUsageSnapshotSetState = serde_json::from_slice(&encoded).expect("set state should deserialize");
|
||||
assert_eq!(decoded, state);
|
||||
|
||||
let snapshot = DataUsageInfo {
|
||||
last_update: Some(SystemTime::UNIX_EPOCH + Duration::from_secs(10)),
|
||||
scanner_cycle: Some(9),
|
||||
scanner_epoch: Some(4),
|
||||
buckets_count: 0,
|
||||
usage_snapshot_converged: Some(false),
|
||||
usage_snapshot_partial: true,
|
||||
usage_snapshot_set_states: vec![
|
||||
DataUsageSnapshotSetState {
|
||||
pool_index: 0,
|
||||
set_index: 0,
|
||||
scanner_cycle: Some(9),
|
||||
scanner_epoch: Some(4),
|
||||
scan_plan_digest: Some(TEST_PLAN_DIGEST.0),
|
||||
complete: true,
|
||||
tombstone: false,
|
||||
},
|
||||
state,
|
||||
],
|
||||
..Default::default()
|
||||
};
|
||||
assert!(snapshot.is_valid_partial_snapshot());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn old_set_completion_cannot_overwrite_new_aggregate() {
|
||||
let all_buckets = vec!["bucket".to_string()];
|
||||
let source = DataUsageCacheSource::new(0, 0);
|
||||
let mut old = completed_root_cache("bucket", 1, 20, source);
|
||||
old.info.next_cycle = 7;
|
||||
old.info.leader_epoch = 2;
|
||||
let expected = HashSet::from([source]);
|
||||
assert!(observational_data_usage_info(&[old], &expected, &all_buckets, TEST_PLAN_DIGEST, 8, 3).is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn usage_aggregate_survives_restart_and_leader_failover() {
|
||||
let all_buckets = vec!["bucket".to_string()];
|
||||
let source = DataUsageCacheSource::new(0, 0);
|
||||
let mut lkg = lkg_root_cache("bucket", 5, source);
|
||||
lkg.info.lkg_leader_epoch = Some(4);
|
||||
lkg.info.lkg_next_cycle = Some(9);
|
||||
let expected = HashSet::from([source]);
|
||||
let (observed, _) = observational_data_usage_info(&[lkg], &expected, &all_buckets, TEST_PLAN_DIGEST, 10, 5)
|
||||
.expect("compatible LKG should survive a leader change");
|
||||
assert_eq!(observed.usage_snapshot_set_states[0].scanner_epoch, Some(4));
|
||||
assert_eq!(observed.objects_total_size, 50);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn usage_aggregate_cost_is_linear_in_set_count() {
|
||||
let all_buckets = vec!["bucket".to_string()];
|
||||
let mut results = Vec::new();
|
||||
let mut expected = HashSet::new();
|
||||
for index in 0..32 {
|
||||
let source = DataUsageCacheSource::new(index, 0);
|
||||
expected.insert(source);
|
||||
let mut cache = completed_root_cache("bucket", 1, 20, source);
|
||||
cache.info.next_cycle = 8;
|
||||
cache.info.leader_epoch = 3;
|
||||
results.push(cache);
|
||||
}
|
||||
let (observed, _) = observational_data_usage_info(&results, &expected, &all_buckets, TEST_PLAN_DIGEST, 8, 3)
|
||||
.expect("all set snapshots should aggregate");
|
||||
assert_eq!(observed.objects_total_count, 32);
|
||||
let reversed = results.iter().rev().cloned().collect::<Vec<_>>();
|
||||
let (reversed_observed, _) = observational_data_usage_info(&reversed, &expected, &all_buckets, TEST_PLAN_DIGEST, 8, 3)
|
||||
.expect("reordered set snapshots should aggregate");
|
||||
assert_eq!(observed.usage_snapshot_set_states, reversed_observed.usage_snapshot_set_states);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn completed_data_usage_info_publishes_tier_stats_across_sets() {
|
||||
let all_buckets = vec!["bucket-a".to_string(), "bucket-b".to_string()];
|
||||
|
||||
@@ -16,8 +16,7 @@
|
||||
|
||||
use super::storage_api::admin_usecase::admin::get_server_info;
|
||||
use super::storage_api::admin_usecase::capacity::{
|
||||
DecommissionUnresolvedEntry, PoolDecommissionInfo, PoolStatus, RebalStatus, get_total_usable_capacity,
|
||||
get_total_usable_capacity_free,
|
||||
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};
|
||||
@@ -108,8 +107,6 @@ 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)]
|
||||
@@ -622,7 +619,6 @@ impl DefaultAdminUsecase {
|
||||
bytes_done: info.bytes_done,
|
||||
bytes_failed: info.bytes_failed,
|
||||
waiting_reason,
|
||||
unresolved_entries: info.unresolved_entries,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -680,7 +676,7 @@ impl DefaultAdminUsecase {
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::super::storage_api::admin_usecase::capacity::{DecommissionUnresolvedEntry, PoolDecommissionInfo, PoolStatus};
|
||||
use super::super::storage_api::admin_usecase::capacity::{PoolDecommissionInfo, PoolStatus};
|
||||
use super::*;
|
||||
use time::OffsetDateTime;
|
||||
use tracing_subscriber::{Layer, Registry, layer::Context, prelude::*};
|
||||
@@ -991,17 +987,6 @@ 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()
|
||||
}),
|
||||
},
|
||||
@@ -1025,13 +1010,6 @@ 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");
|
||||
}
|
||||
|
||||
|
||||
@@ -31,7 +31,6 @@ 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;
|
||||
|
||||
@@ -398,7 +398,7 @@ pub(crate) mod ecstore_bucket {
|
||||
|
||||
pub(crate) mod ecstore_capacity {
|
||||
pub(crate) use rustfs_ecstore::api::capacity::{
|
||||
DecommissionUnresolvedEntry, PoolDecommissionInfo, PoolStatus, get_total_usable_capacity, get_total_usable_capacity_free,
|
||||
PoolDecommissionInfo, PoolStatus, get_total_usable_capacity, get_total_usable_capacity_free,
|
||||
is_reserved_or_invalid_bucket,
|
||||
};
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user