Coverage for scripts / live_release_validation / cleanup / log_groups.py: 100.00%
193 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"""Delete exactly run-owned CloudWatch log groups."""
3from __future__ import annotations
5import copy
6import json
7import re
8import time
9from collections.abc import Mapping
10from typing import Any
12from botocore.exceptions import ClientError
14from ..constants import (
15 _LOG_CLEANUP_TOKEN_TAG,
16 _LOG_GROUP_ABSENCE_OBSERVATIONS,
17 _LOG_GROUP_CLEANUP_MAX_PASSES,
18 _LOG_GROUP_CLEANUP_STABLE_OBSERVATIONS,
19 _LOG_GROUP_OBSERVATION_POLL_SECONDS,
20 _RUN_STACK_TAG,
21 _LogGroupCleanupError,
22)
23from ..models import RunContext, utc_now
24from ..ownership.cleanup_role import (
25 TagConditionedLogDeleter,
26 _delete_log_cleanup_helper,
27)
28from ..ownership.log_groups import (
29 _log_group_generation,
30 _observe_log_group_stability,
31 _record_log_group_observation,
32 _set_log_group_disposition,
33 _validated_owned_log_group_identity,
34)
35from ..ownership.stacks import (
36 _verify_target_stack_absence,
37)
40def _log_group_adoption_blockers(
41 identity: Mapping[str, Any],
42 *,
43 run_id: str,
44 cleanup_token: str,
45) -> list[str]:
46 """Explain why a regenerated same-name log group cannot be adopted.
48 Teardown-time Lambda invocations flush their final events after their log
49 groups were tagged or deleted, recreating untagged generations that belong
50 to this run. Adoption is refused for any generation carrying another
51 owner's markers: a foreign validation run/cleanup token, or CloudFormation
52 stack tags (a real deployment's explicit LogGroup resources are always
53 stack-tagged, while Lambda-recreated groups start with no tags at all).
54 """
55 tags_value = identity.get("tags")
56 tags: dict[str, str] = dict(tags_value) if isinstance(tags_value, Mapping) else {}
57 blockers = []
58 if tags.get(_RUN_STACK_TAG) not in (None, run_id):
59 blockers.append(f"foreign {_RUN_STACK_TAG}={tags.get(_RUN_STACK_TAG)!r}")
60 if tags.get(_LOG_CLEANUP_TOKEN_TAG) not in (None, cleanup_token):
61 blockers.append(f"foreign {_LOG_CLEANUP_TOKEN_TAG}")
62 stack_tags = sorted(key for key in tags if key.startswith("aws:cloudformation:"))
63 if stack_tags:
64 blockers.append("cloudformation-owned generation: " + ", ".join(stack_tags))
65 return blockers
68def _adopt_regenerated_log_group(
69 ctx: RunContext,
70 record: dict[str, Any],
71 logs_client: Any,
72 *,
73 region: str,
74 name: str,
75 observed_generation: Mapping[str, Any],
76 authority_tags: Mapping[str, str],
77) -> dict[str, Any] | None:
78 """Tag and take ownership of a self-regenerated log-group generation.
80 Callers must already hold the invocation-level proof that every exact
81 target stack is absent, so no live deployment can own this name. Returns
82 the stabilized post-tag identity, or ``None`` when the generation did not
83 stabilize under this run's authority tags.
84 """
85 stack_absence = _verify_target_stack_absence(ctx)
86 if not stack_absence["all_absent"]:
87 raise RuntimeError("Log-group adoption requires every exact target stack to be absent")
88 arn = str(observed_generation.get("arn") or "")
89 if not arn:
90 raise RuntimeError(f"Regenerated log group omitted its ARN: {region}:{name}")
91 logs_client.tag_resource(resourceArn=arn, tags=dict(authority_tags))
92 post_tag = _observe_log_group_stability(
93 logs_client,
94 region,
95 name,
96 expected_identity={**observed_generation, "tags": dict(authority_tags)},
97 expected_tags=authority_tags,
98 required_present=_LOG_GROUP_CLEANUP_STABLE_OBSERVATIONS,
99 required_absent=_LOG_GROUP_ABSENCE_OBSERVATIONS,
100 )
101 _record_log_group_observation(
102 ctx,
103 record,
104 phase="cleanup-adoption-post-tag",
105 outcome=post_tag,
106 )
107 if post_tag["status"] != "present":
108 return None
109 identity = post_tag.get("identity")
110 if not isinstance(identity, dict):
111 raise RuntimeError(f"Adopted log group omitted its identity: {region}:{name}")
112 with ctx.state_lock:
113 record["observed_identity"] = copy.deepcopy(identity)
114 adoptions = record.setdefault("adopted_generations", [])
115 if not isinstance(adoptions, list):
116 raise RuntimeError("Log-group adopted_generations must be a list")
117 adoptions.append(
118 {
119 "adopted_at": utc_now(),
120 "generation": _log_group_generation(identity),
121 "stack_absence_proof_at": stack_absence.get("verified_at") or utc_now(),
122 }
123 )
124 ctx.persist_callback(ctx.checkpoint)
125 return identity
128def _blocked_log_group_entry(
129 ctx: RunContext,
130 record: dict[str, Any],
131 region: str,
132 name: str,
133 *,
134 status: str,
135 phase: str,
136 outcome: dict[str, Any],
137 retryable: bool,
138 delete_requested: bool = False,
139) -> dict[str, Any]:
140 """Record why one generation was preserved and build its report entry.
142 ``retryable`` distinguishes "observe again on the next sweep" (a log
143 delivery landed mid-observation) from "never delete this" (the generation
144 carries another owner's markers). Only retryable blockers trigger another
145 sweep; the rest are terminal and fail cleanup with their evidence intact.
146 """
147 disposition = _set_log_group_disposition(
148 ctx,
149 record,
150 status=status,
151 phase=phase,
152 outcome=outcome,
153 )
154 return {
155 "region": region,
156 "name": name,
157 "original_identity": copy.deepcopy(record.get("observed_identity")),
158 "deleted": False,
159 "blocked": True,
160 "retryable": retryable,
161 "delete_requested": delete_requested,
162 "observation": copy.deepcopy(outcome),
163 "replacement_evidence": copy.deepcopy(record.get("replacement_evidence", [])),
164 "original_generation_disposition": disposition,
165 }
168def _converge_one_log_group(
169 ctx: RunContext,
170 record: dict[str, Any],
171 region: str,
172 name: str,
173 *,
174 authority_tags: Mapping[str, str],
175 cleanup_token: str,
176 deleter: TagConditionedLogDeleter,
177) -> tuple[str, dict[str, Any]]:
178 """Drive one checkpointed log group to confirmed absence, or block it.
180 Returns ``("completed", entry)`` once the exact owned generation is gone
181 (or was already absent), and ``("blocked", entry)`` when the group must be
182 preserved. The four phases each re-establish identity before acting:
184 1. **pending-stability** — repeated reads must agree on the checkpointed
185 identity and both authority tags. An untagged same-name regeneration
186 from teardown-time Lambda logging is adopted here; a generation with a
187 foreign owner's markers blocks permanently.
188 2. **immediate-pre-delete** — one final exact read with nothing between it
189 and the tag-conditioned delete request.
190 3. **post-delete-absence** — absence must hold across repeated reads.
191 4. **disposition** — the outcome is checkpointed either way.
192 """
193 observed = record.get("observed_identity")
194 if not isinstance(observed, dict):
195 raise RuntimeError(f"Log-group checkpoint identity is malformed: {region}:{name}")
196 _log_group_generation(observed)
197 normal_logs = ctx.session.client("logs", region_name=region)
199 initial = _observe_log_group_stability(
200 normal_logs,
201 region,
202 name,
203 expected_identity=observed,
204 expected_tags=authority_tags,
205 required_present=_LOG_GROUP_CLEANUP_STABLE_OBSERVATIONS,
206 required_absent=_LOG_GROUP_ABSENCE_OBSERVATIONS,
207 )
208 _record_log_group_observation(
209 ctx,
210 record,
211 phase="cleanup-pending-stability",
212 outcome=initial,
213 )
214 if initial["status"] == "absent":
215 record["deleted"] = True
216 disposition = _set_log_group_disposition(
217 ctx,
218 record,
219 status="already-absent-confirmed",
220 phase="cleanup-pending-stability",
221 outcome=initial,
222 )
223 return (
224 "completed",
225 {
226 "region": region,
227 "name": name,
228 "original_identity": copy.deepcopy(observed),
229 "already_absent": True,
230 "absence_observations": initial["attempt_count"],
231 "original_generation_disposition": disposition,
232 },
233 )
235 identity: dict[str, Any] | None = None
236 if initial["status"] == "present":
237 candidate = initial.get("identity")
238 if not isinstance(candidate, dict):
239 raise RuntimeError(f"Stable log-group observation omitted identity: {region}:{name}")
240 identity = candidate
241 elif initial["status"] == "replacement":
242 replacement_identity = initial.get("identity")
243 if not isinstance(replacement_identity, dict):
244 # The regeneration vanished mid-observation; the next
245 # pass will see stable absence.
246 return (
247 "blocked",
248 _blocked_log_group_entry(
249 ctx,
250 record,
251 region,
252 name,
253 status="replacement-without-identity",
254 phase="cleanup-pending-stability",
255 outcome=initial,
256 retryable=True,
257 ),
258 )
259 blockers = _log_group_adoption_blockers(
260 replacement_identity,
261 run_id=ctx.settings.run_id,
262 cleanup_token=cleanup_token,
263 )
264 if blockers:
265 return (
266 "blocked",
267 _blocked_log_group_entry(
268 ctx,
269 record,
270 region,
271 name,
272 status="replacement-observed-before-delete",
273 phase="cleanup-pending-stability",
274 outcome={**initial, "adoption_blockers": blockers},
275 retryable=False,
276 ),
277 )
278 adopted = _adopt_regenerated_log_group(
279 ctx,
280 record,
281 normal_logs,
282 region=region,
283 name=name,
284 observed_generation=replacement_identity,
285 authority_tags=authority_tags,
286 )
287 if adopted is None:
288 return (
289 "blocked",
290 _blocked_log_group_entry(
291 ctx,
292 record,
293 region,
294 name,
295 status="adoption-did-not-stabilize",
296 phase="cleanup-adoption-post-tag",
297 outcome=initial,
298 retryable=True,
299 ),
300 )
301 identity = adopted
302 else:
303 if initial["status"] == "tag-drift":
304 disposition_status = "authority-tag-drift-before-delete"
305 retryable = False
306 else:
307 disposition_status = "identity-not-stable-before-delete"
308 retryable = True
309 return (
310 "blocked",
311 _blocked_log_group_entry(
312 ctx,
313 record,
314 region,
315 name,
316 status=disposition_status,
317 phase="cleanup-pending-stability",
318 outcome=initial,
319 retryable=retryable,
320 ),
321 )
323 restricted_logs = deleter.client(region)
324 # No persistence, sleep, or unrelated API call is permitted between
325 # this single exact read and the tag-conditioned delete request.
326 pre_delete = _observe_log_group_stability(
327 normal_logs,
328 region,
329 name,
330 expected_identity=identity,
331 expected_tags=authority_tags,
332 required_present=1,
333 required_absent=1,
334 )
335 if pre_delete["status"] == "present":
336 try:
337 restricted_logs.delete_log_group(logGroupName=name)
338 except ClientError as exc:
339 _record_log_group_observation(
340 ctx,
341 record,
342 phase="cleanup-immediate-pre-delete",
343 outcome=pre_delete,
344 )
345 if exc.response.get("Error", {}).get("Code") != "ResourceNotFoundException":
346 raise
347 else:
348 _record_log_group_observation(
349 ctx,
350 record,
351 phase="cleanup-immediate-pre-delete",
352 outcome=pre_delete,
353 )
354 record["delete_requested_at"] = utc_now()
355 ctx.persist()
356 else:
357 _record_log_group_observation(
358 ctx,
359 record,
360 phase="cleanup-immediate-pre-delete",
361 outcome=pre_delete,
362 )
364 if pre_delete["status"] not in {"present", "absent"}:
365 if pre_delete["status"] == "replacement":
366 disposition_status = "replacement-observed-immediately-before-delete"
367 retryable = not _log_group_adoption_blockers(
368 pre_delete.get("identity") or {},
369 run_id=ctx.settings.run_id,
370 cleanup_token=cleanup_token,
371 )
372 elif pre_delete["status"] == "tag-drift":
373 disposition_status = "authority-tag-drift-immediately-before-delete"
374 retryable = False
375 else:
376 disposition_status = "identity-not-stable-immediately-before-delete"
377 retryable = True
378 return (
379 "blocked",
380 _blocked_log_group_entry(
381 ctx,
382 record,
383 region,
384 name,
385 status=disposition_status,
386 phase="cleanup-immediate-pre-delete",
387 outcome=pre_delete,
388 retryable=retryable,
389 ),
390 )
392 absence = _observe_log_group_stability(
393 normal_logs,
394 region,
395 name,
396 expected_identity=identity,
397 expected_tags=authority_tags,
398 required_present=None,
399 required_absent=_LOG_GROUP_ABSENCE_OBSERVATIONS,
400 )
401 _record_log_group_observation(
402 ctx,
403 record,
404 phase="cleanup-post-delete-absence",
405 outcome=absence,
406 )
407 if absence["status"] != "absent":
408 if absence["status"] == "replacement":
409 disposition_status = "replacement-observed-before-confirmed-absence"
410 retryable = not _log_group_adoption_blockers(
411 absence.get("identity") or {},
412 run_id=ctx.settings.run_id,
413 cleanup_token=cleanup_token,
414 )
415 elif absence["status"] == "tag-drift":
416 disposition_status = "authority-tag-drift-after-delete-request"
417 retryable = False
418 else:
419 disposition_status = "absence-not-stable-after-delete-request"
420 retryable = True
421 return (
422 "blocked",
423 _blocked_log_group_entry(
424 ctx,
425 record,
426 region,
427 name,
428 status=disposition_status,
429 phase="cleanup-post-delete-absence",
430 outcome=absence,
431 retryable=retryable,
432 delete_requested=pre_delete["status"] == "present",
433 ),
434 )
436 record["deleted"] = True
437 disposition = _set_log_group_disposition(
438 ctx,
439 record,
440 status="deleted-confirmed-absent",
441 phase="cleanup-post-delete-absence",
442 outcome=absence,
443 )
444 entry = {
445 "region": region,
446 "name": name,
447 "arn": identity["arn"],
448 "creation_time": identity["creation_time"],
449 "stack_id": record["stack_id"],
450 "source_logical_id": record["source_logical_id"],
451 "source_resource_type": record["source_resource_type"],
452 "authority_phase": record["authority_phase"],
453 "atomic_resource_tag_condition": True,
454 "absence_observations": absence["attempt_count"],
455 "deleted": True,
456 "adopted": bool(record.get("adopted_generations")),
457 "original_generation_disposition": disposition,
458 }
459 ctx.persist()
460 return ("completed", entry)
463def _validated_log_group_cleanup_records(
464 ctx: RunContext,
465) -> tuple[list[tuple[dict[str, Any], str, str]], str]:
466 """Validate the cleanup preconditions and return records plus the token.
468 Cleanup may only proceed once every exact target stack is absent, because
469 a live stack can legitimately recreate its own log groups. The token is the
470 second half of the deletion authority (paired with the run tag) and is
471 checked for shape before it can reach an IAM condition.
472 """
473 records = ctx.checkpoint.state.get("owned_log_groups", [])
474 if not isinstance(records, list):
475 raise RuntimeError("Checkpoint owned_log_groups must be a list")
476 cleanup_token = str(ctx.checkpoint.state.get("log_group_cleanup_token") or "")
478 if records:
479 stack_absence = _verify_target_stack_absence(ctx)
480 if not stack_absence["all_absent"]:
481 raise RuntimeError("Log-group cleanup requires every exact target stack to be absent")
482 if records and not re.fullmatch(r"[0-9a-f]{32}", cleanup_token):
483 raise RuntimeError("Checkpoint log-group cleanup token is malformed")
485 validated: list[tuple[dict[str, Any], str, str]] = []
486 for raw_record in records:
487 if not isinstance(raw_record, dict):
488 raise RuntimeError("Checkpoint owned_log_groups must contain objects")
489 region, name = _validated_owned_log_group_identity(ctx, raw_record)
490 validated.append((raw_record, region, name))
491 return validated, cleanup_token
494def _cleanup_owned_log_groups(ctx: RunContext) -> dict[str, Any]:
495 """Converge every checkpointed log group to stable absence.
497 Each record is processed independently by :func:`_converge_one_log_group`,
498 so one blocked generation never strands the rest. Teardown-time Lambda
499 invocations recreate their own groups after tagging, so untagged same-name
500 regenerations are adopted (re-tagged under this run's proven stack-absence
501 authority) and deleted; generations carrying another owner's markers stay
502 strictly preserved. Bounded extra sweeps absorb log deliveries that land
503 mid-cleanup, and only retryable blockers earn another sweep.
505 The delegated cleanup helper stack is always torn down, even when
506 convergence fails, and both independent failures are reported together.
507 """
508 results: list[dict[str, Any]] = []
509 deleter = TagConditionedLogDeleter(ctx)
510 helper_cleanup: dict[str, Any] = {"needed": False, "deleted": True}
511 cleanup_error: Exception | None = None
512 helper_error: Exception | None = None
513 authority_tags: dict[str, str] = {}
514 try:
515 validated, cleanup_token = _validated_log_group_cleanup_records(ctx)
516 authority_tags = {
517 _RUN_STACK_TAG: ctx.settings.run_id,
518 _LOG_CLEANUP_TOKEN_TAG: cleanup_token,
519 }
521 completed: dict[tuple[str, str], dict[str, Any]] = {}
522 blocked: dict[tuple[str, str], dict[str, Any]] = {}
523 for sweep in range(1, _LOG_GROUP_CLEANUP_MAX_PASSES + 1):
524 if sweep > 1:
525 # Absorb straggling teardown log deliveries before re-observing.
526 time.sleep(_LOG_GROUP_OBSERVATION_POLL_SECONDS * sweep)
527 blocked.clear()
528 for record, region, name in validated:
529 key = (region, name)
530 if key in completed:
531 continue
532 outcome, entry = _converge_one_log_group(
533 ctx,
534 record,
535 region,
536 name,
537 authority_tags=authority_tags,
538 cleanup_token=cleanup_token,
539 deleter=deleter,
540 )
541 if outcome == "completed":
542 completed[key] = entry
543 else:
544 blocked[key] = entry
546 if not blocked or not any(entry["retryable"] for entry in blocked.values()):
547 break
549 results = [*completed.values(), *blocked.values()]
550 if blocked:
551 summary = ", ".join(
552 f"{region}:{name} ({entry['original_generation_disposition']['status']})"
553 for (region, name), entry in sorted(blocked.items())
554 )
555 raise RuntimeError(f"Log-group cleanup could not converge for: {summary}")
556 except Exception as exc: # noqa: BLE001 - attach helper cleanup and partial evidence
557 cleanup_error = exc
558 finally:
559 try:
560 helper_cleanup = _delete_log_cleanup_helper(ctx)
561 except Exception as exc: # noqa: BLE001 - preserve both independent failures
562 helper_error = exc
564 errors = []
565 if cleanup_error is not None:
566 errors.append(
567 {"phase": "log-groups", "error": f"{type(cleanup_error).__name__}: {cleanup_error}"}
568 )
569 if helper_error is not None:
570 errors.append(
571 {
572 "phase": "cleanup-helper",
573 "error": f"{type(helper_error).__name__}: {helper_error}",
574 }
575 )
576 details = {
577 "log_groups": results,
578 "authorization": deleter.authorization,
579 "helper_stack_cleanup": helper_cleanup,
580 "errors": errors,
581 }
582 ctx.checkpoint.state["last_log_group_cleanup"] = copy.deepcopy(details)
583 ctx.persist()
584 if errors:
585 message = "Retained CloudWatch log cleanup failed: " + json.dumps(errors, sort_keys=True)
586 primary_error = cleanup_error if cleanup_error is not None else helper_error
587 raise _LogGroupCleanupError(message, details) from primary_error
588 return details