fix(quota): enforce durable hard quota reservations (#6058)

* fix(quota): enforce durable hard quota reservations

* fix(quota): close reservation bypasses

* fix(quota): isolate tests and box object futures

* fix(quota): close legacy and deferred settlement bypasses

* fix(app): keep object futures off caller stacks

* fix(metrics): preserve object operation labels

* fix(logging): retain GET trace guard contract
This commit is contained in:
cxymds
2026-08-14 14:26:00 +08:00
committed by GitHub
parent 307f50ee1b
commit d60a77b750
35 changed files with 4183 additions and 425 deletions
+72 -20
View File
@@ -465,6 +465,16 @@ impl Operation for ImportBucketMetadata {
file_contents.push((file_path, content));
}
let durable_quota_import = imported_quota_requires_fleet_proof(&file_contents)?;
let quota_fleet_proof =
if durable_quota_import {
Some(crate::admin::storage_api::acquire_cross_pool_fence_fleet_proof().ok_or_else(|| {
s3_error!(ServiceUnavailable, "durable quota capability is not confirmed across the cluster")
})?)
} else {
None
};
// Extract bucket names
let mut bucket_names = Vec::new();
for (file_path, _) in &file_contents {
@@ -707,21 +717,6 @@ impl Operation for ImportBucketMetadata {
}
BUCKET_QUOTA_CONFIG_FILE => {
if let Err(e) = serde_json::from_slice::<BucketQuota>(&content) {
warn!(
event = EVENT_ADMIN_BUCKET_META_STATE,
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_BUCKET_META,
action = "import_bucket_metadata",
result = "config_deserialize_failed",
bucket = %bucket_name,
config_name = %conf_name,
error = %e,
"admin bucket meta state"
);
continue;
}
let metadata = match bucket_metadatas.get_mut(bucket_name) {
Some(m) => m,
None => continue,
@@ -830,10 +825,6 @@ impl Operation for ImportBucketMetadata {
}
}
// Persist the assembled metadata to disk. Prior to this, the import only mutated the
// in-memory `bucket_metadatas` map and returned 200, silently dropping every imported
// config. `metadata_sys::update` loads the on-disk metadata, overwrites the given config
// field and saves it, preserving any configs not present in the import archive.
for (bucket_name, metadata) in &bucket_metadatas {
for (config_file, data) in imported_configs_to_persist(metadata) {
let site_replication_item = imported_config_to_site_replication_item(bucket_name, metadata, config_file, &data)?;
@@ -842,7 +833,21 @@ impl Operation for ImportBucketMetadata {
} else {
None
};
if let Err(e) = metadata_sys::update(bucket_name, config_file, data).await {
let persist_result = if config_file == BUCKET_QUOTA_CONFIG_FILE {
let quota: BucketQuota =
serde_json::from_slice(&data).map_err(|e| s3_error!(InvalidRequest, "invalid bucket quota: {e}"))?;
if quota.uses_durable_reservations() {
let proof = quota_fleet_proof.as_ref().ok_or_else(|| {
s3_error!(ServiceUnavailable, "durable quota capability is not confirmed across the cluster")
})?;
metadata_sys::update_quota_if_incarnation(bucket_name, data, metadata.bucket_incarnation_id, proof).await
} else {
metadata_sys::update_if_incarnation(bucket_name, config_file, data, metadata.bucket_incarnation_id).await
}
} else {
metadata_sys::update_if_incarnation(bucket_name, config_file, data, metadata.bucket_incarnation_id).await
};
if let Err(e) = persist_result {
warn!(
event = EVENT_ADMIN_BUCKET_META_STATE,
component = LOG_COMPONENT_ADMIN,
@@ -885,6 +890,26 @@ impl Operation for ImportBucketMetadata {
}
}
fn imported_quota_requires_fleet_proof(file_contents: &[(String, Vec<u8>)]) -> S3Result<bool> {
let mut durable = false;
for (file_path, content) in file_contents {
let mut parts = file_path.split(SLASH_SEPARATOR);
let Some(_bucket) = parts.next() else {
continue;
};
if parts.next() != Some(BUCKET_QUOTA_CONFIG_FILE) {
continue;
}
let quota: BucketQuota =
serde_json::from_slice(content).map_err(|e| s3_error!(InvalidRequest, "invalid bucket quota: {e}"))?;
if quota.has_unsupported_reservation_protocol() {
return Err(s3_error!(InvalidRequest, "unsupported bucket quota reservation protocol"));
}
durable |= quota.uses_durable_reservations();
}
Ok(durable)
}
/// The `(config_file, data)` pairs to persist for an imported bucket's metadata: every non-empty
/// config field keyed by its on-disk config-file name, as owned data ready for
/// `metadata_sys::update`. Empty fields are skipped so an import never overwrites an existing
@@ -1076,4 +1101,31 @@ mod import_persist_tests {
assert!(!has_site_replication_item);
}
#[test]
fn quota_import_preflight_rejects_invalid_and_unknown_protocols() {
let missing_limit = vec![(
format!("bucket/{BUCKET_QUOTA_CONFIG_FILE}"),
br#"{"quota":0,"reservation_protocol":1}"#.to_vec(),
)];
assert!(imported_quota_requires_fleet_proof(&missing_limit).is_err());
let unknown_protocol = vec![(
format!("bucket/{BUCKET_QUOTA_CONFIG_FILE}"),
br#"{"quota":0,"reservation_protocol":2,"reservation_quota":1024}"#.to_vec(),
)];
assert!(imported_quota_requires_fleet_proof(&unknown_protocol).is_err());
}
#[test]
fn quota_import_preflight_requires_proof_only_for_durable_quota() {
let legacy = vec![(format!("bucket/{BUCKET_QUOTA_CONFIG_FILE}"), br#"{"quota":1024}"#.to_vec())];
assert!(!imported_quota_requires_fleet_proof(&legacy).expect("legacy quota should remain compatible"));
let durable = vec![(
format!("bucket/{BUCKET_QUOTA_CONFIG_FILE}"),
serde_json::to_vec(&BucketQuota::new(Some(1024))).expect("durable quota should encode"),
)];
assert!(imported_quota_requires_fleet_proof(&durable).expect("durable quota should pass preflight"));
}
}
+25 -4
View File
@@ -22,6 +22,7 @@ use crate::admin::storage_api::bucket::metadata_sys::{self, BucketMetadataSys};
use crate::admin::storage_api::bucket::quota::checker::QuotaChecker;
use crate::admin::storage_api::bucket::quota::{BucketQuota, QuotaError, QuotaOperation};
use crate::auth::{check_key_valid, get_session_token};
use crate::error::ApiError;
use crate::server::ADMIN_PREFIX;
use hyper::{Method, StatusCode};
use matchit::Params;
@@ -293,16 +294,36 @@ impl Operation for SetBucketQuotaHandler {
return Err(s3_error!(InvalidArgument, "{}", rustfs_config::QUOTA_INVALID_TYPE_ERROR_MSG));
}
let fleet_proof = if request.quota.is_some() {
Some(crate::admin::storage_api::acquire_cross_pool_fence_fleet_proof().ok_or_else(|| {
S3Error::with_message(
s3s::S3ErrorCode::ServiceUnavailable,
"durable quota capability is not confirmed across the cluster".to_string(),
)
})?)
} else {
None
};
let quota = BucketQuota::new(request.quota);
let metadata_sys_lock = bucket_metadata_from_context()
.ok_or_else(|| s3_error!(InternalError, "{}", rustfs_config::QUOTA_METADATA_SYSTEM_ERROR_MSG))?;
let mut quota_checker = QuotaChecker::new(metadata_sys_lock.clone());
let updated_at = quota_checker
.set_quota_config_if_incarnation(&bucket, quota.clone(), expected_incarnation_id)
.await
.map_err(|e| s3_error!(InternalError, "failed to set quota: {}", e))?;
let updated_at = match fleet_proof.as_ref() {
Some(fleet_proof) => {
quota_checker
.set_durable_quota_config_if_incarnation(&bucket, quota.clone(), expected_incarnation_id, fleet_proof)
.await
}
None => {
quota_checker
.set_quota_config_if_incarnation(&bucket, quota.clone(), expected_incarnation_id)
.await
}
}
.map_err(ApiError::from)?;
if let Err(err) = site_replication_bucket_meta_hook(SRBucketMeta {
bucket: bucket.clone(),
+30 -3
View File
@@ -29,6 +29,7 @@ use crate::admin::storage_api::bucket::metadata::{
BUCKET_SSECONFIG, BUCKET_TAGGING_CONFIG, BUCKET_TARGETS_FILE, BUCKET_VERSIONING_CONFIG, OBJECT_LOCK_CONFIG,
};
use crate::admin::storage_api::bucket::metadata_sys;
use crate::admin::storage_api::bucket::quota::BucketQuota;
use crate::admin::storage_api::bucket::replication;
use crate::admin::storage_api::bucket::target::{ARN, BucketTarget, BucketTargetType, BucketTargets, Credentials};
use crate::admin::storage_api::bucket::target_sys::BucketTargetSys;
@@ -7895,9 +7896,35 @@ async fn apply_bucket_meta_item(item: SRBucketMeta) -> S3Result<()> {
if !skip_config_write {
if let Some(data) = data {
metadata_sys::update_if_incarnation(&item.bucket, config_file, data, expected_incarnation_id)
.await
.map_err(ApiError::from)?;
if item.r#type == "quota-config" {
let quota: BucketQuota = serde_json::from_slice(&data)
.map_err(|e| S3Error::with_message(S3ErrorCode::InvalidRequest, format!("invalid bucket quota: {e}")))?;
if quota.has_unsupported_reservation_protocol() {
return Err(S3Error::with_message(
S3ErrorCode::InvalidRequest,
"unsupported bucket quota reservation protocol".to_string(),
));
}
if quota.uses_durable_reservations() {
let proof = crate::admin::storage_api::acquire_cross_pool_fence_fleet_proof().ok_or_else(|| {
S3Error::with_message(
S3ErrorCode::ServiceUnavailable,
"durable quota capability is not confirmed across the cluster".to_string(),
)
})?;
metadata_sys::update_quota_if_incarnation(&item.bucket, data, expected_incarnation_id, &proof)
.await
.map_err(ApiError::from)?;
} else {
metadata_sys::update_if_incarnation(&item.bucket, config_file, data, expected_incarnation_id)
.await
.map_err(ApiError::from)?;
}
} else {
metadata_sys::update_if_incarnation(&item.bucket, config_file, data, expected_incarnation_id)
.await
.map_err(ApiError::from)?;
}
} else {
metadata_sys::delete_if_incarnation(&item.bucket, config_file, expected_incarnation_id)
.await
+1 -1
View File
@@ -2901,7 +2901,7 @@ async fn handle_misc_extension_request(req: &mut S3Request<Body>, route: &MiscEx
MiscExtRoute::ObjectLambda { bucket, object } => {
let get_req = build_object_lambda_get_request(req, bucket, object)?;
let usecase = default_object_usecase();
let get_resp = Box::pin(usecase.execute_get_object(get_req)).await?;
let get_resp = usecase.execute_get_object(get_req).await?;
invoke_object_lambda_target(req, bucket, object, get_resp).await
}
MiscExtRoute::ListenNotification { bucket } => {
+16 -1
View File
@@ -64,7 +64,9 @@ mod ecstore_metrics {
}
mod ecstore_notification {
pub(crate) use crate::storage::storage_api::ecstore_notification::NotificationSys;
pub(crate) use crate::storage::storage_api::ecstore_notification::{
CrossPoolFenceFleetProofToken, NotificationSys, acquire_cross_pool_fence_fleet_proof,
};
}
#[allow(unused_imports)]
@@ -112,6 +114,10 @@ pub(crate) type TierCreds = ecstore_tier::tier_admin::TierCreds;
pub(crate) type TierType = ecstore_tier::tier_config::TierType;
pub(crate) type TierConfigUpdateError = crate::storage::storage_api::TierConfigUpdateError;
pub(crate) fn acquire_cross_pool_fence_fleet_proof() -> Option<ecstore_notification::CrossPoolFenceFleetProofToken> {
ecstore_notification::acquire_cross_pool_fence_fleet_proof()
}
pub(crate) mod runtime_sources {
pub(crate) type DailyAllTierStats = super::DailyAllTierStats;
pub(crate) type ECStore = super::ECStore;
@@ -296,6 +302,15 @@ pub(crate) mod metadata_sys {
super::ecstore_bucket::metadata_sys::update_if_incarnation(bucket, config_file, data, expected_incarnation_id).await
}
pub(crate) async fn update_quota_if_incarnation(
bucket: &str,
data: Vec<u8>,
expected_incarnation_id: uuid::Uuid,
proof: &super::ecstore_notification::CrossPoolFenceFleetProofToken,
) -> Result<OffsetDateTime> {
super::ecstore_bucket::metadata_sys::update_quota_if_incarnation(bucket, data, expected_incarnation_id, proof).await
}
pub(crate) async fn capture_bucket_metadata_incarnation(bucket: &str) -> Result<uuid::Uuid> {
super::ecstore_bucket::metadata_sys::capture_bucket_metadata_incarnation(bucket).await
}
+32
View File
@@ -23,6 +23,9 @@
//! `ECStore` and one metadata-sys initialization exist per test binary.
use super::storage_api::test::bucket::metadata_sys;
use super::storage_api::test::bucket::quota::BucketQuota;
use super::storage_api::test::bucket::quota::checker::QuotaChecker;
use super::storage_api::test::contract::bucket::MakeBucketOptions;
use super::storage_api::test::contract::bucket::{BucketOperations, BucketOptions};
use super::storage_api::test::{ECStore, Endpoint, EndpointServerPools, Endpoints, PoolEndpoints};
use super::{context::AppContext, object_traffic_health::ObjectTrafficHealth};
@@ -87,6 +90,17 @@ pub(crate) async fn shared_gating_ecstore() -> Arc<ECStore> {
crate::storage::storage_api::new_global_notification_sys(endpoint_pools.clone())
.await
.expect("initialize notification system for gating test env");
let topology_fingerprint =
crate::storage::storage_api::heal_control_startup_consumer::heal_topology_fingerprint(&endpoint_pools)
.expect("single-node gating topology should hash");
crate::storage::storage_api::start_remote_version_state_fleet_probe(topology_fingerprint);
tokio::time::timeout(std::time::Duration::from_secs(5), async {
while crate::storage::storage_api::ecstore_notification::acquire_cross_pool_fence_fleet_proof().is_none() {
tokio::task::yield_now().await;
}
})
.await
.expect("single-node cross-pool fence capability proof should publish");
let server_addr: std::net::SocketAddr = "127.0.0.1:0".parse().unwrap();
let ecstore = ECStore::new(server_addr, endpoint_pools, CancellationToken::new())
@@ -107,6 +121,24 @@ pub(crate) async fn shared_gating_ecstore() -> Arc<ECStore> {
ecstore
}
pub(crate) async fn durable_quota_test_bucket(prefix: &str, limit: u64) -> (Arc<ECStore>, String) {
let store = shared_gating_ecstore().await;
crate::app::runtime_sources::install_test_app_context(Arc::clone(&store)).await;
let bucket = format!("{prefix:.30}-{}", uuid::Uuid::new_v4().simple());
store
.make_bucket(&bucket, &MakeBucketOptions::default())
.await
.expect("create durable quota test bucket");
super::storage_api::test::data_usage::seed_bucket_usage_memory_for_test(&bucket, 0).await;
let metadata_sys =
crate::app::storage_api::test::get_global_bucket_metadata_sys().expect("test app context should expose bucket metadata");
QuotaChecker::new(metadata_sys)
.set_quota_config(&bucket, BucketQuota::new(Some(limit)))
.await
.expect("configure durable quota test bucket");
(store, bucket)
}
pub(crate) async fn shared_gating_ambient() -> Arc<AppContext> {
let store = shared_gating_ecstore().await;
if let Some(ambient) = crate::runtime_sources::current_app_context() {
+217 -81
View File
@@ -33,7 +33,7 @@ use super::storage_api::multipart_usecase::contract::multipart::{CompletePart, M
use super::storage_api::multipart_usecase::contract::object::{ObjectIO as _, ObjectOperations as _};
use super::storage_api::multipart_usecase::contract::range::HTTPRangeSpec;
use super::storage_api::multipart_usecase::data_usage::{
record_bucket_object_version_write_memory, record_bucket_object_write_memory,
quota_object_size, record_bucket_object_version_write_memory, record_bucket_object_write_memory,
};
use super::storage_api::multipart_usecase::error::{StorageError, is_err_object_not_found, is_err_version_not_found};
use super::storage_api::multipart_usecase::helper::OperationHelper;
@@ -64,10 +64,10 @@ use super::storage_api::multipart_usecase::{
};
use crate::app::object_data_cache::{
ObjectDataCacheAdapter, invalidate_object_data_cache_after_complete_multipart_success,
invalidate_object_data_cache_after_delete_success, invalidate_object_data_cache_before_mutation,
invalidate_object_data_cache_before_mutation,
};
use crate::app::object_usecase::{
acquire_copy_bucket_lifecycle_locks, build_put_like_object_lock_metadata, map_quota_check_outcome,
acquire_copy_bucket_lifecycle_locks, apply_quota_admission, build_put_like_object_lock_metadata, map_quota_check_outcome,
validate_existing_object_lock_for_write,
};
use crate::app::runtime_sources::{
@@ -83,6 +83,8 @@ use rustfs_io_metrics::record_s3_op;
use rustfs_s3_ops::S3Operation;
use rustfs_targets::EventName;
use rustfs_utils::CompressionAlgorithm;
#[cfg(test)]
use rustfs_utils::http::insert_header;
use rustfs_utils::http::{
SUFFIX_REPLICATION_PRESERVE_CIPHERTEXT, SUFFIX_REPLICATION_STATUS, SUFFIX_REPLICATION_TIMESTAMP,
SUFFIX_SOURCE_REPLICATION_REQUEST, contains_key_str, get_header, get_source_scheme,
@@ -226,12 +228,8 @@ fn internal_object_info_lookup_opts(mut opts: ObjectOptions) -> ObjectOptions {
opts
}
fn logical_object_size(info: &ObjectInfo) -> Result<u64, StorageError> {
u64::try_from(info.get_actual_size()?).map_err(|_| StorageError::PartMissingOrCorrupt)
}
fn quota_accounting_object_size(info: &ObjectInfo, fail_closed: bool) -> S3Result<u64> {
match logical_object_size(info) {
match quota_object_size(info) {
Ok(size) => Ok(size),
Err(err) if fail_closed => Err(ApiError::from(err).into()),
Err(_) => Ok(info.size.max(0) as u64),
@@ -507,11 +505,7 @@ impl DefaultMultipartUsecase {
Ok(existing_obj_info) => {
validate_existing_object_lock_for_write(&existing_obj_info, &current_opts)?;
let physical_size = existing_obj_info.size.max(0) as u64;
let logical_size = if opts.replication_request {
Ok(physical_size)
} else {
logical_object_size(&existing_obj_info)
};
let logical_size = quota_object_size(&existing_obj_info);
Some((physical_size, logical_size))
}
Err(err) => {
@@ -560,32 +554,20 @@ impl DefaultMultipartUsecase {
};
let quota_metadata_sys = self.bucket_metadata_sys();
let quota_tracking = quota_metadata_sys.is_some();
let mut quota_enabled = false;
if let Some(metadata_sys) = quota_metadata_sys.as_ref() {
let quota_checker = QuotaChecker::new(metadata_sys.clone());
let check_result =
map_quota_check_outcome(&bucket, quota_checker.check_quota(&bucket, QuotaOperation::PutObject, 0).await)?;
// Ciphertext-passthrough replication parts use a different size basis and retain
// the existing post-commit accounting path until they carry a trusted logical-size proof.
if !opts.replication_request
&& let Some(quota_limit) = check_result.quota_limit
{
let installed = check_result
.current_usage
.is_some_and(|current_usage| opts.set_quota_admission(current_usage, quota_limit));
if !installed {
return Err(S3Error::with_message(
S3ErrorCode::ServiceUnavailable,
"Bucket quota check temporarily unavailable, please retry".to_string(),
));
}
}
quota_enabled = check_result.quota_limit.is_some();
apply_quota_admission(&mut opts, &check_result)?;
}
let previous_current_size = match previous_current_sizes {
Some((physical_size, _)) if opts.replication_request => Some(physical_size),
Some((_, Ok(logical_size))) => Some(logical_size),
Some((_, Err(err))) if opts.quota_admission.is_some() => return Err(ApiError::from(err).into()),
Some((physical_size, Err(_))) => Some(physical_size),
Some((_, Ok(logical_size))) if quota_enabled => Some(logical_size),
Some((_, Err(err))) if quota_enabled => return Err(ApiError::from(err).into()),
Some((physical_size, _)) => Some(physical_size),
None => None,
};
@@ -595,7 +577,6 @@ impl DefaultMultipartUsecase {
let key = key.clone();
let upload_id = upload_id.clone();
let opts = opts.clone();
let quota_metadata_sys = quota_metadata_sys.clone();
async move {
let obj_info = store
.clone()
@@ -605,37 +586,9 @@ impl DefaultMultipartUsecase {
let _ = invalidate_object_data_cache_after_complete_multipart_success(&cache_adapter, &bucket, &key).await;
record_capacity_write(Some(capacity_scope_token)).await;
if let Some(metadata_sys) = quota_metadata_sys.as_ref() {
if opts.replication_request {
let quota_checker = QuotaChecker::new(metadata_sys.clone());
match quota_checker
.check_quota(&bucket, QuotaOperation::PutObject, obj_info.size.max(0) as u64)
.await
{
Ok(check_result) if !check_result.allowed => {
let _ = store.delete_object(&bucket, &key, ObjectOptions::default()).await;
let _ = invalidate_object_data_cache_after_delete_success(&cache_adapter, &bucket, &key).await;
return Err(S3Error::with_message(
S3ErrorCode::InvalidRequest,
format!(
"Bucket quota exceeded. Current usage: {} bytes, limit: {} bytes",
check_result.current_usage.unwrap_or(0),
check_result.quota_limit.unwrap_or(0)
),
));
}
Err(err) => {
warn!("Quota check failed for bucket {} after multipart completion: {}", bucket, err);
}
Ok(_) => {}
}
}
if quota_tracking {
let committed_size = quota_accounting_object_size(&obj_info, quota_enabled)?;
let committed_size = if opts.replication_request {
obj_info.size.max(0) as u64
} else {
quota_accounting_object_size(&obj_info, opts.quota_admission.is_some())?
};
if versioned {
record_bucket_object_version_write_memory(&bucket, previous_current_size, committed_size).await;
} else {
@@ -845,6 +798,20 @@ impl DefaultMultipartUsecase {
let ciphertext_passthrough = replication_authorized
&& get_header(&req.headers, SUFFIX_SOURCE_REPLICATION_REQUEST).as_deref() == Some("true")
&& rustfs_utils::http::ssec_transport_to_stored_metadata(&req.headers).is_some();
if ciphertext_passthrough && let Some(metadata_sys) = self.bucket_metadata_sys() {
let check_result = map_quota_check_outcome(
&bucket,
QuotaChecker::new(metadata_sys)
.check_quota(&bucket, QuotaOperation::PutObject, 0)
.await,
)?;
if check_result.quota_limit.is_some() {
return Err(S3Error::with_message(
S3ErrorCode::InvalidRequest,
"SSE-C ciphertext replication is unavailable for quota-enabled buckets".to_string(),
));
}
}
if ciphertext_passthrough {
insert_str(&mut metadata, SUFFIX_REPLICATION_PRESERVE_CIPHERTEXT, "true".to_string());
}
@@ -1678,6 +1645,24 @@ mod tests {
assert_eq!(quota_accounting_object_size(&info, true).expect("logical size should resolve"), 8192);
assert_eq!(quota_accounting_object_size(&info, false).expect("logical size should resolve"), 8192);
let mut poisoned_metadata = HashMap::new();
insert_str(&mut poisoned_metadata, rustfs_utils::http::SUFFIX_COMPRESSION, "S2".to_string());
insert_str(&mut poisoned_metadata, rustfs_utils::http::SUFFIX_ACTUAL_SIZE, "1".to_string());
let poisoned = ObjectInfo {
size: 17,
parts: Arc::new(vec![rustfs_filemeta::ObjectPartInfo {
size: 4096,
actual_size: 4096,
..Default::default()
}]),
user_defined: Arc::new(poisoned_metadata),
..Default::default()
};
assert_eq!(
quota_accounting_object_size(&poisoned, true).expect("persisted part size must be charged"),
4096
);
}
#[test]
@@ -2086,30 +2071,14 @@ mod tests {
#[tokio::test]
#[serial_test::serial]
async fn compressed_complete_records_logical_quota_usage_and_overwrite_delta() {
use crate::app::storage_api::multipart_usecase::bucket::quota::BucketQuota;
use crate::app::storage_api::test::contract::bucket::{BucketOperations as _, MakeBucketOptions};
use crate::app::storage_api::test::data_usage::seed_bucket_usage_memory_for_test;
let store = crate::app::gating_test_env::shared_gating_ecstore().await;
crate::app::runtime_sources::install_test_app_context(Arc::clone(&store)).await;
let bucket = format!("compressed-complete-quota-{}", Uuid::new_v4());
let (store, bucket) = crate::app::gating_test_env::durable_quota_test_bucket("compressed-complete-quota", 16_384).await;
let object = "object";
store
.make_bucket(&bucket, &MakeBucketOptions::default())
.await
.expect("create compressed quota bucket");
seed_bucket_usage_memory_for_test(&bucket, 0).await;
let usecase = DefaultMultipartUsecase::from_global();
let metadata_sys = usecase
.bucket_metadata_sys()
.expect("test app context should expose bucket metadata");
let mut quota_checker = QuotaChecker::new(metadata_sys);
quota_checker
.set_quota_config(&bucket, BucketQuota::new(Some(16_384)))
.await
.expect("configure bucket quota");
let quota_checker = QuotaChecker::new(metadata_sys);
for (actual_size, payload_byte) in [(8192_i64, 0x61), (4096_i64, 0x62)] {
let mut create_opts = ObjectOptions::default();
@@ -2154,6 +2123,173 @@ mod tests {
}
}
#[tokio::test]
#[serial_test::serial]
async fn create_multipart_rejects_ciphertext_replication_before_parts_are_staged() {
let (_store, bucket) = crate::app::gating_test_env::durable_quota_test_bucket("ciphertext-multipart-quota", 4096).await;
let usecase = DefaultMultipartUsecase::from_global();
let input = CreateMultipartUploadInput::builder()
.bucket(bucket)
.key("object".to_string())
.build()
.expect("create multipart request should build");
let mut request = build_request(input, Method::POST);
insert_header(&mut request.headers, SUFFIX_SOURCE_REPLICATION_REQUEST, "true");
request
.headers
.insert(rustfs_utils::http::REPLICATION_SSEC_ALGORITHM_HEADER, HeaderValue::from_static("AES256"));
request.extensions.insert(crate::storage::access::ReqInfo {
replication_request_authorized: true,
..Default::default()
});
let err = usecase
.execute_create_multipart_upload(request)
.await
.expect_err("quota-enabled ciphertext multipart replication should fail before upload creation");
assert_eq!(err.code(), &S3ErrorCode::InvalidRequest);
}
#[tokio::test]
#[serial_test::serial]
async fn concurrent_completions_share_durable_bucket_quota_reservations() {
let (store, bucket) = crate::app::gating_test_env::durable_quota_test_bucket("concurrent-complete-quota", 6000).await;
let usecase = DefaultMultipartUsecase::from_global();
let mut inputs = Vec::new();
for object in ["first", "second"] {
let upload = store
.new_multipart_upload(&bucket, object, &ObjectOptions::default())
.await
.expect("create concurrent multipart upload");
let mut reader = PutObjReader::from_vec(vec![0x71; 4096]);
let part = store
.put_object_part(&bucket, object, &upload.upload_id, 1, &mut reader, &ObjectOptions::default())
.await
.expect("stage concurrent multipart part");
inputs.push(
CompleteMultipartUploadInput::builder()
.bucket(bucket.clone())
.key(object.to_string())
.upload_id(upload.upload_id)
.multipart_upload(Some(CompletedMultipartUpload {
parts: Some(vec![CompletedPart {
part_number: Some(1),
e_tag: part.etag.map(|etag| to_s3s_etag(&etag)),
..Default::default()
}]),
}))
.build()
.expect("build concurrent completion input"),
);
}
let first_usecase = usecase.clone();
let first = first_usecase.execute_complete_multipart_upload(build_request(inputs.remove(0), Method::POST));
let second = usecase.execute_complete_multipart_upload(build_request(inputs.remove(0), Method::POST));
let (first, second) = tokio::join!(first, second);
assert_eq!(usize::from(first.is_ok()) + usize::from(second.is_ok()), 1);
let denied = first.err().or_else(|| second.err()).expect("one completion must be denied");
assert_eq!(denied.code(), &S3ErrorCode::InvalidRequest);
}
#[tokio::test]
#[serial_test::serial]
async fn multipart_completion_rejects_rotated_quota_capability_before_rename() {
use crate::app::storage_api::test::set_disk::{MultipartCommitBarrier, MultipartCommitPause};
let (store, bucket) = crate::app::gating_test_env::durable_quota_test_bucket("rotated-proof-mpu-quota", 4096).await;
let object = "object";
let upload = store
.new_multipart_upload(&bucket, object, &ObjectOptions::default())
.await
.expect("create multipart upload");
let mut reader = PutObjReader::from_vec(vec![0x78; 4096]);
let part = store
.put_object_part(&bucket, object, &upload.upload_id, 1, &mut reader, &ObjectOptions::default())
.await
.expect("stage multipart part");
let barrier = MultipartCommitBarrier::install(&bucket, object, MultipartCommitPause::BeforeQuotaRename);
let complete_store = Arc::clone(&store);
let complete_bucket = bucket.clone();
let upload_id = upload.upload_id.clone();
let complete = tokio::spawn(async move {
complete_store
.complete_multipart_upload(
&complete_bucket,
object,
&upload_id,
vec![CompletePart {
part_num: 1,
etag: part.etag,
..Default::default()
}],
&ObjectOptions::default(),
)
.await
});
barrier.wait_until_paused().await;
assert!(
crate::storage::storage_api::ecstore_notification::rotate_cross_pool_fence_fleet_proof_for_test(),
"the gating environment must have a current fleet proof"
);
barrier.release();
let err = complete
.await
.expect("completion task should not panic")
.expect_err("a replaced fleet proof must fence multipart rename");
assert!(matches!(
err,
StorageError::NamespaceLockQuorumUnavailable {
mode: "quota_reservation",
..
}
));
store
.get_multipart_info(&bucket, object, &upload.upload_id, &ObjectOptions::default())
.await
.expect("proof rotation must preserve the multipart upload for retry");
}
#[tokio::test]
#[serial_test::serial]
async fn data_movement_multipart_completion_has_zero_quota_growth() {
let (store, bucket) = crate::app::gating_test_env::durable_quota_test_bucket("data-movement-mpu-quota", 0).await;
let object = "object";
let mut movement_opts = ObjectOptions {
data_movement: true,
..Default::default()
};
let upload = store
.new_multipart_upload(&bucket, object, &movement_opts)
.await
.expect("create data-movement multipart upload");
let mut reader = PutObjReader::from_vec(vec![0x7a; 4096]);
let part = store
.put_object_part(&bucket, object, &upload.upload_id, 1, &mut reader, &movement_opts)
.await
.expect("stage data-movement multipart part");
movement_opts.preserve_etag = Some("movement-etag".to_string());
let completed = store
.complete_multipart_upload(
&bucket,
object,
&upload.upload_id,
vec![CompletePart {
part_num: 1,
etag: part.etag,
..Default::default()
}],
&movement_opts,
)
.await
.expect("moving an already-accounted multipart object between pools must have zero quota growth");
assert_eq!(completed.size, 4096);
}
#[tokio::test]
#[serial_test::serial]
async fn rejected_empty_parts_preserve_existing_object_and_staging() {
File diff suppressed because it is too large Load Diff
+12
View File
@@ -47,6 +47,8 @@ pub(crate) mod capacity {
pub(crate) mod data_usage {
use std::sync::Arc;
pub(crate) use crate::storage::storage_api::ecstore_data_usage::quota_object_size;
pub(crate) async fn apply_bucket_usage_memory_overlay(data_usage_info: &mut rustfs_data_usage::DataUsageInfo) {
crate::storage::storage_api::ecstore_data_usage::apply_bucket_usage_memory_overlay(data_usage_info).await;
}
@@ -1214,4 +1216,14 @@ pub(crate) mod test {
pub(crate) use crate::storage::storage_api::{
ECStore, Endpoint, Endpoints, PoolEndpoints, StorageObjectInfo, StorageObjectOptions, StoragePutObjReader,
};
pub(crate) mod set_disk {
pub(crate) use crate::storage::storage_api::ecstore_set_disk::{
MultipartCommitBarrier, MultipartCommitPause, PutObjectCommitBarrier, PutObjectCommitPause,
fail_next_quota_ledger_save_for_test,
};
}
pub(crate) mod metadata_sys {
pub(crate) use crate::storage::storage_api::ecstore_bucket::metadata_sys::ConfigWriteLockProbe;
}
}
+3 -3
View File
@@ -301,7 +301,7 @@ impl S3 for FS {
#[instrument(level = "debug", skip(self, req))]
async fn copy_object(&self, req: S3Request<CopyObjectInput>) -> S3Result<S3Response<CopyObjectOutput>> {
let usecase = s3_api::object_usecase_for(self);
Box::pin(usecase.execute_copy_object(req)).await
usecase.execute_copy_object(req).await
}
#[instrument(
@@ -704,7 +704,7 @@ impl S3 for FS {
async fn get_object(&self, req: S3Request<GetObjectInput>) -> S3Result<S3Response<GetObjectOutput>> {
crate::hp_guard!("S3::get_object");
let usecase = s3_api::object_usecase_for(self);
Box::pin(usecase.execute_get_object(req)).await
usecase.execute_get_object(req).await
}
async fn get_object_acl(&self, req: S3Request<GetObjectAclInput>) -> S3Result<S3Response<GetObjectAclOutput>> {
@@ -1261,7 +1261,7 @@ impl S3 for FS {
async fn put_object(&self, req: S3Request<PutObjectInput>) -> S3Result<S3Response<PutObjectOutput>> {
crate::hp_guard!("S3::put_object");
let usecase = s3_api::object_usecase_for(self);
Box::pin(usecase.execute_put_object(self, req)).await
usecase.execute_put_object(self, req).await
}
async fn put_object_acl(&self, req: S3Request<PutObjectAclInput>) -> S3Result<S3Response<PutObjectAclOutput>> {
+3 -5
View File
@@ -153,9 +153,7 @@ fn remove_heal_control_replay(
static HEAL_CONTROL_REPLAY_CACHE: OnceLock<tokio::sync::Mutex<HashMap<String, Arc<HealControlReplayEntry>>>> = OnceLock::new();
static NODE_CAPABILITY_SERVER_EPOCH: LazyLock<Uuid> = LazyLock::new(Uuid::new_v4);
// RUSTFS_COMPAT_TODO(cross-pool-fence-v1): advertise unsupported during predeployment. Remove after composite acquisition,
// activation fencing, fleet proof, commit-time proof revalidation, and fail-closed revocation ship together.
const CROSS_POOL_FENCE_SUPPORTED_VERSION: u32 = 0;
const CROSS_POOL_FENCE_SUPPORTED_VERSION: u32 = 1;
fn admit_heal_control_replay(
replay_cache: &mut HashMap<String, Arc<HealControlReplayEntry>>,
@@ -3342,7 +3340,7 @@ mod tests {
}
#[tokio::test]
async fn cross_pool_fence_probe_authenticates_unsupported_rollout_state() {
async fn cross_pool_fence_probe_authenticates_supported_v1_state() {
let _ = rustfs_credentials::set_global_rpc_secret("cross-pool-fence-node-service-test-secret".to_string());
let endpoints = heal_control_test_endpoints_with_coordinator("node-0", true);
assert!(
@@ -3407,7 +3405,7 @@ mod tests {
assert!(response.success);
assert_eq!(response.error_info, None);
assert_eq!(&response.result[..4], &0_u32.to_be_bytes());
assert_eq!(&response.result[..4], &1_u32.to_be_bytes());
let (topology_member, process_epoch) = rustfs_protos::decode_remote_version_state_capability(&response.result[4..])
.expect("capability identity should decode");
assert_eq!(topology_member, "node-a:9000");
+38 -25
View File
@@ -37,19 +37,28 @@ const MSGPACK_ENCODE_CAPACITY_HINT: usize = 512;
const FILE_INFO_MSGPACK_ENCODE_CAPACITY_HINT: usize = 1024;
const SNAPSHOT_LEASE_PROTOCOL_VERSION: u32 = 1;
fn snapshot_lease_response(result: Result<SnapshotLeaseToken, DiskError>) -> Response<SnapshotLeaseResponse> {
match result {
Ok(token) => Response::new(SnapshotLeaseResponse {
success: true,
token: token.as_bytes().to_vec().into(),
protocol_version: SNAPSHOT_LEASE_PROTOCOL_VERSION,
error: None,
}),
Err(err) => Response::new(SnapshotLeaseResponse {
success: false,
token: Bytes::new(),
protocol_version: SNAPSHOT_LEASE_PROTOCOL_VERSION,
error: Some(err.into()),
}),
}
}
struct DecodedRpcPayload<T> {
value: T,
from_msgpack: bool,
}
fn snapshot_lease_disabled_response() -> SnapshotLeaseResponse {
SnapshotLeaseResponse {
success: false,
token: Bytes::new(),
protocol_version: SNAPSHOT_LEASE_PROTOCOL_VERSION,
error: Some(DiskError::UnsupportedDisk.into()),
}
}
fn decode_msgpack_or_json<T: DeserializeOwned>(
binary: &[u8],
json: &str,
@@ -241,7 +250,12 @@ impl NodeService {
rustfs_protos::canonical_snapshot_lease_request_body(request.get_ref()),
"acquire_snapshot_lease",
)?;
Ok(Response::new(snapshot_lease_disabled_response()))
let request = request.into_inner();
let result = match self.find_disk(&request.disk).await {
Some(disk) => disk.acquire_snapshot_lease(&request.volume, &request.path).await,
None => Err(DiskError::other("cannot find disk")),
};
Ok(snapshot_lease_response(result))
}
pub(super) async fn handle_renew_snapshot_lease(
@@ -253,7 +267,14 @@ impl NodeService {
rustfs_protos::canonical_snapshot_lease_renew_request_body(request.get_ref()),
"renew_snapshot_lease",
)?;
Ok(Response::new(snapshot_lease_disabled_response()))
let request = request.into_inner();
let token =
SnapshotLeaseToken::from_slice(&request.token).map_err(|_| Status::invalid_argument("invalid lease token"))?;
let result = match self.find_disk(&request.disk).await {
Some(disk) => disk.renew_snapshot_lease(&request.volume, &request.path, token).await,
None => Err(DiskError::other("cannot find disk")),
};
Ok(snapshot_lease_response(result))
}
pub(super) async fn handle_release_snapshot_lease(
@@ -266,8 +287,11 @@ impl NodeService {
"release_snapshot_lease",
)?;
let request = request.into_inner();
let token =
SnapshotLeaseToken::from_slice(&request.token).map_err(|_| Status::invalid_argument("invalid lease token"))?;
let token = if request.token.as_ref() == SnapshotLeaseToken::revoke_all().as_bytes() {
SnapshotLeaseToken::revoke_all()
} else {
SnapshotLeaseToken::from_slice(&request.token).map_err(|_| Status::invalid_argument("invalid lease token"))?
};
let Some(disk) = self.find_disk(&request.disk).await else {
return Ok(Response::new(SnapshotLeaseMutationResponse {
success: false,
@@ -1494,11 +1518,11 @@ mod tests {
use super::{
compat_response_json, decode_msgpack_or_json, decode_rename_data_request_file_info,
encode_batch_read_version_response_payloads, encode_file_info_msgpack, encode_msgpack, encode_msgpack_named,
encode_read_multiple_response_payloads, encode_rename_data_response_payloads, snapshot_lease_disabled_response,
encode_read_multiple_response_payloads, encode_rename_data_response_payloads,
};
use crate::storage::storage_api::ReadMultipleResp;
use crate::storage::storage_api::RenameDataResp;
use crate::storage::storage_api::rpc_consumer::node_service::BatchReadVersionResp;
use crate::storage::storage_api::{DiskError, RenameDataResp};
use rustfs_filemeta::FileInfo;
use rustfs_io_metrics::internode_metrics::global_internode_metrics;
use serde::{Deserialize, Serialize};
@@ -1509,17 +1533,6 @@ mod tests {
count: u32,
}
#[test]
fn snapshot_lease_acquire_and_renew_fail_closed() {
let response = snapshot_lease_disabled_response();
let expected_error = DiskError::UnsupportedDisk.into();
assert!(!response.success);
assert!(response.token.is_empty());
assert_eq!(response.protocol_version, 1);
assert_eq!(response.error, Some(expected_error));
}
#[test]
fn decode_msgpack_or_json_prefers_binary_payload() {
let payload = SamplePayload {
+21 -2
View File
@@ -430,7 +430,7 @@ pub(crate) mod ecstore_config {
pub(crate) mod ecstore_data_usage {
pub(crate) use rustfs_ecstore::api::data_usage::{
apply_bucket_usage_memory_overlay, init_compression_total_memory_from_backend, load_admin_data_usage_from_backend_cached,
load_data_usage_from_backend, record_bucket_delete_marker_memory, record_bucket_object_delete_memory,
load_data_usage_from_backend, quota_object_size, record_bucket_delete_marker_memory, record_bucket_object_delete_memory,
record_bucket_object_version_write_memory, record_bucket_object_write_memory,
record_bucket_object_write_unknown_previous_memory, store_compression_total_in_backend,
};
@@ -487,8 +487,12 @@ pub(crate) mod ecstore_metrics {
#[allow(unused_imports)]
pub(crate) mod ecstore_notification {
#[cfg(test)]
pub(crate) use rustfs_ecstore::api::notification::rotate_cross_pool_fence_fleet_proof_for_test;
pub(crate) use rustfs_ecstore::api::notification::{
NotificationSys, get_global_notification_sys, new_global_notification_sys, start_remote_version_state_fleet_probe,
CrossPoolFenceFleetProofToken, NotificationSys, acquire_cross_pool_fence_fleet_proof,
cross_pool_fence_fleet_proof_matches, get_global_notification_sys, new_global_notification_sys,
start_remote_version_state_fleet_probe,
};
}
@@ -546,6 +550,11 @@ pub(crate) mod ecstore_test_support {
}
pub(crate) mod ecstore_set_disk {
#[cfg(test)]
pub(crate) use rustfs_ecstore::api::set_disk::test_util::{
MultipartCommitBarrier, MultipartCommitPause, PutObjectCommitBarrier, PutObjectCommitPause,
fail_next_quota_ledger_save_for_test,
};
pub(crate) use rustfs_ecstore::api::set_disk::{
DEFAULT_READ_BUFFER_SIZE, file_info_quorum_hash, get_lock_acquire_timeout, is_valid_storage_class,
};
@@ -1130,7 +1139,9 @@ pub(crate) trait StorageDiskRpcExt {
) -> DiskResult<()>;
async fn read_metadata(&self, volume: &str, path: &str) -> DiskResult<bytes::Bytes>;
async fn delete_paths(&self, volume: &str, paths: &[String]) -> DiskResult<()>;
async fn acquire_snapshot_lease(&self, volume: &str, path: &str) -> DiskResult<SnapshotLeaseToken>;
async fn release_snapshot_lease(&self, volume: &str, path: &str, token: SnapshotLeaseToken) -> DiskResult<()>;
async fn renew_snapshot_lease(&self, volume: &str, path: &str, token: SnapshotLeaseToken) -> DiskResult<SnapshotLeaseToken>;
async fn stat_volume(&self, volume: &str) -> DiskResult<VolumeInfo>;
async fn list_volumes(&self) -> DiskResult<Vec<VolumeInfo>>;
async fn make_volume(&self, volume: &str) -> DiskResult<()>;
@@ -1258,10 +1269,18 @@ where
ecstore_disk::DiskAPI::delete_paths(self, volume, paths).await
}
async fn acquire_snapshot_lease(&self, volume: &str, path: &str) -> DiskResult<SnapshotLeaseToken> {
ecstore_disk::DiskAPI::acquire_snapshot_lease(self, volume, path).await
}
async fn release_snapshot_lease(&self, volume: &str, path: &str, token: SnapshotLeaseToken) -> DiskResult<()> {
ecstore_disk::DiskAPI::release_snapshot_lease(self, volume, path, token).await
}
async fn renew_snapshot_lease(&self, volume: &str, path: &str, token: SnapshotLeaseToken) -> DiskResult<SnapshotLeaseToken> {
ecstore_disk::DiskAPI::renew_snapshot_lease(self, volume, path, token).await
}
async fn stat_volume(&self, volume: &str) -> DiskResult<VolumeInfo> {
ecstore_disk::DiskAPI::stat_volume(self, volume).await
}