Compare commits

...

1 Commits

Author SHA1 Message Date
cxymds b5cdfd496a fix: enforce S3 authorization for recursive deletion 2026-09-11 21:23:50 +08:00
9 changed files with 1026 additions and 266 deletions
@@ -0,0 +1,642 @@
// Copyright 2026 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.
//! Standard S3 deletion permissions and the explicit recursive-delete extension.
use crate::common::{
AdminTransport, RustFSTestEnvironment, admin_add_canned_policy_via, admin_attach_user_policy_via, admin_create_user,
init_logging,
};
use aws_sdk_s3::Client;
use aws_sdk_s3::error::{ProvideErrorMetadata, SdkError};
use aws_sdk_s3::operation::delete_object::{DeleteObjectError, DeleteObjectOutput};
use aws_sdk_s3::primitives::ByteStream;
use aws_sdk_s3::types::{BucketVersioningStatus, Delete, ObjectIdentifier, VersioningConfiguration};
use futures::{StreamExt, TryStreamExt, stream};
use serde_json::{Value, json};
use std::collections::BTreeSet;
use std::error::Error;
use uuid::Uuid;
type TestResult<T = ()> = Result<T, Box<dyn Error + Send + Sync>>;
type VersionSnapshot = BTreeSet<(String, String, bool)>;
async fn set_policy(env: &RustFSTestEnvironment, name: &str, policy: &Value) -> TestResult {
admin_add_canned_policy_via(
AdminTransport::Signed,
&env.url,
&env.access_key,
&env.secret_key,
name,
&policy.to_string(),
)
.await
}
async fn policy_user(env: &RustFSTestEnvironment, policy_name: &str, policy: Option<Value>) -> TestResult<Client> {
let username = Uuid::new_v4().simple().to_string();
let secret = Uuid::new_v4().simple().to_string();
admin_create_user(env, &username, &secret).await?;
if let Some(policy) = policy {
set_policy(env, policy_name, &policy).await?;
}
admin_attach_user_policy_via(AdminTransport::Signed, &env.url, &env.access_key, &env.secret_key, policy_name, &username)
.await?;
Ok(env.create_s3_client_with_credentials(&username, &secret))
}
async fn versioning(client: &Client, bucket: &str, status: BucketVersioningStatus) -> TestResult {
client
.put_bucket_versioning()
.bucket(bucket)
.versioning_configuration(VersioningConfiguration::builder().status(status).build())
.send()
.await?;
Ok(())
}
async fn put(client: &Client, bucket: &str, key: &str) -> TestResult<String> {
let result = client
.put_object()
.bucket(bucket)
.key(key)
.body(ByteStream::from_static(b"delete authorization fixture"))
.send()
.await?;
Ok(result.version_id().unwrap_or("null").to_string())
}
async fn versions(client: &Client, bucket: &str, prefix: &str) -> TestResult<VersionSnapshot> {
let mut result = BTreeSet::new();
let mut markers = (None, None);
loop {
let page = client
.list_object_versions()
.bucket(bucket)
.prefix(prefix)
.set_key_marker(markers.0.clone())
.set_version_id_marker(markers.1.clone())
.send()
.await?;
for version in page.versions() {
result.insert((
version.key().ok_or("listed version missing key")?.to_string(),
version.version_id().ok_or("listed version missing ID")?.to_string(),
false,
));
}
for marker in page.delete_markers() {
result.insert((
marker.key().ok_or("listed delete marker missing key")?.to_string(),
marker.version_id().ok_or("listed delete marker missing ID")?.to_string(),
true,
));
}
if page.is_truncated() != Some(true) {
return Ok(result);
}
let next = (
Some(
page.next_key_marker()
.ok_or("truncated versions page missing next key marker")?
.to_string(),
),
page.next_version_id_marker().map(str::to_string),
);
assert_ne!(markers, next, "ListObjectVersions pagination must advance");
markers = next;
}
}
async fn force_delete(client: &Client, bucket: &str, prefix: &str) -> Result<DeleteObjectOutput, SdkError<DeleteObjectError>> {
client
.delete_object()
.bucket(bucket)
.key(prefix)
.customize()
.mutate_request(|request| {
request.headers_mut().insert("x-rustfs-force-delete", "true");
})
.send()
.await
}
async fn replica_force_delete(
client: &Client,
bucket: &str,
prefix: &str,
) -> Result<DeleteObjectOutput, SdkError<DeleteObjectError>> {
client
.delete_object()
.bucket(bucket)
.key(prefix)
.customize()
.mutate_request(|request| {
request.headers_mut().insert("x-rustfs-force-delete", "true");
request.headers_mut().insert("x-amz-replication-status", "REPLICA");
})
.send()
.await
}
fn assert_denied<T, E>(result: Result<T, SdkError<E>>)
where
T: std::fmt::Debug,
E: ProvideErrorMetadata + std::fmt::Debug,
{
let error = result.expect_err("request must be denied by its S3 permission");
assert_eq!(
error.as_service_error().and_then(ProvideErrorMetadata::code),
Some("AccessDenied"),
"expected an S3 authorization denial, got {error:?}"
);
}
#[tokio::test]
async fn sdk_version_deletion_requires_only_delete_object_version() -> TestResult {
init_logging();
let mut env = RustFSTestEnvironment::new().await?;
env.start_rustfs_server(vec![]).await?;
let root = env.create_s3_client();
let bucket = "delete-version-permissions";
root.create_bucket().bucket(bucket).send().await?;
put(&root, bucket, "single-null.txt").await?;
put(&root, bucket, "batch-null.txt").await?;
versioning(&root, bucket, BucketVersioningStatus::Enabled).await?;
let old = put(&root, bucket, "single.txt").await?;
let current = put(&root, bucket, "single.txt").await?;
let batch_version = put(&root, bucket, "batch.txt").await?;
let ordinary_version = put(&root, bucket, "ordinary.txt").await?;
let user = policy_user(
&env,
"version-deleter",
Some(json!({"Version":"2012-10-17","Statement":[
{"Effect":"Allow","Action":"s3:DeleteObjectVersion","Resource":format!("arn:aws:s3:::{bucket}/*")},
{"Effect":"Deny","Action":"s3:DeleteObject","Resource":format!("arn:aws:s3:::{bucket}/*")}
]})),
)
.await?;
user.delete_object()
.bucket(bucket)
.key("single.txt")
.version_id(&old)
.send()
.await?;
user.delete_object()
.bucket(bucket)
.key("single-null.txt")
.version_id("null")
.send()
.await?;
assert_denied(user.delete_object().bucket(bucket).key("ordinary.txt").send().await);
let batch = user
.delete_objects()
.bucket(bucket)
.delete(
Delete::builder()
.objects(
ObjectIdentifier::builder()
.key("batch.txt")
.version_id(&batch_version)
.build()?,
)
.objects(ObjectIdentifier::builder().key("batch-null.txt").version_id("null").build()?)
.objects(ObjectIdentifier::builder().key("ordinary.txt").build()?)
.build()?,
)
.send()
.await?;
assert_eq!(batch.deleted().len(), 2, "both explicit version items must succeed");
assert_eq!(batch.errors().len(), 1, "only the unversioned item must be denied");
assert_eq!(batch.errors()[0].key(), Some("ordinary.txt"));
assert_eq!(batch.errors()[0].code(), Some("AccessDenied"));
assert_eq!(
versions(&root, bucket, "").await?,
BTreeSet::from([
("single.txt".into(), current, false),
("ordinary.txt".into(), ordinary_version, false)
]),
"version-only deletion must preserve the current single-object version and denied object"
);
Ok(())
}
#[tokio::test]
async fn sdk_list_bucket_and_list_bucket_versions_permissions_are_independent() -> TestResult {
init_logging();
let mut env = RustFSTestEnvironment::new().await?;
env.start_rustfs_server(vec![]).await?;
let root = env.create_s3_client();
let bucket = "list-version-permissions";
root.create_bucket().bucket(bucket).send().await?;
versioning(&root, bucket, BucketVersioningStatus::Enabled).await?;
put(&root, bucket, "visible.txt").await?;
for (action, name) in [
("s3:ListBucket", "object-lister"),
("s3:ListBucketVersions", "version-lister"),
] {
let user = policy_user(
&env,
name,
Some(json!({"Version":"2012-10-17","Statement":[
{"Effect":"Allow","Action":action,"Resource":format!("arn:aws:s3:::{bucket}")}
]})),
)
.await?;
if action == "s3:ListBucket" {
assert_eq!(user.list_objects_v2().bucket(bucket).send().await?.contents().len(), 1);
assert_denied(user.list_object_versions().bucket(bucket).send().await);
} else {
assert_eq!(user.list_object_versions().bucket(bucket).send().await?.versions().len(), 1);
assert_denied(user.list_objects_v2().bucket(bucket).send().await);
}
}
Ok(())
}
#[tokio::test]
async fn console_admin_force_delete_removes_prefix_versions_and_delete_markers() -> TestResult {
init_logging();
let mut env = RustFSTestEnvironment::new().await?;
env.start_rustfs_server(vec![]).await?;
let root = env.create_s3_client();
let bucket = "force-console-admin";
root.create_bucket().bucket(bucket).send().await?;
put(&root, bucket, "folder/null.txt").await?;
versioning(&root, bucket, BucketVersioningStatus::Enabled).await?;
for key in ["folder/a.txt", "folder/deep/b.txt", "single.txt"] {
put(&root, bucket, key).await?;
put(&root, bucket, key).await?;
root.delete_object().bucket(bucket).key(key).send().await?;
}
put(&root, bucket, "keep.txt").await?;
put(&root, bucket, "folder-sibling/keep.txt").await?;
let keep = versions(&root, bucket, "keep.txt").await?;
let sibling = versions(&root, bucket, "folder-sibling/").await?;
let user = policy_user(&env, "consoleAdmin", None).await?;
force_delete(&user, bucket, "folder/").await?;
assert!(
versions(&root, bucket, "folder/").await?.is_empty(),
"force prefix deletion must remove null versions and markers"
);
assert_eq!(versions(&root, bucket, "folder-sibling/").await?, sibling);
assert_eq!(
versions(&root, bucket, "single.txt").await?.len(),
3,
"the separate key must survive folder deletion"
);
force_delete(&user, bucket, "single.txt").await?;
assert!(
versions(&root, bucket, "single.txt").await?.is_empty(),
"explicit force deletion must remove every version of the selected key"
);
assert_eq!(versions(&root, bucket, "keep.txt").await?, keep);
Ok(())
}
#[tokio::test]
async fn force_delete_authorizes_only_its_path_scope_without_list_permissions() -> TestResult {
init_logging();
let mut env = RustFSTestEnvironment::new().await?;
env.start_rustfs_server(vec![]).await?;
let root = env.create_s3_client();
let bucket = "force-delete-only";
root.create_bucket().bucket(bucket).send().await?;
versioning(&root, bucket, BucketVersioningStatus::Enabled).await?;
put(&root, bucket, "selected.txt").await?;
put(&root, bucket, "selected.txt/child.txt").await?;
put(&root, bucket, "selected.txt-sibling").await?;
let sibling = versions(&root, bucket, "selected.txt-sibling").await?;
let user = policy_user(
&env,
"delete-only",
Some(json!({"Version":"2012-10-17","Statement":[
{"Effect":"Allow","Action":["s3:DeleteObject","s3:DeleteObjectVersion"],"Resource":[
format!("arn:aws:s3:::{bucket}/selected.txt"), format!("arn:aws:s3:::{bucket}/selected.txt/*")
]},
{"Effect":"Deny","Action":["s3:DeleteObject","s3:DeleteObjectVersion"],"Resource":format!("arn:aws:s3:::{bucket}/selected.txt-sibling")}
]})),
)
.await?;
assert_denied(user.list_objects_v2().bucket(bucket).send().await);
assert_denied(user.list_object_versions().bucket(bucket).send().await);
force_delete(&user, bucket, "selected.txt").await?;
assert_eq!(
versions(&root, bucket, "selected.txt").await?,
sibling,
"force deletion must remove the selected path and descendants without authorizing or deleting its similarly prefixed sibling"
);
Ok(())
}
#[tokio::test]
async fn force_directory_delete_cannot_remove_an_unauthorized_colliding_parent() -> TestResult {
init_logging();
let mut env = RustFSTestEnvironment::new().await?;
env.start_rustfs_server(vec![]).await?;
let root = env.create_s3_client();
let bucket = "force-directory-collision";
root.create_bucket().bucket(bucket).send().await?;
versioning(&root, bucket, BucketVersioningStatus::Enabled).await?;
let protected_parent_version = put(&root, bucket, "collision.txt").await?;
for key in ["collision.txt/child", "collision.txt-sibling"] {
put(&root, bucket, key).await?;
}
put(&root, bucket, "collision.txt").await?;
root.delete_object().bucket(bucket).key("collision.txt").send().await?;
let user = policy_user(
&env,
"parent-denier",
Some(json!({"Version":"2012-10-17","Statement":[
{"Effect":"Allow","Action":["s3:DeleteObject","s3:DeleteObjectVersion"],"Resource":format!("arn:aws:s3:::{bucket}/*")},
{"Effect":"Deny","Action":"s3:DeleteObjectVersion","Resource":format!("arn:aws:s3:::{bucket}/collision.txt"),
"Condition":{"StringEquals":{"s3:VersionId":protected_parent_version}}}
]})),
)
.await?;
let mut expected = versions(&root, bucket, "").await?;
expected.retain(|(key, _, _)| key != "collision.txt/child");
force_delete(&user, bucket, "collision.txt/").await?;
assert_eq!(
versions(&root, bucket, "").await?,
expected,
"folder deletion must preserve the denied parent's historical versions and delete marker, plus its sibling"
);
Ok(())
}
#[tokio::test]
async fn force_unversioned_directory_requires_only_delete_object() -> TestResult {
init_logging();
let mut env = RustFSTestEnvironment::new().await?;
env.start_rustfs_server(vec![]).await?;
let root = env.create_s3_client();
let bucket = "force-unversioned-permissions";
root.create_bucket().bucket(bucket).send().await?;
for key in ["folder/", "folder/child.txt", "outside.txt"] {
put(&root, bucket, key).await?;
}
let outside = versions(&root, bucket, "outside.txt").await?;
let user = policy_user(
&env,
"unversioned-deleter",
Some(json!({"Version":"2012-10-17","Statement":[
{"Effect":"Allow","Action":"s3:DeleteObject","Resource":format!("arn:aws:s3:::{bucket}/*")},
{"Effect":"Deny","Action":"s3:DeleteObjectVersion","Resource":format!("arn:aws:s3:::{bucket}/*")}
]})),
)
.await?;
force_delete(&user, bucket, "folder/").await?;
assert_eq!(
versions(&root, bucket, "").await?,
outside,
"unversioned force deletion, including a synthetic nil directory marker, must use DeleteObject permission"
);
Ok(())
}
#[tokio::test]
async fn force_delete_denied_child_preserves_every_object_despite_bucket_allow() -> TestResult {
init_logging();
let mut env = RustFSTestEnvironment::new().await?;
env.start_rustfs_server(vec![]).await?;
let root = env.create_s3_client();
let bucket = "force-child-denial";
root.create_bucket().bucket(bucket).send().await?;
for key in ["folder/a-allowed.txt", "folder/z-denied.txt"] {
put(&root, bucket, key).await?;
}
let user = policy_user(
&env,
"child-denier",
Some(json!({"Version":"2012-10-17","Statement":[
{"Effect":"Allow","Action":["s3:DeleteObject","s3:DeleteObjectVersion","s3:ReplicateDelete"],"Resource":format!("arn:aws:s3:::{bucket}/*")},
{"Effect":"Deny","Action":["s3:DeleteObject","s3:DeleteObjectVersion","s3:ReplicateDelete"],"Resource":format!("arn:aws:s3:::{bucket}/folder/z-denied.txt")}
]})),
)
.await?;
root.put_bucket_policy()
.bucket(bucket)
.policy(json!({"Version":"2012-10-17","Statement":[
{"Effect":"Allow","Principal":"*","Action":["s3:DeleteObject","s3:DeleteObjectVersion","s3:ReplicateDelete"],"Resource":format!("arn:aws:s3:::{bucket}/*")}
]}).to_string())
.send().await?;
let before = versions(&root, bucket, "folder/").await?;
assert_denied(force_delete(&user, bucket, "folder/").await);
assert_eq!(
versions(&root, bucket, "folder/").await?,
before,
"a denied descendant must prevent every mutation in the force scope"
);
assert_denied(replica_force_delete(&user, bucket, "folder/").await);
assert_eq!(
versions(&root, bucket, "folder/").await?,
before,
"the REPLICA header must not bypass a descendant's ReplicateDelete denial"
);
root.delete_bucket_policy().bucket(bucket).send().await?;
let replica_user = policy_user(
&env,
"replica-deleter",
Some(json!({"Version":"2012-10-17","Statement":[
{"Effect":"Allow","Action":"s3:DeleteObject","Resource":format!("arn:aws:s3:::{bucket}/*")},
{"Effect":"Allow","Action":"s3:ReplicateDelete","Resource":format!("arn:aws:s3:::{bucket}/*")}
]})),
)
.await?;
replica_force_delete(&replica_user, bucket, "folder/").await?;
assert!(
versions(&root, bucket, "folder/").await?.is_empty(),
"an authorized replica force request must check ReplicateDelete for its descendants"
);
Ok(())
}
#[tokio::test]
async fn force_delete_denied_historical_version_preserves_versions_and_markers() -> TestResult {
init_logging();
let mut env = RustFSTestEnvironment::new().await?;
env.start_rustfs_server(vec![]).await?;
let root = env.create_s3_client();
let bucket = "force-version-denial";
root.create_bucket().bucket(bucket).send().await?;
put(&root, bucket, "folder/null.txt").await?;
versioning(&root, bucket, BucketVersioningStatus::Enabled).await?;
let protected_version = put(&root, bucket, "folder/versioned.txt").await?;
put(&root, bucket, "folder/versioned.txt").await?;
let marker = root.delete_object().bucket(bucket).key("folder/versioned.txt").send().await?;
let marker_version = marker
.version_id()
.ok_or("versioned delete must return a marker version ID")?;
put(&root, bucket, "folder/a-allowed.txt").await?;
let policy = |version: &str| {
json!({"Version":"2012-10-17","Statement":[
{"Effect":"Allow","Action":["s3:DeleteObject","s3:DeleteObjectVersion"],"Resource":format!("arn:aws:s3:::{bucket}/*")},
{"Effect":"Deny","Action":"s3:DeleteObjectVersion","Resource":format!("arn:aws:s3:::{bucket}/folder/*"),
"Condition":{"StringEquals":{"s3:VersionId":version}}}
]})
};
let user = policy_user(&env, "version-denier", Some(policy(&protected_version))).await?;
let before = versions(&root, bucket, "folder/").await?;
for version in [protected_version.as_str(), "null", marker_version] {
set_policy(&env, "version-denier", &policy(version)).await?;
assert_denied(force_delete(&user, bucket, "folder/").await);
assert_eq!(
versions(&root, bucket, "folder/").await?,
before,
"denial of a historical, null, or delete-marker version must prevent recursive deletion"
);
}
Ok(())
}
#[tokio::test]
async fn sdk_ordinary_deletion_preserves_versions_and_directory_children() -> TestResult {
init_logging();
let mut env = RustFSTestEnvironment::new().await?;
env.start_rustfs_server(vec![]).await?;
let root = env.create_s3_client();
let user = policy_user(
&env,
"ordinary-deleter",
Some(json!({"Version":"2012-10-17","Statement":[
{"Effect":"Allow","Action":"s3:DeleteObject","Resource":"arn:aws:s3:::*/*"},
{"Effect":"Deny","Action":"s3:DeleteObjectVersion","Resource":"arn:aws:s3:::*/*"}
]})),
)
.await?;
for state in ["unversioned", "enabled", "suspended"] {
let bucket = format!("ordinary-directory-{state}");
root.create_bucket().bucket(&bucket).send().await?;
if state != "unversioned" {
versioning(&root, &bucket, BucketVersioningStatus::Enabled).await?;
}
let historical = put(&root, &bucket, "object.txt").await?;
if state == "suspended" {
versioning(&root, &bucket, BucketVersioningStatus::Suspended).await?;
put(&root, &bucket, "object.txt").await?;
}
put(&root, &bucket, "folder/").await?;
put(&root, &bucket, "folder/child.txt").await?;
let child_before = versions(&root, &bucket, "folder/child.txt").await?;
user.delete_object().bucket(&bucket).key("folder/").send().await?;
assert_eq!(
versions(&root, &bucket, "folder/").await?,
child_before,
"ordinary {state} directory-key deletion must remove only its synthetic marker and preserve children"
);
let deleted = user.delete_object().bucket(&bucket).key("object.txt").send().await?;
let object_versions = versions(&root, &bucket, "object.txt").await?;
if state == "unversioned" {
assert!(object_versions.is_empty());
} else {
assert_eq!(deleted.delete_marker(), Some(true));
assert_eq!(
object_versions.len(),
2,
"ordinary {state} deletion must retain its historical data version"
);
assert!(object_versions.contains(&("object.txt".into(), historical, false)));
if state == "suspended" {
assert!(object_versions.contains(&("object.txt".into(), "null".into(), true)));
}
}
}
Ok(())
}
#[tokio::test]
async fn sdk_delete_objects_force_header_keeps_explicit_item_scope() -> TestResult {
init_logging();
let mut env = RustFSTestEnvironment::new().await?;
env.start_rustfs_server(vec![]).await?;
let root = env.create_s3_client();
let bucket = "batch-force-explicit-scope";
root.create_bucket().bucket(bucket).send().await?;
put(&root, bucket, "folder/").await?;
put(&root, bucket, "folder/child.txt").await?;
let child = versions(&root, bucket, "folder/child.txt").await?;
let user = policy_user(&env, "consoleAdmin", None).await?;
let result = user
.delete_objects()
.bucket(bucket)
.delete(
Delete::builder()
.objects(ObjectIdentifier::builder().key("folder/").build()?)
.build()?,
)
.customize()
.mutate_request(|request| {
request.headers_mut().insert("x-rustfs-force-delete", "true");
})
.send()
.await?;
assert!(result.errors().is_empty());
assert_eq!(result.deleted().len(), 1);
assert_eq!(
versions(&root, bucket, "folder/").await?,
child,
"batch deletion must remove only the explicit directory marker even with the force header"
);
Ok(())
}
#[tokio::test]
async fn force_delete_checks_every_version_page_before_mutation() -> TestResult {
init_logging();
let mut env = RustFSTestEnvironment::new().await?;
env.start_rustfs_server(vec![]).await?;
let root = env.create_s3_client();
let bucket = "force-delete-pagination";
root.create_bucket().bucket(bucket).send().await?;
versioning(&root, bucket, BucketVersioningStatus::Enabled).await?;
stream::iter(0..1000)
.map(|index| {
let root = &root;
async move { put(root, bucket, &format!("folder/{index:04}.txt")).await.map(|_| ()) }
})
.buffer_unordered(16)
.try_collect::<Vec<_>>()
.await?;
put(&root, bucket, "folder/z-denied.txt").await?;
let allow = json!({"Effect":"Allow","Action":["s3:DeleteObject","s3:DeleteObjectVersion"],"Resource":format!("arn:aws:s3:::{bucket}/*")});
let user = policy_user(
&env,
"paged-deleter",
Some(json!({"Version":"2012-10-17","Statement":[allow.clone(),
{"Effect":"Deny","Action":"s3:DeleteObjectVersion","Resource":format!("arn:aws:s3:::{bucket}/folder/z-denied.txt")}
]})),
)
.await?;
let before = versions(&root, bucket, "folder/").await?;
assert_eq!(before.len(), 1001, "the denied key must be beyond one default versions page");
assert_denied(force_delete(&user, bucket, "folder/").await);
assert_eq!(
versions(&root, bucket, "folder/").await?,
before,
"a denial on the second page must preserve the first page too"
);
set_policy(&env, "paged-deleter", &json!({"Version":"2012-10-17","Statement":[allow]})).await?;
force_delete(&user, bucket, "folder/").await?;
assert!(
versions(&root, bucket, "folder/").await?.is_empty(),
"authorized recursive deletion must cover all pages"
);
Ok(())
}
+3
View File
@@ -179,6 +179,9 @@ mod compression_test;
#[cfg(test)]
mod delete_objects_versioning_test;
#[cfg(test)]
mod delete_authorization_test;
// Regression test for signed DELETE Object?versionId requests without Content-Length.
#[cfg(test)]
mod delete_object_no_content_length_test;
+13
View File
@@ -364,6 +364,19 @@ impl ECStore {
Ok(pieces.into_guard(bucket, registration.token))
}
/// Hold this guard through recursive-delete authorization and mutation so
/// writers cannot introduce an unchecked object into the deletion scope.
pub async fn lock_bucket_for_recursive_delete(&self, bucket: &str) -> Result<rustfs_lock::NamespaceLockGuard> {
if self.ctx.lock_manager().is_disabled() {
return Err(StorageError::InvalidArgument(
bucket.to_owned(),
String::new(),
"Recursive deletion requires namespace locking".to_owned(),
));
}
self.acquire_bucket_lifecycle_write_lock(bucket).await
}
pub(crate) async fn acquire_bucket_lifecycle_write_lock(&self, bucket: &str) -> Result<rustfs_lock::NamespaceLockGuard> {
let lock = self.new_ns_lock(bucket, BUCKET_LIFECYCLE_LOCK_OBJECT).await?;
lock.get_write_lock(get_lock_acquire_timeout())
+77 -7
View File
@@ -859,9 +859,78 @@ async fn delete_recursive_prefix_with_tier_delete_journal(
}
}
}
// A trailing slash selects a directory, not the object at its parent key.
// Raw filesystem recursion would also remove that object's metadata and
// data. Preserve it by purging the selected keys individually when they
// share this physical directory. The bucket write lock covers both scans.
if object.ends_with('/') && !is_meta_bucketname(bucket) {
let parent = object.strip_suffix('/').unwrap_or(object);
for pool in &store.pools {
for set in &pool.disk_set {
let page = set
.clone()
.inner_list_object_versions_for_recursive_delete(bucket, parent, None, None, 1)
.await?;
if page.objects.iter().any(|info| info.name == parent) {
return delete_directory_keys_with_tier_delete_journal(store, bucket, object, opts, tier_journal_api).await;
}
}
}
}
delete_prefix_with_tier_delete_journal(store, bucket, object, opts, tier_journal_api).await
}
async fn delete_directory_keys_with_tier_delete_journal(
store: &ECStore,
bucket: &str,
prefix: &str,
opts: &ObjectOptions,
tier_journal_api: Option<&Arc<ECStore>>,
) -> Result<()> {
for pool in &store.pools {
for set in &pool.disk_set {
let mut previous_keys = std::collections::BTreeSet::new();
loop {
// Restart after each bounded batch: its version markers have
// been deleted, and the bucket write lock excludes new keys.
let page = set
.clone()
.inner_list_object_versions_for_recursive_delete(
bucket,
prefix,
None,
None,
RECURSIVE_DELETE_VERSION_SCAN_PAGE_SIZE,
)
.await?;
let keys = page
.objects
.into_iter()
.map(|info| info.name)
.filter(|key| key.starts_with(prefix))
.collect::<std::collections::BTreeSet<_>>();
if keys.is_empty() {
break;
}
if keys == previous_keys {
return Err(Error::other("directory deletion did not advance"));
}
for key in &keys {
let encoded_key = encode_dir_object(key);
let mut exact_opts = opts.clone();
exact_opts.delete_prefix_object = true;
let _guard = store
.acquire_object_write_lock_if_needed("delete_object", bucket, &encoded_key, &mut exact_opts)
.await?;
delete_prefix_with_tier_delete_journal(store, bucket, &encoded_key, &exact_opts, tier_journal_api).await?;
}
previous_keys = keys;
}
}
}
Ok(())
}
/// A GET whose object identity has been resolved while its namespace read lock
/// remains held, but whose body reader has not been constructed yet.
///
@@ -4682,13 +4751,14 @@ impl ECStore {
return Err(Error::other("lifecycle delete-all requires namespace locking"));
}
let _bucket_lifecycle_guard = if is_meta_bucketname(bucket) {
None
} else if opts.delete_prefix {
Some(self.acquire_bucket_lifecycle_write_lock(bucket).await?)
} else {
Some(self.acquire_bucket_lifecycle_read_lock(bucket).await?)
};
let _bucket_lifecycle_guard =
if is_meta_bucketname(bucket) || (opts.delete_prefix && opts.bucket_lifecycle_lock_fence.is_some()) {
None
} else if opts.delete_prefix {
Some(self.acquire_bucket_lifecycle_write_lock(bucket).await?)
} else {
Some(self.acquire_bucket_lifecycle_read_lock(bucket).await?)
};
let object = if opts.delete_prefix && !opts.delete_prefix_object {
object.to_owned()
} else {
+201 -42
View File
@@ -16,6 +16,67 @@
use super::*;
async fn authorize_recursive_delete<T>(
req: &mut S3Request<T>,
store: &Arc<ECStore>,
bucket: &str,
prefix: &str,
versioned: bool,
replica: bool,
) -> S3Result<()> {
let original_info = req_info_ref(req)?.clone();
let descendant_prefix = if prefix.ends_with('/') {
prefix.to_owned()
} else {
format!("{prefix}/")
};
let result = async {
let mut marker = None;
let mut version_marker = None;
loop {
let page = store
.clone()
.list_object_versions(bucket, prefix, marker.clone(), version_marker.clone(), None, 1000)
.await
.map_err(ApiError::from)?;
for object in page.objects {
// Disk prefix deletion follows path boundaries, while S3's
// string-prefix listing can also return unrelated siblings.
if object.name != prefix && !object.name.starts_with(&descendant_prefix) {
continue;
}
let object_version = object.version_id.filter(|version| !version.is_nil());
let version_id = (versioned || object_version.is_some())
.then(|| object_version.map_or_else(|| "null".to_owned(), |id| id.to_string()));
let info = req_info_mut(req)?;
info.object = Some(object.name.clone());
info.version_id = version_id.clone();
let action = if replica {
Action::S3Action(S3Action::ReplicateDeleteAction)
} else {
delete_object_authorize_action(version_id.as_deref())
};
authorize_request(req, action).await?;
if has_bypass_governance_header(&req.headers) {
authorize_request(req, Action::S3Action(S3Action::BypassGovernanceRetentionAction)).await?;
}
validate_table_catalog_object_mutation(bucket, &object.name).await?;
}
if !page.is_truncated {
return Ok(());
}
if page.next_marker.is_none() || (marker == page.next_marker && version_marker == page.next_version_idmarker) {
return Err(s3_error!(InternalError, "Recursive delete listing did not advance"));
}
marker = page.next_marker;
version_marker = page.next_version_idmarker;
}
}
.await;
req.extensions.insert(original_info);
result
}
fn successful_delete_audit_objects(
delete: &s3s::dto::Delete,
successful_results: impl IntoIterator<Item = bool>,
@@ -411,14 +472,6 @@ impl DefaultObjectUsecase {
));
}
let is_owner = req_info_ref(&req).map(|info| info.is_owner).unwrap_or(false);
if !recursive_force_delete_is_authorized(&req.headers, is_owner, false) {
return Err(S3Error::with_message(
S3ErrorCode::AccessDenied,
"Recursive force-delete is restricted to administrative requests",
));
}
let Some(store) = self.object_store() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
};
@@ -474,7 +527,7 @@ impl DefaultObjectUsecase {
req_info.version_id = version_id.clone();
}
let auth_res = authorize_request(&mut req, Action::S3Action(S3Action::DeleteObjectAction)).await;
let auth_res = authorize_request(&mut req, delete_object_authorize_action(version_id.as_deref())).await;
if auth_res.is_err() {
if !bulk_denial_logged {
bulk_denial_logged = true;
@@ -882,11 +935,11 @@ impl DefaultObjectUsecase {
authorize_request(&mut req, Action::S3Action(S3Action::ReplicateDeleteAction)).await?;
}
let is_owner = req_info_ref(&req).map(|info| info.is_owner).unwrap_or(false);
if !recursive_force_delete_is_authorized(&req.headers, is_owner, replica) {
let authenticated = req_info_ref(&req).is_ok_and(|info| info.is_owner || info.cred.is_some());
if !recursive_force_delete_has_authenticated_caller(&req.headers, authenticated, replica) {
return Err(S3Error::with_message(
S3ErrorCode::AccessDenied,
"Recursive force-delete is restricted to internal or administrative requests",
"Recursive force-delete requires an authenticated caller",
));
}
validate_table_catalog_object_mutation(&bucket, &key).await?;
@@ -901,6 +954,22 @@ impl DefaultObjectUsecase {
};
validate_bucket_exists(&store, &bucket).await?;
// Lock order is bucket lifecycle, then object/commit locks in storage.
// Keep this guard alive through the physical delete: a preflight without
// writer exclusion could authorize one subtree and delete a newer one.
let recursive_delete_guard = if rustfs_utils::http::get_header(&req.headers, rustfs_utils::http::SUFFIX_FORCE_DELETE)
.is_some_and(|value| value == "true")
{
Some(
store
.lock_bucket_for_recursive_delete(&bucket)
.await
.map_err(ApiError::from)?,
)
} else {
None
};
let metadata = extract_metadata(&req.headers);
// Clone version_id before it's moved
let version_id_clone = version_id.clone();
@@ -924,6 +993,12 @@ impl DefaultObjectUsecase {
apply_bucket_generation_guard(&req, &bucket, &mut opts)?;
let force_delete = opts.delete_prefix;
if let Some(guard) = &recursive_delete_guard {
opts.add_bucket_lifecycle_lock_guard(guard);
authorize_recursive_delete(&mut req, &store, &bucket, &key, opts.versioned || opts.version_suspended, replica)
.await?;
}
// let mut vid = opts.version_id.clone();
if replica {
@@ -1026,6 +1101,7 @@ impl DefaultObjectUsecase {
}
}
};
drop(recursive_delete_guard);
if force_delete {
let _ = invalidate_object_data_cache_prefix_after_delete(&cache_adapter, &bucket, &key).await;
@@ -2212,14 +2288,122 @@ mod tests {
}
#[test]
fn recursive_force_delete_requires_administrative_or_replica_context() {
#[serial_test::serial]
fn recursive_delete_holds_writer_exclusion_after_authorization() {
crate::app::gating_test_env::run_large_stack_test("recursive-delete-writer-exclusion", || async {
use crate::app::storage_api::test::contract::bucket::{
BucketOperations as _, DeleteBucketOptions, MakeBucketOptions,
};
use std::time::Duration;
let store = crate::app::gating_test_env::shared_gating_ecstore().await;
if current_app_context().is_none() {
crate::app::runtime_sources::install_test_app_context(Arc::clone(&store)).await;
}
let context = current_app_context().expect("recursive delete test requires an AppContext");
let bucket = format!("recursive-delete-writer-{}", Uuid::new_v4().simple());
store
.make_bucket(&bucket, &MakeBucketOptions::default())
.await
.expect("create test bucket");
let mut reader = PutObjReader::from_vec(b"old".to_vec());
store
.put_object(&bucket, "folder/old", &mut reader, &ObjectOptions::default())
.await
.expect("seed old object");
let policy_json = format!(
r#"{{"Version":"2012-10-17","Statement":[{{"Effect":"Allow","Principal":{{"AWS":"*"}},"Action":["s3:DeleteObject","s3:DeleteObjectVersion"],"Resource":["arn:aws:s3:::{bucket}/*"]}}]}}"#
);
let mut metadata = (*crate::storage::get_bucket_metadata(&bucket)
.await
.expect("load test metadata"))
.clone();
metadata.policy_config = Some(serde_json::from_str(&policy_json).expect("parse test policy"));
metadata.policy_config_json = policy_json.into_bytes();
crate::storage::storage_api::set_bucket_metadata(bucket.clone(), metadata)
.await
.expect("publish test policy");
let input = DeleteObjectInput::builder()
.bucket(bucket.clone())
.key("folder/".to_owned())
.build()
.expect("build force delete");
let mut req = build_request(input, Method::DELETE);
req.headers.insert("x-rustfs-force-delete", HeaderValue::from_static("true"));
req.extensions.insert(crate::storage::access::ReqInfo {
is_owner: true,
bucket: Some(bucket.clone()),
object: Some("folder/".to_owned()),
..Default::default()
});
let loaded = Arc::new(tokio::sync::Barrier::new(2));
let resume = Arc::new(tokio::sync::Barrier::new(2));
install_delete_source_test_hook(bucket.clone(), Arc::clone(&loaded), Arc::clone(&resume));
let usecase = DefaultObjectUsecase::with_context(Some(context));
let delete = tokio::spawn(async move { usecase.execute_delete_object(req).await });
tokio::time::timeout(Duration::from_secs(30), loaded.wait())
.await
.expect("force delete reaches authorized pre-commit pause");
let writer_store = Arc::clone(&store);
let writer_bucket = bucket.clone();
let mut writer = tokio::spawn(async move {
let mut reader = PutObjReader::from_vec(b"new".to_vec());
writer_store
.put_object(&writer_bucket, "folder/new", &mut reader, &ObjectOptions::default())
.await
});
let before_delete = tokio::time::timeout(Duration::from_secs(1), &mut writer).await;
resume.wait().await;
tokio::time::timeout(Duration::from_secs(30), delete)
.await
.expect("force delete completes without lock recursion")
.expect("delete task joins")
.expect("force delete succeeds");
assert!(
before_delete.is_err(),
"a writer must not enter the authorized subtree before deletion commits"
);
tokio::time::timeout(Duration::from_secs(30), writer)
.await
.expect("writer resumes after deletion")
.expect("writer task joins")
.expect("writer succeeds");
store
.get_object_info(&bucket, "folder/new", &ObjectOptions::default())
.await
.expect("post-delete writer's object survives");
assert!(
store
.get_object_info(&bucket, "folder/old", &ObjectOptions::default())
.await
.is_err(),
"authorized old object is removed"
);
store
.delete_bucket(
&bucket,
&DeleteBucketOptions {
force: true,
..Default::default()
},
)
.await
.expect("remove test bucket");
});
}
#[test]
fn recursive_force_delete_requires_authenticated_or_replica_context() {
let mut headers = HeaderMap::new();
headers.insert("x-rustfs-force-delete", HeaderValue::from_static("true"));
assert!(!recursive_force_delete_is_authorized(&headers, false, false));
assert!(recursive_force_delete_is_authorized(&headers, true, false));
assert!(recursive_force_delete_is_authorized(&headers, false, true));
assert!(recursive_force_delete_is_authorized(&HeaderMap::new(), false, false));
assert!(!recursive_force_delete_has_authenticated_caller(&headers, false, false));
assert!(recursive_force_delete_has_authenticated_caller(&headers, true, false));
assert!(recursive_force_delete_has_authenticated_caller(&headers, false, true));
assert!(recursive_force_delete_has_authenticated_caller(&HeaderMap::new(), false, false));
}
#[tokio::test]
@@ -2240,31 +2424,6 @@ mod tests {
assert_eq!(err.code(), &S3ErrorCode::AccessDenied);
}
#[tokio::test]
async fn execute_delete_objects_rejects_untrusted_force_delete_before_store_access() {
let input = DeleteObjectsInput::builder()
.bucket("test-bucket".to_string())
.delete(Delete {
objects: vec![ObjectIdentifier {
key: "prefix/object".to_string(),
version_id: None,
..Default::default()
}],
quiet: None,
})
.build()
.unwrap();
let mut req = build_request(input, Method::POST);
req.headers.insert("x-rustfs-force-delete", HeaderValue::from_static("true"));
req.extensions.insert(crate::storage::access::ReqInfo::default());
let err = DefaultObjectUsecase::without_context()
.execute_delete_objects(req)
.await
.expect_err("untrusted force-delete must be rejected before storage lookup");
assert_eq!(err.code(), &S3ErrorCode::AccessDenied);
}
// 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.
+4 -3
View File
@@ -21,8 +21,9 @@ use crate::storage_api::table::get_bucket_metadata;
use super::storage_api::object_usecase::access::{
PostObjectRequestMarker, apply_bucket_generation_guard, apply_copy_source_bucket_generation_guard, authorize_request,
has_bypass_governance_header, load_bucket_generation_from_store, odm_read_generation, prepare_odm_read_generation,
recursive_force_delete_is_authorized, replication_request_authorized, req_info_mut, req_info_ref,
delete_object_authorize_action, has_bypass_governance_header, load_bucket_generation_from_store, odm_read_generation,
prepare_odm_read_generation, recursive_force_delete_has_authenticated_caller, replication_request_authorized, req_info_mut,
req_info_ref,
};
#[cfg(test)]
use super::storage_api::object_usecase::bucket::quota::BucketQuota;
@@ -64,7 +65,7 @@ pub(crate) use super::storage_api::object_usecase::concurrency::{
#[cfg(test)]
use super::storage_api::object_usecase::contract::http::HTTPPreconditions;
use super::storage_api::object_usecase::contract::namespace::NamespaceLocking;
use super::storage_api::object_usecase::contract::object::{ObjectIO as _, ObjectOperations as _};
use super::storage_api::object_usecase::contract::object::{ListOperations as _, ObjectIO as _, ObjectOperations as _};
use super::storage_api::object_usecase::contract::range::HTTPRangeSpec;
use super::storage_api::object_usecase::data_usage::{
quota_object_size, record_bucket_delete_marker_memory, record_bucket_object_delete_memory,
+5 -5
View File
@@ -266,10 +266,10 @@ pub(crate) mod access {
pub(crate) use crate::storage::storage_api::access_consumer::ReqInfo;
pub(crate) use crate::storage::storage_api::access_consumer::{
PostObjectRequestMarker, apply_bucket_generation_guard, apply_copy_source_bucket_generation_guard, authorize_request,
bucket_config_mutation_incarnation, has_bypass_governance_header, load_bucket_generation_from_store,
log_list_buckets_iam_implicit_deny, odm_read_generation, prepare_list_buckets_iam_authorization,
prepare_odm_read_generation, recursive_force_delete_is_authorized, replication_request_authorized, req_info_mut,
req_info_ref,
bucket_config_mutation_incarnation, delete_object_authorize_action, has_bypass_governance_header,
load_bucket_generation_from_store, log_list_buckets_iam_implicit_deny, odm_read_generation,
prepare_list_buckets_iam_authorization, prepare_odm_read_generation, recursive_force_delete_has_authenticated_caller,
replication_request_authorized, req_info_mut, req_info_ref,
};
}
@@ -1196,7 +1196,7 @@ pub(crate) mod object_usecase {
}
pub(crate) mod object {
pub(crate) use super::super::super::storage_contracts::{ObjectIO, ObjectOperations};
pub(crate) use super::super::super::storage_contracts::{ListOperations, ObjectIO, ObjectOperations};
}
pub(crate) mod range {
+78 -206
View File
@@ -56,6 +56,10 @@ use std::sync::Arc;
use std::sync::OnceLock;
use url::{Url, form_urlencoded};
const EVENT_OBJECT_TAG_AUTHORIZATION: &str = "object_tag_authorization";
const LOG_COMPONENT_ACCESS: &str = "storage_access";
const LOG_SUBSYSTEM_AUTHORIZATION: &str = "authorization";
#[derive(Default, Clone, Debug)]
pub(crate) struct ReqInfo {
pub cred: Option<rustfs_credentials::Credentials>,
@@ -102,9 +106,13 @@ async fn authorize_replication_only_put_headers<T>(req: &mut S3Request<T>) -> S3
Ok(())
}
pub(crate) fn recursive_force_delete_is_authorized(headers: &HeaderMap, is_owner: bool, replica_request: bool) -> bool {
pub(crate) fn recursive_force_delete_has_authenticated_caller(
headers: &HeaderMap,
authenticated: bool,
replica_request: bool,
) -> bool {
!get_header(headers, SUFFIX_FORCE_DELETE).is_some_and(|value| value.eq_ignore_ascii_case("true"))
|| is_owner
|| authenticated
|| replica_request
}
@@ -833,15 +841,12 @@ pub(crate) fn log_list_buckets_iam_implicit_deny<T>(req: &S3Request<T>) -> S3Res
Ok(())
}
/// Extra action that may be evaluated in the same authorization flow and can
/// independently require `ExistingObjectTag` conditions.
fn secondary_tag_hint_action(action: Action, version_id: Option<&str>) -> Option<Action> {
match action {
Action::S3Action(S3Action::DeleteObjectAction) if version_id.is_some() => {
Some(Action::S3Action(S3Action::DeleteObjectVersionAction))
}
_ => None,
}
pub(crate) fn delete_object_authorize_action(version_id: Option<&str>) -> Action {
Action::S3Action(if version_id.is_some() {
S3Action::DeleteObjectVersionAction
} else {
S3Action::DeleteObjectAction
})
}
/// GHSA-3ppv: select the IAM action for an object read by whether the request
@@ -1015,14 +1020,14 @@ pub async fn authorize_request<T>(req: &mut S3Request<T>, action: Action) -> S3R
deny_only: false,
};
let prepared = iam_store.prepare_auth(&action_args).await;
let mut needs_tag_from_iam = prepared.needs_existing_object_tag;
let needs_tag_from_iam = prepared.needs_existing_object_tag;
let bucket_tag_hint = if !bucket.is_empty() && !object.is_empty() {
Some(load_bucket_policy_existing_object_tag_hint(store.as_ref(), bucket.as_str(), action).await)
} else {
None
};
let mut needs_tag_from_bucket = if let Some(hint) = bucket_tag_hint.as_ref() {
let needs_tag_from_bucket = if let Some(hint) = bucket_tag_hint.as_ref() {
let bucket_args = BucketPolicyArgs {
bucket: bucket.as_str(),
action,
@@ -1037,44 +1042,18 @@ pub async fn authorize_request<T>(req: &mut S3Request<T>, action: Action) -> S3R
false
};
let secondary_action = secondary_tag_hint_action(action, version_id.as_deref());
if let Some(extra_action) = secondary_action {
let extra_args = Args {
account: &cred.access_key,
groups: &cred.groups,
action: extra_action,
bucket: bucket.as_str(),
conditions: &conditions,
is_owner,
object: object.as_str(),
claims,
deny_only: false,
};
needs_tag_from_iam |= prepared.needs_existing_object_tag_for_args(&extra_args).await;
if let Some(hint) = bucket_tag_hint.as_ref() {
let extra_bucket_args = BucketPolicyArgs {
bucket: bucket.as_str(),
action: extra_action,
is_owner,
account: cred.access_key.as_str(),
groups: &cred.groups,
conditions: &conditions,
object: object.as_str(),
};
needs_tag_from_bucket |= bucket_policy_needs_existing_object_tag_from_hint(hint, &extra_bucket_args).await;
}
}
let needs_tag = needs_tag_from_iam || needs_tag_from_bucket;
if needs_tag {
tracing::debug!(
event = EVENT_OBJECT_TAG_AUTHORIZATION,
component = LOG_COMPONENT_ACCESS,
subsystem = LOG_SUBSYSTEM_AUTHORIZATION,
anonymous = false,
bucket = %bucket,
?action,
?secondary_action,
needs_tag_from_iam,
needs_tag_from_bucket,
"authorize_request ExistingObjectTag hint requires tag conditions"
"Object tag authorization conditions required"
);
}
maybe_merge_object_tag_conditions(
@@ -1116,41 +1095,8 @@ pub async fn authorize_request<T>(req: &mut S3Request<T>, action: Action) -> S3R
return Err(denial.deny("bucket_policy_explicit_deny", action));
}
if action == Action::S3Action(S3Action::DeleteObjectAction) && version_id.is_some() {
let delete_version_args = Args {
account: &cred.access_key,
groups: &cred.groups,
action: Action::S3Action(S3Action::DeleteObjectVersionAction),
bucket: bucket.as_str(),
conditions: &conditions,
is_owner,
object: object.as_str(),
claims,
deny_only: false,
};
let delete_version_allowed = iam_store.eval_prepared(&prepared, &delete_version_args).await;
if !delete_version_allowed
&& !PolicySys::try_is_allowed_for_store(
store.as_ref(),
&BucketPolicyArgs {
bucket: bucket.as_str(),
action: Action::S3Action(S3Action::DeleteObjectVersionAction),
is_owner,
account: &cred.access_key,
groups: &cred.groups,
conditions: &conditions,
object: object.as_str(),
},
)
.await
.map_err(ApiError::from)?
{
return Err(denial.deny("delete_object_version_denied", Action::S3Action(S3Action::DeleteObjectVersionAction)));
}
}
let iam_allowed = {
let final_args = Args {
let mut final_args = Args {
account: &cred.access_key,
groups: &cred.groups,
action,
@@ -1161,7 +1107,28 @@ pub async fn authorize_request<T>(req: &mut S3Request<T>, action: Action) -> S3R
claims,
deny_only: false,
};
iam_store.eval_prepared(&prepared, &final_args).await
let allowed = iam_store.eval_prepared(&prepared, &final_args).await;
if !allowed
&& matches!(
action,
Action::S3Action(
S3Action::DeleteObjectAction
| S3Action::DeleteObjectVersionAction
| S3Action::ListBucketVersionsAction
| S3Action::BypassGovernanceRetentionAction
| S3Action::ReplicateDeleteAction
)
)
&& prepared.combined_policy_for_view().is_some()
{
// Bucket policy Allow may supplement an implicit IAM denial,
// but must not override an explicit deletion-policy Deny.
final_args.deny_only = true;
if !iam_store.eval_prepared(&prepared, &final_args).await {
return Err(denial.deny("iam_explicit_deny", action));
}
}
allowed
};
if iam_allowed {
@@ -1199,42 +1166,6 @@ pub async fn authorize_request<T>(req: &mut S3Request<T>, action: Action) -> S3R
}
return Ok(());
}
if action == Action::S3Action(S3Action::ListBucketVersionsAction) {
let list_bucket_args = Args {
account: &cred.access_key,
groups: &cred.groups,
action: Action::S3Action(S3Action::ListBucketAction),
bucket: bucket.as_str(),
conditions: &conditions,
is_owner,
object: object.as_str(),
claims,
deny_only: false,
};
let list_bucket_allowed = iam_store.eval_prepared(&prepared, &list_bucket_args).await;
if list_bucket_allowed {
return Ok(());
}
if PolicySys::try_is_allowed_for_store(
store.as_ref(),
&BucketPolicyArgs {
bucket: bucket.as_str(),
action: Action::S3Action(S3Action::ListBucketAction),
is_owner,
account: &cred.access_key,
groups: &cred.groups,
conditions: &conditions,
object: object.as_str(),
},
)
.await
.map_err(ApiError::from)?
{
return Ok(());
}
}
} else {
let default_cred = rustfs_credentials::Credentials::default();
let client_info = req.extensions.get::<ClientInfo>();
@@ -1254,7 +1185,7 @@ pub async fn authorize_request<T>(req: &mut S3Request<T>, action: Action) -> S3R
} else {
None
};
let mut needs_tag_from_bucket = if let Some(hint) = bucket_tag_hint.as_ref() {
let needs_tag_from_bucket = if let Some(hint) = bucket_tag_hint.as_ref() {
let bucket_args = BucketPolicyArgs {
bucket: bucket.as_str(),
action,
@@ -1268,27 +1199,15 @@ pub async fn authorize_request<T>(req: &mut S3Request<T>, action: Action) -> S3R
} else {
false
};
let secondary_action = secondary_tag_hint_action(action, version_id.as_deref());
if let Some(extra_action) = secondary_action
&& let Some(hint) = bucket_tag_hint.as_ref()
{
let extra_bucket_args = BucketPolicyArgs {
bucket: bucket.as_str(),
action: extra_action,
is_owner: false,
account: "",
groups: &no_groups,
conditions: &conditions,
object: object.as_str(),
};
needs_tag_from_bucket |= bucket_policy_needs_existing_object_tag_from_hint(hint, &extra_bucket_args).await;
}
if needs_tag_from_bucket {
tracing::debug!(
event = EVENT_OBJECT_TAG_AUTHORIZATION,
component = LOG_COMPONENT_ACCESS,
subsystem = LOG_SUBSYSTEM_AUTHORIZATION,
anonymous = true,
bucket = %bucket,
?action,
?secondary_action,
"anonymous authorize_request ExistingObjectTag hint requires tag conditions"
"Object tag authorization conditions required"
);
}
maybe_merge_object_tag_conditions(
@@ -1324,28 +1243,6 @@ pub async fn authorize_request<T>(req: &mut S3Request<T>, action: Action) -> S3R
}
if action != Action::S3Action(S3Action::ListAllMyBucketsAction) {
if action == Action::S3Action(S3Action::DeleteObjectAction) && version_id.is_some() {
let delete_version_allowed = PolicySys::try_is_allowed_for_store(
store.as_ref(),
&BucketPolicyArgs {
bucket: bucket.as_str(),
action: Action::S3Action(S3Action::DeleteObjectVersionAction),
is_owner: false,
account: "",
groups: &None,
conditions: &conditions,
object: object.as_str(),
},
)
.await
.map_err(ApiError::from)?;
if !delete_version_allowed {
return Err(
denial.deny("delete_object_version_denied", Action::S3Action(S3Action::DeleteObjectVersionAction))
);
}
}
let policy_allowed = PolicySys::try_is_allowed_for_store(
store.as_ref(),
&BucketPolicyArgs {
@@ -1361,27 +1258,6 @@ pub async fn authorize_request<T>(req: &mut S3Request<T>, action: Action) -> S3R
.await
.map_err(ApiError::from)?;
// A bucket policy granting s3:ListBucket also covers listing versions. This
// fallback has to feed the same post-authorization gates as the direct grant
// below, otherwise a public bucket keeps serving anonymous
// ListObjectVersions after RestrictPublicBuckets is turned on.
let policy_allowed = policy_allowed
|| (action == Action::S3Action(S3Action::ListBucketVersionsAction)
&& PolicySys::try_is_allowed_for_store(
store.as_ref(),
&BucketPolicyArgs {
bucket: bucket.as_str(),
action: Action::S3Action(S3Action::ListBucketAction),
is_owner: false,
account: "",
groups: &None,
conditions: &conditions,
object: "",
},
)
.await
.map_err(ApiError::from)?);
if policy_allowed {
deny_anonymous_table_data_plane_if_needed(req, action, bucket.as_str(), object.as_str()).await?;
// RestrictPublicBuckets: when true, deny public access even if bucket policy allows it.
@@ -2173,20 +2049,18 @@ impl S3Access for FS {
req_info.bucket = Some(req.input.bucket.clone());
req_info.object = Some(req.input.key.clone());
req_info.version_id = req.input.version_id.clone();
let is_owner = req_info.is_owner;
let authenticated = req_info.is_owner || req_info.cred.is_some();
let action = delete_object_authorize_action(req_info.version_id.as_deref());
authorize_request(req, Action::S3Action(S3Action::DeleteObjectAction)).await?;
authorize_request(req, action).await?;
let replica_request = req
.headers
.get(AMZ_BUCKET_REPLICATION_STATUS)
.and_then(|value| value.to_str().ok())
.is_some_and(|value| value == ReplicationStatusType::Replica.as_str());
if !recursive_force_delete_is_authorized(&req.headers, is_owner, replica_request) {
return Err(s3_error!(
AccessDenied,
"Recursive force-delete is restricted to internal or administrative requests"
));
if !recursive_force_delete_has_authenticated_caller(&req.headers, authenticated, replica_request) {
return Err(s3_error!(AccessDenied, "Recursive force-delete requires an authenticated caller"));
}
// S3 Standard: When bypass_governance header is set, must have s3:BypassGovernanceRetention permission
@@ -3108,12 +2982,12 @@ mod tests {
PostObjectRequestMarker, ReqInfo, S3Access, StorageError, TableDataPlanePublicationGuards, apply_bucket_generation_guard,
apply_copy_source_bucket_generation_guard, authorization_conditions, bucket_policy_needs_existing_object_tag_from_hint,
bucket_website_config_authorize_action, classify_bucket_policy_raw_load_error,
complete_multipart_upload_authorize_action, get_bucket_policy_authorize_action, has_write_offset_bytes_header,
install_restore_authorization_test_hook, legal_hold_write_requested, list_parts_authorize_action,
load_bucket_policy_existing_object_tag_hint, maybe_merge_object_tag_conditions, merge_list_bucket_query_conditions,
merge_request_object_tag_conditions, owner_can_bypass_policy_deny, post_object_authorize_action,
put_bucket_policy_authorize_action, request_context_from_req, request_object_store, retention_write_requested,
secondary_tag_hint_action, table_data_plane_admin_action, table_data_plane_content_mutation,
complete_multipart_upload_authorize_action, delete_object_authorize_action, get_bucket_policy_authorize_action,
has_write_offset_bytes_header, install_restore_authorization_test_hook, legal_hold_write_requested,
list_parts_authorize_action, load_bucket_policy_existing_object_tag_hint, maybe_merge_object_tag_conditions,
merge_list_bucket_query_conditions, merge_request_object_tag_conditions, owner_can_bypass_policy_deny,
post_object_authorize_action, put_bucket_policy_authorize_action, request_context_from_req, request_object_store,
retention_write_requested, table_data_plane_admin_action, table_data_plane_content_mutation,
table_data_plane_resource_for_request, table_publication_guard_error, validate_post_object_success_controls,
versioned_read_action,
};
@@ -3946,20 +3820,18 @@ mod tests {
}
#[test]
fn test_secondary_tag_hint_action_for_delete_object_version() {
assert_eq!(
secondary_tag_hint_action(Action::S3Action(S3Action::DeleteObjectAction), Some("v1")),
Some(Action::S3Action(S3Action::DeleteObjectVersionAction))
);
assert_eq!(secondary_tag_hint_action(Action::S3Action(S3Action::DeleteObjectAction), None), None);
assert_eq!(
secondary_tag_hint_action(Action::S3Action(S3Action::ListBucketVersionsAction), None),
None
);
fn delete_authorization_selects_the_addressed_version() {
assert_eq!(delete_object_authorize_action(None), Action::S3Action(S3Action::DeleteObjectAction));
for version_id in ["null", "8f418ad0-f9f4-4458-83b2-cc72bc6f1b70"] {
assert_eq!(
delete_object_authorize_action(Some(version_id)),
Action::S3Action(S3Action::DeleteObjectVersionAction)
);
}
}
#[tokio::test]
async fn test_anonymous_delete_object_with_version_requires_secondary_policy_and_tag_hint() {
async fn test_anonymous_version_delete_uses_version_policy_and_tag_hint() {
let policy: BucketPolicy = serde_json::from_str(
r#"{
"Version":"2012-10-17",
@@ -4013,16 +3885,16 @@ mod tests {
"DeleteObjectVersion should still be denied without matching ExistingObjectTag conditions"
);
let needs_tag_main = bucket_policy_needs_existing_object_tag_from_hint(&hint, &args_delete).await;
let needs_tag_secondary = bucket_policy_needs_existing_object_tag_from_hint(&hint, &args_delete_version).await;
assert!(!needs_tag_main, "DeleteObject statement itself does not require ExistingObjectTag");
let needs_tag_current = bucket_policy_needs_existing_object_tag_from_hint(&hint, &args_delete).await;
let needs_tag_version = bucket_policy_needs_existing_object_tag_from_hint(&hint, &args_delete_version).await;
assert!(!needs_tag_current, "DeleteObject statement itself does not require ExistingObjectTag");
assert!(
needs_tag_secondary,
needs_tag_version,
"DeleteObjectVersion statement requires ExistingObjectTag when version delete is evaluated"
);
assert!(
needs_tag_main || needs_tag_secondary,
"combined primary+secondary check must require tag fetch for DeleteObject(versionId)"
needs_tag_version,
"the selected version action must require tag fetch for DeleteObject(versionId)"
);
}
+3 -3
View File
@@ -119,9 +119,9 @@ pub(crate) use super::sse::{
pub(crate) mod access_consumer {
pub(crate) use super::super::access::{
PostObjectRequestMarker, ReqInfo, apply_bucket_generation_guard, apply_copy_source_bucket_generation_guard,
authorize_internal_object_request, authorize_request, bucket_config_mutation_incarnation, has_bypass_governance_header,
load_bucket_generation_from_store, log_list_buckets_iam_implicit_deny, odm_read_generation,
prepare_list_buckets_iam_authorization, prepare_odm_read_generation, recursive_force_delete_is_authorized,
authorize_internal_object_request, authorize_request, bucket_config_mutation_incarnation, delete_object_authorize_action,
has_bypass_governance_header, load_bucket_generation_from_store, log_list_buckets_iam_implicit_deny, odm_read_generation,
prepare_list_buckets_iam_authorization, prepare_odm_read_generation, recursive_force_delete_has_authenticated_caller,
replication_request_authorized, req_info_mut, req_info_ref,
};
}