diff --git a/.config/e2e-distributed-selection.txt b/.config/e2e-distributed-selection.txt index c034893ad..33ddd59a4 100644 --- a/.config/e2e-distributed-selection.txt +++ b/.config/e2e-distributed-selection.txt @@ -1,2 +1,2 @@ -sha256-linux=6fc377fa1f9f06e065f077c7bd53d04efbccf4453f95dde67afaa47778185377 -sha256-darwin=6fc377fa1f9f06e065f077c7bd53d04efbccf4453f95dde67afaa47778185377 +sha256-linux=233bd7a68777eda1f044b6cb77128056be2b6492583c0dfd114b21e98a6102ef +sha256-darwin=233bd7a68777eda1f044b6cb77128056be2b6492583c0dfd114b21e98a6102ef diff --git a/.config/e2e-full-selection.txt b/.config/e2e-full-selection.txt index 71b9061aa..f8ae61180 100644 --- a/.config/e2e-full-selection.txt +++ b/.config/e2e-full-selection.txt @@ -1,2 +1,2 @@ -sha256-darwin=6b8e35ef69456bd244d6a6b19179e408763763e1e163aaaf90055490ae12da4b -sha256-linux=ff4c40f288ab5e91fd28d533c277538c0e3281fa353328ab52f5a8f88e0dc931 +sha256-darwin=98dff0961c6eb21d6128bfd2d4a6100a624ddf1fb7db05dba360a5d3d0a310c2 +sha256-linux=9ff0da7406de9b89c0e848ff29f2238b30ba9da89c7e9df97c246ce86631350f diff --git a/.config/e2e-smoke-selection.txt b/.config/e2e-smoke-selection.txt index 0938501c1..eb68a79f1 100644 --- a/.config/e2e-smoke-selection.txt +++ b/.config/e2e-smoke-selection.txt @@ -1 +1 @@ -sha256=7e3989625e2e6087c2e7db7ee02d30624b566d5a99832122032fb413ff456662 +sha256=ba8f392b617cf9f2338142233a417e43927b28096635e628233ce8a60cf49d10 diff --git a/.config/nextest.toml b/.config/nextest.toml index 2094b6640..6134ac90b 100644 --- a/.config/nextest.toml +++ b/.config/nextest.toml @@ -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' diff --git a/crates/e2e_test/src/distributed/s3_basic_test.rs b/crates/e2e_test/src/distributed/s3_basic_test.rs index cd0048f10..baf1eaa90 100644 --- a/crates/e2e_test/src/distributed/s3_basic_test.rs +++ b/crates/e2e_test/src/distributed/s3_basic_test.rs @@ -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::, _>>()?; + + // 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(); diff --git a/crates/e2e_test/src/sts_query_compat_test.rs b/crates/e2e_test/src/sts_query_compat_test.rs index 4f96c2e5c..134b3bfc6 100644 --- a/crates/e2e_test/src/sts_query_compat_test.rs +++ b/crates/e2e_test/src/sts_query_compat_test.rs @@ -35,7 +35,7 @@ use tokio::time::{Duration, timeout}; type BoxError = Box; 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, validation_started: mpsc::UnboundedReceiver<()>, task: JoinHandle<()>, } impl OpaMock { - async fn start() -> Result { + pub(crate) async fn start() -> Result { 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(©_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(©_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(©_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(©_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(); diff --git a/crates/ecstore/src/store/bucket.rs b/crates/ecstore/src/store/bucket.rs index 1e67d852f..c835ea76c 100644 --- a/crates/ecstore/src/store/bucket.rs +++ b/crates/ecstore/src/store/bucket.rs @@ -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> { + 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"", Some(HashMap::new())), + ( + b"DepartmentFinancenote", + Some(HashMap::from([("Department".to_string(), "Finance".to_string()), ("note".to_string(), String::new())])), + ), + (b"", None), + (b"finance", None), + (b"department", None), + (b"finance", None), + ( + b"departmentfinancedepartmentengineering", + 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(), diff --git a/crates/iam/src/store.rs b/crates/iam/src/store.rs index 6c5d7b2f6..756e733e3 100644 --- a/crates/iam/src/store.rs +++ b/crates/iam/src/store.rs @@ -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> { + 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 diff --git a/crates/iam/src/store/object.rs b/crates/iam/src/store/object.rs index a8737015c..143a78e62 100644 --- a/crates/iam/src/store/object.rs +++ b/crates/iam/src/store/object.rs @@ -814,6 +814,10 @@ impl ObjectStore { #[async_trait::async_trait] impl Store for ObjectStore { + async fn load_bucket_tags(&self, bucket: &str) -> Result> { + Ok(self.object_api.get_bucket_tags_for_policy(bucket).await?) + } + fn has_watcher(&self) -> bool { false } diff --git a/crates/iam/src/sys.rs b/crates/iam/src/sys.rs index 5b0a58c92..07870a706 100644 --- a/crates/iam/src/sys.rs +++ b/crates/iam/src/sys.rs @@ -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 IamSys { } 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 { + 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 IamSys { } => { 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 IamSys { 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 IamSys { } } +fn opa_bucket_tag_conditions( + conditions: &HashMap>, + tags: HashMap, +) -> Cow<'_, HashMap>> { + // 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::().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::(&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::::policy_plugin_state().await; + IamSys::::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>>>, saved_sts_users: Arc>>, saved_service_account_count: Arc>, fail_delete: Arc, @@ -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> { + self.bucket_tags.lock().expect("tag state").clone() + } + fn has_watcher(&self) -> bool { false } diff --git a/crates/policy/README.md b/crates/policy/README.md index 96a573808..e5683d4f0 100644 --- a/crates/policy/README.md +++ b/crates/policy/README.md @@ -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/`. 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). diff --git a/rustfs/src/storage/access.rs b/rustfs/src/storage/access.rs index 68c781cc2..b566eb2b4 100644 --- a/rustfs/src/storage/access.rs +++ b/rustfs/src/storage/access.rs @@ -1292,7 +1292,10 @@ pub async fn authorize_request(req: &mut S3Request, 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(req: &mut S3Request, 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)); } }