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
« 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."""
3from __future__ import annotations
5import copy
6import re
7import time
8import uuid
9from collections.abc import Mapping
10from datetime import datetime
11from typing import Any, cast
13from botocore.exceptions import ClientError
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)
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 ()
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
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.
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"}
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
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
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 }
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}
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")
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
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 }
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)
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 )
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
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)
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
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
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)
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
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
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 )
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}")
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)
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)