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

1"""Delete exactly run-owned CloudWatch log groups.""" 

2 

3from __future__ import annotations 

4 

5import copy 

6import json 

7import re 

8import time 

9from collections.abc import Mapping 

10from typing import Any 

11 

12from botocore.exceptions import ClientError 

13 

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) 

38 

39 

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. 

47 

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 

66 

67 

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. 

79 

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 

126 

127 

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. 

141 

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 } 

166 

167 

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. 

179 

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: 

183 

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) 

198 

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 ) 

234 

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 ) 

322 

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 ) 

363 

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 ) 

391 

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 ) 

435 

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) 

461 

462 

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. 

467 

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

477 

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

484 

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 

492 

493 

494def _cleanup_owned_log_groups(ctx: RunContext) -> dict[str, Any]: 

495 """Converge every checkpointed log group to stable absence. 

496 

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. 

504 

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 } 

520 

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 

545 

546 if not blocked or not any(entry["retryable"] for entry in blocked.values()): 

547 break 

548 

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 

563 

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