mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-25 13:36:50 +00:00
fix(scanner): account versioned delete markers in usage (#3904)
Co-authored-by: Henry Guo <marshawcoco@users.noreply.github.com>
This commit is contained in:
@@ -69,7 +69,8 @@ pub mod config {
|
||||
pub mod data_usage {
|
||||
pub use crate::data_usage::{
|
||||
DATA_USAGE_CACHE_NAME, apply_bucket_usage_memory_overlay, load_data_usage_from_backend,
|
||||
record_bucket_object_delete_memory, record_bucket_object_write_memory, remove_bucket_usage_from_backend,
|
||||
record_bucket_delete_marker_memory, record_bucket_object_delete_memory, record_bucket_object_version_write_memory,
|
||||
record_bucket_object_write_memory, refresh_versioned_bucket_usage_from_object_layer, remove_bucket_usage_from_backend,
|
||||
replace_bucket_usage_memory_from_info,
|
||||
};
|
||||
}
|
||||
|
||||
@@ -16,7 +16,7 @@ pub mod local_snapshot;
|
||||
|
||||
use crate::storage_api_contracts::{list::ListOperations as _, object::ObjectIO as _};
|
||||
use crate::{
|
||||
bucket::metadata_sys::get_replication_config,
|
||||
bucket::{metadata_sys::get_replication_config, versioning::VersioningApi as _, versioning_sys::BucketVersioningSys},
|
||||
config::com::read_config,
|
||||
disk::DiskAPI,
|
||||
error::{Error, classify_system_path_failure_reason},
|
||||
@@ -398,8 +398,9 @@ pub async fn aggregate_local_snapshots(store: Arc<ECStore>) -> Result<(Vec<DiskU
|
||||
|
||||
/// Calculate accurate bucket usage statistics by enumerating objects through the object layer.
|
||||
pub async fn compute_bucket_usage(store: Arc<ECStore>, bucket_name: &str) -> Result<BucketUsageInfo, Error> {
|
||||
let mut continuation: Option<String> = None;
|
||||
let mut objects_count: u64 = 0;
|
||||
let mut marker: Option<String> = None;
|
||||
let mut version_marker: Option<String> = None;
|
||||
let mut object_names: HashSet<String> = HashSet::new();
|
||||
let mut versions_count: u64 = 0;
|
||||
let mut total_size: u64 = 0;
|
||||
let mut delete_markers: u64 = 0;
|
||||
@@ -407,15 +408,13 @@ pub async fn compute_bucket_usage(store: Arc<ECStore>, bucket_name: &str) -> Res
|
||||
loop {
|
||||
let result = store
|
||||
.clone()
|
||||
.list_objects_v2(
|
||||
.list_object_versions(
|
||||
bucket_name,
|
||||
"", // prefix
|
||||
continuation.clone(),
|
||||
None, // delimiter
|
||||
1000, // max_keys
|
||||
false, // fetch_owner
|
||||
None, // start_after
|
||||
false, // incl_deleted
|
||||
marker.clone(),
|
||||
version_marker.clone(),
|
||||
None, // delimiter
|
||||
1000, // max_keys
|
||||
)
|
||||
.await?;
|
||||
|
||||
@@ -430,34 +429,27 @@ pub async fn compute_bucket_usage(store: Arc<ECStore>, bucket_name: &str) -> Res
|
||||
}
|
||||
|
||||
let object_size = object.size.max(0) as u64;
|
||||
objects_count = objects_count.saturating_add(1);
|
||||
object_names.insert(object.name.clone());
|
||||
total_size = total_size.saturating_add(object_size);
|
||||
|
||||
let detected_versions = if object.num_versions > 0 {
|
||||
object.num_versions as u64
|
||||
} else {
|
||||
1
|
||||
};
|
||||
versions_count = versions_count.saturating_add(detected_versions);
|
||||
versions_count = versions_count.saturating_add(1);
|
||||
}
|
||||
|
||||
if !result.is_truncated {
|
||||
break;
|
||||
}
|
||||
|
||||
continuation = result.next_continuation_token.clone();
|
||||
if continuation.is_none() {
|
||||
marker = result.next_marker.clone();
|
||||
version_marker = result.next_version_idmarker.clone();
|
||||
if marker.is_none() {
|
||||
info!(
|
||||
"Bucket {} listing marked truncated but no continuation token returned; stopping early",
|
||||
"Bucket {} version listing marked truncated but no marker returned; stopping early",
|
||||
bucket_name
|
||||
);
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
if versions_count == 0 {
|
||||
versions_count = objects_count;
|
||||
}
|
||||
let objects_count = object_names.len() as u64;
|
||||
|
||||
let usage = BucketUsageInfo {
|
||||
size: total_size,
|
||||
@@ -470,6 +462,42 @@ pub async fn compute_bucket_usage(store: Arc<ECStore>, bucket_name: &str) -> Res
|
||||
Ok(usage)
|
||||
}
|
||||
|
||||
pub async fn refresh_versioned_bucket_usage_from_object_layer(store: Arc<ECStore>, data_usage_info: &mut DataUsageInfo) {
|
||||
let buckets = data_usage_info.buckets_usage.keys().cloned().collect::<Vec<String>>();
|
||||
let mut changed = false;
|
||||
|
||||
for bucket in buckets {
|
||||
let Ok(versioning) = BucketVersioningSys::get(&bucket).await else {
|
||||
continue;
|
||||
};
|
||||
|
||||
if !versioning.enabled() && !versioning.suspended() {
|
||||
continue;
|
||||
}
|
||||
|
||||
let usage = match compute_bucket_usage(store.clone(), &bucket).await {
|
||||
Ok(usage) => usage,
|
||||
Err(err) => {
|
||||
debug!(
|
||||
bucket = %bucket,
|
||||
error = %err,
|
||||
"failed to refresh versioned bucket usage from object layer"
|
||||
);
|
||||
continue;
|
||||
}
|
||||
};
|
||||
|
||||
data_usage_info.bucket_sizes.insert(bucket.clone(), usage.size);
|
||||
data_usage_info.buckets_usage.insert(bucket, usage);
|
||||
changed = true;
|
||||
}
|
||||
|
||||
if changed {
|
||||
set_buckets_count_from_usage(data_usage_info);
|
||||
data_usage_info.calculate_totals();
|
||||
}
|
||||
}
|
||||
|
||||
async fn ensure_bucket_usage_cached(bucket: &str) {
|
||||
let cache = memory_cache().read().await;
|
||||
if cache.contains_key(bucket) {
|
||||
@@ -514,6 +542,20 @@ fn bucket_usage_counts_match(left: &BucketUsageInfo, right: &BucketUsageInfo) ->
|
||||
|
||||
/// Fast in-memory update for immediate quota and admin usage consistency.
|
||||
pub async fn record_bucket_object_write_memory(bucket: &str, previous_current_size: Option<u64>, new_size: u64) {
|
||||
record_bucket_object_write_memory_inner(bucket, previous_current_size, new_size, false).await;
|
||||
}
|
||||
|
||||
/// Fast in-memory update for versioned object writes.
|
||||
pub async fn record_bucket_object_version_write_memory(bucket: &str, previous_current_size: Option<u64>, new_size: u64) {
|
||||
record_bucket_object_write_memory_inner(bucket, previous_current_size, new_size, true).await;
|
||||
}
|
||||
|
||||
async fn record_bucket_object_write_memory_inner(
|
||||
bucket: &str,
|
||||
previous_current_size: Option<u64>,
|
||||
new_size: u64,
|
||||
creates_new_version: bool,
|
||||
) {
|
||||
ensure_bucket_usage_cached(bucket).await;
|
||||
|
||||
let mut cache = memory_cache().write().await;
|
||||
@@ -521,14 +563,22 @@ pub async fn record_bucket_object_write_memory(bucket: &str, previous_current_si
|
||||
.entry(bucket.to_string())
|
||||
.or_insert_with(|| cached_bucket_usage_now(BucketUsageInfo::default()));
|
||||
|
||||
match previous_current_size {
|
||||
Some(previous_size) => {
|
||||
entry.usage.size = entry.usage.size.saturating_sub(previous_size).saturating_add(new_size);
|
||||
}
|
||||
None => {
|
||||
entry.usage.size = entry.usage.size.saturating_add(new_size);
|
||||
if creates_new_version {
|
||||
entry.usage.size = entry.usage.size.saturating_add(new_size);
|
||||
if previous_current_size.is_none() {
|
||||
entry.usage.objects_count = entry.usage.objects_count.saturating_add(1);
|
||||
entry.usage.versions_count = entry.usage.versions_count.saturating_add(1);
|
||||
}
|
||||
entry.usage.versions_count = entry.usage.versions_count.saturating_add(1);
|
||||
} else {
|
||||
match previous_current_size {
|
||||
Some(previous_size) => {
|
||||
entry.usage.size = entry.usage.size.saturating_sub(previous_size).saturating_add(new_size);
|
||||
}
|
||||
None => {
|
||||
entry.usage.size = entry.usage.size.saturating_add(new_size);
|
||||
entry.usage.objects_count = entry.usage.objects_count.saturating_add(1);
|
||||
entry.usage.versions_count = entry.usage.versions_count.saturating_add(1);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -566,6 +616,24 @@ pub async fn record_bucket_object_delete_memory(bucket: &str, deleted_size: u64,
|
||||
entry.stale_snapshot_pending = false;
|
||||
}
|
||||
|
||||
/// Fast in-memory update for successful delete marker creation.
|
||||
pub async fn record_bucket_delete_marker_memory(bucket: &str) {
|
||||
ensure_bucket_usage_cached(bucket).await;
|
||||
|
||||
let mut cache = memory_cache().write().await;
|
||||
let entry = cache
|
||||
.entry(bucket.to_string())
|
||||
.or_insert_with(|| cached_bucket_usage_now(BucketUsageInfo::default()));
|
||||
|
||||
entry.usage.delete_markers_count = entry.usage.delete_markers_count.saturating_add(1);
|
||||
|
||||
let now = SystemTime::now();
|
||||
entry.refreshed_at = now;
|
||||
entry.usage_updated_at = now;
|
||||
entry.dirty = true;
|
||||
entry.stale_snapshot_pending = false;
|
||||
}
|
||||
|
||||
/// Fast in-memory decrement for immediate quota consistency
|
||||
pub async fn decrement_bucket_usage_memory(bucket: &str, size_decrement: u64) {
|
||||
record_bucket_object_delete_memory(bucket, size_decrement, size_decrement > 0).await;
|
||||
@@ -605,20 +673,12 @@ async fn update_usage_cache_if_needed() {
|
||||
*updating = true;
|
||||
drop(updating);
|
||||
|
||||
let cache_clone = (*memory_cache()).clone();
|
||||
let updating_clone = (*cache_updating()).clone();
|
||||
tokio::spawn(async move {
|
||||
if let Some(store) = runtime_sources::object_store_handle()
|
||||
&& let Ok(data_usage_info) = load_data_usage_from_backend(store.clone()).await
|
||||
{
|
||||
let mut cache = cache_clone.write().await;
|
||||
let usage_updated_at = data_usage_info_updated_at(&data_usage_info);
|
||||
for (bucket_name, bucket_usage) in data_usage_info.buckets_usage.iter() {
|
||||
cache.insert(
|
||||
bucket_name.clone(),
|
||||
cached_bucket_usage_from_backend(bucket_usage.clone(), usage_updated_at),
|
||||
);
|
||||
}
|
||||
replace_bucket_usage_memory_from_info(&data_usage_info).await;
|
||||
}
|
||||
let mut updating = updating_clone.write().await;
|
||||
*updating = false;
|
||||
@@ -642,14 +702,7 @@ async fn update_usage_cache_if_needed() {
|
||||
if let Some(store) = runtime_sources::object_store_handle()
|
||||
&& let Ok(data_usage_info) = load_data_usage_from_backend(store.clone()).await
|
||||
{
|
||||
let mut cache = memory_cache().write().await;
|
||||
let usage_updated_at = data_usage_info_updated_at(&data_usage_info);
|
||||
for (bucket_name, bucket_usage) in data_usage_info.buckets_usage.iter() {
|
||||
cache.insert(
|
||||
bucket_name.clone(),
|
||||
cached_bucket_usage_from_backend(bucket_usage.clone(), usage_updated_at),
|
||||
);
|
||||
}
|
||||
replace_bucket_usage_memory_from_info(&data_usage_info).await;
|
||||
}
|
||||
|
||||
let mut updating = cache_updating().write().await;
|
||||
@@ -1161,6 +1214,83 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn memory_overlay_counts_versioned_overwrite_as_new_version() {
|
||||
clear_usage_memory_cache_for_test().await;
|
||||
|
||||
let persisted = data_usage_info_for_test("bucket-a", 1, 10, SystemTime::now() - Duration::from_secs(10));
|
||||
replace_bucket_usage_memory_from_info(&persisted).await;
|
||||
record_bucket_object_version_write_memory("bucket-a", Some(10), 20).await;
|
||||
|
||||
let mut response = persisted.clone();
|
||||
apply_bucket_usage_memory_overlay(&mut response).await;
|
||||
|
||||
assert_eq!(response.objects_total_count, 1);
|
||||
assert_eq!(response.versions_total_count, 2);
|
||||
assert_eq!(response.objects_total_size, 30);
|
||||
assert_eq!(
|
||||
response
|
||||
.buckets_usage
|
||||
.get("bucket-a")
|
||||
.map(|usage| (usage.objects_count, usage.versions_count, usage.size)),
|
||||
Some((1, 2, 30))
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn memory_overlay_records_delete_marker_without_removing_versions() {
|
||||
clear_usage_memory_cache_for_test().await;
|
||||
|
||||
let persisted = data_usage_info_for_test("bucket-a", 2, 30, SystemTime::now() - Duration::from_secs(10));
|
||||
replace_bucket_usage_memory_from_info(&persisted).await;
|
||||
record_bucket_delete_marker_memory("bucket-a").await;
|
||||
|
||||
let mut response = persisted.clone();
|
||||
apply_bucket_usage_memory_overlay(&mut response).await;
|
||||
|
||||
assert_eq!(response.objects_total_count, 2);
|
||||
assert_eq!(response.versions_total_count, 2);
|
||||
assert_eq!(response.delete_markers_total_count, 1);
|
||||
assert_eq!(response.objects_total_size, 30);
|
||||
assert_eq!(
|
||||
response
|
||||
.buckets_usage
|
||||
.get("bucket-a")
|
||||
.map(|usage| { (usage.objects_count, usage.versions_count, usage.delete_markers_count, usage.size,) }),
|
||||
Some((2, 2, 1, 30))
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn scanner_sync_preserves_dirty_delete_marker_with_later_snapshot() {
|
||||
clear_usage_memory_cache_for_test().await;
|
||||
|
||||
let now = SystemTime::now();
|
||||
let old_persisted = data_usage_info_for_test("bucket-a", 2, 30, now - Duration::from_secs(10));
|
||||
replace_bucket_usage_memory_from_info(&old_persisted).await;
|
||||
record_bucket_delete_marker_memory("bucket-a").await;
|
||||
|
||||
let scanner_without_marker = data_usage_info_for_test("bucket-a", 2, 30, now + Duration::from_secs(10));
|
||||
replace_bucket_usage_memory_from_info(&scanner_without_marker).await;
|
||||
|
||||
let mut response = scanner_without_marker.clone();
|
||||
apply_bucket_usage_memory_overlay(&mut response).await;
|
||||
|
||||
assert_eq!(response.objects_total_count, 2);
|
||||
assert_eq!(response.versions_total_count, 2);
|
||||
assert_eq!(response.delete_markers_total_count, 1);
|
||||
assert_eq!(
|
||||
response
|
||||
.buckets_usage
|
||||
.get("bucket-a")
|
||||
.map(|usage| { (usage.objects_count, usage.versions_count, usage.delete_markers_count, usage.size,) }),
|
||||
Some((2, 2, 1, 30))
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn memory_overlay_does_not_replace_newer_persisted_usage() {
|
||||
|
||||
Reference in New Issue
Block a user