feat(iam): expose stored bucket tags to OPA (#8298)

* feat(iam): expose stored bucket tags to OPA

* fix(ci): lint bucket tag fixtures and register new e2e tests

* fix(ecstore): restore strict Clippy compatibility on Rust 1.99

Use try_update without changing atomic ordering or overflow behavior. Keep
recursive storage futures boxed once at each frame and remove the redundant
async-recursion macro, including its non-recursive SQL planner use. Remove
needless closure borrows and orphaned dependency entries.

Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
(cherry picked from commit 537c15277a)

* fix(deps): replace yanked yoke-derive release

(cherry picked from commit b54b32b8a9)

* perf(iam): borrow unchanged OPA bucket tag conditions

Reuse clean conditions for bucketless and untagged OPA decisions. Keep the existing clone-and-sanitize path for tagged requests, reserve enrichment capacity, and preserve authoritative lookup errors. Cover borrowed input, reserved namespaces, immutable request state, and both paths in the OPA input regression.

Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>

* test(iam): clarify OPA bucket policy fallback semantics

Document that an OPA tag mismatch remains an implicit IAM denial which a bucket-policy Allow can supplement. Extend the existing S3 contract test to contrast a normal OPA denial with corrupt bucket metadata, whose lookup error must abort the same Allow fallback.

Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>

* fix(ci): isolate durable progress fixtures from competing IO

Reserve the nextest execution budget for schedule, streaming lock, retirement, MRF checkpoint/crash, and native scanner progress scenarios. Preserve their internal concurrency, existing deadlines, typed failure assertions, and zero-retry policy in default and CI profiles.

Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>

---------

Co-authored-by: rjregenold <214054+rjregenold@users.noreply.github.com>
Co-authored-by: Hauser <housemecn@gmail.com>
Co-authored-by: heihutu <heihutu@gmail.com>
Co-authored-by: zhi22915 <qiuzgang@gmail.com>
Co-authored-by: overtrue <anzhengchao@gmail.com>
This commit is contained in:
RJ Regenold
2026-10-03 04:18:04 -05:00
committed by GitHub
parent 4804162be8
commit 48a6104aad
12 changed files with 783 additions and 23 deletions
+2 -2
View File
@@ -1,2 +1,2 @@
sha256-linux=6fc377fa1f9f06e065f077c7bd53d04efbccf4453f95dde67afaa47778185377
sha256-darwin=6fc377fa1f9f06e065f077c7bd53d04efbccf4453f95dde67afaa47778185377
sha256-linux=233bd7a68777eda1f044b6cb77128056be2b6492583c0dfd114b21e98a6102ef
sha256-darwin=233bd7a68777eda1f044b6cb77128056be2b6492583c0dfd114b21e98a6102ef
+2 -2
View File
@@ -1,2 +1,2 @@
sha256-darwin=6b8e35ef69456bd244d6a6b19179e408763763e1e163aaaf90055490ae12da4b
sha256-linux=ff4c40f288ab5e91fd28d533c277538c0e3281fa353328ab52f5a8f88e0dc931
sha256-darwin=98dff0961c6eb21d6128bfd2d4a6100a624ddf1fb7db05dba360a5d3d0a310c2
sha256-linux=9ff0da7406de9b89c0e848ff29f2238b30ba9da89c7e9df97c246ce86631350f
+1 -1
View File
@@ -1 +1 @@
sha256=7e3989625e2e6087c2e7db7ee02d30624b566d5a99832122032fb413ff456662
sha256=ba8f392b617cf9f2338142233a417e43927b28096635e628233ce8a60cf49d10
+38
View File
@@ -64,6 +64,26 @@ command = ['sh', '-c', 'echo RUST_MIN_STACK=33554432 >> "$NEXTEST_ENV"']
command = ['sh', '-c', 'echo RUST_MIN_STACK=33554432 >> "$NEXTEST_ENV"']
# --- default profile (local): serialize the flaky groups, never retry --------
# These real-disk progress oracles exercise fsync, quorum IO, and crash/replay
# boundaries. Reserve the run's capacity so unrelated fixtures cannot starve
# their existing deadlines and typed-failure assertions (rustfs/backlog#2706).
# Keep each test's internal concurrency and zero-retry policy unchanged.
[[profile.default.overrides]]
filter = 'package(rustfs) & test(/^connect::diagnostics::schedule::tests::/)'
threads-required = "num-test-threads"
[[profile.default.overrides]]
filter = 'package(rustfs-ecstore) & (test(/^set_disk::tests::streaming_get_blocks_concurrent_/) | test(=store::object::tests::retired_marker_converges_after_bucket_recreation_on_sixteen_disks))'
threads-required = "num-test-threads"
[[profile.default.overrides]]
filter = 'package(rustfs-heal) & (binary(mrf_partial_write_test) | (binary(mrf_pipeline_test) & test(/^journal_replay_survives_/)) | test(=admin_selector_survives_sigkill_with_original_scope_and_receipts))'
threads-required = "num-test-threads"
[[profile.default.overrides]]
filter = 'package(rustfs-scanner) & (test(/^scanner::backlog::tests::native_(writer|retirement)_/) | test(=scanner_folder::tests::mrf_ownership::mrf_ownership_manager_completion_preserves_scanner_pending) | test(/^scanner_io::tests::service_cohort::/))'
threads-required = "num-test-threads"
[[profile.default.scripts]]
filter = 'package(rustfs-ecstore) & test(/^(bucket::lifecycle::bucket_lifecycle_ops::tests::manual_transition_worker_result_recovery_marks_unknown_for_corrupt_marker|services::rebalance::entry::tests::real_rebalance_run_fence_loss_blocks_multipart_publication|store::init::tests::(batch_transitioned_delete_uses_free_version_per_item|decommission_entry_(allows_free_version_consumed_before_source_lock|rejects_subquorum_free_version_conflict_and_retains_source|skips_cleanup_only_marker_when_free_version_is_present)|dispatched_tier_delete_recovery_(checks_later_pool_then_commits_after_source_removal|finds_directory_source_on_encoded_set|retains_journal_on_source_metadata_error)|force_tier_remove_blocks_on_physical_free_version_hidden_by_other_pool|legacy_unknown_transition_delete_falls_back_for_single_batch_and_blocks_prefix|multi_pool_(recursive_prefix_rejects_legacy_or_hidden_merge_loser_before_delete|same_remote_tuple_(batch|single)_delete_waits_for_all_sources|same_tuple_recursive_prefix_uses_one_journal_owner|transitioned_delete_persists_one_free_version_per_remote_tuple)|recursive_prefix_partial_(pool|set)_failure_keeps_prepared_cleanup_owners|restored_transitioned_delete_uses_free_version_as_cleanup_owner|stable_transitioned_recursive_prefix_delete_uses_journal_owners|suspended_null_transition_delete_uses_free_version_as_sole_owner|tier_mutation_peer_handler_applies_prepare_commit_and_abort_idempotently|transition_response_loss_persists_unknown_outcome_for_provider_recovery|transition_transaction_recovery_(drops_record_after_confirmed_local_commit|keeps_cleanup_pending_local_commit)|transitioned_delete_(free_version_replays_after_store_restart|local_quorum_failure_rolls_back_without_cleanup_owner|uses_free_version_as_cleanup_owner)|versioned_delete_marker_keeps_transitioned_source_and_remote_object|versioned_explicit_transition_delete_preserves_other_version_then_allows_bucket_delete))$/)'
setup = 'ecstore-large-stack'
@@ -290,6 +310,24 @@ fail-fast = false
# marker is the observable signal the flake policy is built around.
path = "junit.xml"
# Match the local durable-progress reservations. All members, deadlines,
# assertions, and internal concurrency remain enabled without retries.
[[profile.ci.overrides]]
filter = 'package(rustfs) & test(/^connect::diagnostics::schedule::tests::/)'
threads-required = "num-test-threads"
[[profile.ci.overrides]]
filter = 'package(rustfs-ecstore) & (test(/^set_disk::tests::streaming_get_blocks_concurrent_/) | test(=store::object::tests::retired_marker_converges_after_bucket_recreation_on_sixteen_disks))'
threads-required = "num-test-threads"
[[profile.ci.overrides]]
filter = 'package(rustfs-heal) & (binary(mrf_partial_write_test) | (binary(mrf_pipeline_test) & test(/^journal_replay_survives_/)) | test(=admin_selector_survives_sigkill_with_original_scope_and_receipts))'
threads-required = "num-test-threads"
[[profile.ci.overrides]]
filter = 'package(rustfs-scanner) & (test(/^scanner::backlog::tests::native_(writer|retirement)_/) | test(=scanner_folder::tests::mrf_ownership::mrf_ownership_manager_completion_preserves_scanner_pending) | test(/^scanner_io::tests::service_cohort::/))'
threads-required = "num-test-threads"
[[profile.ci.scripts]]
filter = 'package(rustfs-ecstore) & test(/^(bucket::lifecycle::bucket_lifecycle_ops::tests::manual_transition_worker_result_recovery_marks_unknown_for_corrupt_marker|services::rebalance::entry::tests::real_rebalance_run_fence_loss_blocks_multipart_publication|store::init::tests::(batch_transitioned_delete_uses_free_version_per_item|decommission_entry_(allows_free_version_consumed_before_source_lock|rejects_subquorum_free_version_conflict_and_retains_source|skips_cleanup_only_marker_when_free_version_is_present)|dispatched_tier_delete_recovery_(checks_later_pool_then_commits_after_source_removal|finds_directory_source_on_encoded_set|retains_journal_on_source_metadata_error)|force_tier_remove_blocks_on_physical_free_version_hidden_by_other_pool|legacy_unknown_transition_delete_falls_back_for_single_batch_and_blocks_prefix|multi_pool_(recursive_prefix_rejects_legacy_or_hidden_merge_loser_before_delete|same_remote_tuple_(batch|single)_delete_waits_for_all_sources|same_tuple_recursive_prefix_uses_one_journal_owner|transitioned_delete_persists_one_free_version_per_remote_tuple)|recursive_prefix_partial_(pool|set)_failure_keeps_prepared_cleanup_owners|restored_transitioned_delete_uses_free_version_as_cleanup_owner|stable_transitioned_recursive_prefix_delete_uses_journal_owners|suspended_null_transition_delete_uses_free_version_as_sole_owner|tier_mutation_peer_handler_applies_prepare_commit_and_abort_idempotently|transition_response_loss_persists_unknown_outcome_for_provider_recovery|transition_transaction_recovery_(drops_record_after_confirmed_local_commit|keeps_cleanup_pending_local_commit)|transitioned_delete_(free_version_replays_after_store_restart|local_quorum_failure_rolls_back_without_cleanup_owner|uses_free_version_as_cleanup_owner)|versioned_delete_marker_keeps_transitioned_source_and_remote_object|versioned_explicit_transition_delete_preserves_other_version_then_allows_bucket_delete))$/)'
setup = 'ecstore-large-stack'
@@ -20,6 +20,60 @@ use aws_sdk_s3::primitives::ByteStream;
use aws_sdk_s3::types::{Delete, MetadataDirective, ObjectIdentifier};
use std::time::Duration;
#[tokio::test]
async fn four_node_bucket_tags_opa_refreshes_warm_peer_caches() -> TestResult {
use super::harness::cluster_admin_ok;
use crate::sts_query_compat_test::{OPA_AUTH_TOKEN, OpaMock, set_department};
init_logging();
let mut opa = OpaMock::start().await?;
let dist = DistCluster::start_with_env(
DistLayout::FourNodeFourDisk,
&[
("RUSTFS_POLICY_PLUGIN_URL", opa.url.as_str()),
("RUSTFS_POLICY_PLUGIN_AUTH_TOKEN", OPA_AUTH_TOKEN),
],
)
.await?;
let secret = uuid::Uuid::new_v4().to_string();
cluster_admin_ok(
&dist.cluster,
http::Method::PUT,
"/rustfs/admin/v3/add-user?accessKey=opabuckettags",
Some(serde_json::json!({"secretKey": secret, "status": "enabled"}).to_string()),
)
.await?;
let bucket = unique_bucket("opa-tags");
dist.create_bucket(&bucket).await?;
let admin = dist.client(0)?;
put_object(&admin, &bucket, "report", b"report".to_vec()).await?;
let clients = (0..4)
.map(|node| dist.client_with_credentials(node, "opabuckettags", &secret))
.collect::<Result<Vec<_>, _>>()?;
// Warm every peer, then verify the acknowledged updates/removal without
// polling or restarting nodes. This exercises healthy-peer metadata reload.
for department in [Some("finance"), Some("engineering"), None, Some("finance")] {
match department {
Some(value) => set_department(&admin, &bucket, value).await?,
None => {
admin.delete_bucket_tagging().bucket(&bucket).send().await?;
}
}
for client in &clients {
let result = client.get_object().bucket(&bucket).key("report").send().await;
if department == Some("finance") {
assert_eq!(result?.body.collect().await?.into_bytes().as_ref(), b"report");
} else {
let error = result.expect_err("updated or removed bucket tags must revoke this policy's read access");
assert_eq!(error.as_service_error().and_then(ProvideErrorMetadata::code), Some("AccessDenied"));
}
opa.expect_bucket_tags("s3:GetObject", &bucket, department).await?;
}
}
Ok(())
}
#[tokio::test]
async fn four_node_four_drive_s3_put_get_head_list_copy_rename_delete_and_presign() -> TestResult {
init_logging();
+281 -9
View File
@@ -35,7 +35,7 @@ use tokio::time::{Duration, timeout};
type BoxError = Box<dyn Error + Send + Sync>;
type TestResult = Result<(), BoxError>;
const OPA_AUTH_TOKEN: &str = "sts-opa-token";
pub(crate) const OPA_AUTH_TOKEN: &str = "sts-opa-token";
fn sts_client(url: &str, access_key: &str, secret_key: &str, session_token: Option<&str>) -> Client {
build_test_sts_client(url, access_key, secret_key, session_token, "e2e-sts-query-compat")
@@ -103,6 +103,25 @@ async fn create_user(env: &RustFSTestEnvironment, user: &str, secret: &str) -> T
Ok(())
}
pub(crate) async fn set_department(client: &aws_sdk_s3::Client, bucket: &str, department: &str) -> TestResult {
client
.put_bucket_tagging()
.bucket(bucket)
.tagging(
aws_sdk_s3::types::Tagging::builder()
.tag_set(
aws_sdk_s3::types::Tag::builder()
.key("department")
.value(department)
.build()?,
)
.build()?,
)
.send()
.await?;
Ok(())
}
async fn assert_access_denied(client: &Client, context: &str) -> TestResult {
let error = client
.assume_role()
@@ -240,15 +259,19 @@ async fn handle_opa_request(
.as_ref()
.and_then(|value| value.pointer("/input/action"))
.and_then(Value::as_str);
let bucket = payload
let department = payload
.as_ref()
.and_then(|value| value.pointer("/input/resource/bucket"))
.and_then(|value| value.pointer("/input/context/conditions/ExistingBucketTag~1department/0"))
.and_then(Value::as_str);
matches!(
(action, bucket),
(Some("s3:ListBucket"), Some("opa-list-visible")) | (Some("s3:GetBucketLocation"), Some("opa-list-location"))
(action, department),
(Some("s3:ListBucket"), Some("finance")) | (Some("s3:GetBucketLocation"), Some("legal"))
)
}
Some(Value::String(account)) if account == "opabuckettags" => payload.as_ref().is_some_and(|value| {
value["input"]["action"] == "s3:CreateBucket"
|| value["input"]["context"]["conditions"]["ExistingBucketTag/department"] == serde_json::json!(["finance"])
}),
None => true,
_ => false,
};
@@ -270,15 +293,15 @@ enum OpaValidationMode {
Unavailable,
}
struct OpaMock {
url: String,
pub(crate) struct OpaMock {
pub(crate) url: String,
requests: mpsc::UnboundedReceiver<Value>,
validation_started: mpsc::UnboundedReceiver<()>,
task: JoinHandle<()>,
}
impl OpaMock {
async fn start() -> Result<Self, BoxError> {
pub(crate) async fn start() -> Result<Self, BoxError> {
Self::start_with_mode(OpaValidationMode::Ready, Some(OPA_AUTH_TOKEN)).await
}
@@ -335,6 +358,16 @@ impl OpaMock {
.ok_or_else(|| "OPA request channel closed".into())
}
pub(crate) async fn expect_bucket_tags(&mut self, action: &str, bucket: &str, department: Option<&str>) -> TestResult {
let payload = self.next_request().await?;
let input = &payload["input"];
assert_eq!(input["action"], action);
assert_eq!(input["resource"]["bucket"], bucket);
let expected = department.map(|value| serde_json::json!([value]));
assert_eq!(input["context"]["conditions"].get("ExistingBucketTag/department"), expected.as_ref());
Ok(())
}
async fn wait_for_validation(&mut self) -> TestResult {
timeout(Duration::from_secs(5), self.validation_started.recv())
.await?
@@ -571,8 +604,13 @@ async fn test_list_buckets_opa_contract() -> TestResult {
.await?;
let admin_client = env.create_s3_client();
for bucket in ["opa-list-hidden", "opa-list-location", "opa-list-visible"] {
for (bucket, department) in [
("opa-list-hidden", "engineering"),
("opa-list-location", "legal"),
("opa-list-visible", "finance"),
] {
admin_client.create_bucket().bucket(bucket).send().await?;
set_department(&admin_client, bucket, department).await?;
}
let user = "opalistbuckets";
@@ -640,6 +678,240 @@ async fn test_list_buckets_opa_contract() -> TestResult {
Ok(())
}
#[tokio::test]
async fn test_bucket_tags_opa_contract() -> TestResult {
use aws_sdk_s3::primitives::ByteStream;
use aws_sdk_s3::types::{
BucketVersioningStatus, CompletedMultipartUpload, CompletedPart, Tag, Tagging, VersioningConfiguration,
};
init_logging();
let mut opa = OpaMock::start().await?;
let mut env = RustFSTestEnvironment::new().await?;
env.start_rustfs_server_without_cleanup_with_env(&[
("RUSTFS_POLICY_PLUGIN_URL", opa.url.as_str()),
("RUSTFS_POLICY_PLUGIN_AUTH_TOKEN", OPA_AUTH_TOKEN),
("NO_PROXY", "127.0.0.1,localhost"),
])
.await?;
let secret = uuid::Uuid::new_v4().to_string();
create_user(&env, "opabuckettags", &secret).await?;
let client = aws_sdk_s3::Client::from_conf(build_test_s3_config(&env.url, "opabuckettags", &secret, None, "e2e-bucket-tags"));
let admin = env.create_s3_client();
let source = "opa-tags-source";
let destination = "opa-tags-destination";
client.create_bucket().bucket(source).send().await?;
opa.expect_bucket_tags("s3:CreateBucket", source, None).await?;
set_department(&admin, source, "finance").await?;
admin
.put_bucket_versioning()
.bucket(source)
.versioning_configuration(
VersioningConfiguration::builder()
.status(BucketVersioningStatus::Enabled)
.build(),
)
.send()
.await?;
let written = client
.put_object()
.bucket(source)
.key("report")
.tagging("department=engineering")
.body(ByteStream::from_static(b"report"))
.send()
.await?;
opa.expect_bucket_tags("s3:PutObject", source, Some("finance")).await?;
let version = written.version_id().ok_or("versioned PUT must return a version ID")?;
let read = client
.get_object()
.bucket(source)
.key("report")
.version_id(version)
.send()
.await?;
assert_eq!(read.body.collect().await?.into_bytes().as_ref(), b"report");
opa.expect_bucket_tags("s3:GetObjectVersion", source, Some("finance")).await?;
client
.head_object()
.bucket(source)
.key("report")
.version_id(version)
.send()
.await?;
opa.expect_bucket_tags("s3:GetObject", source, Some("finance")).await?;
client.list_objects_v2().bucket(source).send().await?;
opa.expect_bucket_tags("s3:ListBucket", source, Some("finance")).await?;
client.list_object_versions().bucket(source).send().await?;
opa.expect_bucket_tags("s3:ListBucketVersions", source, Some("finance"))
.await?;
admin.create_bucket().bucket(destination).send().await?;
set_department(&admin, destination, "engineering").await?;
let copy_source = format!("{source}/report?versionId={version}");
let error = client
.copy_object()
.bucket(destination)
.key("copied")
.copy_source(&copy_source)
.send()
.await
.expect_err("source tags must not authorize a different destination");
assert_eq!(error.as_service_error().and_then(ProvideErrorMetadata::code), Some("AccessDenied"));
opa.expect_bucket_tags("s3:GetObjectVersion", source, Some("finance")).await?;
opa.expect_bucket_tags("s3:PutObject", destination, Some("engineering"))
.await?;
set_department(&admin, destination, "finance").await?;
client
.copy_object()
.bucket(destination)
.key("copied")
.copy_source(&copy_source)
.send()
.await?;
opa.expect_bucket_tags("s3:GetObjectVersion", source, Some("finance")).await?;
opa.expect_bucket_tags("s3:PutObject", destination, Some("finance")).await?;
let upload = client
.create_multipart_upload()
.bucket(destination)
.key("multipart")
.send()
.await?;
opa.expect_bucket_tags("s3:PutObject", destination, Some("finance")).await?;
let upload_id = upload.upload_id().ok_or("multipart upload ID")?;
client
.upload_part()
.bucket(destination)
.key("multipart")
.upload_id(upload_id)
.part_number(1)
.body(ByteStream::from_static(b"part"))
.send()
.await?;
opa.expect_bucket_tags("s3:PutObject", destination, Some("finance")).await?;
let copied_part = client
.upload_part_copy()
.bucket(destination)
.key("multipart")
.upload_id(upload_id)
.part_number(1)
.copy_source(&copy_source)
.send()
.await?;
opa.expect_bucket_tags("s3:GetObjectVersion", source, Some("finance")).await?;
opa.expect_bucket_tags("s3:PutObject", destination, Some("finance")).await?;
let part = copied_part.copy_part_result().ok_or("copied part result")?;
client
.complete_multipart_upload()
.bucket(destination)
.key("multipart")
.upload_id(upload_id)
.multipart_upload(
CompletedMultipartUpload::builder()
.parts(
CompletedPart::builder()
.part_number(1)
.e_tag(part.e_tag().ok_or("part ETag")?)
.build(),
)
.build(),
)
.send()
.await?;
opa.expect_bucket_tags("s3:PutObject", destination, Some("finance")).await?;
let read = client.get_object().bucket(destination).key("multipart").send().await?;
assert_eq!(read.body.collect().await?.into_bytes().as_ref(), b"report");
opa.expect_bucket_tags("s3:GetObject", destination, Some("finance")).await?;
// Mutation authorization uses the old tags, not the proposed replacement.
set_department(&client, source, "engineering").await?;
opa.expect_bucket_tags("s3:PutBucketTagging", source, Some("finance")).await?;
let error = client
.list_objects_v2()
.bucket(source)
.send()
.await
.expect_err("updated tags revoke access");
assert_eq!(error.as_service_error().and_then(ProvideErrorMetadata::code), Some("AccessDenied"));
opa.expect_bucket_tags("s3:ListBucket", source, Some("engineering")).await?;
let error = client
.copy_object()
.bucket(destination)
.key("denied-source")
.copy_source(&copy_source)
.send()
.await
.expect_err("destination tags must not authorize a denied source");
assert_eq!(error.as_service_error().and_then(ProvideErrorMetadata::code), Some("AccessDenied"));
opa.expect_bucket_tags("s3:GetObjectVersion", source, Some("engineering"))
.await?;
set_department(&client, source, "finance")
.await
.expect_err("requested tags must not authorize their own mutation");
opa.expect_bucket_tags("s3:PutBucketTagging", source, Some("engineering"))
.await?;
set_department(&admin, source, "finance").await?;
client.delete_bucket_tagging().bucket(source).send().await?;
opa.expect_bucket_tags("s3:PutBucketTagging", source, Some("finance")).await?;
let error = client
.list_objects_v2()
.bucket(source)
.send()
.await
.expect_err("untagged bucket is denied by this policy");
assert_eq!(error.as_service_error().and_then(ProvideErrorMetadata::code), Some("AccessDenied"));
opa.expect_bucket_tags("s3:ListBucket", source, None).await?;
// Ambiguous stored tags must not turn into an implicit IAM denial that a
// separate bucket-policy Allow can override.
admin
.put_bucket_policy()
.bucket(destination)
.policy(
serde_json::json!({
"Version": "2012-10-17", "Statement": [{"Effect": "Allow", "Principal": {"AWS": "*"},
"Action": "s3:GetObject", "Resource": format!("arn:aws:s3:::{destination}/*")}]
})
.to_string(),
)
.send()
.await?;
// A normal OPA denial retains the existing bucket-policy Allow fallback;
// a metadata lookup error must abort that same authorization path.
set_department(&admin, destination, "engineering").await?;
let read = client.get_object().bucket(destination).key("multipart").send().await?;
assert_eq!(read.body.collect().await?.into_bytes().as_ref(), b"report");
opa.expect_bucket_tags("s3:GetObject", destination, Some("engineering"))
.await?;
admin
.put_bucket_tagging()
.bucket(destination)
.tagging(
Tagging::builder()
.tag_set(Tag::builder().key("department").value("finance").build()?)
.tag_set(Tag::builder().key("department").value("engineering").build()?)
.build()?,
)
.send()
.await?;
let error = client
.get_object()
.bucket(destination)
.key("multipart")
.send()
.await
.expect_err("metadata failure must not fall back to a bucket-policy Allow");
assert_eq!(error.raw_response().map(|response| response.status().as_u16()), Some(500));
assert!(
matches!(opa.requests.try_recv(), Err(mpsc::error::TryRecvError::Empty)),
"invalid tags must not be sent to OPA as absent"
);
env.stop_server();
Ok(())
}
#[tokio::test]
async fn test_sts_and_list_buckets_fail_closed_while_opa_is_initializing() -> TestResult {
init_logging();
+81 -1
View File
@@ -14,7 +14,7 @@
use super::*;
use crate::bucket::{
metadata::{BUCKET_TABLE_RESERVED_PREFIX, table_bucket_catalog_metadata_prefix},
metadata::{BUCKET_TABLE_RESERVED_PREFIX, BUCKET_TAGGING_CONFIG, ConfigState, table_bucket_catalog_metadata_prefix},
utils::is_meta_bucketname,
};
use crate::error::is_err_bucket_not_found;
@@ -272,6 +272,32 @@ impl ECStore {
sys.read().await.get(bucket).await
}
/// Resolve stored tags for authorization. Only confirmed absence returns an
/// empty map; unavailable, fabricated or malformed metadata must not grant access.
pub async fn get_bucket_tags_for_policy(&self, bucket: &str) -> Result<HashMap<String, String>> {
let sys = metadata_sys::require_bucket_metadata_sys_in(&self.ctx)?;
let metadata = match sys.read().await.get_authoritative_metadata(bucket).await {
Ok(metadata) => metadata,
Err(Error::ConfigNotFound) => return Ok(HashMap::new()),
Err(err) => return Err(err),
};
let Some(tagging) =
ConfigState::of(&metadata.tagging_config_xml, &metadata.tagging_config).require(bucket, BUCKET_TAGGING_CONFIG)?
else {
return Ok(HashMap::new());
};
let mut tags = HashMap::with_capacity(tagging.tag_set.len());
for tag in &tagging.tag_set {
let (Some(key), Some(value)) = (&tag.key, &tag.value) else {
return Err(Error::other("stored bucket tag is missing a key or value"));
};
if key.is_empty() || tags.insert(key.clone(), value.clone()).is_some() {
return Err(Error::other("stored bucket tags contain empty or duplicate keys"));
}
}
Ok(tags)
}
pub async fn get_bucket_policy(&self, bucket: &str) -> Result<(BucketPolicy, OffsetDateTime)> {
let sys = metadata_sys::require_bucket_metadata_sys_in(&self.ctx)?;
sys.read().await.get_bucket_policy(bucket).await
@@ -1212,6 +1238,7 @@ mod tests {
use rustfs_filemeta::{FileInfo, FileMeta, TRANSITION_COMPLETE};
use rustfs_lock::{LocalClient, LockRequest, LockType, NamespaceLock, ObjectKey};
use serial_test::serial;
use std::collections::HashMap;
use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
@@ -1619,6 +1646,55 @@ mod tests {
}
}
#[tokio::test]
async fn bucket_tags_for_policy_distinguish_absence_from_corrupt_metadata() {
let (_temp_dir, store) = setup_bucket_quorum_test_env(&[4], None).await;
metadata_sys::init_bucket_metadata_sys(Arc::clone(&store), Vec::new()).await;
assert!(
store
.get_bucket_tags_for_policy("not-created")
.await
.expect("confirmed missing bucket")
.is_empty()
);
let bucket = "policy-bucket-tags";
store
.make_bucket(bucket, &MakeBucketOptions::default())
.await
.expect("create tagged bucket");
let original = store.get_bucket_metadata(bucket).await.expect("read bucket metadata");
let cases = &[
(b"".as_slice(), Some(HashMap::new())),
(b"<Tagging><TagSet/></Tagging>", Some(HashMap::new())),
(
b"<Tagging><TagSet><Tag><Key>Department</Key><Value>Finance</Value></Tag><Tag><Key>note</Key><Value></Value></Tag></TagSet></Tagging>",
Some(HashMap::from([("Department".to_string(), "Finance".to_string()), ("note".to_string(), String::new())])),
),
(b"<Tagging><TagSet>", None),
(b"<Tagging><TagSet><Tag><Value>finance</Value></Tag></TagSet></Tagging>", None),
(b"<Tagging><TagSet><Tag><Key>department</Key></Tag></TagSet></Tagging>", None),
(b"<Tagging><TagSet><Tag><Key></Key><Value>finance</Value></Tag></TagSet></Tagging>", None),
(
b"<Tagging><TagSet><Tag><Key>department</Key><Value>finance</Value></Tag><Tag><Key>department</Key><Value>engineering</Value></Tag></TagSet></Tagging>",
None,
),
];
for (xml, expected) in cases {
let mut metadata = (*original).clone();
metadata.tagging_config_xml = xml.to_vec();
metadata.tagging_config = crate::bucket::utils::deserialize(xml).ok();
metadata_sys::set_bucket_metadata_in(&store.ctx, metadata)
.await
.expect("inject persisted tag state");
let result = store.get_bucket_tags_for_policy(bucket).await;
match expected {
Some(tags) => assert_eq!(&result.expect("valid tag configuration"), tags),
None => assert!(result.is_err(), "corrupt tags must not become an untagged bucket: {xml:?}"),
}
}
}
#[tokio::test]
async fn request_metadata_methods_fail_closed_before_instance_initialization() {
let (_temp_dir, store) = setup_multi_pool_bucket_test_env().await;
@@ -1626,6 +1702,10 @@ mod tests {
let expected = "bucket metadata sys not initialized for this instance";
let errors = [
store.get_bucket_metadata("bucket").await.unwrap_err(),
store
.get_bucket_tags_for_policy("bucket")
.await
.expect_err("uninitialized metadata must deny tag lookup"),
store.get_bucket_policy("bucket").await.unwrap_err(),
store.get_bucket_policy_raw("bucket").await.unwrap_err(),
store.restricts_public_bucket_access("bucket").await.unwrap_err(),
+6
View File
@@ -72,6 +72,12 @@ pub trait Store: Clone + Send + Sync + 'static {
async fn load_all(&self, cache: &Cache) -> Result<()>;
/// Stored, case-sensitive bucket tags for OPA. A confirmed missing or untagged
/// bucket has an empty map; unavailable metadata must return an error.
async fn load_bucket_tags(&self, _bucket: &str) -> Result<HashMap<String, String>> {
Err(crate::error::Error::other("bucket tag lookup is unavailable for this IAM store"))
}
// Lock-free variants used by the cross-node notification handlers.
//
// Notification-path cache refreshes are asynchronous, best-effort, and
+4
View File
@@ -814,6 +814,10 @@ impl ObjectStore {
#[async_trait::async_trait]
impl Store for ObjectStore {
async fn load_bucket_tags(&self, bucket: &str) -> Result<HashMap<String, String>> {
Ok(self.object_api.get_bucket_tags_for_policy(bucket).await?)
}
fn has_watcher(&self) -> bool {
false
}
+230 -6
View File
@@ -41,6 +41,7 @@ use rustfs_policy::policy::Args;
use rustfs_policy::policy::opa;
use rustfs_policy::policy::{Policy, PolicyDoc, iam_policy_claim_name_sa, policy_needs_existing_object_tag_for_args};
use serde_json::Value;
use std::borrow::Cow;
use std::collections::{HashMap, HashSet};
use std::sync::Arc;
use std::sync::OnceLock;
@@ -1412,13 +1413,30 @@ impl<T: Store> IamSys<T> {
}
pub async fn eval_prepared(&self, prepared: &PreparedIamAuth, args: &Args<'_>) -> bool {
match &prepared.mode {
self.try_eval_prepared(prepared, args).await.unwrap_or(false)
}
/// Preserve resource lookup errors for callers that otherwise fall back to
/// bucket policies on an IAM denial. An error must not become an implicit deny.
pub async fn try_eval_prepared(&self, prepared: &PreparedIamAuth, args: &Args<'_>) -> Result<bool> {
Ok(match &prepared.mode {
PreparedIamMode::Opa => {
let Some(opa_enable) = Self::get_policy_plugin_client().await else {
tracing::warn!("eval_prepared: OPA mode requested but plugin is unavailable");
return false;
return Ok(false);
};
opa_enable.is_allowed(args).await
let tags = if args.bucket.is_empty() {
HashMap::new()
} else {
self.store.api.load_bucket_tags(args.bucket).await?
};
let conditions = opa_bucket_tag_conditions(args.conditions, tags);
opa_enable
.is_allowed(&Args {
conditions: conditions.as_ref(),
..args.clone()
})
.await
}
PreparedIamMode::Owner => true,
PreparedIamMode::Deny => false,
@@ -1430,7 +1448,7 @@ impl<T: Store> IamSys<T> {
} => {
let session_ok = evaluate_prepared_session_policy(session_policy, args).await;
if let Some(ok) = session_ok {
return ok && (*is_owner || combined_policy.is_allowed(args).await);
return Ok(ok && (*is_owner || combined_policy.is_allowed(args).await));
}
*is_owner || combined_policy.is_allowed(args).await
}
@@ -1460,13 +1478,13 @@ impl<T: Store> IamSys<T> {
PreparedServicePolicyMode::SessionBound => {
let session_ok = evaluate_prepared_session_policy(session_policy, args).await;
if let Some(ok) = session_ok {
return ok && parent_allowed;
return Ok(ok && parent_allowed);
}
parent_allowed
}
}
}
}
})
}
async fn prepare_regular_auth(&self, args: &Args<'_>) -> PreparedIamAuth {
@@ -1796,6 +1814,29 @@ impl<T: Store> IamSys<T> {
}
}
fn opa_bucket_tag_conditions(
conditions: &HashMap<String, Vec<String>>,
tags: HashMap<String, String>,
) -> Cow<'_, HashMap<String, Vec<String>>> {
// Conditions may originate in headers, claims or another resource's
// authorization. Only the addressed bucket supplies this namespace.
let is_bucket_tag = |key: &str| {
let namespace = key.split('/').next().unwrap_or_default();
namespace.eq_ignore_ascii_case("ExistingBucketTag") || namespace.eq_ignore_ascii_case("s3:ExistingBucketTag")
};
if tags.is_empty() && !conditions.keys().any(|key| is_bucket_tag(key)) {
return Cow::Borrowed(conditions);
}
let mut conditions = conditions.clone();
conditions.retain(|key, _| !is_bucket_tag(key));
conditions.reserve(tags.len());
for (key, value) in tags {
conditions.insert(format!("ExistingBucketTag/{key}"), vec![value]);
}
Cow::Owned(conditions)
}
async fn prepared_session_policy_needs_existing_object_tag_for_args(policy: &PreparedSessionPolicy, args: &Args<'_>) -> bool {
match policy {
PreparedSessionPolicy::Policy(p) => policy_needs_existing_object_tag_for_args(p, args).await,
@@ -2185,6 +2226,183 @@ mod tests {
assert!(matches!(state, PolicyPluginState::Failed));
}
#[test]
fn opa_bucket_tag_conditions_reuse_clean_input() {
let empty = HashMap::new();
let clean = HashMap::from([
("userid".to_string(), vec!["tag-user".to_string()]),
("ExistingBucketTagger/role".to_string(), vec!["unrelated".to_string()]),
]);
for conditions in [&empty, &clean] {
let result = opa_bucket_tag_conditions(conditions, HashMap::new());
assert_eq!(result.as_ref(), conditions);
assert!(
matches!(result, Cow::Borrowed(borrowed) if std::ptr::eq(borrowed, conditions)),
"clean, untagged authorization must reuse the existing conditions"
);
}
let result = opa_bucket_tag_conditions(&clean, HashMap::from([("department".to_string(), "finance".to_string())]));
assert_eq!(result.get("ExistingBucketTag/department"), Some(&vec!["finance".to_string()]));
assert_eq!(result.get("userid"), clean.get("userid"));
assert!(!clean.contains_key("ExistingBucketTag/department"));
}
#[test]
fn opa_bucket_tag_conditions_sanitize_before_enrichment() {
let mut untrusted = HashMap::from([("userid".to_string(), vec!["tag-user".to_string()])]);
for key in [
"ExistingBucketTag",
"s3:ExistingBucketTag",
"ExistingBucketTag/Department",
"existingbuckettag/department",
"S3:EXISTINGBUCKETTAG/role",
] {
untrusted.insert(key.to_string(), vec!["forged".to_string()]);
}
let original = untrusted.clone();
for tags in [
HashMap::new(),
HashMap::from([
("Department".to_string(), "Finance".to_string()),
("note".to_string(), String::new()),
("team/name".to_string(), "reporting".to_string()),
]),
] {
let mut expected = HashMap::from([("userid".to_string(), vec!["tag-user".to_string()])]);
for (key, value) in &tags {
expected.insert(format!("ExistingBucketTag/{key}"), vec![value.clone()]);
}
let result = opa_bucket_tag_conditions(&untrusted, tags);
assert_eq!(result.as_ref(), &expected);
assert_eq!(untrusted, original, "enrichment must not alter reusable request conditions");
}
}
#[tokio::test]
#[serial]
async fn opa_bucket_tags_replace_untrusted_conditions_and_lookup_errors_deny() {
use tokio::io::{AsyncBufReadExt, BufReader};
let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind OPA input receiver");
let url = format!("http://{}/decision", listener.local_addr().expect("receiver address"));
let receiver = tokio::spawn(async move {
let mut payloads = Vec::new();
for _ in 0..5 {
let (stream, _) = listener.accept().await.expect("accept OPA request");
let mut stream = BufReader::new(stream);
let mut length = None;
loop {
let mut line = String::new();
assert!(stream.read_line(&mut line).await.expect("read HTTP header") > 0);
if line == "\r\n" {
break;
}
if let Some((name, value)) = line.split_once(':')
&& name.eq_ignore_ascii_case("content-length")
{
length = Some(value.trim().parse::<usize>().expect("HTTP body length"));
}
}
let mut body = vec![0; length.expect("OPA JSON has a content length")];
stream.read_exact(&mut body).await.expect("read complete OPA body");
payloads.push(serde_json::from_slice::<serde_json::Value>(&body).expect("OPA JSON"));
stream
.get_mut()
.write_all(b"HTTP/1.1 200 OK\r\nContent-Length: 15\r\nConnection: close\r\n\r\n{\"result\":true}")
.await
.expect("send decision");
}
payloads
});
let (outcomes, error, denied) = temp_env::async_with_vars(
[
("NO_PROXY", Some("127.0.0.1,localhost")),
("no_proxy", Some("127.0.0.1,localhost")),
],
async {
let store = StsTestMockStore::new(true);
let iam = IamSys::new(IamCache::new(store.clone()).await.expect("initialize IAM cache"));
let previous = IamSys::<StsTestMockStore>::policy_plugin_state().await;
IamSys::<StsTestMockStore>::set_policy_plugin_client(opa::AuthZPlugin::new(opa::Args {
url,
auth_token: String::new(),
}))
.await;
let groups = None;
let claims = HashMap::new();
let conditions = HashMap::from([
("userid".to_string(), vec!["tag-user".to_string()]),
("ExistingBucketTag/Department".to_string(), vec!["forged".to_string()]),
("existingbuckettag/department".to_string(), vec!["forged".to_string()]),
("S3:EXISTINGBUCKETTAG/role".to_string(), vec!["forged".to_string()]),
]);
let mut args = Args {
account: "tag-user",
groups: &groups,
claims: &claims,
conditions: &conditions,
action: Action::S3Action(S3Action::GetBucketLocationAction),
bucket: "finance",
object: "",
is_owner: false,
deny_only: false,
};
let mut outcomes = Vec::new();
for (bucket, tags) in [
(
"finance",
Ok(HashMap::from([
("Department".to_string(), "Finance".to_string()),
("note".to_string(), String::new()),
])),
),
("untagged", Ok(HashMap::new())),
("", Err(Error::other("bucketless decisions must not read metadata"))),
] {
args.bucket = bucket;
*store.bucket_tags.lock().expect("tag state") = tags;
outcomes.push(iam.is_allowed(&args).await);
}
let clean_conditions = HashMap::from([("userid".to_string(), vec!["tag-user".to_string()])]);
args.conditions = &clean_conditions;
for (bucket, tags) in [
("untagged", Ok(HashMap::new())),
("", Err(Error::other("bucketless decisions must not read metadata"))),
] {
args.bucket = bucket;
*store.bucket_tags.lock().expect("tag state") = tags;
outcomes.push(iam.is_allowed(&args).await);
}
args.bucket = "unavailable";
let prepared = iam.prepare_auth(&args).await;
let error = iam.try_eval_prepared(&prepared, &args).await;
let denied = !iam.eval_prepared(&prepared, &args).await;
*get_policy_plugin_state().write().await = previous;
(outcomes, error, denied)
},
)
.await;
assert_eq!(outcomes, [true; 5]);
assert!(
matches!(error, Err(Error::Io(_))),
"lookup failure must remain an error, not an implicit denial"
);
assert!(denied, "boolean IAM callers must fail closed on lookup errors");
let payloads = tokio::time::timeout(std::time::Duration::from_secs(10), receiver)
.await
.expect("OPA input deadline")
.expect("input receiver");
assert_eq!(
payloads[0]["input"]["context"]["conditions"],
serde_json::json!({
"userid": ["tag-user"], "ExistingBucketTag/Department": ["Finance"], "ExistingBucketTag/note": [""]
})
);
for payload in &payloads[1..] {
assert_eq!(payload["input"]["context"]["conditions"], serde_json::json!({"userid": ["tag-user"]}));
}
}
const CUSTOM_STS_CLAIM_POLICY: &str = "custom-sts-claim-getobject";
const CUSTOM_STS_CLAIM_BUCKET: &str = "claim-bucket";
const CUSTOM_STS_CLAIM_POLICY_JSON: &str = r#"{
@@ -2203,6 +2421,7 @@ mod tests {
struct StsTestMockStore {
/// When true, parent user has no groups and no mapped policies (empty `policy_db_get`).
empty_policies: bool,
bucket_tags: Arc<Mutex<Result<HashMap<String, String>>>>,
saved_sts_users: Arc<Mutex<HashMap<String, UserIdentity>>>,
saved_service_account_count: Arc<Mutex<usize>>,
fail_delete: Arc<std::sync::atomic::AtomicBool>,
@@ -2219,6 +2438,7 @@ mod tests {
fn new(empty_policies: bool) -> Self {
Self {
empty_policies,
bucket_tags: Arc::new(Mutex::new(Err(Error::other("bucket metadata unavailable")))),
saved_sts_users: Arc::new(Mutex::new(HashMap::new())),
saved_service_account_count: Arc::new(Mutex::new(0)),
fail_delete: Arc::new(std::sync::atomic::AtomicBool::new(false)),
@@ -2242,6 +2462,10 @@ mod tests {
#[async_trait::async_trait]
impl Store for StsTestMockStore {
async fn load_bucket_tags(&self, _bucket: &str) -> Result<HashMap<String, String>> {
self.bucket_tags.lock().expect("tag state").clone()
}
fn has_watcher(&self) -> bool {
false
}
+75
View File
@@ -28,6 +28,81 @@
- Role-based access control integration
- Dynamic policy evaluation with context
## Stored bucket tags in OPA input
When the OPA authorization plugin is enabled, IAM resolves the addressed bucket's
stored tags before sending its decision request. Tags appear in the existing
conditions map; no additional S3 calls or credentials are needed by OPA:
```json
{
"input": {
"resource": { "bucket": "financial-reports", "object": "annual/report.parquet" },
"action": "s3:GetObject",
"context": {
"conditions": {
"ExistingBucketTag/department": ["finance"],
"ExistingBucketTag/environment": ["production"]
}
}
}
}
```
The example omits unchanged identity and context fields. Tag names and values
are case-sensitive. Values use single-element string arrays, including `[""]`
for an empty tag value. A simple resource predicate is:
```rego
package rustfs.example
import rego.v1
finance_bucket if {
input.context.conditions["ExistingBucketTag/department"] == ["finance"]
}
```
Combine this predicate with the policy's identity and action checks; tags do not
establish caller roles or grant access by themselves. Restrict `PutBucketTagging`
and `DeleteBucketTagging` separately when tags control access.
- Only server-resolved, stored tags populate `ExistingBucketTag/<key>`. Incoming
conditions in that namespace (including case variants and the `s3:` alias)
are discarded before enrichment. Headers, claims, requested replacement tags
and object tags cannot override it.
- A confirmed missing bucket, or a bucket with no stored tags, contributes no
bucket-tag conditions. New bucket creation therefore has no existing tags;
a create request naming an existing bucket uses its stored tags. Bucketless
requests, including STS and the global ListBuckets check, have no bucket tags.
- Bucket metadata lookup failures, unreadable tag XML, and ambiguous/incomplete
tags fail closed, even if a separate bucket policy would allow the request.
This applies to requests evaluated by OPA even if the policy does not inspect tags.
Boolean IAM callers deny on a lookup error; S3 access checks propagate it.
- Object reads/writes, versions, listing and multipart checks use their addressed
bucket. CopyObject and UploadPartCopy authorize source and destination with
each bucket's own tags. Tag replacement/removal uses the **existing** stored state.
- ListBuckets retains its existing contract: a global `s3:ListAllMyBuckets` allow
lists all buckets. Otherwise, per-bucket `s3:ListBucket`/`s3:GetBucketLocation`
checks carry each candidate's tags and can filter the result.
- Resolution reuses the IAM storage instance's authoritative metadata reader and
existing metadata cache/reload machinery, not a separate OPA cache. Local tag
updates/removal update the local cache, and tag mutations wait for healthy-peer
reload attempts before responding. Failed peer notification can leave that
peer's cached tags stale until a later reload or distributed metadata refresh
(normally every 15 minutes). This adds no stronger consistency barrier or
partition-time revocation guarantee, and does not revoke in-flight requests.
This is an OPA input extension, not a new native IAM condition key. Native policy
evaluation, owner/anonymous handling, and existing authorization combination rules
remain unchanged. For S3 authorization, an OPA `false` decision is an implicit IAM
denial that an applicable bucket-policy Allow may supplement. A tag mismatch alone
therefore does not revoke access granted by such a bucket policy; policy authors
must account for those grants when relying on tags to restrict access. ListBuckets
filtering remains IAM-only. Lookup errors on the OPA path abort authorization and
never fall through to a bucket-policy Allow. Custom IAM `Store` implementations
must implement `load_bucket_tags` to support bucket-scoped OPA requests; the default
fails closed.
## 📚 Documentation
For comprehensive documentation, examples, and usage guides, please visit the main [RustFS repository](https://github.com/rustfs/rustfs).
+9 -2
View File
@@ -1292,7 +1292,10 @@ pub async fn authorize_request<T>(req: &mut S3Request<T>, action: Action) -> S3R
claims,
deny_only: false,
};
let allowed = iam_store.eval_prepared(&prepared, &final_args).await;
let allowed = iam_store
.try_eval_prepared(&prepared, &final_args)
.await
.map_err(ApiError::from)?;
if !allowed
&& matches!(
action,
@@ -1311,7 +1314,11 @@ pub async fn authorize_request<T>(req: &mut S3Request<T>, action: Action) -> S3R
// 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 {
if !iam_store
.try_eval_prepared(&prepared, &final_args)
.await
.map_err(ApiError::from)?
{
return Err(denial.deny("iam_explicit_deny", action));
}
}