#!/usr/bin/env python3 """ cleanup_event_backlog.py — safe backlog cleanup for /opt/homelab/events/. Removes ONLY processed noise events (service_healthy, node_health by default) that are both: 1. Strictly older than the observer's per-node checkpoint timestamp (observer_checkpoint.json) — i.e. guaranteed already processed into world_state, per the same ordering guarantee observer.py itself relies on (oldest-first, monotonically-advancing checkpoint). 2. Older than a minimum safety age (default 3600s / 1h) — defense in depth against clock skew or a checkpoint write racing this script, mirroring the dual-condition pattern already used in node_agent.py's _cleanup_control_plane_fs. Everything else — healthcheck_failed, containers_not_running, incidents, ha_* events, node liveness transitions, disk/memory/cpu pressure, and any event newer than the checkpoint — is never touched. Defaults to DRY RUN: prints counts per node/type and does not delete anything. Pass --apply to actually delete. Usage: python3 cleanup_event_backlog.py # dry run, defaults python3 cleanup_event_backlog.py --apply # actually delete python3 cleanup_event_backlog.py --types service_healthy,node_health python3 cleanup_event_backlog.py --min-age-secs 86400 # extra-cautious """ import argparse import json import os import re import sys import time from pathlib import Path RUNTIME_PATH = Path(os.getenv("RUNTIME_PATH", "/opt/homelab")) EVENTS_DIR = RUNTIME_PATH / "events" CHECKPOINT_FILE = RUNTIME_PATH / "state" / "observer_checkpoint.json" # Matches observer.py's _EVENT_TS_RE: evt----.json _EVENT_TS_RE = re.compile(r"-(\d{9,11})-") DEFAULT_TYPES = {"service_healthy", "node_health"} def _ts_from_name(name: str): m = _EVENT_TS_RE.search(Path(name).stem) return int(m.group(1)) if m else None def _type_from_name(name: str): """Best-effort event type extraction from evt----.json. Mirrors the informal convention node-agent/observer use when building filenames; falls back to 'unknown' rather than guessing wrong. """ stem = Path(name).stem m = re.match(r"evt-.+?-\d{9,11}-(.+)$", stem) if not m: return "unknown" rest = m.group(1) # rest is "-"; type itself may contain underscores but not # dashes in this codebase's event vocabulary, so split on first dash after # the known type token is ambiguous — instead check against DEFAULT_TYPES # and a couple of other well-known multi-word-free types directly. for t in KNOWN_TYPES: if rest == t or rest.startswith(t + "-"): return t return rest.split("-")[0] # Every event type this repo's emitters are known to produce (node_agent.py, # ha_diag/event_emitter.py, observer.py's own _emit_node_transition), longest # first so startswith() matching in _type_from_name doesn't shadow a longer # type with a shorter prefix (e.g. "node_health" vs "node_offline" share no # prefix, but keep this list explicit rather than relying on luck). KNOWN_TYPES = sorted( { "service_healthy", "service_unhealthy", "service_recovered", "node_health", "node_online", "node_offline", "node_stale", "healthcheck_failed", "containers_not_running", "container_restarting", "container_state_unexpected", "disk_pressure", "high_memory", "high_cpu", "mqtt_unreachable", "ha_websocket_dead", "ha_websocket_recovered", "ha_integration_failed", "ha_entity_unavailable_long", "ha_automation_failing", "ha_update_available", "ha_recorder_lag", "ha_system_health_degraded", "deployment_started", "deployment_completed", "deployment_failed", "remediation_started", "remediation_completed", }, key=len, reverse=True, ) def load_checkpoints(checkpoint_file: Path = CHECKPOINT_FILE): if not checkpoint_file.exists(): print(f"WARNING: no checkpoint file at {checkpoint_file} — refusing to " f"guess; nothing is safe to delete without it.", file=sys.stderr) return {} data = json.loads(checkpoint_file.read_text()) raw = data.get("node_checkpoints", {}) out = {} for node, val in raw.items(): if isinstance(val, bool): continue if isinstance(val, (int, float)): out[node] = int(val) return out def scan(types_to_clean, min_age_secs, checkpoints, now, events_dir: Path = EVENTS_DIR): """Return (candidates: list[Path], counts: dict[node][type] -> (delete, keep)).""" candidates = [] counts = {} if not events_dir.exists(): return candidates, counts for node_dir in sorted(p for p in events_dir.iterdir() if p.is_dir()): node = node_dir.name checkpoint = checkpoints.get(node) counts.setdefault(node, {}) for f in node_dir.glob("*.json"): etype = _type_from_name(f.name) ts = _ts_from_name(f.name) if ts is None: try: ts = int(f.stat().st_mtime) except OSError: continue bucket = counts[node].setdefault(etype, {"delete": 0, "keep": 0}) if etype not in types_to_clean: bucket["keep"] += 1 continue if checkpoint is None: # No checkpoint for this node dir at all (e.g. a ghost/foreign # dir) — never guess; keep everything. bucket["keep"] += 1 continue if ts >= checkpoint: # Not guaranteed processed yet — never delete. bucket["keep"] += 1 continue if now - ts < min_age_secs: # Within the safety buffer even though past checkpoint. bucket["keep"] += 1 continue bucket["delete"] += 1 candidates.append(f) return candidates, counts def print_report(counts, min_age_secs, checkpoints): total_delete = 0 total_keep = 0 for node in sorted(counts): cp = checkpoints.get(node) cp_str = f"checkpoint={cp}" if cp is not None else "NO CHECKPOINT (nothing eligible)" print(f"\n=== {node} ({cp_str}) ===") for etype, c in sorted(counts[node].items(), key=lambda kv: -sum(kv[1].values())): if c["delete"] or c["keep"]: print(f" {etype:32s} delete={c['delete']:>7d} keep={c['keep']:>7d}") total_delete += c["delete"] total_keep += c["keep"] print(f"\nTOTAL: delete={total_delete} keep={total_keep}") print(f"(min_age_secs={min_age_secs}, i.e. only files older than checkpoint AND " f"more than {min_age_secs}s old are ever deleted)") def main(): ap = argparse.ArgumentParser(description=__doc__) ap.add_argument("--apply", action="store_true", help="Actually delete. Without this flag, dry-run only.") ap.add_argument("--types", default=",".join(sorted(DEFAULT_TYPES)), help=f"Comma-separated event types to clean (default: {DEFAULT_TYPES})") ap.add_argument("--min-age-secs", type=int, default=3600, help="Extra safety buffer past the checkpoint (default 3600s = 1h)") args = ap.parse_args() types_to_clean = {t.strip() for t in args.types.split(",") if t.strip()} checkpoints = load_checkpoints(CHECKPOINT_FILE) now = int(time.time()) candidates, counts = scan(types_to_clean, args.min_age_secs, checkpoints, now, EVENTS_DIR) print_report(counts, args.min_age_secs, checkpoints) if not args.apply: print(f"\nDRY RUN — no files deleted. Re-run with --apply to delete " f"the {len(candidates)} files listed above as 'delete'.") return deleted = 0 for f in candidates: try: f.unlink(missing_ok=True) deleted += 1 except Exception as exc: print(f"Failed to remove {f}: {exc}", file=sys.stderr) print(f"\nDeleted {deleted}/{len(candidates)} files.") if __name__ == "__main__": main()