mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-24 13:16:28 +00:00
fix(scanner): publish partial usage observations
This commit is contained in:
@@ -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,142 @@ 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,
|
||||
|
||||
Reference in New Issue
Block a user