Compare commits

..

4 Commits

Author SHA1 Message Date
马登山 82fb0a8843 fix(scanner): correct observational usage arguments 2026-08-23 10:31:42 +08:00
马登山 307510749e style: format usage freshness changes 2026-08-23 10:28:47 +08:00
马登山 5496e14960 fix(ecstore): preserve quota baseline across restart 2026-08-23 10:20:39 +08:00
马登山 ec3b7a7dc6 fix(scanner): publish partial usage observations 2026-08-23 09:50:16 +08:00
15 changed files with 746 additions and 469 deletions
+108 -1
View File
@@ -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]
+2 -2
View File
@@ -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;
}
+66 -396
View File
@@ -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());
}
+109 -3
View File
@@ -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() {
+21 -3
View File
@@ -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()
}
}
+2 -3
View File
@@ -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,
}
}
+1 -1
View File
@@ -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,
+13 -2
View File
@@ -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,
+145 -2
View File
@@ -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,
+89 -30
View File
@@ -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",
+33
View File
@@ -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()];
+2 -24
View File
@@ -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");
}
-1
View File
@@ -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;
+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::{
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,
};
}