Coverage for scripts / live_release_validation / ownership / ecr.py: 100.00%
283 statements
« prev ^ index » next coverage.py v7.13.5, created at 2026-09-14 22:07 +0000
« prev ^ index » next coverage.py v7.13.5, created at 2026-09-14 22:07 +0000
1"""ECR expectation, ownership checkpointing, and baseline stripping."""
3from __future__ import annotations
5import copy
6import json
7import re
8from collections.abc import Mapping
9from typing import Any
11from ..constants import (
12 _RUN_STACK_TAG,
13)
14from ..context import (
15 _topology_regions,
16)
17from ..inventory import (
18 collect_ecr_inventory,
19)
20from ..models import RunContext
21from ..protected import (
22 _EC2_TAGGED_RESOURCE_IDENTITIES,
23 _PROTECTED_GLOBAL_RESOURCE_CATEGORIES,
24 _PROTECTED_REGIONAL_RESOURCE_CATEGORIES,
25 _baseline_protected_identities,
26 _ec2_tagged_resource_identity,
27 _eks_pod_parent_cluster,
28 _matches_protected_physical_identity,
29 _tagged_resource_is_protected,
30)
33def _strip_baseline_ecr(
34 project_inventory: dict[str, Any], baseline: dict[str, Any]
35) -> dict[str, Any]:
36 """Strip only exact protected identities and authoritatively absent tag records."""
37 protected_stack_ids, protected_resource_ids = _baseline_protected_identities(baseline)
38 baseline_ecr_names: dict[str, set[str]] = {}
39 baseline_ecr_arns: dict[str, set[str]] = {}
40 for raw_region, repositories in (baseline.get("ecr_repositories") or {}).items():
41 region = str(raw_region)
42 for repository in repositories:
43 name = str(repository.get("name") or "")
44 arn = str(repository.get("arn") or "")
45 if name:
46 baseline_ecr_names.setdefault(region, set()).add(name)
47 if arn:
48 baseline_ecr_arns.setdefault(region, set()).add(arn)
50 inventory = copy.deepcopy(project_inventory)
51 stacks_by_region = inventory.get("cloudformation_stacks")
52 if isinstance(stacks_by_region, dict):
53 for region, stacks in list(stacks_by_region.items()):
54 exact_stack_ids = protected_stack_ids.get(str(region), set())
55 remaining = [
56 stack for stack in stacks if str(stack.get("stack_id") or "") not in exact_stack_ids
57 ]
58 if remaining:
59 stacks_by_region[region] = remaining
60 else:
61 stacks_by_region.pop(region)
63 authoritative_clusters = inventory.get("authoritative_eks_clusters")
64 authoritative_ec2_resources = inventory.get("authoritative_ec2_resources")
65 authority_scope = inventory.get("authority_scope")
66 expected_partition = (
67 str(authority_scope.get("partition") or "") if isinstance(authority_scope, Mapping) else ""
68 )
69 expected_account = (
70 str(authority_scope.get("account") or "") if isinstance(authority_scope, Mapping) else ""
71 )
72 has_exact_authority_scope = bool(
73 expected_partition and re.fullmatch(r"\d{12}", expected_account)
74 )
75 coverage = inventory.get("coverage")
76 coverage_complete = isinstance(coverage, Mapping) and coverage.get("complete") is True
77 completed_scanners = (
78 {str(scanner) for scanner in coverage.get("completed_scanners", [])}
79 if isinstance(coverage, Mapping) and isinstance(coverage.get("completed_scanners"), list)
80 else set()
81 )
82 scanner_regions = coverage.get("scanner_regions") if isinstance(coverage, Mapping) else None
83 eks_scanner_regions = (
84 {str(region) for region in scanner_regions.get("eks_clusters", [])}
85 if isinstance(scanner_regions, Mapping)
86 and isinstance(scanner_regions.get("eks_clusters"), list)
87 else set()
88 )
89 instance_scanner_regions = (
90 {str(region) for region in scanner_regions.get("ec2_instances", [])}
91 if isinstance(scanner_regions, Mapping)
92 and isinstance(scanner_regions.get("ec2_instances"), list)
93 else set()
94 )
95 network_scanner_regions = (
96 {str(region) for region in scanner_regions.get("ec2_networking", [])}
97 if isinstance(scanner_regions, Mapping)
98 and isinstance(scanner_regions.get("ec2_networking"), list)
99 else set()
100 )
101 has_complete_eks_authority = coverage_complete and "eks_clusters" in completed_scanners
102 has_complete_ec2_authority = coverage_complete and {
103 "ec2_instances",
104 "ec2_networking",
105 }.issubset(completed_scanners)
106 ec2_scanner_regions = instance_scanner_regions & network_scanner_regions
107 for region, resources in list(inventory.get("regional", {}).items()):
108 region_key = str(region)
109 region_stack_ids = protected_stack_ids.get(region_key, set())
110 region_resource_ids = protected_resource_ids.get(region_key, {})
111 exact_tagged_arns = {
112 physical_id
113 for physical_ids in region_resource_ids.values()
114 for physical_id in physical_ids
115 if physical_id.startswith("arn:")
116 }
117 exact_tagged_arns.update(region_stack_ids)
118 exact_tagged_arns.update(baseline_ecr_arns.get(region_key, set()))
119 if "tagged_resources" in resources:
120 resources["tagged_resources"] = [
121 record
122 for record in resources.get("tagged_resources", [])
123 if not _tagged_resource_is_protected(
124 record,
125 protected_stack_ids=region_stack_ids,
126 protected_resource_ids=region_resource_ids,
127 exact_arns=exact_tagged_arns,
128 expected_partition=expected_partition,
129 expected_region=region_key,
130 expected_account=expected_account,
131 )
132 ]
134 protected_backup_plan_ids = region_resource_ids.get(
135 "AWS::Backup::BackupPlan",
136 set(),
137 )
138 for resource_type, category in _PROTECTED_REGIONAL_RESOURCE_CATEGORIES.items():
139 if category not in resources:
140 continue
141 physical_ids = region_resource_ids.get(resource_type, set())
142 if physical_ids:
143 resources[category] = [
144 candidate
145 for candidate in resources.get(category, [])
146 if not any(
147 _matches_protected_physical_identity(
148 resource_type,
149 category,
150 candidate,
151 physical_id,
152 protected_backup_plan_ids=protected_backup_plan_ids,
153 )
154 for physical_id in physical_ids
155 )
156 ]
158 if "ecr_repositories" in resources:
159 resources["ecr_repositories"] = [
160 name
161 for name in resources.get("ecr_repositories", [])
162 if name not in baseline_ecr_names.get(str(region), set())
163 ]
165 region_ec2_authority = (
166 authoritative_ec2_resources.get(region_key)
167 if isinstance(authoritative_ec2_resources, dict)
168 else None
169 )
170 if (
171 "tagged_resources" in resources
172 and has_exact_authority_scope
173 and has_complete_ec2_authority
174 and region_key in ec2_scanner_regions
175 and isinstance(region_ec2_authority, dict)
176 ):
177 authoritative_ec2_ids = {
178 category: {str(candidate) for candidate in candidates}
179 for category, candidates in region_ec2_authority.items()
180 if category in {item[0] for item in _EC2_TAGGED_RESOURCE_IDENTITIES.values()}
181 and isinstance(candidates, list)
182 }
183 resources["tagged_resources"] = [
184 record
185 for record in resources["tagged_resources"]
186 if (
187 (
188 identity := _ec2_tagged_resource_identity(
189 str(record.get("arn") or ""),
190 region_key,
191 expected_partition,
192 expected_account,
193 )
194 )
195 is None
196 or identity[0] not in authoritative_ec2_ids
197 or identity[1] in authoritative_ec2_ids[identity[0]]
198 )
199 ]
201 if (
202 "tagged_resources" in resources
203 and has_exact_authority_scope
204 and has_complete_eks_authority
205 and region_key in eks_scanner_regions
206 and isinstance(authoritative_clusters, dict)
207 and region in authoritative_clusters
208 and isinstance(authoritative_clusters[region], list)
209 ):
210 existing_clusters = {str(name) for name in authoritative_clusters[region]}
211 resources["tagged_resources"] = [
212 record
213 for record in resources["tagged_resources"]
214 if (
215 (
216 parent := _eks_pod_parent_cluster(
217 str(record.get("arn") or ""),
218 region_key,
219 expected_partition,
220 expected_account,
221 )
222 )
223 is None
224 or parent in existing_clusters
225 )
226 ]
228 if not any(resources.values()):
229 inventory["regional"].pop(region)
231 for resource_type, category in _PROTECTED_GLOBAL_RESOURCE_CATEGORIES.items():
232 if category not in inventory:
233 continue
234 physical_ids = {
235 physical_id
236 for resources_by_type in protected_resource_ids.values()
237 for physical_id in resources_by_type.get(resource_type, set())
238 }
239 if physical_ids:
240 inventory[category] = [
241 candidate
242 for candidate in inventory.get(category, [])
243 if not any(
244 _matches_protected_physical_identity(
245 resource_type,
246 category,
247 candidate,
248 physical_id,
249 )
250 for physical_id in physical_ids
251 )
252 ]
253 return inventory
256def _strip_expected_retained_ecr(
257 ctx: RunContext,
258 final_baseline: dict[str, Any],
259) -> tuple[dict[str, Any], dict[str, list[dict[str, Any]]]]:
260 """Remove only exact checkpointed ECR residuals from comparison inventory."""
261 sanitized = copy.deepcopy(final_baseline)
262 repositories_by_region = sanitized.setdefault("ecr_repositories", {})
263 baseline_repositories = {
264 (str(region), str(repository["name"])): repository
265 for region, repositories in (ctx.checkpoint.baseline or {})
266 .get("ecr_repositories", {})
267 .items()
268 for repository in repositories
269 }
270 accepted: dict[str, list[dict[str, Any]]] = {
271 "repositories": [],
272 "image_deltas": [],
273 }
274 created_keys: set[tuple[str, str]] = set()
276 for record in ctx.checkpoint.state.get("created_ecr_repositories", []):
277 region = str(record["region"])
278 name = str(record["name"])
279 repository_key = (region, name)
280 if repository_key in created_keys or repository_key in baseline_repositories:
281 raise RuntimeError(f"Invalid retained ECR repository authority for {region}:{name}")
282 created_keys.add(repository_key)
283 repositories = repositories_by_region.get(region, [])
284 matches = [repository for repository in repositories if repository.get("name") == name]
285 if len(matches) > 1:
286 raise RuntimeError(f"Final ECR inventory duplicated {region}:{name}")
287 if not matches:
288 accepted["repositories"].append(
289 {"region": region, "name": name, "arn": record["arn"], "already_absent": True}
290 )
291 continue
292 repository = matches[0]
293 if _ecr_creation_identity(repository) != record.get("creation_identity"):
294 raise RuntimeError(f"Retained ECR repository identity changed for {region}:{name}")
295 if (repository.get("tags") or {}).get(_RUN_STACK_TAG) != record.get("run_tag"):
296 raise RuntimeError(f"Retained ECR repository run tag changed for {record['arn']}")
297 repositories.remove(repository)
298 accepted["repositories"].append(
299 {
300 "region": region,
301 "name": name,
302 "arn": record["arn"],
303 "retained": True,
304 "inventory": repository,
305 }
306 )
308 observed_delta_keys: set[tuple[str, str, str]] = set()
309 for record in ctx.checkpoint.state.get("retained_ecr_image_deltas", []):
310 region = str(record["region"])
311 name = str(record["repository"])
312 tag = str(record["tag"])
313 image_key = (region, name, tag)
314 if image_key in observed_delta_keys or (region, name) in created_keys:
315 raise RuntimeError(f"Invalid retained ECR image-delta authority for {image_key}")
316 observed_delta_keys.add(image_key)
317 baseline_repository = baseline_repositories.get((region, name))
318 if baseline_repository is None or _image_with_tag(baseline_repository, tag) is not None:
319 raise RuntimeError(
320 f"Retained ECR image delta is not absent from the baseline: {image_key}"
321 )
322 repositories = repositories_by_region.get(region, [])
323 matches = [repository for repository in repositories if repository.get("name") == name]
324 if len(matches) != 1:
325 raise RuntimeError(
326 f"Baseline ECR repository changed before final comparison: {image_key}"
327 )
328 repository = matches[0]
329 images = repository.get("images", [])
330 tagged_images = [image for image in images if tag in image.get("tags", [])]
331 if len(tagged_images) > 1:
332 raise RuntimeError(f"Retained ECR tag resolves to multiple images: {image_key}")
333 if not tagged_images:
334 accepted["image_deltas"].append(
335 {"region": region, "repository": name, "tag": tag, "already_absent": True}
336 )
337 continue
338 image = tagged_images[0]
339 if _ecr_image_identity(image) != record.get("identity"):
340 raise RuntimeError(f"Retained ECR image identity changed for {image_key}")
341 remaining_tags = sorted(value for value in image.get("tags", []) if value != tag)
342 baseline_digest_matches = [
343 baseline_image
344 for baseline_image in baseline_repository.get("images", [])
345 if baseline_image.get("digest") == image.get("digest")
346 ]
347 if len(baseline_digest_matches) > 1:
348 raise RuntimeError(f"Baseline ECR digest is ambiguous for {image_key}")
349 if remaining_tags or baseline_digest_matches:
350 image["tags"] = remaining_tags
351 else:
352 images.remove(image)
353 accepted["image_deltas"].append(
354 {
355 "region": region,
356 "repository": name,
357 "tag": tag,
358 "digest": record["identity"]["digest"],
359 "retained": True,
360 }
361 )
363 return sanitized, accepted
366def _strip_accepted_retained_ecr(
367 project_inventory: dict[str, Any],
368 accepted: dict[str, list[dict[str, Any]]],
369) -> dict[str, Any]:
370 """Exclude exact retained repositories after final identity revalidation."""
371 allowed = {
372 (str(item["region"]), str(item["name"]))
373 for item in accepted.get("repositories", [])
374 if item.get("retained")
375 }
376 inventory = copy.deepcopy(project_inventory)
377 for region, resources in list(inventory.get("regional", {}).items()):
378 resources["ecr_repositories"] = [
379 name
380 for name in resources.get("ecr_repositories", [])
381 if (str(region), str(name)) not in allowed
382 ]
383 if not any(resources.values()):
384 inventory["regional"].pop(region)
385 return inventory
388def _merge_expected_ecr_target(
389 targets: dict[tuple[str, str, str], dict[str, Any]],
390 *,
391 region: str,
392 repository: str,
393 tag: str,
394 source: dict[str, str],
395) -> None:
396 if not region or not repository or not tag:
397 raise RuntimeError(f"Invalid expected ECR target: {region}:{repository}:{tag}")
398 key = (region, repository, tag)
399 target = targets.setdefault(
400 key,
401 {"region": region, "repository": repository, "tag": tag, "sources": []},
402 )
403 if source not in target["sources"]:
404 target["sources"].append(source)
407def _expected_ecr_images(ctx: RunContext, stack_names: list[str]) -> list[dict[str, Any]]:
408 """Derive exact CDK-asset and configured mirror tags without AWS writes."""
409 targets: dict[tuple[str, str, str], dict[str, Any]] = {}
410 assembly = ctx.settings.repo_root / "cdk.out"
411 for stack_name in stack_names:
412 path = assembly / f"{stack_name}.assets.json"
413 try:
414 document = json.loads(path.read_text(encoding="utf-8"))
415 except (OSError, UnicodeError, json.JSONDecodeError) as exc:
416 raise RuntimeError(f"Could not read cloud assembly assets {path}: {exc}") from exc
417 docker_images = document.get("dockerImages") if isinstance(document, dict) else None
418 if not isinstance(docker_images, dict):
419 raise RuntimeError(f"Cloud assembly {path} omitted dockerImages")
420 for asset_id, asset in docker_images.items():
421 destinations = asset.get("destinations") if isinstance(asset, dict) else None
422 if not isinstance(destinations, dict):
423 raise RuntimeError(f"Docker asset {stack_name}:{asset_id} has no destinations")
424 for destination in destinations.values():
425 if not isinstance(destination, dict):
426 raise RuntimeError(f"Docker asset {stack_name}:{asset_id} is malformed")
427 _merge_expected_ecr_target(
428 targets,
429 region=str(destination.get("region") or ""),
430 repository=str(destination.get("repositoryName") or ""),
431 tag=str(destination.get("imageTag") or ""),
432 source={"kind": "cdk-asset", "stack": stack_name, "asset_id": str(asset_id)},
433 )
435 from cli import _image_mirror
437 mirror_config = _image_mirror.read_mirror_config(ctx.settings.repo_root / "cdk.json")
438 if mirror_config["enabled"]:
439 source_refs = _image_mirror.collect_source_refs()
440 for region in ctx.deployment_regions:
441 plan = _image_mirror.plan_from_sources(
442 source_refs,
443 f"validation.invalid.{region}",
444 mirror_config["ecr_namespace"],
445 )
446 for item in plan:
447 _merge_expected_ecr_target(
448 targets,
449 region=region,
450 repository=item.dest_repo,
451 tag=item.tag,
452 source={"kind": "configured-mirror", "source_ref": item.source_ref},
453 )
455 return [
456 {
457 **target,
458 "sources": sorted(target["sources"], key=lambda item: json.dumps(item, sort_keys=True)),
459 }
460 for _key, target in sorted(targets.items())
461 ]
464def _ecr_creation_identity(repository: dict[str, Any]) -> dict[str, Any]:
465 """Return immutable-enough ECR creation fields for delete authorization."""
466 return {
467 "name": str(repository.get("name") or ""),
468 "arn": str(repository.get("arn") or ""),
469 "registry_id": str(repository.get("registry_id") or ""),
470 "created_at": repository.get("created_at"),
471 }
474def _record_ecr_repository_creation(
475 ctx: RunContext,
476 region: str,
477 repository: Mapping[str, Any],
478) -> None:
479 """Persist the synchronous create_repository acknowledgement before any copy."""
480 name = str(repository.get("repositoryName") or "")
481 arn = str(repository.get("repositoryArn") or "")
482 registry_id = str(repository.get("registryId") or "")
483 created_at_raw = repository.get("createdAt")
484 if created_at_raw is None:
485 created_at = ""
486 elif hasattr(created_at_raw, "isoformat"):
487 created_at = str(created_at_raw.isoformat())
488 else:
489 created_at = str(created_at_raw)
490 expected = {
491 (str(item["region"]), str(item["repository"]))
492 for item in ctx.checkpoint.state.get("expected_ecr_images", [])
493 }
494 baseline_names = {
495 (str(baseline_region), str(item["name"]))
496 for baseline_region, repositories in (ctx.checkpoint.baseline or {})
497 .get("ecr_repositories", {})
498 .items()
499 for item in repositories
500 }
501 key = (region, name)
502 if key not in expected or key in baseline_names:
503 raise RuntimeError(f"Unexpected ECR repository creation acknowledgement: {region}:{name}")
504 if not arn or not registry_id or not created_at:
505 raise RuntimeError(f"ECR creation acknowledgement is incomplete for {region}:{name}")
506 creation_identity = {
507 "name": name,
508 "arn": arn,
509 "registry_id": registry_id,
510 "created_at": created_at,
511 }
512 with ctx.state_lock:
513 records = ctx.checkpoint.state.setdefault("created_ecr_repositories", [])
514 matches = [
515 item for item in records if item.get("region") == region and item.get("name") == name
516 ]
517 if len(matches) > 1:
518 raise RuntimeError(f"Duplicate ECR creation acknowledgements for {region}:{name}")
519 candidate = {
520 "region": region,
521 "name": name,
522 "arn": arn,
523 "creation_identity": creation_identity,
524 "run_tag": ctx.settings.run_id,
525 "cleanup_policy": "retain-no-conditional-delete",
526 }
527 if matches and matches[0] != candidate:
528 raise RuntimeError(f"ECR creation acknowledgement changed for {region}:{name}")
529 if not matches:
530 records.append(candidate)
531 ctx.persist_callback(ctx.checkpoint)
534def _checkpoint_new_ecr_repositories(ctx: RunContext) -> list[dict[str, Any]]:
535 """Reconcile only repositories backed by persisted create acknowledgements."""
536 records = ctx.checkpoint.state.get("created_ecr_repositories", [])
537 if not records:
538 return []
539 current = collect_ecr_inventory(
540 ctx.session,
541 {str(item["region"]) for item in records},
542 )
543 current_by_key = {
544 (region, str(repository["name"])): repository
545 for region, repositories in current.items()
546 for repository in repositories
547 }
548 for record in records:
549 key = (str(record["region"]), str(record["name"]))
550 repository = current_by_key.get(key)
551 if repository is None:
552 record["observed_absent"] = True
553 continue
554 if _ecr_creation_identity(repository) != record.get("creation_identity"):
555 raise RuntimeError(f"ECR repository identity changed for {key[0]}:{key[1]}")
556 if (repository.get("tags") or {}).get(_RUN_STACK_TAG) != record.get("run_tag"):
557 raise RuntimeError(f"ECR repository run tag changed for {record['arn']}")
558 record["last_observed"] = _ecr_creation_identity(repository)
559 ctx.persist()
560 return copy.deepcopy(records)
563def _ecr_image_identity(image: dict[str, Any]) -> dict[str, Any]:
564 return {
565 "digest": str(image.get("digest") or ""),
566 "manifest_media_type": str(image.get("manifest_media_type") or ""),
567 "artifact_media_type": str(image.get("artifact_media_type") or ""),
568 "manifest": image.get("manifest"),
569 }
572def _image_with_tag(repository: dict[str, Any], tag: str) -> dict[str, Any] | None:
573 matches = [image for image in repository.get("images", []) if tag in image.get("tags", [])]
574 if len(matches) > 1:
575 raise RuntimeError(f"ECR tag {repository.get('name')}:{tag} resolved to multiple digests")
576 return matches[0] if matches else None
579def _checkpoint_new_ecr_images(ctx: RunContext) -> list[dict[str, Any]]:
580 """Observe ECR deltas without converting mutable tags into delete authority."""
581 baseline = ctx.checkpoint.baseline
582 if baseline is None:
583 raise RuntimeError("Cannot reconcile ECR images without a baseline")
584 baseline_repositories = {
585 (region, str(repository["name"])): repository
586 for region, repositories in baseline.get("ecr_repositories", {}).items()
587 for repository in repositories
588 }
589 current = collect_ecr_inventory(
590 ctx.session,
591 baseline.get("ecr_regions") or _topology_regions(ctx),
592 )
593 current_repositories = {
594 (region, str(repository["name"])): repository
595 for region, repositories in current.items()
596 for repository in repositories
597 }
599 deltas: list[dict[str, Any]] = []
600 for expected in ctx.checkpoint.state.get("expected_ecr_images", []):
601 key = (
602 str(expected["region"]),
603 str(expected["repository"]),
604 str(expected["tag"]),
605 )
606 baseline_repository = baseline_repositories.get(key[:2])
607 if baseline_repository is None:
608 continue
609 current_repository = current_repositories.get(key[:2])
610 if current_repository is None:
611 raise RuntimeError(f"Baseline ECR repository disappeared: {key[0]}:{key[1]}")
612 before = _image_with_tag(baseline_repository, key[2])
613 now = _image_with_tag(current_repository, key[2])
614 if before is not None:
615 if now is None or _ecr_image_identity(now) != _ecr_image_identity(before):
616 raise RuntimeError(f"Baseline ECR tag changed during validation: {key}")
617 continue
618 if now is None:
619 continue
620 identity = _ecr_image_identity(now)
621 if not identity["digest"]:
622 raise RuntimeError(f"Expected ECR tag lacks a digest: {key}")
623 deltas.append(
624 {
625 "region": key[0],
626 "repository": key[1],
627 "tag": key[2],
628 "identity": identity,
629 "sources": expected.get("sources", []),
630 "cleanup_policy": "retain-no-conditional-delete",
631 }
632 )
633 with ctx.state_lock:
634 previous = ctx.checkpoint.state.get("retained_ecr_image_deltas")
635 if previous is not None and previous != deltas:
636 raise RuntimeError("Observed ECR image deltas changed during validation")
637 ctx.checkpoint.state["retained_ecr_image_deltas"] = deltas
638 ctx.checkpoint.state["owned_ecr_images"] = []
639 ctx.persist_callback(ctx.checkpoint)
640 return copy.deepcopy(deltas)