diff --git a/fluid_build/_error_catalog.py b/fluid_build/_error_catalog.py index a3ce0340..ebf1e23e 100644 --- a/fluid_build/_error_catalog.py +++ b/fluid_build/_error_catalog.py @@ -307,6 +307,16 @@ 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 a958f9ef..7a26cf1e 100644 --- a/fluid_build/iac/providers/aws.py +++ b/fluid_build/iac/providers/aws.py @@ -410,6 +410,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) @@ -796,8 +801,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"} @@ -1597,6 +1601,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 @@ -1609,9 +1641,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 @@ -1619,10 +1663,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.""" @@ -1679,7 +1726,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. @@ -1725,47 +1773,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 @@ -1775,32 +1907,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 @@ -1815,61 +1958,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} } @@ -1987,7 +2119,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. @@ -2081,11 +2213,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 {} @@ -2480,6 +2610,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"): @@ -2512,9 +2643,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 382d991c..ac43741d 100644 --- a/tests/iac/test_iac_packaging_default_pin.py +++ b/tests/iac/test_iac_packaging_default_pin.py @@ -180,7 +180,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: @@ -210,7 +210,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): @@ -247,7 +247,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", ""))