From 078774a7e809472890e4f62931db012cb436a162 Mon Sep 17 00:00:00 2001 From: Lucas Vieira Date: Mon, 14 Sep 2026 08:56:47 -0300 Subject: [PATCH 1/6] feat(dynamodb): enforce IAM policies on DynamoDB and DynamoDB Streams DynamoDB was not iam_enforceable, so FAKECLOUD_IAM=strict let any signed caller do anything to any table. - Every DynamoDB and DynamoDB Streams operation maps to the actions and resources in the AWS Service Authorization Reference: table, index (Query/Scan/contributor insights with IndexName), stream, backup, export, import and global-table ARNs; `*` for account-level listings. - A request can need several authorizations. AwsService gains iam_actions_for (defaulting to the single iam_action_for), and dispatch evaluates each one, denying if any is denied. Batches need the batch action on every table; transactions need GetItem / PutItem / UpdateItem / DeleteItem / ConditionCheckItem on each item's table; PartiQL needs PartiQLSelect/Insert/Update/Delete per statement; CreateTable with Tags or ResourcePolicy also needs TagResource / PutResourcePolicy; restores need the data-plane actions on the target. - A TableName given as an ARN authorizes against that ARN. - aws:ResourceTag reads the table's tags (also for its indexes and streams); aws:RequestTag / aws:TagKeys read CreateTable, TagResource and UntagResource tags. --- README.md | 2 +- crates/fakecloud-core/src/dispatch.rs | 435 +++++----- crates/fakecloud-core/src/service.rs | 14 + crates/fakecloud-dynamodb/src/service/iam.rs | 788 ++++++++++++++++++ crates/fakecloud-dynamodb/src/service/mod.rs | 25 + .../src/streams_dataplane.rs | 21 + .../tests/iam_enforcement_dynamodb.rs | 527 ++++++++++++ website/content/docs/reference/security.md | 8 +- website/content/docs/services/dynamodb.md | 1 + website/content/docs/services/iam.md | 4 +- website/static/llms-full.txt | 2 +- website/static/llms.txt | 2 +- 12 files changed, 1609 insertions(+), 220 deletions(-) create mode 100644 crates/fakecloud-dynamodb/src/service/iam.rs create mode 100644 crates/fakecloud-e2e/tests/iam_enforcement_dynamodb.rs diff --git a/README.md b/README.md index d67498ee3..8722ee501 100644 --- a/README.md +++ b/README.md @@ -54,7 +54,7 @@ Works as a drop-in for LocalStack in CI, with Terraform (`endpoints` block), CDK - **Single binary.** ~19 MB, ~10 MiB idle, ~300ms startup. No Docker needed to run fakecloud itself. - **Full Bedrock surface.** 216 ops across 4 APIs with real `InvokeModel`/`Converse` streaming, guardrails, agents, and flows. Configurable responses + fault injection for deterministic tests. See [`/bedrock-emulator/`](https://fakecloud.dev/bedrock-emulator/). - **First-party test SDKs** for TypeScript, Python, Go, PHP, Java, and Rust. Assert on what your code called without raw HTTP. -- **Opt-in SigV4 verification and IAM enforcement.** Off by default so tests just work; `--verify-sigv4` for real signature checking and `--iam soft|strict` for policy evaluation across IAM/STS/SQS/SNS/S3. See [security docs](https://fakecloud.dev/docs/reference/security/). +- **Opt-in SigV4 verification and IAM enforcement.** Off by default so tests just work; `--verify-sigv4` for real signature checking and `--iam soft|strict` for policy evaluation across IAM, STS, SQS, SNS, S3, KMS, Lambda, DynamoDB (and Streams), ELBv2 and Scheduler. See [security docs](https://fakecloud.dev/docs/reference/security/). - **Run your app unmodified.** An app that expects an instance/task role resolves the AWS SDK default credential chain against fakecloud with no static keys and no code change: point `AWS_CONTAINER_CREDENTIALS_FULL_URI` at `/_fakecloud/credentials`. See [Run an app unmodified](https://fakecloud.dev/docs/guides/instance-credentials/). - **LocalStack and real-AWS URL compatibility.** Both `*.localhost.localstack.cloud` and `*.amazonaws.com` Host headers route correctly, including every S3 virtual-hosted variant. Persisted URLs and dev scripts from either system replay unchanged. diff --git a/crates/fakecloud-core/src/dispatch.rs b/crates/fakecloud-core/src/dispatch.rs index 83f66a821..d72fb5855 100644 --- a/crates/fakecloud-core/src/dispatch.rs +++ b/crates/fakecloud-core/src/dispatch.rs @@ -592,166 +592,173 @@ pub async fn dispatch( if let Some(evaluator) = config.policy_evaluator.as_ref() { if let Some(principal) = aws_request.principal.as_ref() { if !principal.is_root() { - if let Some(iam_action) = service.iam_action_for(&aws_request) { - let mut condition_context = build_condition_context( - principal, - remote_addr, - &aws_request.region, - is_secure_transport(&aws_request.headers), - ); - // F3 keys riding on the resolved credential. STS - // populates these at mint time so subsequent - // requests under the credential can be evaluated - // against `aws:MultiFactorAuthPresent`, - // `aws:MultiFactorAuthAge`, `aws:TokenIssueTime`, - // and `aws:FederatedProvider`. IAM user access - // keys carry none of these, matching AWS. - if let Some(rc) = resolved.as_ref() { - condition_context.aws_mfa_present = Some(rc.mfa_present); - condition_context.aws_token_issue_time = rc.token_issued_at; - condition_context.aws_federated_provider = - rc.federated_provider.clone(); - // `aws:MultiFactorAuthAge` is "seconds since - // MFA was asserted" — computed at evaluation - // time from the token issue moment so the - // value increases monotonically as the session - // ages. Only set when the session was actually - // minted with MFA; otherwise the key is - // absent, matching AWS. - if rc.mfa_present { - if let Some(issued) = rc.token_issued_at { - let age = chrono::Utc::now() - .signed_duration_since(issued) - .num_seconds() - .max(0); - condition_context.aws_mfa_age_seconds = Some(age); + // A request can need several authorizations -- one per + // table in a batch, say -- and every one must allow it. + let iam_actions = service.iam_actions_for(&aws_request); + if !iam_actions.is_empty() { + for iam_action in &iam_actions { + let mut condition_context = build_condition_context( + principal, + remote_addr, + &aws_request.region, + is_secure_transport(&aws_request.headers), + ); + // F3 keys riding on the resolved credential. STS + // populates these at mint time so subsequent + // requests under the credential can be evaluated + // against `aws:MultiFactorAuthPresent`, + // `aws:MultiFactorAuthAge`, `aws:TokenIssueTime`, + // and `aws:FederatedProvider`. IAM user access + // keys carry none of these, matching AWS. + if let Some(rc) = resolved.as_ref() { + condition_context.aws_mfa_present = Some(rc.mfa_present); + condition_context.aws_token_issue_time = rc.token_issued_at; + condition_context.aws_federated_provider = + rc.federated_provider.clone(); + // `aws:MultiFactorAuthAge` is "seconds since + // MFA was asserted" — computed at evaluation + // time from the token issue moment so the + // value increases monotonically as the session + // ages. Only set when the session was actually + // minted with MFA; otherwise the key is + // absent, matching AWS. + if rc.mfa_present { + if let Some(issued) = rc.token_issued_at { + let age = chrono::Utc::now() + .signed_duration_since(issued) + .num_seconds() + .max(0); + condition_context.aws_mfa_age_seconds = Some(age); + } } } - } - condition_context.service_keys = - service.iam_condition_keys_for(&aws_request, &iam_action); - - // ABAC: populate tag-based condition keys. - // aws:ResourceTag/* - match service.resource_tags_for(&iam_action.resource) { - Some(tags) => condition_context.resource_tags = Some(tags), - None => tracing::debug!( - target: "fakecloud::iam::audit", - service = %detected.service, - resource = %iam_action.resource, - "service does not expose resource tags for ABAC; skipping aws:ResourceTag/* evaluation" - ), - } - // aws:RequestTag/* + aws:TagKeys - match service.request_tags_from(&aws_request, iam_action.action) { - Some(tags) => condition_context.request_tags = Some(tags), - None => tracing::debug!( - target: "fakecloud::iam::audit", - service = %detected.service, - action = %iam_action.action_string(), - "service does not expose request tags for ABAC; skipping aws:RequestTag/* / aws:TagKeys evaluation" - ), - } - // aws:PrincipalTag/* - condition_context.principal_tags = principal.tags.clone(); - - // Phase 2: fetch the resource-based policy (if - // any) attached to the target resource and - // pass it to the evaluator alongside the - // principal's identity policies. The resource's - // owning account is parsed from the ARN (#381 - // multi-account alignment); S3 ARNs have an - // empty account field, so we fall back to the - // server's configured account ID in that case. - let resource_policy_json = - config.resource_policy_provider.as_ref().and_then(|p| { - p.resource_policy(&detected.service, &iam_action.resource) - }); - // Derive the resource-owning account. Prefer a provider - // lookup (S3 ARNs carry no account, so the bucket's - // owner is resolved from state — without this, account - // A reaching account B's bucket would be mis-read as - // same-account and skip B's bucket-policy requirement, - // bug-audit 2026-05-28, 5.3), then fall back to the - // account embedded in the ARN (SQS/SNS/Lambda/…), then - // to the caller's account for wildcard / unscoped - // actions (ListQueues, GetCallerIdentity). - let resource_account_id = config - .resource_policy_provider - .as_ref() - .and_then(|p| { - p.resource_owner_account(&detected.service, &iam_action.resource) - }) - .or_else(|| parse_account_from_arn(&iam_action.resource)) - .unwrap_or_else(|| principal.account_id.clone()); - // SCP ceiling: resolve the inherited SCP chain - // for this principal (management accounts and - // service-linked roles come back as `None`, in - // which case the evaluator treats the layer as - // absent). Audit breadcrumbs emitted by the - // resolver itself, not here. - let scps = config - .scp_resolver - .as_ref() - .and_then(|r| r.scps_for(principal)); - let decision = evaluator.evaluate_with_resource_policy( - principal, - &iam_action, - &condition_context, - resource_policy_json.as_deref(), - &resource_account_id, - &caller_session_policies, - scps.as_deref(), - ); - if !decision.is_allow() { - tracing::warn!( - target: "fakecloud::iam::audit", - service = %detected.service, - action = %iam_action.action_string(), - resource = %iam_action.resource, - principal = %principal.arn, - resource_policy_present = resource_policy_json.is_some(), - decision = ?decision, - mode = %config.iam_mode, - request_id = %request_id, - "IAM policy evaluation denied request" - ); - if config.iam_mode.is_strict() { - // Real AWS includes an "Encoded - // authorization failure message" suffix - // on AccessDeniedException — an opaque - // base64+zlib JSON blob that the caller - // can pass to STS - // `DecodeAuthorizationMessage` to - // recover the structured deny reason - // (action, principal, matched - // statements, condition context). We - // produce the same blob inline so - // existing tooling that decodes deny - // reasons works against fakecloud. - let context_summary = serde_json::json!({ - "aws:PrincipalArn": principal.arn, - "aws:PrincipalAccount": principal.account_id, - "aws:RequestedRegion": condition_context - .aws_requested_region - .clone() - .unwrap_or_default(), - "aws:SecureTransport": condition_context - .aws_secure_transport - .unwrap_or(false), - "aws:Action": iam_action.action_string(), - "aws:Resource": iam_action.resource, - "decision": format!("{:?}", decision), + condition_context.service_keys = + service.iam_condition_keys_for(&aws_request, iam_action); + + // ABAC: populate tag-based condition keys. + // aws:ResourceTag/* + match service.resource_tags_for(&iam_action.resource) { + Some(tags) => condition_context.resource_tags = Some(tags), + None => tracing::debug!( + target: "fakecloud::iam::audit", + service = %detected.service, + resource = %iam_action.resource, + "service does not expose resource tags for ABAC; skipping aws:ResourceTag/* evaluation" + ), + } + // aws:RequestTag/* + aws:TagKeys + match service.request_tags_from(&aws_request, iam_action.action) { + Some(tags) => condition_context.request_tags = Some(tags), + None => tracing::debug!( + target: "fakecloud::iam::audit", + service = %detected.service, + action = %iam_action.action_string(), + "service does not expose request tags for ABAC; skipping aws:RequestTag/* / aws:TagKeys evaluation" + ), + } + // aws:PrincipalTag/* + condition_context.principal_tags = principal.tags.clone(); + + // Phase 2: fetch the resource-based policy (if + // any) attached to the target resource and + // pass it to the evaluator alongside the + // principal's identity policies. The resource's + // owning account is parsed from the ARN (#381 + // multi-account alignment); S3 ARNs have an + // empty account field, so we fall back to the + // server's configured account ID in that case. + let resource_policy_json = + config.resource_policy_provider.as_ref().and_then(|p| { + p.resource_policy(&detected.service, &iam_action.resource) }); - let action_string = iam_action.action_string(); - let encoded = crate::auth_message::encode_deny( - matches!(decision, crate::auth::IamDecision::ExplicitDeny), - Some(&action_string), - Some(&principal.arn), - Vec::new(), - Some(context_summary), + // Derive the resource-owning account. Prefer a provider + // lookup (S3 ARNs carry no account, so the bucket's + // owner is resolved from state — without this, account + // A reaching account B's bucket would be mis-read as + // same-account and skip B's bucket-policy requirement, + // bug-audit 2026-05-28, 5.3), then fall back to the + // account embedded in the ARN (SQS/SNS/Lambda/…), then + // to the caller's account for wildcard / unscoped + // actions (ListQueues, GetCallerIdentity). + let resource_account_id = config + .resource_policy_provider + .as_ref() + .and_then(|p| { + p.resource_owner_account( + &detected.service, + &iam_action.resource, + ) + }) + .or_else(|| parse_account_from_arn(&iam_action.resource)) + .unwrap_or_else(|| principal.account_id.clone()); + // SCP ceiling: resolve the inherited SCP chain + // for this principal (management accounts and + // service-linked roles come back as `None`, in + // which case the evaluator treats the layer as + // absent). Audit breadcrumbs emitted by the + // resolver itself, not here. + let scps = config + .scp_resolver + .as_ref() + .and_then(|r| r.scps_for(principal)); + let decision = evaluator.evaluate_with_resource_policy( + principal, + iam_action, + &condition_context, + resource_policy_json.as_deref(), + &resource_account_id, + &caller_session_policies, + scps.as_deref(), + ); + if !decision.is_allow() { + tracing::warn!( + target: "fakecloud::iam::audit", + service = %detected.service, + action = %iam_action.action_string(), + resource = %iam_action.resource, + principal = %principal.arn, + resource_policy_present = resource_policy_json.is_some(), + decision = ?decision, + mode = %config.iam_mode, + request_id = %request_id, + "IAM policy evaluation denied request" ); - return build_error_response( + if config.iam_mode.is_strict() { + // Real AWS includes an "Encoded + // authorization failure message" suffix + // on AccessDeniedException — an opaque + // base64+zlib JSON blob that the caller + // can pass to STS + // `DecodeAuthorizationMessage` to + // recover the structured deny reason + // (action, principal, matched + // statements, condition context). We + // produce the same blob inline so + // existing tooling that decodes deny + // reasons works against fakecloud. + let context_summary = serde_json::json!({ + "aws:PrincipalArn": principal.arn, + "aws:PrincipalAccount": principal.account_id, + "aws:RequestedRegion": condition_context + .aws_requested_region + .clone() + .unwrap_or_default(), + "aws:SecureTransport": condition_context + .aws_secure_transport + .unwrap_or(false), + "aws:Action": iam_action.action_string(), + "aws:Resource": iam_action.resource, + "decision": format!("{:?}", decision), + }); + let action_string = iam_action.action_string(); + let encoded = crate::auth_message::encode_deny( + matches!(decision, crate::auth::IamDecision::ExplicitDeny), + Some(&action_string), + Some(&principal.arn), + Vec::new(), + Some(context_summary), + ); + return build_error_response( StatusCode::FORBIDDEN, "AccessDeniedException", &format!( @@ -764,9 +771,10 @@ pub async fn dispatch( &request_id, detected.protocol, ); + } + // Soft mode: audit log already emitted; fall + // through to the handler. } - // Soft mode: audit log already emitted; fall - // through to the handler. } } else { // Service opted in via `iam_enforceable()` but its @@ -820,63 +828,66 @@ pub async fn dispatch( // SigV4 verification off, fakecloud does not reject unverified // signed requests, and turning them into anonymous denials would // change long-standing behavior. - if let Some(iam_action) = service.iam_action_for(&aws_request) { - let now = chrono::Utc::now(); - let mut condition_context = ConditionContext { - aws_source_ip: remote_addr.map(|sa| sa.ip()), - aws_current_time: Some(now), - aws_epoch_time: Some(now.timestamp()), - aws_secure_transport: Some(is_secure_transport(&aws_request.headers)), - aws_requested_region: Some(aws_request.region.clone()), - ..Default::default() - }; - condition_context.service_keys = - service.iam_condition_keys_for(&aws_request, &iam_action); - let resource_policy_json = config - .resource_policy_provider - .as_ref() - .and_then(|p| p.resource_policy(&detected.service, &iam_action.resource)); - let policy_decision = evaluator.evaluate_anonymous( - &iam_action, - &condition_context, - resource_policy_json.as_deref(), - ); - let policy_allows = policy_decision.is_allow(); - // An explicit Deny in the resource policy always wins, even - // over a public-read ACL — matching AWS's Deny-overrides - // precedence. Collapsing the decision to a bool and ORing the - // ACL let a public ACL override an explicit anonymous Deny. - let policy_explicit_deny = - matches!(policy_decision, crate::auth::IamDecision::ExplicitDeny); - let acl_allows = !policy_explicit_deny - && config.resource_policy_provider.as_ref().is_some_and(|p| { - p.public_acl_allows( - &detected.service, - &iam_action.resource, - iam_action.action, - ) - }); - if !policy_allows && !acl_allows { - tracing::warn!( - target: "fakecloud::iam::audit", - service = %detected.service, - action = %iam_action.action_string(), - resource = %iam_action.resource, - resource_policy_present = resource_policy_json.is_some(), - mode = %config.iam_mode, - request_id = %request_id, - "anonymous request denied: no public bucket policy or ACL grants the action" + let iam_actions = service.iam_actions_for(&aws_request); + if !iam_actions.is_empty() { + for iam_action in &iam_actions { + let now = chrono::Utc::now(); + let mut condition_context = ConditionContext { + aws_source_ip: remote_addr.map(|sa| sa.ip()), + aws_current_time: Some(now), + aws_epoch_time: Some(now.timestamp()), + aws_secure_transport: Some(is_secure_transport(&aws_request.headers)), + aws_requested_region: Some(aws_request.region.clone()), + ..Default::default() + }; + condition_context.service_keys = + service.iam_condition_keys_for(&aws_request, iam_action); + let resource_policy_json = + config.resource_policy_provider.as_ref().and_then(|p| { + p.resource_policy(&detected.service, &iam_action.resource) + }); + let policy_decision = evaluator.evaluate_anonymous( + iam_action, + &condition_context, + resource_policy_json.as_deref(), ); - if config.iam_mode.is_strict() { - return build_error_response( - StatusCode::FORBIDDEN, - "AccessDenied", - "Access Denied", - &request_id, - detected.protocol, + let policy_allows = policy_decision.is_allow(); + // An explicit Deny in the resource policy always wins, even + // over a public-read ACL — matching AWS's Deny-overrides + // precedence. Collapsing the decision to a bool and ORing the + // ACL let a public ACL override an explicit anonymous Deny. + let policy_explicit_deny = + matches!(policy_decision, crate::auth::IamDecision::ExplicitDeny); + let acl_allows = !policy_explicit_deny + && config.resource_policy_provider.as_ref().is_some_and(|p| { + p.public_acl_allows( + &detected.service, + &iam_action.resource, + iam_action.action, + ) + }); + if !policy_allows && !acl_allows { + tracing::warn!( + target: "fakecloud::iam::audit", + service = %detected.service, + action = %iam_action.action_string(), + resource = %iam_action.resource, + resource_policy_present = resource_policy_json.is_some(), + mode = %config.iam_mode, + request_id = %request_id, + "anonymous request denied: no public bucket policy or ACL grants the action" ); + if config.iam_mode.is_strict() { + return build_error_response( + StatusCode::FORBIDDEN, + "AccessDenied", + "Access Denied", + &request_id, + detected.protocol, + ); + } + // Soft mode: audit log emitted; fall through to the handler. } - // Soft mode: audit log emitted; fall through to the handler. } } else { // Anonymous request to an iam_enforceable service whose diff --git a/crates/fakecloud-core/src/service.rs b/crates/fakecloud-core/src/service.rs index ff002620a..584296153 100644 --- a/crates/fakecloud-core/src/service.rs +++ b/crates/fakecloud-core/src/service.rs @@ -747,6 +747,20 @@ pub trait AwsService: Send + Sync { None } + /// Every IAM authorization an incoming request needs. + /// + /// Most operations act on one resource and need one action, which is + /// what the default returns ([`AwsService::iam_action_for`]). Some need + /// several: a batch or transaction naming several resources needs the + /// action on each, and an operation can require more than one action + /// (DynamoDB's `CreateTable` with `Tags` also needs `TagResource`). + /// Dispatch evaluates every action returned and denies the request if + /// any of them is denied. An empty list means the operation has no + /// mapping, which strict enforcement denies. + fn iam_actions_for(&self, request: &AwsRequest) -> Vec { + self.iam_action_for(request).into_iter().collect() + } + /// Derive service-specific IAM condition keys for an incoming request. /// /// Called right after [`AwsService::iam_action_for`] when IAM diff --git a/crates/fakecloud-dynamodb/src/service/iam.rs b/crates/fakecloud-dynamodb/src/service/iam.rs new file mode 100644 index 000000000..11f79d0e6 --- /dev/null +++ b/crates/fakecloud-dynamodb/src/service/iam.rs @@ -0,0 +1,788 @@ +//! IAM authorization for DynamoDB and DynamoDB Streams requests: which +//! `dynamodb:*` actions each operation needs, on which resources. +//! +//! Follows the AWS Service Authorization Reference for DynamoDB. Most +//! operations need their namesake action on one table. The exceptions: +//! +//! - A batch needs `BatchGetItem` / `BatchWriteItem` on every table it names. +//! - A transaction needs the per-item action (`GetItem`, `PutItem`, +//! `UpdateItem`, `DeleteItem`, `ConditionCheckItem`) on each item's table; +//! there is no `dynamodb:TransactWriteItems` action. +//! - PartiQL statements need `PartiQLSelect` / `PartiQLInsert` / +//! `PartiQLUpdate` / `PartiQLDelete` on the table (or, for a SELECT, the +//! index) each statement names. +//! - `Query`, `Scan` and the contributor-insights operations target the +//! index when `IndexName` is given. +//! - `CreateTable` also needs `TagResource` when it carries `Tags` and +//! `PutResourcePolicy` when it carries `ResourcePolicy`. +//! - A restore needs its own action on the source plus the data-plane +//! actions DynamoDB uses to write the target table. +//! +//! A `TableName` given as an ARN authorizes against that ARN, so the +//! resource carries the table's own account and region. + +use std::collections::HashMap; + +use fakecloud_core::auth::IamAction; +use fakecloud_core::service::AwsRequest; +use serde_json::Value; + +use super::helpers::partiql::{find_outside_quotes, parse_partiql_table_name}; + +const SERVICE: &str = "dynamodb"; + +/// The data-plane actions DynamoDB performs on a restore's target table. +const RESTORE_TARGET_ACTIONS: [&str; 7] = [ + "BatchWriteItem", + "DeleteItem", + "GetItem", + "PutItem", + "Query", + "Scan", + "UpdateItem", +]; + +/// Where a resource ARN built from a bare name lives: the caller's account +/// and the request's region. +struct Scope<'a> { + account: &'a str, + region: &'a str, +} + +impl Scope<'_> { + /// The table ARN for a `TableName` value: kept as given when it already + /// is an ARN (normalized to the table itself), built from the scope + /// otherwise. + fn table(&self, name_or_arn: &str) -> String { + if let Some(arn) = table_arn_of(name_or_arn) { + return arn; + } + format!( + "arn:aws:dynamodb:{}:{}:table/{name_or_arn}", + self.region, self.account + ) + } + + fn index(&self, table: &str, index: &str) -> String { + format!("{}/index/{index}", self.table(table)) + } + + fn global_table(&self, name: &str) -> String { + format!("arn:aws:dynamodb::{}:global-table/{name}", self.account) + } +} + +/// `arn:aws:dynamodb:REGION:ACCOUNT:table/NAME` for an ARN naming a table or +/// one of its sub-resources, or `None` for anything else. +fn table_arn_of(arn: &str) -> Option { + let rest = arn.strip_prefix("arn:aws:dynamodb:")?; + let (scope, resource) = rest.split_once(":table/")?; + let name = resource.split('/').next().filter(|n| !n.is_empty())?; + Some(format!("arn:aws:dynamodb:{scope}:table/{name}")) +} + +fn action(name: &'static str, resource: String) -> IamAction { + IamAction { + service: SERVICE, + action: name, + resource, + } +} + +/// A body string field, or `None` when absent or empty. +fn field<'a>(body: &'a Value, name: &str) -> Option<&'a str> { + body[name].as_str().filter(|s| !s.is_empty()) +} + +/// The `dynamodb:*` authorizations a DynamoDB request needs. Empty only for +/// an operation this service does not implement. +pub(crate) fn actions_for(request: &AwsRequest) -> Vec { + let body: Value = serde_json::from_slice(&request.body).unwrap_or(Value::Null); + let account = request + .principal + .as_ref() + .map(|p| p.account_id.as_str()) + .unwrap_or(request.account_id.as_str()); + let scope = Scope { + account, + region: request.region.as_str(), + }; + let op: &'static str = match DYNAMODB_ACTIONS + .iter() + .find(|a| **a == request.action.as_str()) + { + Some(op) => op, + None => return Vec::new(), + }; + // A required name that is missing still maps, to `*`: the request is + // malformed, and the handler rejects it with its own validation error + // once the caller is authorized for the operation at all. + let table = |key: &str| field(&body, key).map_or_else(|| "*".to_string(), |t| scope.table(t)); + let table_or_index = || match (field(&body, "TableName"), field(&body, "IndexName")) { + (Some(t), Some(i)) => scope.index(t, i), + (Some(t), None) => scope.table(t), + _ => "*".to_string(), + }; + let arn_field = |key: &str| field(&body, key).unwrap_or("*").to_string(); + + match op { + "GetItem" + | "PutItem" + | "UpdateItem" + | "DeleteItem" + | "CreateBackup" + | "DeleteTable" + | "DescribeTable" + | "UpdateTable" + | "DescribeTimeToLive" + | "UpdateTimeToLive" + | "DescribeContinuousBackups" + | "UpdateContinuousBackups" + | "DescribeKinesisStreamingDestination" + | "EnableKinesisStreamingDestination" + | "DisableKinesisStreamingDestination" + | "UpdateKinesisStreamingDestination" + | "DescribeTableReplicaAutoScaling" + | "UpdateTableReplicaAutoScaling" => { + vec![action(op, table("TableName"))] + } + "Query" + | "Scan" + | "DescribeContributorInsights" + | "UpdateContributorInsights" + | "SearchVectors" => vec![action(op, table_or_index())], + "ListTables" + | "DescribeLimits" + | "DescribeEndpoints" + | "ListBackups" + | "ListContributorInsights" + | "ListGlobalTables" => vec![action(op, "*".to_string())], + "ListExports" | "ListImports" => vec![action(op, table("TableArn"))], + "ExportTableToPointInTime" => vec![action(op, table("TableArn"))], + "DescribeBackup" | "DeleteBackup" => vec![action(op, arn_field("BackupArn"))], + "DescribeExport" => vec![action(op, arn_field("ExportArn"))], + "DescribeImport" => vec![action(op, arn_field("ImportArn"))], + "TagResource" + | "UntagResource" + | "ListTagsOfResource" + | "GetResourcePolicy" + | "PutResourcePolicy" + | "DeleteResourcePolicy" => { + vec![action(op, arn_field("ResourceArn"))] + } + "CreateTable" => { + let resource = table("TableName"); + let mut out = vec![action("CreateTable", resource.clone())]; + if body["Tags"].as_array().is_some_and(|t| !t.is_empty()) { + out.push(action("TagResource", resource.clone())); + } + if field(&body, "ResourcePolicy").is_some() { + out.push(action("PutResourcePolicy", resource)); + } + out + } + "ImportTable" => { + let name = body["TableCreationParameters"]["TableName"] + .as_str() + .filter(|s| !s.is_empty()); + vec![action( + "ImportTable", + name.map_or_else(|| "*".to_string(), |n| scope.table(n)), + )] + } + "RestoreTableFromBackup" => { + let target = table("TargetTableName"); + let mut out = vec![ + action("RestoreTableFromBackup", arn_field("BackupArn")), + action("RestoreTableFromBackup", target.clone()), + ]; + out.extend( + RESTORE_TARGET_ACTIONS + .iter() + .map(|a| action(a, target.clone())), + ); + out + } + "RestoreTableToPointInTime" => { + let source = field(&body, "SourceTableArn") + .or_else(|| field(&body, "SourceTableName")) + .map_or_else(|| "*".to_string(), |t| scope.table(t)); + let target = table("TargetTableName"); + let mut out = vec![action("RestoreTableToPointInTime", source)]; + out.extend( + RESTORE_TARGET_ACTIONS + .iter() + .map(|a| action(a, target.clone())), + ); + out + } + "CreateGlobalTable" | "UpdateGlobalTable" | "UpdateGlobalTableSettings" => { + let name = field(&body, "GlobalTableName").unwrap_or("*"); + vec![ + action(op, scope.global_table(name)), + action(op, scope.table(name)), + ] + } + "DescribeGlobalTable" | "DescribeGlobalTableSettings" => { + let name = field(&body, "GlobalTableName").unwrap_or("*"); + vec![action(op, scope.global_table(name))] + } + "BatchGetItem" | "BatchWriteItem" => { + let tables = batch_table_names(&body["RequestItems"]); + if tables.is_empty() { + return vec![action(op, "*".to_string())]; + } + tables + .into_iter() + .map(|t| action(op, scope.table(t))) + .collect() + } + "TransactGetItems" | "TransactWriteItems" => { + let mut out = Vec::new(); + for item in body["TransactItems"].as_array().into_iter().flatten() { + for (member, item_action) in [ + ("Get", "GetItem"), + ("Put", "PutItem"), + ("Update", "UpdateItem"), + ("Delete", "DeleteItem"), + ("ConditionCheck", "ConditionCheckItem"), + ] { + if let Some(name) = item[member]["TableName"].as_str() { + push_unique(&mut out, action(item_action, scope.table(name))); + } + } + } + if out.is_empty() { + let fallback = if op == "TransactGetItems" { + "GetItem" + } else { + "PutItem" + }; + out.push(action(fallback, "*".to_string())); + } + out + } + "ExecuteStatement" => { + let statement = field(&body, "Statement").unwrap_or(""); + vec![partiql_action(&scope, statement)] + } + "BatchExecuteStatement" | "ExecuteTransaction" => { + let list = if op == "BatchExecuteStatement" { + &body["Statements"] + } else { + &body["TransactStatements"] + }; + let mut out = Vec::new(); + for statement in list.as_array().into_iter().flatten() { + let text = statement["Statement"].as_str().unwrap_or(""); + push_unique(&mut out, partiql_action(&scope, text)); + } + if out.is_empty() { + out.push(action("PartiQLSelect", "*".to_string())); + } + out + } + _ => Vec::new(), + } +} + +/// The DynamoDB Streams operations' authorizations. +pub(crate) fn streams_actions_for(request: &AwsRequest) -> Vec { + let body: Value = serde_json::from_slice(&request.body).unwrap_or(Value::Null); + match request.action.as_str() { + "ListStreams" => vec![action("ListStreams", "*".to_string())], + "DescribeStream" => vec![action( + "DescribeStream", + field(&body, "StreamArn").unwrap_or("*").to_string(), + )], + "GetShardIterator" => vec![action( + "GetShardIterator", + field(&body, "StreamArn").unwrap_or("*").to_string(), + )], + // An iterator is `STREAM_ARN|SHARD|SEQUENCE`; the stream it reads is + // the resource. + "GetRecords" => { + let stream = field(&body, "ShardIterator") + .and_then(|it| it.split('|').next()) + .filter(|arn| arn.starts_with("arn:")) + .unwrap_or("*"); + vec![action("GetRecords", stream.to_string())] + } + _ => Vec::new(), + } +} + +/// Every DynamoDB control- and data-plane operation this service serves. +const DYNAMODB_ACTIONS: &[&str] = &[ + "BatchExecuteStatement", + "BatchGetItem", + "BatchWriteItem", + "CreateBackup", + "CreateGlobalTable", + "CreateTable", + "DeleteBackup", + "DeleteItem", + "DeleteResourcePolicy", + "DeleteTable", + "DescribeBackup", + "DescribeContinuousBackups", + "DescribeContributorInsights", + "DescribeEndpoints", + "DescribeExport", + "DescribeGlobalTable", + "DescribeGlobalTableSettings", + "DescribeImport", + "DescribeKinesisStreamingDestination", + "DescribeLimits", + "DescribeTable", + "DescribeTableReplicaAutoScaling", + "DescribeTimeToLive", + "DisableKinesisStreamingDestination", + "EnableKinesisStreamingDestination", + "ExecuteStatement", + "ExecuteTransaction", + "ExportTableToPointInTime", + "GetItem", + "GetResourcePolicy", + "ImportTable", + "ListBackups", + "ListContributorInsights", + "ListExports", + "ListGlobalTables", + "ListImports", + "ListTables", + "ListTagsOfResource", + "PutItem", + "PutResourcePolicy", + "Query", + "RestoreTableFromBackup", + "RestoreTableToPointInTime", + "Scan", + "SearchVectors", + "TagResource", + "TransactGetItems", + "TransactWriteItems", + "UntagResource", + "UpdateContinuousBackups", + "UpdateContributorInsights", + "UpdateGlobalTable", + "UpdateGlobalTableSettings", + "UpdateItem", + "UpdateKinesisStreamingDestination", + "UpdateTable", + "UpdateTableReplicaAutoScaling", + "UpdateTimeToLive", +]; + +fn batch_table_names(request_items: &Value) -> Vec<&str> { + request_items + .as_object() + .map(|m| m.keys().map(String::as_str).collect()) + .unwrap_or_default() +} + +fn push_unique(out: &mut Vec, a: IamAction) { + if !out.contains(&a) { + out.push(a); + } +} + +/// The PartiQL action a statement needs, on the table (or, for a SELECT +/// from `"table"."index"`, the index) it names. A statement too malformed to +/// name a table maps to `PartiQLSelect` on `*`, leaving the syntax error to +/// the handler. +fn partiql_action(scope: &Scope<'_>, statement: &str) -> IamAction { + let trimmed = statement.trim(); + let upper = trimmed.to_ascii_uppercase(); + let (verb, keyword) = if upper.starts_with("SELECT") { + ("PartiQLSelect", Some("FROM")) + } else if upper.starts_with("INSERT") { + ("PartiQLInsert", Some("INTO")) + } else if upper.starts_with("UPDATE") { + ("PartiQLUpdate", None) + } else if upper.starts_with("DELETE") { + ("PartiQLDelete", Some("FROM")) + } else { + return action("PartiQLSelect", "*".to_string()); + }; + let after = match keyword { + Some(kw) => match find_outside_quotes(&upper, kw) { + Some(pos) => &trimmed[pos + kw.len()..], + None => return action(verb, "*".to_string()), + }, + None => &trimmed["UPDATE".len()..], + }; + let (table, rest) = parse_partiql_table_name(after); + if table.is_empty() { + return action(verb, "*".to_string()); + } + if verb == "PartiQLSelect" { + if let Some(index_part) = rest.strip_prefix('.') { + let (index, _) = parse_partiql_table_name(index_part); + if !index.is_empty() { + return action(verb, scope.index(&table, &index)); + } + } + } + action(verb, scope.table(&table)) +} + +/// Tags on the table a resource ARN names (a table, or its index or +/// stream), for `aws:ResourceTag/*`. `Some(empty)` for `*`; `None` when the +/// ARN names no table this state holds. +pub(crate) fn resource_tags( + state: &crate::state::SharedDynamoDbState, + resource_arn: &str, +) -> Option> { + if resource_arn == "*" { + return Some(HashMap::new()); + } + let table_arn = table_arn_of(resource_arn)?; + let account = table_arn.split(':').nth(4)?; + let name = table_arn.rsplit("table/").next()?; + let accounts = state.read(); + let table = accounts.get(account)?.tables.get(name)?; + Some( + table + .tags + .iter() + .map(|(k, v)| (k.clone(), v.clone())) + .collect(), + ) +} + +/// Tags a request writes, for `aws:RequestTag/*` and `aws:TagKeys`: +/// `CreateTable` / `TagResource` carry `Tags: [{Key, Value}]`, and +/// `UntagResource` names keys only. +pub(crate) fn request_tags(request: &AwsRequest, action: &str) -> Option> { + let body: Value = serde_json::from_slice(&request.body).unwrap_or(Value::Null); + match action { + "CreateTable" | "TagResource" => Some( + body["Tags"] + .as_array() + .into_iter() + .flatten() + .filter_map(|t| { + Some(( + t["Key"].as_str()?.to_string(), + t["Value"].as_str().unwrap_or_default().to_string(), + )) + }) + .collect(), + ), + "UntagResource" => Some( + body["TagKeys"] + .as_array() + .into_iter() + .flatten() + .filter_map(|k| Some((k.as_str()?.to_string(), String::new()))) + .collect(), + ), + _ => None, + } +} + +#[cfg(test)] +mod tests { + use super::*; + use fakecloud_core::service::AwsService; + use serde_json::json; + + const ACCOUNT: &str = "111122223333"; + const TABLE: &str = "arn:aws:dynamodb:eu-west-1:111122223333:table/Orders"; + + fn req(action: &str, body: Value) -> AwsRequest { + AwsRequest { + service: "dynamodb".to_string(), + action: action.to_string(), + region: "eu-west-1".to_string(), + account_id: ACCOUNT.to_string(), + request_id: "test-id".to_string(), + headers: http::HeaderMap::new(), + query_params: HashMap::new(), + body: serde_json::to_vec(&body).unwrap().into(), + body_stream: parking_lot::Mutex::new(None), + path_segments: vec![], + raw_path: "/".to_string(), + raw_query: String::new(), + method: http::Method::POST, + is_query_protocol: false, + access_key_id: None, + principal: None, + } + } + + fn pairs(actions: Vec) -> Vec<(String, String)> { + actions + .into_iter() + .map(|a| (a.action_string(), a.resource)) + .collect() + } + + fn one(action: &str, resource: &str) -> Vec<(String, String)> { + vec![(format!("dynamodb:{action}"), resource.to_string())] + } + + /// Strict enforcement denies an operation with no mapping, so every + /// operation either service serves must map to at least one action. + #[test] + fn every_served_operation_maps_to_an_action() { + let state: crate::state::SharedDynamoDbState = + std::sync::Arc::new(parking_lot::RwLock::new( + fakecloud_core::multi_account::MultiAccountState::new(ACCOUNT, "eu-west-1", ""), + )); + let service = crate::DynamoDbService::new(state.clone()); + for op in service.supported_actions() { + assert!( + !actions_for(&req(op, json!({}))).is_empty(), + "DynamoDB {op} has no IAM mapping" + ); + } + let streams = crate::DynamoDbStreamsService::new(state); + for op in streams.supported_actions() { + assert!( + !streams_actions_for(&req(op, json!({}))).is_empty(), + "DynamoDB Streams {op} has no IAM mapping" + ); + } + } + + #[test] + fn item_and_table_operations_target_the_table() { + for op in [ + "GetItem", + "PutItem", + "UpdateItem", + "DeleteItem", + "DescribeTable", + ] { + assert_eq!( + pairs(actions_for(&req(op, json!({"TableName": "Orders"})))), + one(op, TABLE) + ); + } + // A table ARN is authorized as given, in its own account and region. + let foreign = "arn:aws:dynamodb:us-east-2:444455556666:table/Shared"; + assert_eq!( + pairs(actions_for(&req("GetItem", json!({"TableName": foreign})))), + one("GetItem", foreign) + ); + assert_eq!( + pairs(actions_for(&req("ListTables", json!({})))), + one("ListTables", "*") + ); + } + + #[test] + fn query_and_scan_target_the_index_when_named() { + let body = json!({"TableName": "Orders", "IndexName": "by-customer"}); + let index = format!("{TABLE}/index/by-customer"); + assert_eq!( + pairs(actions_for(&req("Query", body.clone()))), + one("Query", &index) + ); + assert_eq!(pairs(actions_for(&req("Scan", body))), one("Scan", &index)); + assert_eq!( + pairs(actions_for(&req("Scan", json!({"TableName": "Orders"})))), + one("Scan", TABLE) + ); + } + + /// Batches need the batch action on every table; transactions need the + /// per-item action on each item's table, once per distinct pair. + #[test] + fn batches_and_transactions_authorize_every_table() { + let batch = json!({"RequestItems": {"Orders": [], "Customers": []}}); + let mut got = pairs(actions_for(&req("BatchWriteItem", batch))); + got.sort(); + assert_eq!( + got, + vec![ + ( + "dynamodb:BatchWriteItem".to_string(), + "arn:aws:dynamodb:eu-west-1:111122223333:table/Customers".to_string() + ), + ("dynamodb:BatchWriteItem".to_string(), TABLE.to_string()), + ] + ); + + let transact = json!({"TransactItems": [ + {"Put": {"TableName": "Orders"}}, + {"Put": {"TableName": "Orders"}}, + {"ConditionCheck": {"TableName": "Customers"}}, + {"Delete": {"TableName": "Orders"}}, + {"Update": {"TableName": "Orders"}} + ]}); + assert_eq!( + pairs(actions_for(&req("TransactWriteItems", transact))), + vec![ + ("dynamodb:PutItem".to_string(), TABLE.to_string()), + ( + "dynamodb:ConditionCheckItem".to_string(), + "arn:aws:dynamodb:eu-west-1:111122223333:table/Customers".to_string() + ), + ("dynamodb:DeleteItem".to_string(), TABLE.to_string()), + ("dynamodb:UpdateItem".to_string(), TABLE.to_string()), + ] + ); + assert_eq!( + pairs(actions_for(&req( + "TransactGetItems", + json!({"TransactItems": [{"Get": {"TableName": "Orders"}}]}) + ))), + one("GetItem", TABLE) + ); + } + + #[test] + fn partiql_statements_map_to_partiql_actions() { + let cases = [ + ( + "SELECT * FROM \"Orders\" WHERE pk = 'a'", + "PartiQLSelect", + TABLE.to_string(), + ), + ("select * from Orders", "PartiQLSelect", TABLE.to_string()), + ( + "SELECT * FROM \"Orders\".\"by-customer\"", + "PartiQLSelect", + format!("{TABLE}/index/by-customer"), + ), + ( + "INSERT INTO \"Orders\" VALUE {'pk': 'a'}", + "PartiQLInsert", + TABLE.to_string(), + ), + ( + "UPDATE \"Orders\" SET x = 1 WHERE pk = 'a'", + "PartiQLUpdate", + TABLE.to_string(), + ), + ( + "DELETE FROM \"Orders\" WHERE pk = 'a'", + "PartiQLDelete", + TABLE.to_string(), + ), + ("EXPLAIN nonsense", "PartiQLSelect", "*".to_string()), + ]; + for (statement, action_name, resource) in cases { + assert_eq!( + pairs(actions_for(&req( + "ExecuteStatement", + json!({"Statement": statement}) + ))), + one(action_name, &resource), + "{statement}" + ); + } + let batch = json!({"Statements": [ + {"Statement": "INSERT INTO \"Orders\" VALUE {'pk': 'a'}"}, + {"Statement": "INSERT INTO \"Orders\" VALUE {'pk': 'b'}"}, + {"Statement": "DELETE FROM \"Orders\" WHERE pk = 'c'"} + ]}); + assert_eq!( + pairs(actions_for(&req("BatchExecuteStatement", batch))), + vec![ + ("dynamodb:PartiQLInsert".to_string(), TABLE.to_string()), + ("dynamodb:PartiQLDelete".to_string(), TABLE.to_string()), + ] + ); + } + + #[test] + fn create_table_and_restores_need_their_companion_actions() { + let create = json!({ + "TableName": "Orders", + "Tags": [{"Key": "team", "Value": "x"}], + "ResourcePolicy": "{}" + }); + assert_eq!( + pairs(actions_for(&req("CreateTable", create))), + vec![ + ("dynamodb:CreateTable".to_string(), TABLE.to_string()), + ("dynamodb:TagResource".to_string(), TABLE.to_string()), + ("dynamodb:PutResourcePolicy".to_string(), TABLE.to_string()), + ] + ); + assert_eq!( + pairs(actions_for(&req( + "CreateTable", + json!({"TableName": "Orders"}) + ))), + one("CreateTable", TABLE) + ); + + let backup = format!("{TABLE}/backup/01700000000000-abcd"); + let restored = pairs(actions_for(&req( + "RestoreTableFromBackup", + json!({"BackupArn": backup, "TargetTableName": "Copy"}), + ))); + let copy = "arn:aws:dynamodb:eu-west-1:111122223333:table/Copy"; + assert_eq!( + restored[0], + ( + "dynamodb:RestoreTableFromBackup".to_string(), + backup.clone() + ) + ); + assert!(restored.contains(&("dynamodb:PutItem".to_string(), copy.to_string()))); + assert!(restored.contains(&("dynamodb:BatchWriteItem".to_string(), copy.to_string()))); + + let pitr = pairs(actions_for(&req( + "RestoreTableToPointInTime", + json!({"SourceTableName": "Orders", "TargetTableName": "Copy"}), + ))); + assert_eq!( + pitr[0], + ( + "dynamodb:RestoreTableToPointInTime".to_string(), + TABLE.to_string() + ) + ); + assert!(pitr.contains(&("dynamodb:UpdateItem".to_string(), copy.to_string()))); + } + + #[test] + fn streams_operations_target_the_stream() { + let stream = format!("{TABLE}/stream/2026-01-01T00:00:00.000"); + assert_eq!( + pairs(streams_actions_for(&req( + "DescribeStream", + json!({"StreamArn": stream}) + ))), + one("DescribeStream", &stream) + ); + assert_eq!( + pairs(streams_actions_for(&req( + "GetRecords", + json!({"ShardIterator": format!("{stream}|shardId-1|0")}) + ))), + one("GetRecords", &stream) + ); + assert_eq!( + pairs(streams_actions_for(&req("ListStreams", json!({})))), + one("ListStreams", "*") + ); + } + + #[test] + fn request_tags_cover_create_tag_and_untag() { + let create = req( + "CreateTable", + json!({"Tags": [{"Key": "team", "Value": "payments"}]}), + ); + assert_eq!( + request_tags(&create, "CreateTable"), + Some(HashMap::from([( + "team".to_string(), + "payments".to_string() + )])) + ); + let untag = req("UntagResource", json!({"TagKeys": ["team"]})); + assert_eq!( + request_tags(&untag, "UntagResource").map(|t| t.into_keys().collect::>()), + Some(vec!["team".to_string()]) + ); + assert_eq!(request_tags(&req("GetItem", json!({})), "GetItem"), None); + } +} diff --git a/crates/fakecloud-dynamodb/src/service/mod.rs b/crates/fakecloud-dynamodb/src/service/mod.rs index 30145eb4b..a9087c9a9 100644 --- a/crates/fakecloud-dynamodb/src/service/mod.rs +++ b/crates/fakecloud-dynamodb/src/service/mod.rs @@ -2,6 +2,7 @@ mod batch; #[cfg(test)] mod expression_corpus_tests; mod global_tables; +pub(crate) mod iam; mod items; mod queries; mod streams; @@ -537,6 +538,30 @@ impl AwsService for DynamoDbService { result } + fn iam_enforceable(&self) -> bool { + true + } + + fn iam_actions_for(&self, request: &AwsRequest) -> Vec { + iam::actions_for(request) + } + + fn iam_action_for(&self, request: &AwsRequest) -> Option { + iam::actions_for(request).into_iter().next() + } + + fn resource_tags_for(&self, resource_arn: &str) -> Option> { + iam::resource_tags(&self.state, resource_arn) + } + + fn request_tags_from( + &self, + request: &AwsRequest, + action: &str, + ) -> Option> { + iam::request_tags(request, action) + } + fn supported_actions(&self) -> &[&str] { &[ "CreateTable", diff --git a/crates/fakecloud-dynamodb/src/streams_dataplane.rs b/crates/fakecloud-dynamodb/src/streams_dataplane.rs index 3a7947890..4a44ff622 100644 --- a/crates/fakecloud-dynamodb/src/streams_dataplane.rs +++ b/crates/fakecloud-dynamodb/src/streams_dataplane.rs @@ -46,6 +46,27 @@ impl AwsService for DynamoDbStreamsService { } } + fn iam_enforceable(&self) -> bool { + true + } + + fn iam_actions_for(&self, request: &AwsRequest) -> Vec { + crate::service::iam::streams_actions_for(request) + } + + fn iam_action_for(&self, request: &AwsRequest) -> Option { + crate::service::iam::streams_actions_for(request) + .into_iter() + .next() + } + + fn resource_tags_for( + &self, + resource_arn: &str, + ) -> Option> { + crate::service::iam::resource_tags(&self.state, resource_arn) + } + fn supported_actions(&self) -> &[&str] { &[ "ListStreams", diff --git a/crates/fakecloud-e2e/tests/iam_enforcement_dynamodb.rs b/crates/fakecloud-e2e/tests/iam_enforcement_dynamodb.rs new file mode 100644 index 000000000..eb5a04982 --- /dev/null +++ b/crates/fakecloud-e2e/tests/iam_enforcement_dynamodb.rs @@ -0,0 +1,527 @@ +//! IAM enforcement for DynamoDB and DynamoDB Streams. +//! +//! Each test starts fakecloud with `FAKECLOUD_IAM=strict`, seeds tables with +//! the root-bypass `test` credentials, gives a user an inline policy, and +//! checks what that user's own credentials may do. + +mod helpers; + +use aws_credential_types::Credentials; +use aws_sdk_dynamodb::types::{ + AttributeDefinition, AttributeValue, BillingMode, GlobalSecondaryIndex, KeySchemaElement, + KeyType, Projection, ProjectionType, Put, PutRequest, ScalarAttributeType, StreamSpecification, + StreamViewType, Tag, TransactWriteItem, WriteRequest, +}; +use aws_sdk_dynamodb::Client as DynamoClient; +use aws_sdk_iam::Client as IamClient; +use helpers::TestServer; + +const ACCOUNT: &str = "123456789012"; +const REGION: &str = "us-east-1"; + +async fn start_strict() -> TestServer { + TestServer::start_with_env(&[ + ("FAKECLOUD_IAM", "strict"), + ("FAKECLOUD_VERIFY_SIGV4", "true"), + ]) + .await +} + +async fn sdk_config_with(server: &TestServer, akid: &str, secret: &str) -> aws_config::SdkConfig { + aws_config::defaults(aws_config::BehaviorVersion::latest()) + .endpoint_url(server.endpoint()) + .region(aws_config::Region::new(REGION)) + .credentials_provider(Credentials::new( + akid, + secret, + None, + None, + "fakecloud-dynamodb-iam", + )) + .load() + .await +} + +async fn admin(server: &TestServer) -> DynamoClient { + DynamoClient::new(&sdk_config_with(server, "test", "test").await) +} + +/// A user whose only permissions are `policy`, and a DynamoDB client signed +/// with that user's credentials. +async fn user_with_policy(server: &TestServer, name: &str, policy: &str) -> DynamoClient { + let boot = sdk_config_with(server, "test", "test").await; + let iam = IamClient::new(&boot); + iam.create_user().user_name(name).send().await.unwrap(); + let key = iam + .create_access_key() + .user_name(name) + .send() + .await + .unwrap(); + let key = key.access_key().unwrap(); + iam.put_user_policy() + .user_name(name) + .policy_name("inline") + .policy_document(policy) + .send() + .await + .unwrap(); + DynamoClient::new(&sdk_config_with(server, key.access_key_id(), key.secret_access_key()).await) +} + +fn table_arn(name: &str) -> String { + format!("arn:aws:dynamodb:{REGION}:{ACCOUNT}:table/{name}") +} + +fn allow(actions: &[&str], resources: &[String]) -> String { + serde_json::json!({ + "Version": "2012-10-17", + "Statement": [{"Effect": "Allow", "Action": actions, "Resource": resources}] + }) + .to_string() +} + +async fn create_table(client: &DynamoClient, name: &str) { + client + .create_table() + .table_name(name) + .key_schema( + KeySchemaElement::builder() + .attribute_name("pk") + .key_type(KeyType::Hash) + .build() + .unwrap(), + ) + .attribute_definitions( + AttributeDefinition::builder() + .attribute_name("pk") + .attribute_type(ScalarAttributeType::S) + .build() + .unwrap(), + ) + .attribute_definitions( + AttributeDefinition::builder() + .attribute_name("g") + .attribute_type(ScalarAttributeType::S) + .build() + .unwrap(), + ) + .global_secondary_indexes( + GlobalSecondaryIndex::builder() + .index_name("by-g") + .key_schema( + KeySchemaElement::builder() + .attribute_name("g") + .key_type(KeyType::Hash) + .build() + .unwrap(), + ) + .projection( + Projection::builder() + .projection_type(ProjectionType::All) + .build(), + ) + .build() + .unwrap(), + ) + .billing_mode(BillingMode::PayPerRequest) + .stream_specification( + StreamSpecification::builder() + .stream_enabled(true) + .stream_view_type(StreamViewType::NewImage) + .build() + .unwrap(), + ) + .send() + .await + .unwrap(); +} + +fn denied(result: Result) -> bool { + match result { + Ok(_) => false, + Err(e) => format!("{e:?}").contains("AccessDenied"), + } +} + +#[tokio::test] +async fn dynamodb_requires_a_policy_under_strict_enforcement() { + let server = start_strict().await; + create_table(&admin(&server).await, "Orders").await; + let nobody = user_with_policy( + &server, + "nobody", + &allow(&["sqs:ListQueues"], &["*".to_string()]), + ) + .await; + + assert!(denied(nobody.list_tables().send().await)); + assert!(denied( + nobody + .get_item() + .table_name("Orders") + .key("pk", AttributeValue::S("a".into())) + .send() + .await + )); +} + +/// A policy scoped to one action on one table allows exactly that. +#[tokio::test] +async fn item_actions_are_scoped_to_the_table_and_action() { + let server = start_strict().await; + let admin = admin(&server).await; + create_table(&admin, "Orders").await; + create_table(&admin, "Customers").await; + let reader = user_with_policy( + &server, + "reader", + &allow(&["dynamodb:GetItem"], &[table_arn("Orders")]), + ) + .await; + + reader + .get_item() + .table_name("Orders") + .key("pk", AttributeValue::S("a".into())) + .send() + .await + .expect("GetItem on the allowed table"); + assert!(denied( + reader + .put_item() + .table_name("Orders") + .item("pk", AttributeValue::S("a".into())) + .send() + .await + )); + assert!(denied( + reader + .get_item() + .table_name("Customers") + .key("pk", AttributeValue::S("a".into())) + .send() + .await + )); + // The table's ARN authorizes the same as its name. + reader + .get_item() + .table_name(table_arn("Orders")) + .key("pk", AttributeValue::S("a".into())) + .send() + .await + .expect("GetItem by table ARN"); +} + +/// A Query on an index is authorized against the index's ARN, not the +/// table's. +#[tokio::test] +async fn index_queries_are_authorized_against_the_index() { + let server = start_strict().await; + create_table(&admin(&server).await, "Orders").await; + let table_only = user_with_policy( + &server, + "table-only", + &allow(&["dynamodb:Query"], &[table_arn("Orders")]), + ) + .await; + let query_index = |client: DynamoClient| async move { + client + .query() + .table_name("Orders") + .index_name("by-g") + .key_condition_expression("g = :g") + .expression_attribute_values(":g", AttributeValue::S("x".into())) + .send() + .await + }; + assert!(denied(query_index(table_only).await)); + + let with_index = user_with_policy( + &server, + "with-index", + &allow( + &["dynamodb:Query"], + &[ + table_arn("Orders"), + format!("{}/index/*", table_arn("Orders")), + ], + ), + ) + .await; + query_index(with_index) + .await + .expect("Query on an allowed index"); +} + +/// A batch or transaction needs the permission on every table it touches: +/// being allowed on one of them is not enough. +#[tokio::test] +async fn batches_and_transactions_need_every_table() { + let server = start_strict().await; + let admin = admin(&server).await; + create_table(&admin, "Orders").await; + create_table(&admin, "Customers").await; + + let orders_only = user_with_policy( + &server, + "orders-only", + &allow( + &["dynamodb:PutItem", "dynamodb:BatchWriteItem"], + &[table_arn("Orders")], + ), + ) + .await; + let both = user_with_policy( + &server, + "both", + &allow( + &["dynamodb:PutItem", "dynamodb:BatchWriteItem"], + &[table_arn("Orders"), table_arn("Customers")], + ), + ) + .await; + + let put = |table: &str, pk: &str| { + TransactWriteItem::builder() + .put( + Put::builder() + .table_name(table) + .item("pk", AttributeValue::S(pk.into())) + .build() + .unwrap(), + ) + .build() + }; + assert!(denied( + orders_only + .transact_write_items() + .transact_items(put("Orders", "a")) + .transact_items(put("Customers", "a")) + .send() + .await + )); + both.transact_write_items() + .transact_items(put("Orders", "a")) + .transact_items(put("Customers", "a")) + .send() + .await + .expect("transaction allowed on both tables"); + + let write = |pk: &str| { + WriteRequest::builder() + .put_request( + PutRequest::builder() + .item("pk", AttributeValue::S(pk.into())) + .build() + .unwrap(), + ) + .build() + }; + assert!(denied( + orders_only + .batch_write_item() + .request_items("Orders", vec![write("b")]) + .request_items("Customers", vec![write("b")]) + .send() + .await + )); + orders_only + .batch_write_item() + .request_items("Orders", vec![write("b")]) + .send() + .await + .expect("batch on the allowed table only"); + // Nothing the denied requests carried was written. + let scan = admin.scan().table_name("Customers").send().await.unwrap(); + let pks: Vec<&str> = scan + .items() + .iter() + .map(|i| i["pk"].as_s().unwrap().as_str()) + .collect(); + assert_eq!(pks, ["a"]); +} + +/// PartiQL statements are authorized by their verb. +#[tokio::test] +async fn partiql_statements_need_their_partiql_action() { + let server = start_strict().await; + create_table(&admin(&server).await, "Orders").await; + let selector = user_with_policy( + &server, + "selector", + &allow(&["dynamodb:PartiQLSelect"], &[table_arn("Orders")]), + ) + .await; + + selector + .execute_statement() + .statement("SELECT * FROM \"Orders\"") + .send() + .await + .expect("SELECT allowed"); + assert!(denied( + selector + .execute_statement() + .statement("INSERT INTO \"Orders\" VALUE {'pk': 'a'}") + .send() + .await + )); + // The data-plane action does not grant the PartiQL one, or the reverse. + assert!(denied(selector.scan().table_name("Orders").send().await)); +} + +/// `aws:ResourceTag/*` conditions see the table's tags, and CreateTable with +/// tags also needs `dynamodb:TagResource`. +#[tokio::test] +async fn tag_conditions_and_create_table_tags() { + let server = start_strict().await; + let admin = admin(&server).await; + create_table(&admin, "Tagged").await; + create_table(&admin, "Untagged").await; + admin + .tag_resource() + .resource_arn(table_arn("Tagged")) + .tags( + Tag::builder() + .key("team") + .value("payments") + .build() + .unwrap(), + ) + .send() + .await + .unwrap(); + + let policy = serde_json::json!({ + "Version": "2012-10-17", + "Statement": [ + { + "Effect": "Allow", + "Action": "dynamodb:GetItem", + "Resource": "*", + "Condition": {"StringEquals": {"aws:ResourceTag/team": "payments"}} + }, + { + "Effect": "Allow", + "Action": "dynamodb:CreateTable", + "Resource": "*" + } + ] + }) + .to_string(); + let payments = user_with_policy(&server, "payments", &policy).await; + let get = |client: &DynamoClient, table: &str| { + client + .get_item() + .table_name(table) + .key("pk", AttributeValue::S("a".into())) + .send() + }; + get(&payments, "Tagged").await.expect("tag matches"); + assert!(denied(get(&payments, "Untagged").await)); + + let create = |name: &str, tagged: bool| { + let mut req = payments + .create_table() + .table_name(name) + .key_schema( + KeySchemaElement::builder() + .attribute_name("pk") + .key_type(KeyType::Hash) + .build() + .unwrap(), + ) + .attribute_definitions( + AttributeDefinition::builder() + .attribute_name("pk") + .attribute_type(ScalarAttributeType::S) + .build() + .unwrap(), + ) + .billing_mode(BillingMode::PayPerRequest); + if tagged { + req = req.tags( + Tag::builder() + .key("team") + .value("payments") + .build() + .unwrap(), + ); + } + req.send() + }; + create("Plain", false).await.expect("CreateTable allowed"); + assert!( + denied(create("WithTags", true).await), + "CreateTable with Tags also needs dynamodb:TagResource" + ); +} + +/// DynamoDB Streams operations are authorized against the stream ARN. +#[tokio::test] +async fn streams_are_authorized_against_the_stream() { + let server = start_strict().await; + let admin = admin(&server).await; + create_table(&admin, "Orders").await; + let stream_arn = admin + .describe_table() + .table_name("Orders") + .send() + .await + .unwrap() + .table() + .unwrap() + .latest_stream_arn() + .unwrap() + .to_string(); + + let boot = sdk_config_with(&server, "test", "test").await; + let iam = IamClient::new(&boot); + iam.create_user() + .user_name("streamer") + .send() + .await + .unwrap(); + let key = iam + .create_access_key() + .user_name("streamer") + .send() + .await + .unwrap(); + let key = key.access_key().unwrap(); + iam.put_user_policy() + .user_name("streamer") + .policy_name("inline") + .policy_document(allow( + &["dynamodb:DescribeStream"], + std::slice::from_ref(&stream_arn), + )) + .send() + .await + .unwrap(); + let streams = aws_sdk_dynamodbstreams::Client::new( + &sdk_config_with(&server, key.access_key_id(), key.secret_access_key()).await, + ); + + let described = streams + .describe_stream() + .stream_arn(&stream_arn) + .send() + .await + .expect("DescribeStream allowed"); + let shard = described.stream_description().unwrap().shards()[0] + .shard_id() + .unwrap() + .to_string(); + assert!(denied( + streams + .get_shard_iterator() + .stream_arn(&stream_arn) + .shard_id(shard) + .shard_iterator_type(aws_sdk_dynamodbstreams::types::ShardIteratorType::TrimHorizon) + .send() + .await + )); + assert!(denied(streams.list_streams().send().await)); +} diff --git a/website/content/docs/reference/security.md b/website/content/docs/reference/security.md index f16d5a76b..045c448ac 100644 --- a/website/content/docs/reference/security.md +++ b/website/content/docs/reference/security.md @@ -67,6 +67,8 @@ Opt-in enforcement covers the services most commonly subject to real IAM policie | **SNS** | All 34 supported actions | Topic / subscription / platform-app / endpoint ARNs | | **S3** | All 74 supported actions | `arn:aws:s3:::[/]` (object actions include the key; bucket actions don't) | | **KMS** | All 47 supported actions | `arn:aws:kms:::key/` (key-targeted actions) or `*` (account-level actions like CreateKey, ListKeys) | +| **DynamoDB** | All 58 supported operations | `arn:aws:dynamodb:::table/`, with `/index/` for `Query`, `Scan` and contributor insights on an index, `/backup/...`, `/export/...` and `/import/...` for those operations, `arn:aws:dynamodb:::global-table/` for legacy global tables, and `*` for account-level listings. Batches need the batch action on every table they name; transactions need `GetItem` / `PutItem` / `UpdateItem` / `DeleteItem` / `ConditionCheckItem` on each item's table; PartiQL statements need `PartiQLSelect` / `PartiQLInsert` / `PartiQLUpdate` / `PartiQLDelete`; `CreateTable` with `Tags` or `ResourcePolicy` also needs `TagResource` / `PutResourcePolicy`; restores also need the data-plane actions on the target table | +| **DynamoDB Streams** | All 4 supported operations (`dynamodb:` prefix) | `arn:aws:dynamodb:::table//stream/