perf(delete): gate and parallelize DeleteObjects per-object stat fanout (#4398)

This commit is contained in:
houseme
2026-07-08 08:45:47 +08:00
committed by GitHub
parent eaff17cade
commit a413729b16
5 changed files with 814 additions and 85 deletions
@@ -0,0 +1,322 @@
// Copyright 2024 RustFS Team
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
//! Regression coverage for backlog#929 (HP-8): the DeleteObjects batch path
//! gates its two per-object metadata stat fanouts on the bucket configuration.
//! These tests run against a real 4-disk `ECStore` with the bucket metadata
//! sys initialized, so both gate branches are exercised with production
//! metadata resolution:
//!
//! - buckets created with Object Lock keep the held-lock stat and the #4297
//! delete protection (explicit-version deletes of retained objects are
//! rejected);
//! - buckets without Object Lock take the gated (stat-skipping) path and must
//! behave exactly as before: unversioned batch deletes remove objects and
//! report per-key results, versioned batch deletes still create delete
//! markers and preserve the underlying version.
use super::storage_api::test::bucket::metadata_sys;
use super::storage_api::test::contract::bucket::{BucketOperations, BucketOptions, MakeBucketOptions};
use super::storage_api::test::contract::object::{ObjectIO as _, ObjectOperations as _};
use super::storage_api::test::{ECStore, Endpoint, EndpointServerPools, Endpoints, PoolEndpoints};
use super::storage_api::test::{StorageObjectOptions as ObjectOptions, StoragePutObjReader as PutObjReader};
use crate::storage::storage_api::{StorageObjectLockDeleteOptions, StorageObjectToDelete as ObjectToDelete};
use serial_test::serial;
use std::path::PathBuf;
use std::sync::{Arc, OnceLock};
use tempfile::TempDir;
use tokio::fs;
use tokio_util::sync::CancellationToken;
use uuid::Uuid;
static DELETE_GATING_ENV: OnceLock<(Vec<PathBuf>, Arc<ECStore>, TempDir)> = OnceLock::new();
async fn setup_delete_gating_env() -> Arc<ECStore> {
if let Some((_paths, store, _)) = DELETE_GATING_ENV.get() {
return store.clone();
}
let temp_dir = TempDir::new().expect("create temp dir for delete gating test");
let temp_path = temp_dir.path().to_path_buf();
let disk_paths = vec![
temp_path.join("disk1"),
temp_path.join("disk2"),
temp_path.join("disk3"),
temp_path.join("disk4"),
];
for disk_path in &disk_paths {
fs::create_dir_all(disk_path).await.unwrap();
}
let mut endpoints = Vec::new();
for (i, disk_path) in disk_paths.iter().enumerate() {
let mut endpoint = Endpoint::try_from(disk_path.to_str().unwrap()).unwrap();
endpoint.set_pool_index(0);
endpoint.set_set_index(0);
endpoint.set_disk_index(i);
endpoints.push(endpoint);
}
let pool_endpoints = PoolEndpoints {
legacy: false,
set_count: 1,
drives_per_set: 4,
endpoints: Endpoints::from(endpoints),
cmd_line: "delete-objects-stat-gating-test".to_string(),
platform: format!("OS: {} | Arch: {}", std::env::consts::OS, std::env::consts::ARCH),
};
let endpoint_pools = EndpointServerPools(vec![pool_endpoints]);
super::storage_api::test::runtime::init_local_disks(endpoint_pools.clone())
.await
.unwrap();
let server_addr: std::net::SocketAddr = "127.0.0.1:0".parse().unwrap();
let ecstore = ECStore::new(server_addr, endpoint_pools, CancellationToken::new())
.await
.unwrap();
let buckets_list = ecstore
.list_bucket(&BucketOptions {
no_metadata: true,
..Default::default()
})
.await
.unwrap();
let buckets = buckets_list.into_iter().map(|v| v.name).collect();
metadata_sys::init_bucket_metadata_sys(ecstore.clone(), buckets).await;
let _ = DELETE_GATING_ENV.set((disk_paths, ecstore.clone(), temp_dir));
ecstore
}
fn compliance_retention_metadata() -> std::collections::HashMap<String, String> {
let retain_until = time::OffsetDateTime::now_utc() + time::Duration::days(30);
let mut user_defined = std::collections::HashMap::new();
user_defined.insert("x-amz-object-lock-mode".to_string(), "COMPLIANCE".to_string());
user_defined.insert(
"x-amz-object-lock-retain-until-date".to_string(),
retain_until
.format(&time::format_description::well_known::Rfc3339)
.expect("retain-until date should format"),
);
user_defined
}
#[tokio::test]
#[serial]
async fn object_lock_bucket_batch_delete_keeps_held_lock_protection() {
let ecstore = setup_delete_gating_env().await;
let bucket = format!("hp8-lock-{}", Uuid::new_v4());
ecstore
.make_bucket(
&bucket,
&MakeBucketOptions {
lock_enabled: true,
..Default::default()
},
)
.await
.expect("create object-lock bucket");
let mut reader = PutObjReader::from_vec(b"retained payload".to_vec());
let put_info = ecstore
.put_object(
&bucket,
"retained.bin",
&mut reader,
&ObjectOptions {
versioned: true,
user_defined: compliance_retention_metadata(),
..Default::default()
},
)
.await
.expect("put retained object");
let version_id = put_info.version_id.expect("lock bucket writes must be versioned");
let (_deleted, errs) = ecstore
.delete_objects(
&bucket,
vec![ObjectToDelete {
object_name: "retained.bin".to_string(),
version_id: Some(version_id),
..Default::default()
}],
ObjectOptions {
versioned: true,
object_lock_delete: Some(StorageObjectLockDeleteOptions {
bypass_governance: false,
}),
..Default::default()
},
)
.await;
assert!(
errs[0].is_some(),
"explicit-version delete of a COMPLIANCE-retained object must be rejected on lock buckets"
);
ecstore
.get_object_info(
&bucket,
"retained.bin",
&ObjectOptions {
version_id: Some(version_id.to_string()),
versioned: true,
..Default::default()
},
)
.await
.expect("retained version must survive the batch delete");
}
#[tokio::test]
#[serial]
async fn non_lock_versioned_bucket_batch_delete_still_creates_delete_marker() {
let ecstore = setup_delete_gating_env().await;
let bucket = format!("hp8-versioned-{}", Uuid::new_v4());
ecstore
.make_bucket(
&bucket,
&MakeBucketOptions {
versioning_enabled: true,
..Default::default()
},
)
.await
.expect("create versioned bucket");
let mut reader = PutObjReader::from_vec(b"versioned payload".to_vec());
let put_info = ecstore
.put_object(
&bucket,
"versioned.bin",
&mut reader,
&ObjectOptions {
versioned: true,
..Default::default()
},
)
.await
.expect("put versioned object");
let version_id = put_info.version_id.expect("versioned write must return a version id");
// No explicit version id: this is the delete-marker-creating shape that
// skips both stat fanouts on a non-lock, non-replicated bucket.
let (deleted, errs) = ecstore
.delete_objects(
&bucket,
vec![ObjectToDelete {
object_name: "versioned.bin".to_string(),
..Default::default()
}],
ObjectOptions {
versioned: true,
object_lock_delete: Some(StorageObjectLockDeleteOptions {
bypass_governance: false,
}),
..Default::default()
},
)
.await;
assert!(errs[0].is_none(), "delete-marker creation must succeed: {:?}", errs[0]);
assert!(
deleted[0].delete_marker,
"versioned delete without version id must create a delete marker"
);
assert!(
deleted[0].delete_marker_version_id.is_some(),
"delete marker must carry its own version id"
);
ecstore
.get_object_info(
&bucket,
"versioned.bin",
&ObjectOptions {
version_id: Some(version_id.to_string()),
versioned: true,
..Default::default()
},
)
.await
.expect("original version must survive delete-marker creation");
}
#[tokio::test]
#[serial]
async fn non_lock_unversioned_bucket_batch_delete_reports_per_key_results() {
let ecstore = setup_delete_gating_env().await;
let bucket = format!("hp8-plain-{}", Uuid::new_v4());
ecstore
.make_bucket(&bucket, &MakeBucketOptions::default())
.await
.expect("create plain bucket");
for object in ["keep-a.bin", "keep-b.bin"] {
let mut reader = PutObjReader::from_vec(b"plain payload".to_vec());
ecstore
.put_object(&bucket, object, &mut reader, &ObjectOptions::default())
.await
.expect("put plain object");
}
let (deleted, errs) = ecstore
.delete_objects(
&bucket,
vec![
ObjectToDelete {
object_name: "keep-a.bin".to_string(),
..Default::default()
},
ObjectToDelete {
object_name: "missing.bin".to_string(),
..Default::default()
},
ObjectToDelete {
object_name: "keep-b.bin".to_string(),
..Default::default()
},
],
ObjectOptions {
object_lock_delete: Some(StorageObjectLockDeleteOptions {
bypass_governance: false,
}),
..Default::default()
},
)
.await;
assert!(
errs.iter().all(Option::is_none),
"batch delete on the gated (stat-skipping) path must keep S3 per-key semantics: {errs:?}"
);
assert_eq!(deleted[0].object_name, "keep-a.bin");
assert_eq!(deleted[1].object_name, "missing.bin");
assert_eq!(deleted[2].object_name, "keep-b.bin");
for object in ["keep-a.bin", "keep-b.bin"] {
ecstore
.get_object_info(&bucket, object, &ObjectOptions::default())
.await
.expect_err("deleted object must be gone");
}
}
+2
View File
@@ -29,4 +29,6 @@ pub(crate) mod storage_api;
#[cfg(test)]
mod capacity_dirty_scope_test;
#[cfg(test)]
mod delete_objects_stat_gating_test;
#[cfg(test)]
mod lifecycle_transition_api_test;
+247 -52
View File
@@ -2160,6 +2160,39 @@ fn delete_creates_delete_marker(opts: &ObjectOptions) -> bool {
opts.version_id.is_none() && opts.versioned && !opts.version_suspended
}
/// Bounded concurrency for the per-object pre-delete stat fanout in
/// `execute_delete_objects` (backlog#929 / HP-8). Keeps the metadata reads for
/// a 1000-key batch from serializing while capping the disk fanout pressure.
const DELETE_OBJECTS_PRE_STAT_CONCURRENCY: usize = 16;
/// backlog#929 (HP-8): whether the pre-delete `get_object_info` for one entry
/// of a DeleteObjects batch can be skipped without changing behavior.
///
/// The stat result feeds four consumers, and each must be provably idle:
/// - the app-layer object-lock admission check never runs for deletes that
/// create a delete marker, and non-lock buckets cannot hold retention or
/// legal-hold metadata (`bucket_lock_enabled == false`);
/// - the replication delete decision is only consulted when the bucket has
/// active replication rules for the batch (`replicate_deletes == false`);
/// - usage accounting for delete-marker creation goes through
/// `record_bucket_delete_marker_memory` and never reads the object size
/// (`accounting_creates_delete_marker` is computed from the same versioning
/// snapshot the accounting branch uses);
/// - transitioned-object (ILM tier) cleanup journaling is a no-op for
/// delete-marker creation because no version is removed, so `ObjSweeper`
/// produces no journal entry regardless of the stat result.
///
/// Object-lock enabled buckets always keep the stat, so their delete path is
/// byte-for-byte the pre-#929 one (see PR #4297).
fn can_skip_delete_objects_pre_stat(
bucket_lock_enabled: bool,
replicate_deletes: bool,
opts: &ObjectOptions,
accounting_creates_delete_marker: bool,
) -> bool {
!bucket_lock_enabled && !replicate_deletes && delete_creates_delete_marker(opts) && accounting_creates_delete_marker
}
fn resolve_put_object_extract_options(headers: &HeaderMap) -> S3Result<PutObjectExtractOptions> {
let prefix = snowball_meta_value(headers, SNOWBALL_PREFIX_HEADER_KEYS, SNOWBALL_PREFIX_SUFFIX_LOWER)
.map(|value| normalize_snowball_prefix(&value))
@@ -4858,6 +4891,7 @@ impl DefaultObjectUsecase {
let version_cfg = BucketVersioningSys::get(&bucket).await.unwrap_or_default();
let bypass_governance = has_bypass_governance_header(&req.headers);
let bucket_lock_enabled = bucket_object_locking_enabled(&bucket).await;
#[derive(Default, Clone)]
struct DeleteResult {
@@ -4867,10 +4901,18 @@ impl DefaultObjectUsecase {
let mut delete_results = vec![DeleteResult::default(); delete.objects.len()];
let mut object_to_delete = Vec::new();
let mut object_to_delete_idx = Vec::new();
let mut object_sizes = Vec::new();
let mut existing_object_infos = Vec::new();
struct PreparedDelete {
idx: usize,
object: ObjectToDelete,
opts: ObjectOptions,
version_id: Option<String>,
skip_stat: bool,
}
// Phase 1 (serial): request-scoped validation and authorization. These
// steps mutate the request info between authorization calls, so they
// stay sequential; they perform no per-object disk I/O.
let mut prepared_deletes: Vec<PreparedDelete> = Vec::with_capacity(delete.objects.len());
for (idx, obj_id) in delete.objects.iter().enumerate() {
let raw_version_id = obj_id.version_id.clone();
let (version_id, version_uuid) = match normalize_delete_objects_version_id(raw_version_id.clone()) {
@@ -4927,7 +4969,7 @@ impl DefaultObjectUsecase {
continue;
}
let mut object = ObjectToDelete {
let object = ObjectToDelete {
object_name: obj_id.key.clone(),
version_id: version_uuid,
..Default::default()
@@ -4944,57 +4986,140 @@ impl DefaultObjectUsecase {
.await
.map_err(ApiError::from)?;
let (goi, gerr) = match store.get_object_info(&bucket, &object.object_name, &opts).await {
Ok(res) => (res, None),
Err(e) => (ObjectInfo::default(), Some(e.to_string())),
};
// backlog#929 (HP-8): the accounting branch after the store delete
// decides delete-marker vs object-delete from this exact snapshot,
// so evaluate it here with the same inputs to keep the stat-skip
// decision and the accounting path provably consistent.
let accounting_creates_delete_marker = object.version_id.is_none()
&& version_cfg.prefix_enabled(object.object_name.as_str())
&& !version_cfg.suspended();
let skip_stat =
can_skip_delete_objects_pre_stat(bucket_lock_enabled, replicate_deletes, &opts, accounting_creates_delete_marker);
if gerr.is_none()
&& !delete_creates_delete_marker(&opts)
&& let Some(block_reason) = check_object_lock_for_deletion(&bucket, &goi, bypass_governance).await
{
delete_results[idx].error = Some(s3s::dto::Error {
code: Some("AccessDenied".to_string()),
key: Some(obj_id.key.clone()),
message: Some(block_reason.error_message()),
version_id: version_id.clone(),
});
prepared_deletes.push(PreparedDelete {
idx,
object,
opts,
version_id,
skip_stat,
});
}
struct AdmittedDelete {
idx: usize,
object: ObjectToDelete,
size: i64,
existing: Option<ObjectInfo>,
blocked: Option<s3s::dto::Error>,
}
// Phase 2 (bounded concurrency, backlog#929 / HP-8): the per-object
// pre-delete stat plus the admission checks that consume it. Entries
// are independent per key, and `buffered` preserves input order so the
// per-key result mapping below is identical to the previous serial
// loop. The authoritative object-lock enforcement stays in the
// set_disk layer under the held write lock (#4297); the check here is
// the same early, advisory rejection as before.
let store_ref = &store;
let bucket_ref = bucket.as_str();
let admitted_deletes: Vec<AdmittedDelete> =
futures::stream::iter(prepared_deletes.into_iter().map(|prepared| async move {
let PreparedDelete {
idx,
mut object,
opts,
version_id,
skip_stat,
} = prepared;
let (goi, gerr) = if skip_stat {
(ObjectInfo::default(), None)
} else {
match store_ref.get_object_info(bucket_ref, &object.object_name, &opts).await {
Ok(res) => (res, None),
Err(e) => (ObjectInfo::default(), Some(e.to_string())),
}
};
if !skip_stat
&& gerr.is_none()
&& !delete_creates_delete_marker(&opts)
&& let Some(block_reason) = check_object_lock_for_deletion(bucket_ref, &goi, bypass_governance).await
{
let blocked_key = object.object_name.clone();
return AdmittedDelete {
idx,
object,
size: 0,
existing: None,
blocked: Some(s3s::dto::Error {
code: Some("AccessDenied".to_string()),
key: Some(blocked_key),
message: Some(block_reason.error_message()),
version_id,
}),
};
}
let size = goi.size;
if is_dir_object(&object.object_name) && object.version_id.is_none() {
object.version_id = Some(Uuid::nil());
}
if replicate_deletes {
let dsc = check_replicate_delete(
bucket_ref,
&ObjectToDelete {
object_name: object.object_name.clone(),
version_id: object.version_id,
..Default::default()
},
&goi,
&opts,
gerr.clone(),
)
.await;
if dsc.replicate_any() {
if object.version_id.is_some() {
set_object_to_delete_version_purge_status(&mut object, VersionPurgeStatusType::Pending);
object.version_purge_statuses = dsc.pending_status();
} else {
object.delete_marker_replication_status = dsc.pending_status();
}
object.replicate_decision_str = Some(dsc.to_string());
}
}
let existing = (!skip_stat && gerr.is_none()).then_some(goi);
AdmittedDelete {
idx,
object,
size,
existing,
blocked: None,
}
}))
.buffered(DELETE_OBJECTS_PRE_STAT_CONCURRENCY)
.collect()
.await;
// Phase 3 (serial): apply outcomes in the original request order so
// per-key success/failure reporting is unchanged.
let mut object_to_delete = Vec::new();
let mut object_to_delete_idx = Vec::new();
let mut object_sizes = Vec::new();
let mut existing_object_infos = Vec::new();
for admitted in admitted_deletes {
if let Some(err) = admitted.blocked {
delete_results[admitted.idx].error = Some(err);
continue;
}
object_sizes.push(goi.size);
if is_dir_object(&object.object_name) && object.version_id.is_none() {
object.version_id = Some(Uuid::nil());
}
if replicate_deletes {
let dsc = check_replicate_delete(
&bucket,
&ObjectToDelete {
object_name: object.object_name.clone(),
version_id: object.version_id,
..Default::default()
},
&goi,
&opts,
gerr.clone(),
)
.await;
if dsc.replicate_any() {
if object.version_id.is_some() {
set_object_to_delete_version_purge_status(&mut object, VersionPurgeStatusType::Pending);
object.version_purge_statuses = dsc.pending_status();
} else {
object.delete_marker_replication_status = dsc.pending_status();
}
object.replicate_decision_str = Some(dsc.to_string());
}
}
object_to_delete_idx.push(idx);
object_to_delete.push(object);
existing_object_infos.push(gerr.is_none().then_some(goi));
object_sizes.push(admitted.size);
object_to_delete_idx.push(admitted.idx);
object_to_delete.push(admitted.object);
existing_object_infos.push(admitted.existing);
}
let cache_adapter = self.object_data_cache();
@@ -8621,6 +8746,76 @@ mod tests {
assert_eq!(internal_version_id, Some(Uuid::nil()));
}
// backlog#929 (HP-8): the pre-delete stat may only be skipped when every
// consumer of its result is provably idle. Each guard flips one condition
// to prove the skip is fenced on all four data dependencies.
fn delete_marker_creating_opts() -> ObjectOptions {
ObjectOptions {
version_id: None,
versioned: true,
version_suspended: false,
..Default::default()
}
}
#[test]
fn delete_objects_pre_stat_skippable_for_delete_marker_on_plain_bucket() {
assert!(can_skip_delete_objects_pre_stat(false, false, &delete_marker_creating_opts(), true));
}
#[test]
fn delete_objects_pre_stat_kept_for_object_lock_buckets() {
assert!(!can_skip_delete_objects_pre_stat(true, false, &delete_marker_creating_opts(), true));
}
#[test]
fn delete_objects_pre_stat_kept_when_replication_rules_match() {
assert!(!can_skip_delete_objects_pre_stat(false, true, &delete_marker_creating_opts(), true));
}
#[test]
fn delete_objects_pre_stat_kept_for_explicit_version_deletes() {
let opts = ObjectOptions {
version_id: Some(Uuid::new_v4().to_string()),
versioned: true,
version_suspended: false,
..Default::default()
};
assert!(!can_skip_delete_objects_pre_stat(false, false, &opts, true));
}
#[test]
fn delete_objects_pre_stat_kept_for_unversioned_buckets() {
// Unversioned deletes remove the current object: usage accounting needs
// the object size and ILM tier cleanup needs the transition metadata.
let opts = ObjectOptions {
version_id: None,
versioned: false,
version_suspended: false,
..Default::default()
};
assert!(!can_skip_delete_objects_pre_stat(false, false, &opts, false));
}
#[test]
fn delete_objects_pre_stat_kept_for_suspended_versioning() {
let opts = ObjectOptions {
version_id: None,
versioned: true,
version_suspended: true,
..Default::default()
};
assert!(!can_skip_delete_objects_pre_stat(false, false, &opts, false));
}
#[test]
fn delete_objects_pre_stat_kept_when_accounting_snapshot_disagrees() {
// If the accounting-side versioning snapshot does not also classify the
// delete as a delete-marker creation, the stat must stay so usage
// accounting keeps its size input.
assert!(!can_skip_delete_objects_pre_stat(false, false, &delete_marker_creating_opts(), false));
}
#[test]
fn should_schedule_delete_replication_skips_replica_requests() {
let opts = ObjectOptions {