Coverage for scripts / live_release_validation / ownership / log_groups.py: 100.00%

347 statements  

« prev     ^ index     » next       coverage.py v7.13.5, created at 2026-09-14 22:07 +0000

1"""CloudWatch log-group identity, stability observation, and ownership checkpointing.""" 

2 

3from __future__ import annotations 

4 

5import copy 

6import re 

7import time 

8import uuid 

9from collections.abc import Mapping 

10from datetime import datetime 

11from typing import Any, cast 

12 

13from botocore.exceptions import ClientError 

14 

15from ..constants import ( 

16 _EKS_LOG_GROUP_SUFFIXES, 

17 _LOG_CLEANUP_TOKEN_TAG, 

18 _LOG_GROUP_CHECKPOINT_STABLE_OBSERVATIONS, 

19 _LOG_GROUP_CLEANUP_STABLE_OBSERVATIONS, 

20 _LOG_GROUP_OBSERVATION_ATTEMPTS, 

21 _LOG_GROUP_OBSERVATION_HISTORY_LIMIT, 

22 _LOG_GROUP_OBSERVATION_POLL_SECONDS, 

23 _LOG_GROUP_RETRYABLE_OBSERVATION_CODES, 

24 _LOG_GROUP_SOURCE_TYPES, 

25 _RUN_STACK_TAG, 

26) 

27from ..inventory import ( 

28 describe_stack, 

29) 

30from ..models import RunContext, utc_now 

31from ..ownership.stacks import ( 

32 _owned_stack_record, 

33) 

34 

35 

36def _derived_log_group_names(resource_type: str, physical_id: str) -> tuple[str, ...]: 

37 if resource_type == "AWS::Logs::LogGroup": 

38 return (physical_id,) 

39 if resource_type == "AWS::Lambda::Function": 

40 return (f"/aws/lambda/{physical_id}",) 

41 if resource_type == "AWS::EKS::Cluster": 

42 return ( 

43 f"/aws/eks/{physical_id}/cluster", 

44 *( 

45 f"/aws/containerinsights/{physical_id}/{suffix}" 

46 for suffix in _EKS_LOG_GROUP_SUFFIXES 

47 ), 

48 ) 

49 return () 

50 

51 

52def _live_eks_cluster_identity( 

53 ctx: RunContext, 

54 region: str, 

55 cluster_name: str, 

56) -> dict[str, str]: 

57 """Require the exact ACTIVE service-side cluster before deriving log authority.""" 

58 partition = ctx.session.get_partition_for_region(region) 

59 if not partition: 

60 raise RuntimeError(f"Could not resolve AWS partition for EKS cluster in {region}") 

61 expected_arn = ( 

62 f"arn:{partition}:eks:{region}:{ctx.settings.expected_account}:cluster/{cluster_name}" 

63 ) 

64 response = ctx.session.client("eks", region_name=region).describe_cluster(name=cluster_name) 

65 cluster = response.get("cluster") 

66 if not isinstance(cluster, dict): 

67 raise RuntimeError(f"EKS omitted cluster identity for {region}:{cluster_name}") 

68 identity = { 

69 "name": str(cluster.get("name") or ""), 

70 "arn": str(cluster.get("arn") or ""), 

71 "status": str(cluster.get("status") or ""), 

72 } 

73 if identity != {"name": cluster_name, "arn": expected_arn, "status": "ACTIVE"}: 

74 raise RuntimeError( 

75 f"EKS cluster identity is not exact and ACTIVE for {region}:{cluster_name}" 

76 ) 

77 return identity 

78 

79 

80def _eks_cluster_log_authority_identity( 

81 ctx: RunContext, 

82 region: str, 

83 cluster_name: str, 

84 *, 

85 allow_deleted: bool, 

86) -> dict[str, str]: 

87 """Resolve EKS log authority, tolerating a rolled-back (deleted) cluster. 

88 

89 A create rollback deletes the cluster itself while its control-plane and 

90 Container Insights log groups survive. The DELETED tombstone identity is 

91 only ever derived from this run's own stack resource record. 

92 """ 

93 if not allow_deleted: 

94 return _live_eks_cluster_identity(ctx, region, cluster_name) 

95 partition = ctx.session.get_partition_for_region(region) 

96 if not partition: 

97 raise RuntimeError(f"Could not resolve AWS partition for EKS cluster in {region}") 

98 expected_arn = ( 

99 f"arn:{partition}:eks:{region}:{ctx.settings.expected_account}:cluster/{cluster_name}" 

100 ) 

101 try: 

102 return _live_eks_cluster_identity(ctx, region, cluster_name) 

103 except ClientError as exc: 

104 if exc.response.get("Error", {}).get("Code", "") != "ResourceNotFoundException": 

105 raise 

106 return {"name": cluster_name, "arn": expected_arn, "status": "DELETED"} 

107 

108 

109def _validated_owned_log_group_identity( 

110 ctx: RunContext, 

111 record: Mapping[str, Any], 

112) -> tuple[str, str]: 

113 """Validate an exact log-group name derived from a checkpointed stack resource.""" 

114 region = str(record.get("region") or "") 

115 name = str(record.get("name") or "") 

116 stack_name = str(record.get("stack_name") or "") 

117 stack_id = str(record.get("stack_id") or "") 

118 resource_type = str(record.get("source_resource_type") or "") 

119 logical_id = str(record.get("source_logical_id") or "") 

120 physical_id = str(record.get("source_physical_id") or "") 

121 target_regions = ctx.checkpoint.state.get("target_stack_regions") 

122 if not isinstance(target_regions, dict) or str(target_regions.get(stack_name) or "") != region: 

123 raise RuntimeError(f"Log-group checkpoint target stack is invalid for {region}:{name}") 

124 partition = ctx.session.get_partition_for_region(region) 

125 if not partition: 

126 raise RuntimeError(f"Could not resolve AWS partition for log group in {region}") 

127 expected_stack_prefix = ( 

128 f"arn:{partition}:cloudformation:{region}:{ctx.settings.expected_account}:" 

129 f"stack/{stack_name}/" 

130 ) 

131 owned_stack_record = _owned_stack_record(ctx, region, stack_name) 

132 expected_stack_id = str((owned_stack_record or {}).get("stack_id") or "") 

133 if ( 

134 not name 

135 or not logical_id 

136 or resource_type not in _LOG_GROUP_SOURCE_TYPES 

137 or name not in _derived_log_group_names(resource_type, physical_id) 

138 ): 

139 raise RuntimeError(f"Log-group checkpoint source is invalid for {region}:{name}") 

140 source_service_identity = record.get("source_service_identity") 

141 if resource_type == "AWS::EKS::Cluster": 

142 expected_arn = ( 

143 f"arn:{partition}:eks:{region}:{ctx.settings.expected_account}:cluster/{physical_id}" 

144 ) 

145 # ACTIVE is the pre-destroy authority; DELETED is the exact tombstone 

146 # recorded when a rolled-back create removed the cluster but left its 

147 # control-plane and Container Insights log groups behind. 

148 accepted_source_identities = tuple( 

149 {"name": physical_id, "arn": expected_arn, "status": status} 

150 for status in ("ACTIVE", "DELETED") 

151 ) 

152 if source_service_identity not in accepted_source_identities: 

153 raise RuntimeError( 

154 f"Log-group checkpoint lacks exact live EKS identity for {region}:{name}" 

155 ) 

156 elif source_service_identity not in (None, {}): 

157 raise RuntimeError(f"Unexpected service identity for log group {region}:{name}") 

158 cleanup_token = str(record.get("cleanup_token") or "") 

159 expected_cleanup_token = str(ctx.checkpoint.state.get("log_group_cleanup_token") or "") 

160 if ( 

161 not expected_stack_id.startswith(expected_stack_prefix) 

162 or stack_id != expected_stack_id 

163 or (owned_stack_record or {}).get("run_tag") != ctx.settings.run_id 

164 or record.get("run_tag") != ctx.settings.run_id 

165 or record.get("ownership_authority") != "cloudformation-stack-resource-derived" 

166 or record.get("authority_phase") != "pre-destroy" 

167 or not cleanup_token 

168 or cleanup_token != expected_cleanup_token 

169 ): 

170 raise RuntimeError(f"Log-group checkpoint authority is invalid for {region}:{name}") 

171 return region, name 

172 

173 

174def _describe_exact_log_group(client: Any, name: str) -> dict[str, Any] | None: 

175 kwargs: dict[str, Any] = {"logGroupNamePrefix": name, "limit": 50} 

176 while True: 

177 response = client.describe_log_groups(**kwargs) 

178 for log_group in response.get("logGroups", []): 

179 if not isinstance(log_group, Mapping): 

180 raise RuntimeError("CloudWatch Logs returned a non-object log-group record") 

181 candidate = cast(Mapping[str, Any], log_group) 

182 if str(candidate.get("logGroupName") or "") == name: 

183 return {str(key): value for key, value in candidate.items()} 

184 token = response.get("nextToken") 

185 if not token: 

186 return None 

187 kwargs["nextToken"] = token 

188 

189 

190def _log_group_identity(client: Any, region: str, name: str) -> dict[str, Any] | None: 

191 log_group = _describe_exact_log_group(client, name) 

192 if log_group is None: 

193 return None 

194 arn = str(log_group.get("logGroupArn") or log_group.get("arn") or "").removesuffix(":*") 

195 creation_time = log_group.get("creationTime") 

196 if not arn or not isinstance(creation_time, int): 

197 raise RuntimeError(f"CloudWatch Logs omitted identity for {region}:{name}") 

198 tags = client.list_tags_for_resource(resourceArn=arn).get("tags") or {} 

199 return { 

200 "arn": arn, 

201 "creation_time": creation_time, 

202 "tags": {str(key): str(value) for key, value in tags.items()}, 

203 } 

204 

205 

206def _log_group_generation(identity: Mapping[str, Any]) -> dict[str, Any]: 

207 """Return the immutable fields that distinguish same-name log generations.""" 

208 arn = str(identity.get("arn") or "") 

209 creation_time = identity.get("creation_time") 

210 if not arn or not isinstance(creation_time, int): 

211 raise RuntimeError("Checkpointed CloudWatch log-group identity is malformed") 

212 return {"arn": arn, "creation_time": creation_time} 

213 

214 

215def _observe_log_group_stability( 

216 client: Any, 

217 region: str, 

218 name: str, 

219 *, 

220 expected_identity: Mapping[str, Any] | None = None, 

221 expected_tags: Mapping[str, str] | None = None, 

222 required_present: int | None, 

223 required_absent: int | None, 

224 attempts: int = _LOG_GROUP_OBSERVATION_ATTEMPTS, 

225 poll_seconds: float = _LOG_GROUP_OBSERVATION_POLL_SECONDS, 

226) -> dict[str, Any]: 

227 """Bound identity reads until presence/absence is stable or a fence is crossed.""" 

228 for label, value in ( 

229 ("attempts", attempts), 

230 ("required_present", required_present), 

231 ("required_absent", required_absent), 

232 ): 

233 if value is not None and ( 

234 isinstance(value, bool) or not isinstance(value, int) or value <= 0 

235 ): 

236 raise ValueError(f"{label} must be a positive integer or None") 

237 if required_present is None and required_absent is None: 

238 raise ValueError("At least one stable log-group outcome must be requested") 

239 if poll_seconds < 0: 

240 raise ValueError("poll_seconds must be non-negative") 

241 

242 expected_generation = ( 

243 _log_group_generation(expected_identity) if expected_identity is not None else None 

244 ) 

245 expected_authority_tags = {str(key): str(value) for key, value in (expected_tags or {}).items()} 

246 observations: list[dict[str, Any]] = [] 

247 seen_generations: list[dict[str, Any]] = [] 

248 present_streak = 0 

249 absent_streak = 0 

250 replacement_streak = 0 

251 replacement_generation: dict[str, Any] | None = None 

252 last_identity: dict[str, Any] | None = None 

253 

254 def result(status: str, **extra: Any) -> dict[str, Any]: 

255 return { 

256 "status": status, 

257 "region": region, 

258 "name": name, 

259 "attempt_count": len(observations), 

260 "observations": observations, 

261 "identity": copy.deepcopy(last_identity), 

262 **extra, 

263 } 

264 

265 for attempt in range(1, attempts + 1): 

266 observed_at = utc_now() 

267 try: 

268 identity = _log_group_identity(client, region, name) 

269 except ClientError as exc: 

270 code = str(exc.response.get("Error", {}).get("Code") or "") 

271 if code not in _LOG_GROUP_RETRYABLE_OBSERVATION_CODES: 

272 raise 

273 observations.append( 

274 { 

275 "attempt": attempt, 

276 "observed_at": observed_at, 

277 "status": "retryable-error", 

278 "error_code": code, 

279 "error": str(exc), 

280 } 

281 ) 

282 present_streak = 0 

283 absent_streak = 0 

284 replacement_streak = 0 

285 replacement_generation = None 

286 else: 

287 last_identity = copy.deepcopy(identity) 

288 if identity is None: 

289 if expected_generation is None and seen_generations: 

290 observations.append( 

291 { 

292 "attempt": attempt, 

293 "observed_at": observed_at, 

294 "status": "replacement", 

295 "observed_generation": None, 

296 } 

297 ) 

298 return result( 

299 "replacement", 

300 expected_generation=seen_generations[-1], 

301 observed_generation=None, 

302 ) 

303 replacement_streak = 0 

304 replacement_generation = None 

305 absent_streak += 1 

306 present_streak = 0 

307 observations.append( 

308 { 

309 "attempt": attempt, 

310 "observed_at": observed_at, 

311 "status": "absent", 

312 "consecutive": absent_streak, 

313 } 

314 ) 

315 if required_absent is not None and absent_streak >= required_absent: 

316 return result("absent", consecutive=absent_streak) 

317 else: 

318 generation = _log_group_generation(identity) 

319 if generation not in seen_generations: 

320 seen_generations.append(generation) 

321 observations.append( 

322 { 

323 "attempt": attempt, 

324 "observed_at": observed_at, 

325 "status": "present", 

326 "generation": copy.deepcopy(generation), 

327 } 

328 ) 

329 if expected_generation is not None and generation != expected_generation: 

330 if generation == replacement_generation: 

331 replacement_streak += 1 

332 else: 

333 replacement_generation = generation 

334 replacement_streak = 1 

335 observations[-1]["status"] = "replacement-candidate" 

336 observations[-1]["consecutive"] = replacement_streak 

337 present_streak = 0 

338 absent_streak = 0 

339 if replacement_streak >= _LOG_GROUP_CLEANUP_STABLE_OBSERVATIONS: 

340 observations[-1]["status"] = "replacement" 

341 return result( 

342 "replacement", 

343 expected_generation=expected_generation, 

344 observed_generation=generation, 

345 consecutive=replacement_streak, 

346 ) 

347 if attempt < attempts: 

348 time.sleep(poll_seconds) 

349 continue 

350 replacement_streak = 0 

351 replacement_generation = None 

352 if expected_generation is None and len(seen_generations) > 1: 

353 observations[-1]["status"] = "replacement" 

354 return result( 

355 "replacement", 

356 expected_generation=seen_generations[0], 

357 observed_generation=generation, 

358 ) 

359 tags = identity.get("tags") or {} 

360 tag_drift = { 

361 key: {"expected": value, "observed": tags.get(key)} 

362 for key, value in expected_authority_tags.items() 

363 if tags.get(key) != value 

364 } 

365 if tag_drift: 

366 observations[-1]["status"] = "tag-drift" 

367 observations[-1]["tag_drift"] = copy.deepcopy(tag_drift) 

368 return result("tag-drift", tag_drift=tag_drift) 

369 present_streak += 1 

370 absent_streak = 0 

371 observations[-1]["consecutive"] = present_streak 

372 if required_present is not None and present_streak >= required_present: 

373 return result("present", consecutive=present_streak) 

374 if attempt < attempts: 

375 time.sleep(poll_seconds) 

376 

377 return result( 

378 "unsettled", 

379 required_present=required_present, 

380 required_absent=required_absent, 

381 present_streak=present_streak, 

382 absent_streak=absent_streak, 

383 ) 

384 

385 

386def _record_log_group_observation( 

387 ctx: RunContext, 

388 record: dict[str, Any], 

389 *, 

390 phase: str, 

391 outcome: Mapping[str, Any], 

392) -> dict[str, Any]: 

393 """Persist bounded identity evidence and any confirmed replacement generation.""" 

394 entry = {"phase": phase, "recorded_at": utc_now(), **copy.deepcopy(dict(outcome))} 

395 with ctx.state_lock: 

396 history = record.setdefault("identity_observation_history", []) 

397 if not isinstance(history, list): 

398 raise RuntimeError("Log-group identity_observation_history must be a list") 

399 history.append(entry) 

400 del history[:-_LOG_GROUP_OBSERVATION_HISTORY_LIMIT] 

401 if outcome.get("status") == "replacement": 

402 replacements = record.setdefault("replacement_evidence", []) 

403 if not isinstance(replacements, list): 

404 raise RuntimeError("Log-group replacement_evidence must be a list") 

405 replacements.append(copy.deepcopy(entry)) 

406 ctx.persist_callback(ctx.checkpoint) 

407 return entry 

408 

409 

410def _record_log_group_checkpoint_incident( 

411 ctx: RunContext, 

412 candidate: Mapping[str, Any], 

413 *, 

414 phase: str, 

415 outcome: Mapping[str, Any], 

416) -> None: 

417 """Preserve failed pre-authority observations without adopting the generation.""" 

418 with ctx.state_lock: 

419 incidents = ctx.checkpoint.state.setdefault("log_group_checkpoint_incidents", []) 

420 if not isinstance(incidents, list): 

421 raise RuntimeError("Checkpoint log_group_checkpoint_incidents must be a list") 

422 incidents.append( 

423 { 

424 "phase": phase, 

425 "recorded_at": utc_now(), 

426 "candidate": copy.deepcopy(dict(candidate)), 

427 "outcome": copy.deepcopy(dict(outcome)), 

428 } 

429 ) 

430 ctx.persist_callback(ctx.checkpoint) 

431 

432 

433def _set_log_group_disposition( 

434 ctx: RunContext, 

435 record: dict[str, Any], 

436 *, 

437 status: str, 

438 phase: str, 

439 outcome: Mapping[str, Any], 

440) -> dict[str, Any]: 

441 disposition = { 

442 "status": status, 

443 "phase": phase, 

444 "recorded_at": utc_now(), 

445 "original_identity": copy.deepcopy(record.get("observed_identity")), 

446 "last_observation_status": str(outcome.get("status") or ""), 

447 } 

448 with ctx.state_lock: 

449 record["original_generation_disposition"] = disposition 

450 ctx.persist_callback(ctx.checkpoint) 

451 return disposition 

452 

453 

454def _checkpoint_owned_log_groups(ctx: RunContext) -> list[dict[str, Any]]: 

455 """Fence, tag, and checkpoint exact generations while source stacks are live.""" 

456 target_regions = ctx.checkpoint.state.get("target_stack_regions") 

457 if not isinstance(target_regions, dict): 

458 raise RuntimeError("Checkpoint target_stack_regions must be an object") 

459 try: 

460 run_started_ms = int(datetime.fromisoformat(ctx.checkpoint.created_at).timestamp() * 1000) 

461 except ValueError as exc: 

462 raise RuntimeError("Checkpoint created_at is not a valid timestamp") from exc 

463 

464 with ctx.state_lock: 

465 cleanup_token = str(ctx.checkpoint.state.get("log_group_cleanup_token") or "") 

466 if not cleanup_token: 

467 cleanup_token = uuid.uuid4().hex 

468 ctx.checkpoint.state["log_group_cleanup_token"] = cleanup_token 

469 if not re.fullmatch(r"[0-9a-f]{32}", cleanup_token): 

470 raise RuntimeError("Checkpoint log-group cleanup token is malformed") 

471 records = ctx.checkpoint.state.setdefault("owned_log_groups", []) 

472 if not isinstance(records, list): 

473 raise RuntimeError("Checkpoint owned_log_groups must be a list") 

474 by_identity = { 

475 (str(item.get("region") or ""), str(item.get("name") or "")): item 

476 for item in records 

477 if isinstance(item, dict) 

478 } 

479 authority_tags = { 

480 _RUN_STACK_TAG: ctx.settings.run_id, 

481 _LOG_CLEANUP_TOKEN_TAG: cleanup_token, 

482 } 

483 for stack_name, raw_region in sorted(target_regions.items()): 

484 region = str(raw_region) 

485 stack_record = _owned_stack_record(ctx, region, str(stack_name)) 

486 if stack_record is None: 

487 continue 

488 live_stack = describe_stack(ctx.session, region, stack_record["stack_id"]) 

489 if ( 

490 live_stack is None 

491 or str(live_stack.get("status") or "").startswith("DELETE") 

492 or (live_stack.get("tags") or {}).get(_RUN_STACK_TAG) != ctx.settings.run_id 

493 ): 

494 # Destructive authority is never created or completed from a 

495 # deleted-stack tombstone. Existing pre-destroy records remain usable. 

496 continue 

497 # A rolled-back create leaves log groups behind: retained LogGroup 

498 # resources, Lambda-created default groups, and EKS control-plane 

499 # groups all survive resource deletion. Their stack resources read 

500 # DELETE_COMPLETE while the stack itself is still describable, so 

501 # rollback statuses widen the resource filter to those tombstones. 

502 rolled_back = str(live_stack.get("status") or "") in { 

503 "ROLLBACK_COMPLETE", 

504 "ROLLBACK_FAILED", 

505 "UPDATE_ROLLBACK_COMPLETE", 

506 "UPDATE_ROLLBACK_FAILED", 

507 } 

508 allowed_resource_statuses = {"CREATE_COMPLETE", "UPDATE_COMPLETE"} 

509 if rolled_back: 

510 allowed_resource_statuses |= {"DELETE_COMPLETE", "DELETE_FAILED", "DELETE_SKIPPED"} 

511 cfn = ctx.session.client("cloudformation", region_name=region) 

512 pages = cfn.get_paginator("list_stack_resources").paginate( 

513 StackName=stack_record["stack_id"] 

514 ) 

515 resources = [ 

516 item 

517 for page in pages 

518 for item in page.get("StackResourceSummaries", []) 

519 if str(item.get("ResourceType") or "") in _LOG_GROUP_SOURCE_TYPES 

520 and item.get("LogicalResourceId") 

521 and item.get("PhysicalResourceId") 

522 and str(item.get("ResourceStatus") or "") in allowed_resource_statuses 

523 ] 

524 logs = ctx.session.client("logs", region_name=region) 

525 lambda_client = ctx.session.client("lambda", region_name=region) 

526 for resource in resources: 

527 resource_type = str(resource["ResourceType"]) 

528 physical_id = str(resource["PhysicalResourceId"]) 

529 source_service_identity = None 

530 if resource_type == "AWS::EKS::Cluster": 

531 source_service_identity = _eks_cluster_log_authority_identity( 

532 ctx, 

533 region, 

534 physical_id, 

535 allow_deleted=rolled_back, 

536 ) 

537 names = _derived_log_group_names(resource_type, physical_id) 

538 if resource_type == "AWS::Lambda::Function": 

539 default_name = f"/aws/lambda/{physical_id}" 

540 try: 

541 function = lambda_client.get_function_configuration( 

542 FunctionName=physical_id 

543 ) 

544 except ClientError as exc: 

545 error_code = exc.response.get("Error", {}).get("Code", "") 

546 if not (rolled_back and error_code == "ResourceNotFoundException"): 

547 raise 

548 # The rolled-back function is gone; only its default 

549 # log group can remain. 

550 names = (default_name,) 

551 else: 

552 configured_name = str( 

553 (function.get("LoggingConfig") or {}).get("LogGroup") or default_name 

554 ) 

555 names = (default_name,) if configured_name == default_name else () 

556 for name in names: 

557 key = (region, name) 

558 candidate = { 

559 "region": region, 

560 "name": name, 

561 "stack_name": str(stack_name), 

562 "stack_id": stack_record["stack_id"], 

563 "source_resource_type": resource_type, 

564 "source_logical_id": str(resource["LogicalResourceId"]), 

565 "source_physical_id": physical_id, 

566 "ownership_authority": "cloudformation-stack-resource-derived", 

567 "authority_phase": "pre-destroy", 

568 "run_tag": ctx.settings.run_id, 

569 "cleanup_token": cleanup_token, 

570 } 

571 if source_service_identity is not None: 

572 candidate["source_service_identity"] = source_service_identity 

573 _validated_owned_log_group_identity(ctx, candidate) 

574 

575 previous = by_identity.get(key) 

576 immutable = tuple(candidate) 

577 expected_identity: Mapping[str, Any] | None = None 

578 if previous is not None: 

579 if any(previous.get(field) != candidate[field] for field in immutable): 

580 raise RuntimeError(f"Log-group ownership changed for {region}:{name}") 

581 observed = previous.get("observed_identity") 

582 if not isinstance(observed, dict): 

583 raise RuntimeError( 

584 f"Log-group checkpoint identity is malformed for {region}:{name}" 

585 ) 

586 expected_identity = observed 

587 

588 initial = _observe_log_group_stability( 

589 logs, 

590 region, 

591 name, 

592 expected_identity=expected_identity, 

593 expected_tags=authority_tags if previous is not None else None, 

594 required_present=_LOG_GROUP_CHECKPOINT_STABLE_OBSERVATIONS, 

595 required_absent=_LOG_GROUP_CHECKPOINT_STABLE_OBSERVATIONS, 

596 ) 

597 if previous is not None: 

598 _record_log_group_observation( 

599 ctx, 

600 previous, 

601 phase="checkpoint-revalidation", 

602 outcome=initial, 

603 ) 

604 if initial["status"] != "present": 

605 status = ( 

606 "replacement-observed-during-checkpoint" 

607 if initial["status"] == "replacement" 

608 else "checkpoint-generation-not-stable" 

609 ) 

610 _set_log_group_disposition( 

611 ctx, 

612 previous, 

613 status=status, 

614 phase="checkpoint-revalidation", 

615 outcome=initial, 

616 ) 

617 raise RuntimeError( 

618 f"Log-group checkpoint generation is not stable for " 

619 f"{region}:{name}: {initial['status']}" 

620 ) 

621 continue 

622 

623 if initial["status"] == "absent": 

624 if resource_type == "AWS::Logs::LogGroup": 

625 _record_log_group_checkpoint_incident( 

626 ctx, 

627 candidate, 

628 phase="checkpoint-explicit-group-absence", 

629 outcome=initial, 

630 ) 

631 if rolled_back: 

632 # A rolled-back create deletes non-retained 

633 # LogGroup resources; the stack resource is a 

634 # DELETE_COMPLETE tombstone widened in only to 

635 # catch groups that DID survive (retained or 

636 # recreated by late log deliveries). A group 

637 # that is genuinely gone is the expected 

638 # rollback outcome, not missing authority — 

639 # aborting here strands every remaining stack 

640 # (observed live: example-job validation run 

641 # ex241-2913b044, where the Valkey subnet 

642 # failure rolled back gco-us-east-1 and this 

643 # guard then blocked the entire teardown). 

644 continue 

645 raise RuntimeError( 

646 f"CloudFormation log group is absent before teardown: " 

647 f"{region}:{name}" 

648 ) 

649 try: 

650 logs.create_log_group(logGroupName=name, tags=authority_tags) 

651 except ClientError as exc: 

652 if ( 

653 exc.response.get("Error", {}).get("Code") 

654 != "ResourceAlreadyExistsException" 

655 ): 

656 raise 

657 raced = _observe_log_group_stability( 

658 logs, 

659 region, 

660 name, 

661 expected_identity=None, 

662 expected_tags=None, 

663 required_present=_LOG_GROUP_CHECKPOINT_STABLE_OBSERVATIONS, 

664 required_absent=_LOG_GROUP_CHECKPOINT_STABLE_OBSERVATIONS, 

665 ) 

666 _record_log_group_checkpoint_incident( 

667 ctx, 

668 candidate, 

669 phase="checkpoint-create-race", 

670 outcome=raced, 

671 ) 

672 raise RuntimeError( 

673 f"Log group appeared during checkpoint creation; refusing to " 

674 f"adopt or tag it: {region}:{name}" 

675 ) from exc 

676 initial = _observe_log_group_stability( 

677 logs, 

678 region, 

679 name, 

680 expected_identity=None, 

681 expected_tags=authority_tags, 

682 required_present=_LOG_GROUP_CHECKPOINT_STABLE_OBSERVATIONS, 

683 required_absent=_LOG_GROUP_CHECKPOINT_STABLE_OBSERVATIONS, 

684 ) 

685 if initial["status"] != "present": 

686 _record_log_group_checkpoint_incident( 

687 ctx, 

688 candidate, 

689 phase="checkpoint-initial-stability", 

690 outcome=initial, 

691 ) 

692 raise RuntimeError( 

693 f"Log group could not be stably checkpointed: " 

694 f"{region}:{name}: {initial['status']}" 

695 ) 

696 

697 identity = initial.get("identity") 

698 if not isinstance(identity, dict): 

699 raise RuntimeError(f"Log group omitted identity: {region}:{name}") 

700 if identity["creation_time"] < run_started_ms: 

701 raise RuntimeError( 

702 f"Log group predates this validation run: {region}:{name}" 

703 ) 

704 tags = identity.get("tags") or {} 

705 conflicting_tags = { 

706 key: {"expected": value, "observed": tags.get(key)} 

707 for key, value in authority_tags.items() 

708 if key in tags and tags.get(key) != value 

709 } 

710 if conflicting_tags: 

711 conflict = { 

712 **copy.deepcopy(initial), 

713 "status": "tag-drift", 

714 "tag_drift": conflicting_tags, 

715 } 

716 _record_log_group_checkpoint_incident( 

717 ctx, 

718 candidate, 

719 phase="checkpoint-authority-tag-conflict", 

720 outcome=conflict, 

721 ) 

722 raise RuntimeError(f"Log-group authority tags conflict for {region}:{name}") 

723 

724 # This read is intentionally adjacent to tag_resource. A generation 

725 # change after the stable reads is fenced before any authority tags 

726 # can be applied. The post-tag stable reads catch the unavoidable 

727 # service-side TOCTOU without ever adopting that replacement. 

728 pre_tag = _observe_log_group_stability( 

729 logs, 

730 region, 

731 name, 

732 expected_identity=identity, 

733 expected_tags=None, 

734 required_present=1, 

735 required_absent=1, 

736 ) 

737 if pre_tag["status"] != "present": 

738 _record_log_group_checkpoint_incident( 

739 ctx, 

740 candidate, 

741 phase="checkpoint-immediate-pre-tag", 

742 outcome=pre_tag, 

743 ) 

744 raise RuntimeError( 

745 f"Log-group generation changed immediately before tagging: " 

746 f"{region}:{name}: {pre_tag['status']}" 

747 ) 

748 if any(tags.get(key) != value for key, value in authority_tags.items()): 

749 logs.tag_resource(resourceArn=identity["arn"], tags=authority_tags) 

750 

751 post_tag = _observe_log_group_stability( 

752 logs, 

753 region, 

754 name, 

755 expected_identity=identity, 

756 expected_tags=authority_tags, 

757 required_present=_LOG_GROUP_CHECKPOINT_STABLE_OBSERVATIONS, 

758 required_absent=_LOG_GROUP_CHECKPOINT_STABLE_OBSERVATIONS, 

759 ) 

760 if post_tag["status"] != "present": 

761 _record_log_group_checkpoint_incident( 

762 ctx, 

763 candidate, 

764 phase="checkpoint-post-tag-stability", 

765 outcome=post_tag, 

766 ) 

767 raise RuntimeError( 

768 f"Log-group generation or authority changed while checkpointing: " 

769 f"{region}:{name}: {post_tag['status']}" 

770 ) 

771 final_identity = post_tag.get("identity") 

772 if not isinstance(final_identity, dict): 

773 raise RuntimeError( 

774 f"Log group omitted its post-tag identity: {region}:{name}" 

775 ) 

776 candidate["observed_identity"] = final_identity 

777 candidate["checkpoint_observations"] = { 

778 "initial": initial, 

779 "immediate_pre_tag": pre_tag, 

780 "post_tag": post_tag, 

781 } 

782 candidate["original_generation_disposition"] = { 

783 "status": "checkpointed-present", 

784 "phase": "pre-destroy", 

785 "recorded_at": utc_now(), 

786 "original_identity": copy.deepcopy(final_identity), 

787 } 

788 records.append(candidate) 

789 by_identity[key] = candidate 

790 ctx.persist_callback(ctx.checkpoint) 

791 return copy.deepcopy(records)