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

1"""ECR expectation, ownership checkpointing, and baseline stripping.""" 

2 

3from __future__ import annotations 

4 

5import copy 

6import json 

7import re 

8from collections.abc import Mapping 

9from typing import Any 

10 

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) 

31 

32 

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) 

49 

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) 

62 

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 ] 

133 

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 ] 

157 

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 ] 

164 

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 ] 

200 

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 ] 

227 

228 if not any(resources.values()): 

229 inventory["regional"].pop(region) 

230 

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 

254 

255 

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() 

275 

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 ) 

307 

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 ) 

362 

363 return sanitized, accepted 

364 

365 

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 

386 

387 

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) 

405 

406 

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 ) 

434 

435 from cli import _image_mirror 

436 

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 ) 

454 

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 ] 

462 

463 

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 } 

472 

473 

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) 

532 

533 

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) 

561 

562 

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 } 

570 

571 

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 

577 

578 

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 } 

598 

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)