From 9864cb1cc4c9fcfc0d5b73f44694f9839e7eb605 Mon Sep 17 00:00:00 2001 From: Speculator55005 Date: Mon, 28 Sep 2026 21:33:55 +0200 Subject: [PATCH 1/2] fix(iac): one bucket policy per bucket, fallback LF ARNs name the account, generate iac calls no AWS API Three defects found on main while fixing the examples for #674. 1. On the {account}-fluid-data fallback bucket (an unresolved {{ env.DATA_LAKE_BUCKET }}), the Lake Formation location and the bucket-policy ARNs were plain strings, so the renderer escaped the account token to $${...}, a literal: tofu plan registered arn:aws:s3:::${data.aws_caller_identity...}-fluid-data/... as written. The bucket policy also referenced an aws_s3_bucket the module never declares, so a plan with an ARN grantee failed "Reference to undeclared resource". The fallback ARNs are now TofuExpr with the contract-derived suffix escaped (the way _emit_glue builds the table location), decided on the raw input, and the policy addresses the fallback bucket by name. 2. Every expose emitted the bucket's aws_s3_bucket_policy under the same key, so a second expose on the bucket replaced the first: a grantee in another account lost s3:GetObject on the first prefix. The exposes on a bucket are now merged into one policy, grantees numbered across it so Sids and data.aws_arn keys stay unique. A one-expose bucket gets the same policy as before, byte for byte. Two exposes asking one bucket for different bucketPolicy modes are refused. 3. fluid generate iac built the native AWS provider, which resolves a missing account with sts:GetCallerIdentity over whatever credentials the machine has. generate iac now gives the planner a placeholder account when AWS_ACCOUNT_ID is unset, so no boto3 client is built and the sovereignty check still runs, and refuses a module that would carry the placeholder (generate_iac_aws_account_required, naming AWS_ACCOUNT_ID). fluid apply keeps resolving the account as before. Proved by tofu plan against moto (tests/iac/test_iac_lakeformation_bucket_policy_plan.py) and by a guard on botocore's Session.create_client (tests/iac/test_iac_generate_offline.py). --- fluid_build/_error_catalog.py | 8 + fluid_build/cli/generate_iac.py | 57 ++- fluid_build/iac/providers/aws.py | 336 ++++++++++++------ tests/iac/_lf_bucket_policy.py | 59 +-- tests/iac/test_iac_generate_offline.py | 212 +++++++++++ .../test_iac_lakeformation_bucket_policy.py | 216 +++++++++++ ...st_iac_lakeformation_bucket_policy_plan.py | 111 +++++- tests/iac/test_iac_packaging_default_pin.py | 6 +- 8 files changed, 869 insertions(+), 136 deletions(-) create mode 100644 tests/iac/test_iac_generate_offline.py diff --git a/fluid_build/_error_catalog.py b/fluid_build/_error_catalog.py index a3ce0340..254c1f59 100644 --- a/fluid_build/_error_catalog.py +++ b/fluid_build/_error_catalog.py @@ -307,6 +307,14 @@ def slug_for(event: str) -> str: ], None, ), + "generate_iac_aws_account_required": ( + [ + "Set AWS_ACCOUNT_ID to the account the module will be applied in, then re-run " + "'fluid generate iac' (it never looks the account up in AWS)", + "'fluid apply' resolves the account itself and needs no AWS_ACCOUNT_ID", + ], + None, + ), "generate_ci_failed": ( [ "Check the --system value is a supported CI provider and the contract validates", diff --git a/fluid_build/cli/generate_iac.py b/fluid_build/cli/generate_iac.py index 73c16a9e..1fa6e732 100644 --- a/fluid_build/cli/generate_iac.py +++ b/fluid_build/cli/generate_iac.py @@ -44,6 +44,12 @@ _PROVIDER_CHOICES = ["auto", *sorted(IAC_PLUGINS)] +#: The AWS account the native planner is given when ``fluid generate iac`` runs +#: without ``AWS_ACCOUNT_ID`` (see :func:`native_actions`). Shaped to pass the +#: planner's own name checks (it becomes part of bucket names), and never an +#: account id, so a module that would carry it is refused rather than written. +UNSET_AWS_ACCOUNT = "fluid-aws-account-id-unset" + def register_subcommand(subparsers: argparse._SubParsersAction): """Register as a subcommand of ``fluid generate``.""" @@ -116,7 +122,7 @@ def run(args, logger: logging.Logger) -> int: logger.debug("packaging resolved: legacy=%s pool=%s", packaging.is_legacy, packaging.pool) provider = _resolve_provider(contract, getattr(args, "provider", "auto")) plugin = get_iac_plugin(provider) - actions = native_actions(contract, logger) + actions = native_actions(contract, logger, offline=True) resources = plugin.emit(contract, actions) count = sum(len(items) for items in resources.values()) provider_cfg = provider_config(plugin, contract) @@ -127,11 +133,13 @@ def run(args, logger: logging.Logger) -> int: # `.tf.json` keys the provider block by the provider's local name. provider={plugin.name: provider_cfg} if provider_cfg else None, ) + rendered = render_tofu_json(document) + _refuse_an_unset_aws_account(rendered) out_dir = getattr(args, "out", None) or "runtime/iac" os.makedirs(out_dir, exist_ok=True) out_path = os.path.join(out_dir, "main.tf.json") with open(out_path, "w", encoding="utf-8") as f: - f.write(render_tofu_json(document)) + f.write(rendered) except CLIError: raise except UnsupportedBindingError as exc: @@ -188,6 +196,32 @@ def run(args, logger: logging.Logger) -> int: return 0 +def _refuse_an_unset_aws_account(rendered: str) -> None: + """Refuse a module that would carry :data:`UNSET_AWS_ACCOUNT`. + + The emitter's own account references are ``data.aws_caller_identity`` + lookups that ``tofu`` resolves at plan time; only the native planner's + actions (Lambda, EventBridge and Step Functions ARNs, for instance) spell + the account out. When they reach the module and no account was given, + the module is refused, naming ``AWS_ACCOUNT_ID``, rather than written with + a placeholder where the account belongs. + """ + if UNSET_AWS_ACCOUNT not in rendered: + return + raise CLIError( + 1, + "generate_iac_aws_account_required", + { + "error": ( + "the module names the AWS account (in the ARNs of the resources the " + "contract's orchestration plans), and AWS_ACCOUNT_ID is not set.\n" + " `fluid generate iac` never asks AWS for the account. Set " + "AWS_ACCOUNT_ID to the account the module will be applied in, then re-run." + ) + }, + ) + + def _validate_with_tofu(out_dir: str) -> None: """Run ``tofu init -backend=false`` + ``tofu validate`` on the emitted module. @@ -318,7 +352,7 @@ def _resolve_provider(contract, requested: str) -> str: return cloud -def native_actions(contract, logger: logging.Logger) -> list: +def native_actions(contract, logger: logging.Logger, *, offline: bool = False) -> list: """Best-effort native ``provider.plan()`` actions for the contract. The OpenTofu emitter consumes these to translate the schedule / @@ -334,12 +368,25 @@ def native_actions(contract, logger: logging.Logger) -> list: apply path kept announcing three actions for a contract that owns exactly one leaf table — while ``tofu`` correctly planned ``+1``. This is the apply-side chokepoint, matching ``cli/plan.py``'s. + + ``offline`` is set by ``fluid generate iac``, which must reach no cloud + API. The AWS provider resolves a missing account with + ``sts:GetCallerIdentity`` over whatever credentials the machine has, so + without ``AWS_ACCOUNT_ID`` it is given :data:`UNSET_AWS_ACCOUNT` instead: + the planner (and its sovereignty check) still runs, and :func:`run` + refuses a module that would carry the placeholder. ``fluid apply`` leaves + it unset and resolves the account as before. """ try: from ._common import build_provider, resolve_provider_from_contract name, loc = resolve_provider_from_contract(contract) - native = build_provider(name, loc.get("project"), loc.get("region"), logger) + account = loc.get("project") + # The provider name as ``build_provider`` resolves it. + provider = (name or os.getenv("FLUID_PROVIDER") or "").strip().lower().replace("-", "_") + if offline and provider == "aws" and not account and not os.getenv("AWS_ACCOUNT_ID"): + account = UNSET_AWS_ACCOUNT + native = build_provider(name, account, loc.get("region"), logger) if hasattr(native, "plan"): return _drop_referenced_container_actions(contract, list(native.plan(contract)), logger) except Exception as exc: # noqa: BLE001 — native planner is best-effort @@ -421,7 +468,7 @@ def _print_shadow_report(contract, plugin, logger: logging.Logger) -> None: """Run shadow-compare and print the native↔OpenTofu parity report.""" from fluid_build.iac import shadow_compare - actions = native_actions(contract, logger) + actions = native_actions(contract, logger, offline=True) if not actions: cprint( "\nShadow-compare: native planner produced no actions " diff --git a/fluid_build/iac/providers/aws.py b/fluid_build/iac/providers/aws.py index 7ac76dc3..11dc33d0 100644 --- a/fluid_build/iac/providers/aws.py +++ b/fluid_build/iac/providers/aws.py @@ -396,6 +396,11 @@ def emit( # (binding.encryption), per bucket this product owns. Nothing is # added for a contract that declares neither. _emit_storage_policies(resources, contract, cid) + # The S3 bucket policy of the Lake Formation grants: one per bucket, + # with the statements of every exposure on it (a second policy on the + # bucket would replace the first). + for bucket_policy in _lf_bucket_policies(contract, cid).values(): + _emit_lf_bucket_policy(resources, bucket_policy) # Glue ETL jobs / Step Functions / the Lambda schedule path — # the planner's build & orchestration ops. _emit_from_actions(resources, actions, cid) @@ -782,8 +787,7 @@ def _emit_glue( # (``bucket_uses_fallback``), which a contract cannot spoof. bucket, path = _warehouse.normalize_location(loc, account_ref=_CALLER_ACCOUNT_TOKEN) if _warehouse.bucket_uses_fallback(loc): - safe_path = path.replace("${", "$${").replace("%{", "%%{") - storage["location"] = TofuExpr(f"s3://{bucket}/{safe_path}") + storage["location"] = TofuExpr(f"s3://{bucket}/{_literal(path)}") else: storage["location"] = f"s3://{bucket}/{path}" parameters: Dict[str, str] = {"classification": fmt, "managed_by": "fluid"} @@ -1583,6 +1587,34 @@ def _lf_location(loc: Mapping[str, Any], placement: _Placement) -> Tuple[Optiona return bucket, path +def _literal(text: str) -> str: + """``text`` escaped as the renderer escapes a plain string (``${`` and ``%{``). + + For contract-derived text that has to sit inside an emitter-built + ``TofuExpr``, which the renderer leaves as it is. + """ + return text.replace("${", "$${").replace("%{", "%%{") + + +def _lf_s3_arn(loc: Mapping[str, Any], bucket: str, suffix: str = "") -> str: + """``arn:aws:s3:::``, for the Lake Formation location and bucket policy. + + On the ``{account}-fluid-data`` fallback (:func:`_lf_location`) the bucket + carries ``_CALLER_ACCOUNT_TOKEN``, and the ARN has to reach OpenTofu as an + interpolation. As a plain string the renderer escaped it to ``$${...}``, + a literal (the OpenTofu string docs: "Use ``$${`` to produce a literal + ``${``"), so ``registerLocation`` registered an ARN whose account was + never filled in. So the fallback ARN is a ``TofuExpr``, with the + contract-derived ``suffix`` escaped first, the way :func:`_emit_glue` + builds the table's location. The fallback is decided on the raw input + (``bucket_uses_fallback``), which a contract cannot spoof by writing the + token itself; every other ARN stays a plain string the renderer escapes. + """ + if not _warehouse.bucket_uses_fallback(loc): + return f"arn:aws:s3:::{bucket}{suffix}" + return TofuExpr(f"arn:aws:s3:::{bucket}{_literal(suffix)}") + + # ``binding.governance.lakeFormation.bucketPolicy`` (fluid-schema 0.7.6): which # ``grants[]`` principals get a statement in the companion ``aws_s3_bucket_policy``. _LF_BUCKET_POLICY_CROSS_ACCOUNT = "cross-account" # the default @@ -1595,9 +1627,21 @@ def _lf_location(loc: Mapping[str, Any], placement: _Placement) -> Tuple[Optiona ) +@dataclass(frozen=True) +class _LfBucketGrants: + """One exposure's statements in its bucket's policy (see :func:`_lf_bucket_policies`).""" + + object_arn: str + #: ``s3:prefix`` values ListBucket is narrowed to on a shared pool, else ``()``. + list_prefixes: Tuple[str, ...] + principals: Tuple[str, ...] + #: Index of the first of ``principals`` among every grantee of the policy. + first: int = 0 + + @dataclass(frozen=True) class _LfBucketPolicy: - """The bucket-policy half of one exposure's LF grants (see :func:`_lf_bucket_policy`).""" + """The bucket policy of one bucket's LF grants (see :func:`_lf_bucket_policy`).""" mode: str #: Key of the ``aws_s3_bucket_policy`` resource, and of its policy-document @@ -1605,10 +1649,13 @@ class _LfBucketPolicy: policy_key: str bucket_ref: TofuExpr bucket_arn: str - object_arn: str - #: ``s3:prefix`` values ListBucket is narrowed to on a shared pool, else ``()``. - list_prefixes: Tuple[str, ...] - principals: Tuple[str, ...] + #: One entry per exposure whose grants land in the bucket, in contract order. + grants: Tuple[_LfBucketGrants, ...] + + @property + def principals(self) -> Tuple[str, ...]: + """Every grantee of the policy, in index order.""" + return tuple(p for grants in self.grants for p in grants.principals) def grantee_key(self, index: int) -> str: """Key of the ``data.aws_arn`` that parses grantee ``index``'s ARN.""" @@ -1665,7 +1712,8 @@ def _lf_bucket_policy( ) -> Optional[_LfBucketPolicy]: """The ``aws_s3_bucket_policy`` one exposure's LF grants need, or ``None``. - THE one derivation, shared by :meth:`AwsIacPlugin.emit` (the resource) and + THE one derivation, which :func:`_lf_bucket_policies` merges per bucket + for :meth:`AwsIacPlugin.emit` (the resource) and :meth:`AwsIacPlugin.emit_data` (its data sources), so the two halves can never disagree about keys, grantees or ARNs. @@ -1711,47 +1759,131 @@ def _lf_bucket_policy( # statements degrade to the whole bucket and this authoritative policy also # replaces the platform team's own. Fail closed. _require_pool_prefix(placement, path, what="governance.lakeFormation.grants[]") + if _warehouse.bucket_uses_fallback(loc): + # The ``{account}-fluid-data`` fallback is no ``aws_s3_bucket`` of this + # module, so it is addressed by its name, interpolated at plan time. + bucket_ref = TofuExpr(bucket) + else: + bucket_ref = _s3_bucket_ref( + safe_ident(f"{cid}_{bucket}"), referenced=placement.bucket_referenced + ) return _LfBucketPolicy( mode=mode, policy_key=safe_ident(f"{cid}_lf_bucket_policy_{bucket}"), - bucket_ref=_s3_bucket_ref( - safe_ident(f"{cid}_{bucket}"), referenced=placement.bucket_referenced + bucket_ref=bucket_ref, + bucket_arn=_lf_s3_arn(loc, bucket), + grants=( + _LfBucketGrants( + # ListBucket targets the bucket ARN itself; GetObject targets the + # per-object ARN under the configured path prefix (or everything if + # no path). + object_arn=_lf_s3_arn(loc, bucket, f"/{path}*" if path else "/*"), + # ListBucket is inherently bucket-scoped, so on a shared pool it is + # narrowed with the standard ``s3:prefix`` condition; otherwise this + # product's consumers could enumerate every other tenant's keys in the + # pool. ``path`` is guaranteed non-empty here by ``_require_pool_prefix``. + list_prefixes=(f"{path}*",) if placement.bucket_referenced else (), + principals=tuple(principals), + ), ), - bucket_arn=f"arn:aws:s3:::{bucket}", - # ListBucket targets the bucket ARN itself; GetObject targets the - # per-object ARN under the configured path prefix (or everything if - # no path). - object_arn=f"arn:aws:s3:::{bucket}/{path}*" if path else f"arn:aws:s3:::{bucket}/*", - # ListBucket is inherently bucket-scoped, so on a shared pool it is - # narrowed with the standard ``s3:prefix`` condition — otherwise this - # product's consumers could enumerate every other tenant's keys in the - # pool. ``path`` is guaranteed non-empty here by ``_require_pool_prefix``. - list_prefixes=(f"{path}*",) if placement.bucket_referenced else (), - principals=tuple(principals), ) -def _lf_other_account_grantees(bp: _LfBucketPolicy) -> str: +def _lf_bucket_policies(contract: Mapping[str, Any], cid: str) -> Dict[str, _LfBucketPolicy]: + """Every bucket's one ``aws_s3_bucket_policy``, by resource key, over all its exposures. + + S3 keeps one policy per bucket, and the hashicorp/aws docs warn that more + than one ``aws_s3_bucket_policy`` on a bucket means "the policy applied + last will silently override any previously applied policy". Each exposure + used to emit its own under the bucket's key, so a second exposure on the + bucket replaced the first one's statements: a grantee in another account + lost ``s3:GetObject`` on the first prefix. So the exposures on a bucket are + merged, in contract order, into one policy whose grantees are numbered + across all of them (Sids and ``data.aws_arn`` keys stay unique). A bucket + with one exposure gets exactly the policy it got before. + + Walks the exposures as :meth:`AwsIacPlugin.emit` does. Two exposures that + ask the bucket for different ``bucketPolicy`` modes are refused, because + the one policy cannot be both; ``none`` contributes no statement. + """ + packaging = resolve_packaging(contract) + policies: Dict[str, _LfBucketPolicy] = {} + for exposure in contract.get("exposes") or []: + binding = exposure.get("binding") or {} + if not is_cloud(binding, "aws"): + continue + bp = _lf_bucket_policy( + binding, + binding.get("location") or {}, + binding.get("format") or "parquet", + cid, + placement=_placement(packaging, exposure), + ) + if bp is None: + continue + merged = policies.get(bp.policy_key) + if merged is None: + policies[bp.policy_key] = bp + continue + if merged.mode != bp.mode: + raise UnsupportedBindingError( + "lakeformation-bucket-policy", + f"two exposes grant through Lake Formation on the bucket {bp.bucket_arn!r} " + f"with different bucketPolicy values ({merged.mode!r} and {bp.mode!r}); a " + "bucket has one bucket policy.", + ("Declare the same bucketPolicy on every expose whose grants land in the bucket.",), + ) + policies[bp.policy_key] = _LfBucketPolicy( + mode=merged.mode, + policy_key=merged.policy_key, + bucket_ref=merged.bucket_ref, + bucket_arn=merged.bucket_arn, + grants=merged.grants + + tuple( + _LfBucketGrants( + object_arn=grants.object_arn, + list_prefixes=grants.list_prefixes, + principals=grants.principals, + first=len(merged.principals), + ) + for grants in bp.grants + ), + ) + return policies + + +def _lf_other_account_grantees( + bp: _LfBucketPolicy, grants: Optional[_LfBucketGrants] = None +) -> str: """HCL map ``{"" = }`` of the grantees NOT in the applying account. + Of one exposure's ``grants``, or of the whole policy: a ``merge`` of every + exposure's map, whose keys never collide because the grantees are numbered + across the policy. A one-exposure policy's map is written as before. + Only emitter-controlled text goes in: ``safe_ident`` keys and literal attribute names. The ARNs themselves stay in the ``data.aws_arn`` blocks, as plain strings the renderer escapes, so a contract cannot inject an interpolation through this expression. """ - grantees = ", ".join(f"data.aws_arn.{bp.grantee_key(i)}" for i in range(len(bp.principals))) + if grants is None: + maps = [_lf_other_account_grantees(bp, each) for each in bp.grants] + return maps[0] if len(maps) == 1 else f"merge({', '.join(maps)})" + grantees = ", ".join( + f"data.aws_arn.{bp.grantee_key(grants.first + i)}" for i in range(len(grants.principals)) + ) + index = f"i + {grants.first}" if grants.first else "i" return ( - f"{{for i, g in [{grantees}] : tostring(i) => g.arn " + f"{{for i, g in [{grantees}] : tostring({index}) => g.arn " "if g.account != data.aws_caller_identity.fluid_lf_caller.account_id}" ) def _emit_lf_bucket_policy(resources: Dict[str, Any], bp: _LfBucketPolicy) -> None: - """Emit the ``aws_s3_bucket_policy`` resource for one exposure's LF grants. + """Emit the ``aws_s3_bucket_policy`` resource for one bucket's LF grants. - One resource per bucket; a later exposure on the same bucket overwrites an - earlier one's, because there is exactly one ``aws_s3_bucket_policy`` slot - per bucket. + One resource per bucket, carrying the statements of every exposure whose + grants land in it (:func:`_lf_bucket_policies`). NOTE (v2): ``aws_s3_bucket_policy`` is authoritative for the whole bucket, so two products sharing one pool would each rewrite the other's policy on @@ -1761,32 +1893,43 @@ def _emit_lf_bucket_policy(resources: Dict[str, Any], bp: _LfBucketPolicy) -> No """ body: Dict[str, Any] = {"bucket": bp.bucket_ref} if bp.mode == _LF_BUCKET_POLICY_ALL_GRANTEES: + # On the ``{account}-fluid-data`` fallback the ARNs interpolate the + # account (:func:`_lf_s3_arn`), so the document is a ``TofuExpr`` and + # the contract-derived text in it is escaped here instead of by the + # renderer. + live = isinstance(bp.bucket_arn, TofuExpr) + text = _literal if live else str statements: List[Dict[str, Any]] = [] - for sid_idx, principal in enumerate(bp.principals): - list_statement: Dict[str, Any] = { - "Sid": f"FluidLfBucketList{sid_idx}", - "Effect": "Allow", - "Principal": {"AWS": principal}, - "Action": ["s3:ListBucket", "s3:GetBucketLocation"], - "Resource": bp.bucket_arn, - } - if bp.list_prefixes: - list_statement["Condition"] = {"StringLike": {"s3:prefix": list(bp.list_prefixes)}} - statements.append(list_statement) - statements.append( - { - "Sid": f"FluidLfBucketGet{sid_idx}", + for grants in bp.grants: + for offset, principal in enumerate(grants.principals): + sid_idx = grants.first + offset + list_statement: Dict[str, Any] = { + "Sid": f"FluidLfBucketList{sid_idx}", "Effect": "Allow", - "Principal": {"AWS": principal}, - "Action": ["s3:GetObject"], - "Resource": bp.object_arn, + "Principal": {"AWS": text(principal)}, + "Action": ["s3:ListBucket", "s3:GetBucketLocation"], + "Resource": bp.bucket_arn, } - ) - body["policy"] = json.dumps( + if grants.list_prefixes: + list_statement["Condition"] = { + "StringLike": {"s3:prefix": [text(p) for p in grants.list_prefixes]} + } + statements.append(list_statement) + statements.append( + { + "Sid": f"FluidLfBucketGet{sid_idx}", + "Effect": "Allow", + "Principal": {"AWS": text(principal)}, + "Action": ["s3:GetObject"], + "Resource": grants.object_arn, + } + ) + policy = json.dumps( {"Version": "2012-10-17", "Statement": statements}, sort_keys=True, separators=(",", ":"), ) + body["policy"] = TofuExpr(policy) if live else policy else: # cross-account: the statements are built at plan time by the policy # document :func:`_emit_lf_bucket_policy_data` declares; no instance @@ -1801,61 +1944,50 @@ def _emit_lf_bucket_policy_data( ) -> None: """Declare the data sources a ``cross-account`` LF bucket policy reads. - Walks the exposures exactly as :meth:`AwsIacPlugin.emit` does and asks - :func:`_lf_bucket_policy` the same question, so every ``data.`` reference + Asks :func:`_lf_bucket_policies` the same question as + :meth:`AwsIacPlugin.emit`, so every ``data.`` reference :func:`_emit_lf_bucket_policy` writes has its declaration here. Per bucket policy: one ``aws_arn`` per grantee (the provider's own ARN parser supplies ``.account``) and one ``aws_iam_policy_document`` whose ``dynamic`` - statements iterate only over the grantees in other accounts. Nothing is - added for ``all-grantees`` or ``none``. + statements iterate only over the grantees in other accounts, a pair per + exposure on the bucket, each over its own grantees and scoped to its own + prefix. Nothing is added for ``all-grantees`` or ``none``. """ - packaging = resolve_packaging(contract) - for exposure in contract.get("exposes") or []: - binding = exposure.get("binding") or {} - if not is_cloud(binding, "aws"): - continue - bp = _lf_bucket_policy( - binding, - binding.get("location") or {}, - binding.get("format") or "parquet", - cid, - placement=_placement(packaging, exposure), - ) - if bp is None or bp.mode != _LF_BUCKET_POLICY_CROSS_ACCOUNT: + for bp in _lf_bucket_policies(contract, cid).values(): + if bp.mode != _LF_BUCKET_POLICY_CROSS_ACCOUNT: continue arns = data.setdefault("aws_arn", {}) for index, principal in enumerate(bp.principals): arns[bp.grantee_key(index)] = {"arn": principal} - grantees = tofu_ref(_lf_other_account_grantees(bp)) - list_content: Dict[str, Any] = { - # ``statement.key`` is the grantee's index, so a Sid names the - # same grantee as the all-grantees policy does. - "sid": TofuExpr("FluidLfBucketList" + tofu_ref("statement.key")), - "effect": "Allow", - "principals": {"type": "AWS", "identifiers": [tofu_ref("statement.value")]}, - "actions": ["s3:ListBucket", "s3:GetBucketLocation"], - "resources": [bp.bucket_arn], - } - if bp.list_prefixes: - list_content["condition"] = { - "test": "StringLike", - "variable": "s3:prefix", - "values": list(bp.list_prefixes), + statements: List[Dict[str, Any]] = [] + for grants in bp.grants: + grantees = tofu_ref(_lf_other_account_grantees(bp, grants)) + list_content: Dict[str, Any] = { + # ``statement.key`` is the grantee's index, so a Sid names the + # same grantee as the all-grantees policy does. + "sid": TofuExpr("FluidLfBucketList" + tofu_ref("statement.key")), + "effect": "Allow", + "principals": {"type": "AWS", "identifiers": [tofu_ref("statement.value")]}, + "actions": ["s3:ListBucket", "s3:GetBucketLocation"], + "resources": [bp.bucket_arn], } - get_content: Dict[str, Any] = { - "sid": TofuExpr("FluidLfBucketGet" + tofu_ref("statement.key")), - "effect": "Allow", - "principals": {"type": "AWS", "identifiers": [tofu_ref("statement.value")]}, - "actions": ["s3:GetObject"], - "resources": [bp.object_arn], - } - data.setdefault("aws_iam_policy_document", {})[bp.policy_key] = { - "dynamic": { - "statement": [ - {"for_each": grantees, "content": list_content}, - {"for_each": grantees, "content": get_content}, - ] + if grants.list_prefixes: + list_content["condition"] = { + "test": "StringLike", + "variable": "s3:prefix", + "values": list(grants.list_prefixes), + } + get_content: Dict[str, Any] = { + "sid": TofuExpr("FluidLfBucketGet" + tofu_ref("statement.key")), + "effect": "Allow", + "principals": {"type": "AWS", "identifiers": [tofu_ref("statement.value")]}, + "actions": ["s3:GetObject"], + "resources": [grants.object_arn], } + statements.append({"for_each": grantees, "content": list_content}) + statements.append({"for_each": grantees, "content": get_content}) + data.setdefault("aws_iam_policy_document", {})[bp.policy_key] = { + "dynamic": {"statement": statements} } @@ -1900,7 +2032,7 @@ def _emit_lakeformation( _require_pool_prefix(placement, path, what="governance.lakeFormation.registerLocation") loc_key = safe_ident(f"{cid}_lf_loc_{bucket}_{path or 'root'}") resources.setdefault("aws_lakeformation_resource", {})[loc_key] = { - "arn": f"arn:aws:s3:::{bucket}/{path}" if path else f"arn:aws:s3:::{bucket}", + "arn": _lf_s3_arn(loc, bucket, f"/{path}" if path else ""), # ``use_service_linked_role: true`` is the default safe path # — LF uses the AWSServiceRoleForLakeFormationDataAccess SLR # to access objects under the registered location. @@ -1994,11 +2126,9 @@ def _emit_lakeformation( body_key = safe_ident(f"{cid}_lf_grant_{table or database}_{idx}") resources.setdefault("aws_lakeformation_permissions", {})[body_key] = body - # 2b. The S3 bucket policy for the grantees above — by default only those - # in another AWS account. See :func:`_lf_bucket_policy`. - bucket_policy = _lf_bucket_policy(binding, loc, fmt, cid, placement=placement) - if bucket_policy is not None: - _emit_lf_bucket_policy(resources, bucket_policy) + # 2b. The S3 bucket policy for the grantees above (by default only those in + # another AWS account) is emitted once per bucket, over every exposure + # on it, by :meth:`AwsIacPlugin.emit`. See :func:`_lf_bucket_policies`. # 3. LF-tag associations on the table (LF-TBAC). tag_assoc = gov.get("tags") or {} @@ -2390,6 +2520,7 @@ def _storage_by_bucket( """ packaging = resolve_packaging(contract) owned: Dict[str, _BucketStorage] = {} + policies: Optional[Dict[str, _LfBucketPolicy]] = None for index, exposure in enumerate(contract.get("exposes") or []): binding = exposure.get("binding") or {} if not is_cloud(binding, "aws"): @@ -2422,9 +2553,12 @@ def _storage_by_bucket( binding, loc, binding.get("format") or "parquet", cid, placement=placement ) if bucket_policy is not None: - # The same slot :func:`_emit_lf_bucket_policy` fills: one per bucket, - # the last exposure's wins. - state.bucket_policy = bucket_policy + # The bucket's one policy, over every exposure on it: the policy + # :func:`_emit_lf_bucket_policy` writes, whose readers the key lets + # decrypt. + if policies is None: + policies = _lf_bucket_policies(contract, cid) + state.bucket_policy = policies[bucket_policy.policy_key] for state in owned.values(): _check_retention_overlaps(state) _check_key_usable_by_lakeformation(state) diff --git a/tests/iac/_lf_bucket_policy.py b/tests/iac/_lf_bucket_policy.py index 9960e7ce..0f04e202 100644 --- a/tests/iac/_lf_bucket_policy.py +++ b/tests/iac/_lf_bucket_policy.py @@ -19,10 +19,11 @@ * ``all-grantees`` puts a static JSON document on the ``aws_s3_bucket_policy``. * ``cross-account`` (the default) points the resource at a ``data.aws_iam_policy_document`` whose ``dynamic "statement"`` blocks - iterate over the grantees in other accounts, decided at plan time. The - helper expands those blocks for EVERY grantee, i.e. the candidate - statements before the plan-time filter runs, which is the set the - prefix-scoping assertions must hold for. + iterate over the grantees in other accounts, decided at plan time: a + pair of blocks per exposure on the bucket, each over that exposure's + grantees. The helper expands those blocks for EVERY grantee, i.e. the + candidate statements before the plan-time filter runs, which is the set + the prefix-scoping assertions must hold for. Both come back in the static document's form and order, so a test can assert one property across both modes. @@ -36,6 +37,8 @@ _DOC_REF = re.compile(r"^\$\{data\.aws_iam_policy_document\.([A-Za-z0-9_]+)\.json\}$") _ARN_REF = re.compile(r"data\.aws_arn\.([A-Za-z0-9_]+)") +#: ``tostring(i)`` for an exposure's grantees numbered from 0, ``tostring(i + N)`` from N. +_FIRST = re.compile(r"tostring\(i(?: \+ (\d+))?\)") def policy_statements( @@ -47,27 +50,37 @@ def policy_statements( return list(json.loads(policy["policy"])["Statement"]) assert data is not None, "a cross-account policy needs emit_data() to be read" document = data["aws_iam_policy_document"][match.group(1)] - blocks = document["dynamic"]["statement"] - grantee_keys = _ARN_REF.findall(blocks[0]["for_each"]) - arns = [data["aws_arn"][key]["arn"] for key in grantee_keys] + # Consecutive blocks over the same grantees are one exposure's pair. + groups: List[List[Mapping[str, Any]]] = [] + for block in document["dynamic"]["statement"]: + if groups and groups[-1][0]["for_each"] == block["for_each"]: + groups[-1].append(block) + else: + groups.append([block]) statements: List[Dict[str, Any]] = [] - for index, arn in enumerate(arns): - for block in blocks: - content = block["content"] - resources = list(content["resources"]) - statement: Dict[str, Any] = { - "Sid": content["sid"].replace("${statement.key}", str(index)), - "Effect": content["effect"], - "Principal": {content["principals"]["type"]: arn}, - "Action": list(content["actions"]), - "Resource": resources[0] if len(resources) == 1 else resources, - } - condition = content.get("condition") - if condition: - statement["Condition"] = { - condition["test"]: {condition["variable"]: list(condition["values"])} + for group in groups: + for_each = group[0]["for_each"] + first = _FIRST.search(for_each) + assert first is not None, f"no grantee index in {for_each}" + offset = int(first.group(1) or 0) + arns = [data["aws_arn"][key]["arn"] for key in _ARN_REF.findall(for_each)] + for position, arn in enumerate(arns): + for block in group: + content = block["content"] + resources = list(content["resources"]) + statement: Dict[str, Any] = { + "Sid": content["sid"].replace("${statement.key}", str(offset + position)), + "Effect": content["effect"], + "Principal": {content["principals"]["type"]: arn}, + "Action": list(content["actions"]), + "Resource": resources[0] if len(resources) == 1 else resources, } - statements.append(statement) + condition = content.get("condition") + if condition: + statement["Condition"] = { + condition["test"]: {condition["variable"]: list(condition["values"])} + } + statements.append(statement) return statements diff --git a/tests/iac/test_iac_generate_offline.py b/tests/iac/test_iac_generate_offline.py new file mode 100644 index 00000000..d432c5a9 --- /dev/null +++ b/tests/iac/test_iac_generate_offline.py @@ -0,0 +1,212 @@ +# Copyright 2024-2026 Agentics Transformation Ltd +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +"""``fluid generate iac`` calls no cloud API. + +The native planner's AWS provider resolved a missing account with +``sts:GetCallerIdentity`` over whatever credentials the machine had, so +``fluid generate iac`` called AWS whenever ``AWS_ACCOUNT_ID`` was unset, and +wrote the account it found into the ARNs of the planned orchestration. The +examples' READMEs say the command needs no AWS account and no credentials. + +Every boto3 client and resource is made by botocore's +``Session.create_client`` (boto3's ``Session.client`` returns +``self._session.create_client(...)`` and ``Session.resource`` calls +``client``), so the guard replaces it and records each call. The provider +swallows any error from the lookup, so the record, not an exception, is what +fails a test. ``fluid apply`` still resolves the account; the last class pins +that. +""" + +from __future__ import annotations + +import argparse +import json +import logging +import os +from pathlib import Path +from typing import Any, Dict, List + +import pytest +import yaml + +from fluid_build.cli import generate_iac, main +from fluid_build.cli._common import CLIError + +pytestmark = [pytest.mark.unit, pytest.mark.provider] + +_EXAMPLES = Path(__file__).resolve().parents[2] / "examples" +AWS_EXAMPLES = sorted(_EXAMPLES.glob("aws-*/contract*.fluid.yaml")) +ACCOUNT = "123456789012" + + +@pytest.fixture +def no_aws(monkeypatch: pytest.MonkeyPatch) -> None: + """No account, no credential source: env vars, shared files, profile, metadata.""" + for key in ( + "AWS_ACCOUNT_ID", + "AWS_ACCESS_KEY_ID", + "AWS_SECRET_ACCESS_KEY", + "AWS_SESSION_TOKEN", + "AWS_PROFILE", + "AWS_DEFAULT_PROFILE", + "AWS_REGION", + "AWS_DEFAULT_REGION", + "FLUID_PROVIDER", + ): + monkeypatch.delenv(key, raising=False) + monkeypatch.setenv("AWS_SHARED_CREDENTIALS_FILE", os.devnull) + monkeypatch.setenv("AWS_CONFIG_FILE", os.devnull) + monkeypatch.setenv("AWS_EC2_METADATA_DISABLED", "true") + + +@pytest.fixture +def clients(monkeypatch: pytest.MonkeyPatch) -> List[str]: + """The service of every boto3 client constructed while the test runs.""" + session = pytest.importorskip("botocore.session") + made: List[str] = [] + + def refuse(self: Any, service_name: str, *args: Any, **kwargs: Any) -> Any: + made.append(service_name) + raise AssertionError(f"a boto3 {service_name} client was constructed") + + monkeypatch.setattr(session.Session, "create_client", refuse) + return made + + +def _write(tmp_path: Path, contract: Dict[str, Any]) -> Path: + path = tmp_path / "contract.fluid.yaml" + path.write_text(yaml.safe_dump(contract, sort_keys=False), encoding="utf-8") + return path + + +def _orders(**extra: Any) -> Dict[str, Any]: + contract: Dict[str, Any] = { + "fluidVersion": "0.7.6", + "kind": "DataProduct", + "id": "offline_orders", + "name": "Offline orders", + "metadata": {"layer": "Silver", "owner": {"team": "data", "email": "d@example.com"}}, + "exposes": [ + { + "exposeId": "orders", + "kind": "table", + "binding": { + "platform": "aws", + "format": "parquet", + "location": { + "bucket": "offline-lake", + "path": "orders/", + "database": "sales", + "table": "orders", + "region": "us-east-1", + }, + }, + "contract": {"schema": [{"name": "id", "type": "string"}]}, + } + ], + } + contract.update(extra) + return contract + + +def _lake_formation() -> Dict[str, Any]: + """Lake Formation grants to another account, on the ``{account}-fluid-data`` fallback.""" + contract = _orders() + binding = contract["exposes"][0]["binding"] + binding["location"]["bucket"] = "{{ env.FLUID_TEST_OFFLINE_BUCKET }}" + binding["governance"] = { + "lakeFormation": { + "registerLocation": True, + "grants": [ + { + "principal": "arn:aws:iam::222222222222:role/reader", + "permissions": ["SELECT", "DESCRIBE"], + } + ], + } + } + return contract + + +def _scheduled() -> Dict[str, Any]: + """An EventBridge schedule, whose planned Lambda and schedule ARNs name the account.""" + return _orders(orchestration={"schedule": "rate(1 hour)"}) + + +def _args(contract: Path, out: Path) -> argparse.Namespace: + return argparse.Namespace(contract=str(contract), provider="auto", out=str(out), env=None) + + +@pytest.mark.usefixtures("no_aws") +class TestGenerateIacCallsNoCloudApi: + def test_the_examples_exist(self): + assert len(AWS_EXAMPLES) >= 3, "the parametrisation below would prove nothing" + + @pytest.mark.parametrize("flags", [[], ["--shadow"]], ids=["plain", "shadow"]) + @pytest.mark.parametrize("contract", AWS_EXAMPLES, ids=lambda p: f"{p.parent.name}/{p.name}") + def test_an_example_constructs_no_boto3_client(self, contract, flags, tmp_path, clients): + rc = main(["generate", "iac", str(contract), "--out", str(tmp_path), *flags]) + assert rc == 0 + assert clients == [] + text = (tmp_path / "main.tf.json").read_text(encoding="utf-8") + assert generate_iac.UNSET_AWS_ACCOUNT not in text + + def test_lake_formation_on_the_fallback_bucket_constructs_no_client( + self, tmp_path, clients, monkeypatch + ): + monkeypatch.delenv("FLUID_TEST_OFFLINE_BUCKET", raising=False) + out = tmp_path / "out" + rc = main(["generate", "iac", str(_write(tmp_path, _lake_formation())), "--out", str(out)]) + assert rc == 0 + assert clients == [] + # The emitter's account references stay lookups ``tofu`` makes at plan time. + text = (out / "main.tf.json").read_text(encoding="utf-8") + assert "${data.aws_caller_identity.fluid_lf_caller.account_id}-fluid-data" in text + + def test_an_orchestration_that_names_the_account_is_refused(self, tmp_path, clients): + out = tmp_path / "out" + with pytest.raises(CLIError) as excinfo: + generate_iac.run(_args(_write(tmp_path, _scheduled()), out), logging.getLogger("t")) + assert excinfo.value.event == "generate_iac_aws_account_required" + assert "AWS_ACCOUNT_ID" in excinfo.value.context["error"] + assert excinfo.value.suggestions + assert not (out / "main.tf.json").exists(), "nothing is written with a placeholder" + assert clients == [] + + def test_with_aws_account_id_the_orchestration_is_emitted(self, tmp_path, clients, monkeypatch): + monkeypatch.setenv("AWS_ACCOUNT_ID", ACCOUNT) + out = tmp_path / "out" + rc = generate_iac.run(_args(_write(tmp_path, _scheduled()), out), logging.getLogger("t")) + assert rc == 0 + module = json.loads((out / "main.tf.json").read_text(encoding="utf-8")) + (function,) = module["resource"]["aws_lambda_function"].values() + assert function["role"] == f"arn:aws:iam::{ACCOUNT}:role/fluid-workflow-execution" + assert clients == [] + + +@pytest.mark.usefixtures("no_aws") +class TestApplyStillResolvesTheAccount: + """``fluid apply`` builds its native plan with ``native_actions(contract, logger)`` + (``cli/_apply_opentofu_engine.py``); only ``generate iac`` passes ``offline``.""" + + def test_the_apply_call_asks_sts_for_a_missing_account(self, clients): + generate_iac.native_actions(_scheduled(), logging.getLogger("t")) + assert clients == ["sts"] + + def test_the_offline_call_plans_without_asking(self, clients): + actions = generate_iac.native_actions(_scheduled(), logging.getLogger("t"), offline=True) + assert clients == [] + # The planner still ran, sovereignty check included, on the placeholder. + assert any(a.get("op") == "lambda.ensure_function" for a in actions) diff --git a/tests/iac/test_iac_lakeformation_bucket_policy.py b/tests/iac/test_iac_lakeformation_bucket_policy.py index 22bc8665..b661a886 100644 --- a/tests/iac/test_iac_lakeformation_bucket_policy.py +++ b/tests/iac/test_iac_lakeformation_bucket_policy.py @@ -286,6 +286,222 @@ def test_a_pool_without_a_path_is_fine_when_nothing_bucket_level_is_emitted(self assert "aws_lakeformation_permissions" in resources +# --------------------------------------------------------------------------- +# Two exposes on one bucket: one policy +# --------------------------------------------------------------------------- + +POLICY_KEY = "lf_bucket_policy_lf_bucket_policy_acme_lake" + + +def _two_zones( + raw: List[str], curated: List[str], *, bucket_policy: Optional[str] = None +) -> Dict[str, Any]: + """Two exposes on the bucket ``acme-lake``, as in examples/aws-medallion-lake.""" + contract = _contract(raw, bucket_policy=bucket_policy, path="raw/") + zone = _contract(curated, bucket_policy=bucket_policy, path="curated/")["exposes"][0] + zone["exposeId"] = "curated" + zone["binding"]["location"]["database"] = "silver" + contract["exposes"].append(zone) + return contract + + +def _other_account(grantees: List[str], first: int = 0) -> str: + """The HCL map one expose's grantees in other accounts are filtered into.""" + index = f"i + {first}" if first else "i" + arns = ", ".join(f"data.aws_arn.{POLICY_KEY}_grantee_{first + n}" for n in range(len(grantees))) + return f"{{for i, g in [{arns}] : tostring({index}) => g.arn if g.account != {CALLER}}}" + + +def _gets(statements: List[Dict[str, Any]]) -> Dict[str, Any]: + return {s["Sid"]: s["Resource"] for s in statements if s["Action"] == ["s3:GetObject"]} + + +class TestTwoExposesOnOneBucketShareOnePolicy: + """S3 keeps one policy per bucket. Each expose used to write its own under + the bucket's key, so the second replaced the first: a grantee in another + account lost ``s3:GetObject`` on the first expose's prefix.""" + + def test_one_policy_carries_both_prefixes(self): + resources, data = _emit(_two_zones([OTHER], [OTHER])) + assert _gets(policy_statements(only_policy(resources), data)) == { + "FluidLfBucketGet0": "arn:aws:s3:::acme-lake/raw/*", + "FluidLfBucketGet1": "arn:aws:s3:::acme-lake/curated/*", + } + + def test_the_grantees_are_numbered_across_the_bucket(self): + resources, data = _emit(_two_zones([SAME, OTHER], [OTHER])) + assert list(data["aws_arn"]) == [f"{POLICY_KEY}_grantee_{i}" for i in range(3)] + assert [body["arn"] for body in data["aws_arn"].values()] == [SAME, OTHER, OTHER] + raw, curated = _other_account([SAME, OTHER]), _other_account([OTHER], first=2) + # One instance while any expose keeps a grantee; the keys never collide. + assert only_policy(resources)["count"] == ( + "${length(merge(" + raw + ", " + curated + ")) > 0 ? 1 : 0}" + ) + blocks = data["aws_iam_policy_document"][POLICY_KEY]["dynamic"]["statement"] + assert [block["for_each"] for block in blocks] == ["${" + raw + "}"] * 2 + [ + "${" + curated + "}" + ] * 2 + sids = [s["Sid"] for s in policy_statements(only_policy(resources), data)] + assert len(sids) == len(set(sids)) == 6 + + def test_all_grantees_writes_every_exposes_statements(self): + resources, _ = _emit(_two_zones([SAME], [OTHER], bucket_policy="all-grantees")) + statements = json.loads(only_policy(resources)["policy"])["Statement"] + assert [(s["Sid"], s["Principal"]["AWS"], s["Resource"]) for s in statements] == [ + ("FluidLfBucketList0", SAME, "arn:aws:s3:::acme-lake"), + ("FluidLfBucketGet0", SAME, "arn:aws:s3:::acme-lake/raw/*"), + ("FluidLfBucketList1", OTHER, "arn:aws:s3:::acme-lake"), + ("FluidLfBucketGet1", OTHER, "arn:aws:s3:::acme-lake/curated/*"), + ] + + def test_the_candidates_are_exactly_the_all_grantees_statements(self): + resources, data = _emit(_two_zones([SAME, OTHER], [OTHER])) + legacy, legacy_data = _emit( + _two_zones([SAME, OTHER], [OTHER], bucket_policy="all-grantees") + ) + assert policy_statements(only_policy(resources), data) == policy_statements( + only_policy(legacy), legacy_data + ) + + def test_a_none_expose_leaves_the_other_exposes_statements(self): + contract = _two_zones([OTHER], [OTHER]) + contract["exposes"][1]["binding"]["governance"]["lakeFormation"]["bucketPolicy"] = "none" + resources, data = _emit(contract) + assert _gets(policy_statements(only_policy(resources), data)) == { + "FluidLfBucketGet0": "arn:aws:s3:::acme-lake/raw/*" + } + + @pytest.mark.parametrize("method", ["emit", "emit_data"]) + def test_two_modes_on_one_bucket_are_refused(self, method): + contract = _two_zones([OTHER], [OTHER]) + gov = contract["exposes"][1]["binding"]["governance"]["lakeFormation"] + gov["bucketPolicy"] = "all-grantees" + with pytest.raises(UnsupportedBindingError) as excinfo: + getattr(AwsIacPlugin(), method)(contract) + assert excinfo.value.kind == "lakeformation-bucket-policy" + assert "one bucket policy" in str(excinfo.value) + assert excinfo.value.remediation + + def test_two_buckets_keep_a_policy_each(self): + contract = _two_zones([OTHER], [OTHER]) + contract["exposes"][1]["binding"]["location"]["bucket"] = "acme-curated" + resources, data = _emit(contract) + policies = resources["aws_s3_bucket_policy"] + assert sorted(policies) == [ + "lf_bucket_policy_lf_bucket_policy_acme_curated", + POLICY_KEY, + ] + assert _gets(policy_statements(policies[POLICY_KEY], data)) == { + "FluidLfBucketGet0": "arn:aws:s3:::acme-lake/raw/*" + } + + def test_every_reference_resolves_and_every_declaration_is_read(self): + contract = _two_zones([SAME, OTHER], [OTHER]) + _, data = _emit(contract) + rendered = build_module(AwsIacPlugin(), contract) + referenced = set(re.findall(r"data\.(aws_[a-z_]+)\.([A-Za-z0-9_]+)", rendered)) + declared = {(dtype, name) for dtype, block in data.items() for name in block} + assert referenced == declared + + def test_the_product_key_lets_every_exposes_readers_decrypt(self): + contract = _two_zones([OTHER], [OTHER]) + for exposure in contract["exposes"]: + exposure["binding"]["encryption"] = {"kms": "product"} + _, data = _emit(contract) + key_policy = data["aws_iam_policy_document"]["lf_bucket_policy_acme_lake_kms_policy"] + (reader,) = key_policy["dynamic"]["statement"] + both = "merge(" + _other_account([OTHER]) + ", " + _other_account([OTHER], 1) + ")" + assert reader["for_each"] == "${" + both + "}" + + +# --------------------------------------------------------------------------- +# The {account}-fluid-data fallback bucket +# --------------------------------------------------------------------------- + +UNSET = "{{ env.FLUID_TEST_LF_UNSET_BUCKET }}" +ACCOUNT_BUCKET = "${" + CALLER + "}-fluid-data" + + +@pytest.fixture +def unset_bucket(monkeypatch): + monkeypatch.delenv("FLUID_TEST_LF_UNSET_BUCKET", raising=False) + + +def _fallback(grantees: List[str], **kwargs: Any) -> Dict[str, Any]: + contract = _contract(grantees, **kwargs) + contract["exposes"][0]["binding"]["location"]["bucket"] = UNSET + return contract + + +def _rendered(contract: Dict[str, Any]) -> Dict[str, Any]: + return json.loads(build_module(AwsIacPlugin(), contract)) + + +@pytest.mark.usefixtures("unset_bucket") +class TestTheFallbackBucketInterpolatesTheAccount: + """An unresolved ``location.bucket`` falls back to ``{account}-fluid-data``. + The Lake Formation location and bucket-policy ARNs were plain strings, so + the renderer escaped the account token to ``$${...}``, a literal, and the + policy referenced an ``aws_s3_bucket`` the module never declared.""" + + def test_the_registered_location_names_the_applying_account(self): + (resource,) = _rendered(_fallback([OTHER]))["resource"][ + "aws_lakeformation_resource" + ].values() + assert resource["arn"] == f"arn:aws:s3:::{ACCOUNT_BUCKET}/orders/" + + def test_the_bucket_policy_names_the_applying_account(self): + rendered = _rendered(_fallback([OTHER])) + (policy,) = rendered["resource"]["aws_s3_bucket_policy"].values() + assert policy["bucket"] == ACCOUNT_BUCKET + (document,) = rendered["data"]["aws_iam_policy_document"].values() + assert [block["content"]["resources"] for block in document["dynamic"]["statement"]] == [ + [f"arn:aws:s3:::{ACCOUNT_BUCKET}"], + [f"arn:aws:s3:::{ACCOUNT_BUCKET}/orders/*"], + ] + assert "$${" + CALLER not in json.dumps(rendered) + + def test_all_grantees_names_the_applying_account(self): + rendered = _rendered(_fallback([SAME, OTHER], bucket_policy="all-grantees")) + (policy,) = rendered["resource"]["aws_s3_bucket_policy"].values() + assert policy["bucket"] == ACCOUNT_BUCKET + statements = json.loads(policy["policy"])["Statement"] + assert {s["Resource"] for s in statements} == { + f"arn:aws:s3:::{ACCOUNT_BUCKET}", + f"arn:aws:s3:::{ACCOUNT_BUCKET}/orders/*", + } + assert "$${" + CALLER not in json.dumps(rendered) + + @pytest.mark.parametrize("bucket_policy", [None, "all-grantees"]) + def test_every_resource_reference_is_declared(self, bucket_policy): + rendered = _rendered(_fallback([OTHER], bucket_policy=bucket_policy)) + text = json.dumps(rendered).replace("$${", "") + referenced = set(re.findall(r"\$\{(aws_[a-z0-9_]+)\.([A-Za-z0-9_]+)\.", text)) + declared = { + (rtype, name) for rtype, block in rendered["resource"].items() for name in block + } + assert referenced <= declared, f"dangling: {sorted(referenced - declared)}" + + @pytest.mark.parametrize("bucket_policy", [None, "all-grantees"]) + def test_contract_text_stays_inert_beside_the_account(self, bucket_policy): + # SECURITY: only the emitter's own account token is interpolated; the + # contract's path and principal stay escaped inside the same values. + evil = 'arn:aws:iam::222222222222:role/${file("/etc/passwd")}' + contract = _fallback([evil], bucket_policy=bucket_policy, path='o/${file("/etc/hosts")}/') + rendered = build_module(AwsIacPlugin(), contract) + assert "${file(" not in rendered.replace("$${", "") + (resource,) = json.loads(rendered)["resource"]["aws_lakeformation_resource"].values() + assert resource["arn"] == f'arn:aws:s3:::{ACCOUNT_BUCKET}/o/$${{file("/etc/hosts")}}/' + + def test_a_bucket_that_spells_the_token_is_not_interpolated(self): + # The fallback is decided on the raw contract value, so a contract that + # writes the token itself gets a literal. + contract = _contract([OTHER]) + contract["exposes"][0]["binding"]["location"]["bucket"] = ACCOUNT_BUCKET + (resource,) = _rendered(contract)["resource"]["aws_lakeformation_resource"].values() + assert resource["arn"].startswith("arn:aws:s3:::$${") + + # --------------------------------------------------------------------------- # Validation # --------------------------------------------------------------------------- diff --git a/tests/iac/test_iac_lakeformation_bucket_policy_plan.py b/tests/iac/test_iac_lakeformation_bucket_policy_plan.py index 27cbcfe3..3656439a 100644 --- a/tests/iac/test_iac_lakeformation_bucket_policy_plan.py +++ b/tests/iac/test_iac_lakeformation_bucket_policy_plan.py @@ -19,6 +19,11 @@ ``tofu`` runs with. So the emit carries the filter as HCL, and the only proof that it keeps and drops the right statements is a plan that evaluates it. +The same plans prove the two other halves of the bucket policy that only a +plan evaluates: two exposes on one bucket share one policy, with a statement +for each prefix, and on the ``{account}-fluid-data`` fallback bucket the Lake +Formation location and the policy name the applying account. + A moto ``ThreadedMotoServer`` answers ``sts:GetCallerIdentity`` for ``data.aws_caller_identity`` (moto's account is ``123456789012``); the ``aws_arn`` and ``aws_iam_policy_document`` data sources are computed by the @@ -141,8 +146,10 @@ def _contract(grantees: List[str], bucket_policy: Optional[str] = None) -> Dict[ } -def _plan(contract: Dict[str, Any], workdir: Path, endpoint: str, env: Dict[str, str]): - """``tofu init`` + ``plan``; the planned ``aws_s3_bucket_policy`` instances by address.""" +def _planned_changes( + contract: Dict[str, Any], workdir: Path, endpoint: str, env: Dict[str, str] +) -> List[Dict[str, Any]]: + """``tofu init`` + ``plan``; every planned resource change.""" workdir.mkdir(parents=True, exist_ok=True) (workdir / "main.tf.json").write_text(build_module(get_iac_plugin("aws"), contract)) (workdir / "provider.tf.json").write_text(json.dumps(_provider_override(endpoint))) @@ -174,10 +181,14 @@ def run(*args: str) -> "subprocess.CompletedProcess[str]": timeout=600, check=True, ) - changes = json.loads(shown.stdout)["resource_changes"] + return list(json.loads(shown.stdout)["resource_changes"]) + + +def _plan(contract: Dict[str, Any], workdir: Path, endpoint: str, env: Dict[str, str]): + """The planned ``aws_s3_bucket_policy`` instances by address, policies parsed.""" return { change["address"]: json.loads(change["change"]["after"]["policy"]) - for change in changes + for change in _planned_changes(contract, workdir, endpoint, env) if change["type"] == "aws_s3_bucket_policy" and change["change"]["actions"] == ["create"] } @@ -219,3 +230,95 @@ def test_all_grantees_restores_the_same_account_statement( def test_none_plans_no_bucket_policy(self, tmp_path, moto_endpoint, tofu_env): planned = _plan(_contract([SAME, OTHER], "none"), tmp_path, moto_endpoint, tofu_env) assert planned == {} + + +def _two_zones( + raw: List[str], curated: List[str], bucket_policy: Optional[str] = None +) -> Dict[str, Any]: + """Two exposes on the bucket ``lf-plan-lake``, as in examples/aws-medallion-lake.""" + contract = _contract(raw, bucket_policy) + contract["exposes"][0]["binding"]["location"]["path"] = "raw/" + zone = _contract(curated, bucket_policy)["exposes"][0] + zone["exposeId"] = "curated" + zone["binding"]["location"].update({"path": "curated/", "database": "silver"}) + contract["exposes"].append(zone) + return contract + + +def _statements(policy: Dict[str, Any]) -> Dict[str, Any]: + return {s["Sid"]: (s["Principal"]["AWS"], s["Resource"]) for s in policy["Statement"]} + + +@pytest.mark.skipif(_SKIP, reason=_SKIP_REASON) +class TestOneBucketTwoExposesAtPlanTime: + """Each expose used to emit the bucket's policy under the same key, so the + second replaced the first and ``raw/`` lost its cross-account statements.""" + + def test_both_prefixes_keep_their_cross_account_statements( + self, tmp_path, moto_endpoint, tofu_env + ): + planned = _plan(_two_zones([SAME, OTHER], [OTHER]), tmp_path, moto_endpoint, tofu_env) + assert list(planned) == [f"{POLICY}[0]"], "one policy for the one bucket" + assert _statements(planned[f"{POLICY}[0]"]) == { + "FluidLfBucketList1": (OTHER, "arn:aws:s3:::lf-plan-lake"), + "FluidLfBucketGet1": (OTHER, "arn:aws:s3:::lf-plan-lake/raw/*"), + "FluidLfBucketList2": (OTHER, "arn:aws:s3:::lf-plan-lake"), + "FluidLfBucketGet2": (OTHER, "arn:aws:s3:::lf-plan-lake/curated/*"), + } + + def test_same_account_grantees_on_both_exposes_plan_no_policy( + self, tmp_path, moto_endpoint, tofu_env + ): + planned = _plan(_two_zones([SAME], [SAME]), tmp_path, moto_endpoint, tofu_env) + assert planned == {} + + def test_all_grantees_writes_both_prefixes(self, tmp_path, moto_endpoint, tofu_env): + contract = _two_zones([SAME], [OTHER], "all-grantees") + planned = _plan(contract, tmp_path, moto_endpoint, tofu_env) + assert list(planned) == [POLICY] + assert _statements(planned[POLICY]) == { + "FluidLfBucketList0": (SAME, "arn:aws:s3:::lf-plan-lake"), + "FluidLfBucketGet0": (SAME, "arn:aws:s3:::lf-plan-lake/raw/*"), + "FluidLfBucketList1": (OTHER, "arn:aws:s3:::lf-plan-lake"), + "FluidLfBucketGet1": (OTHER, "arn:aws:s3:::lf-plan-lake/curated/*"), + } + + +@pytest.mark.skipif(_SKIP, reason=_SKIP_REASON) +class TestTheFallbackBucketAtPlanTime: + """An unresolved ``location.bucket`` falls back to ``{account}-fluid-data``. + The Lake Formation location was planned as the literal + ``arn:aws:s3:::${data.aws_caller_identity...}-fluid-data/...``, and a + bucket policy failed the plan with "Reference to undeclared resource".""" + + ACCOUNT_BUCKET = f"{APPLYING_ACCOUNT}-fluid-data" + + @pytest.fixture(autouse=True) + def _unset(self, monkeypatch): + monkeypatch.delenv("FLUID_TEST_LF_PLAN_BUCKET", raising=False) + + def _fallback(self, bucket_policy: Optional[str] = None) -> Dict[str, Any]: + contract = _contract([OTHER], bucket_policy) + contract["exposes"][0]["binding"]["location"][ + "bucket" + ] = "{{ env.FLUID_TEST_LF_PLAN_BUCKET }}" + return contract + + @pytest.mark.parametrize("bucket_policy", [None, "all-grantees"]) + def test_the_location_and_the_policy_name_the_applying_account( + self, bucket_policy, tmp_path, moto_endpoint, tofu_env + ): + changes = _planned_changes(self._fallback(bucket_policy), tmp_path, moto_endpoint, tofu_env) + by_type: Dict[str, List[Dict[str, Any]]] = {} + for change in changes: + by_type.setdefault(change["type"], []).append(change["change"]["after"]) + assert [r["arn"] for r in by_type["aws_lakeformation_resource"]] == [ + f"arn:aws:s3:::{self.ACCOUNT_BUCKET}/orders/" + ] + (policy,) = by_type["aws_s3_bucket_policy"] + assert policy["bucket"] == self.ACCOUNT_BUCKET + resources = {s["Resource"] for s in json.loads(policy["policy"])["Statement"]} + assert resources == { + f"arn:aws:s3:::{self.ACCOUNT_BUCKET}", + f"arn:aws:s3:::{self.ACCOUNT_BUCKET}/orders/*", + } diff --git a/tests/iac/test_iac_packaging_default_pin.py b/tests/iac/test_iac_packaging_default_pin.py index 933155e6..41382fc9 100644 --- a/tests/iac/test_iac_packaging_default_pin.py +++ b/tests/iac/test_iac_packaging_default_pin.py @@ -179,7 +179,7 @@ def test_wired_cli_emit_is_byte_identical( # Pin determinism: the native planner is best-effort (credentials- # dependent) — force the no-actions path on BOTH sides so the pin # never depends on the developer's cloud environment. - monkeypatch.setattr(generate_iac, "native_actions", lambda contract, logger: []) + monkeypatch.setattr(generate_iac, "native_actions", lambda contract, logger, **_: []) try: _, expected = _expected_module_bytes(path) except CLIError as exc: @@ -199,7 +199,7 @@ def test_resolver_called_once_and_returns_legacy( ): example = REPO_ROOT / "examples" / "aws-s3-glue-athena" / "contract.fluid.yaml" assert example.is_file() - monkeypatch.setattr(generate_iac, "native_actions", lambda contract, logger: []) + monkeypatch.setattr(generate_iac, "native_actions", lambda contract, logger, **_: []) seen = [] def spy(contract): @@ -236,7 +236,7 @@ def test_invalid_packaging_fails_generate_with_typed_error( } path = tmp_path / "contract.fluid.yaml" path.write_text(yaml.safe_dump(contract), encoding="utf-8") - monkeypatch.setattr(generate_iac, "native_actions", lambda contract, logger: []) + monkeypatch.setattr(generate_iac, "native_actions", lambda contract, logger, **_: []) with pytest.raises(CLIError) as excinfo: _run_generate_iac(path, tmp_path / "out") assert "pool" in str((excinfo.value.context or {}).get("error", "")) From 8f9a58b9a3a93a24d740cf0fd33af49c3fbb7ba0 Mon Sep 17 00:00:00 2001 From: Speculator55005 Date: Mon, 28 Sep 2026 22:50:57 +0200 Subject: [PATCH 2/2] fix(errors): the generate-iac account hint is one explicit item --- fluid_build/_error_catalog.py | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/fluid_build/_error_catalog.py b/fluid_build/_error_catalog.py index 254c1f59..ebf1e23e 100644 --- a/fluid_build/_error_catalog.py +++ b/fluid_build/_error_catalog.py @@ -309,8 +309,10 @@ def slug_for(event: str) -> str: ), "generate_iac_aws_account_required": ( [ - "Set AWS_ACCOUNT_ID to the account the module will be applied in, then re-run " - "'fluid generate iac' (it never looks the account up in AWS)", + ( + "Set AWS_ACCOUNT_ID to the account the module will be applied in, then re-run " + "'fluid generate iac' (it never looks the account up in AWS)" + ), "'fluid apply' resolves the account itself and needs no AWS_ACCOUNT_ID", ], None,