#!/usr/bin/env python3 """Produce a chronological timeline report from one experiment's JSONL streams. Usage: ./scripts/analyze.py results/experiment-1-saturated-deploy/ Reads every *.jsonl file in the directory, normalizes timestamps, merges into a single timeline, and produces: - A narrative timeline (events in order, abbreviated) - Per-revision summaries (how many enable/disable/exit events each task-def revision saw) - Convergence markers (when deployments dropped to 1 entry, etc.) - Task lifecycle summary (when each task was created, stopped, why) - DEPLOYMENT_BLOCKED tally Writes `timeline.txt` and `summary.md` into the same directory. Pure stdlib so you don't need to install anything to run it. """ from __future__ import annotations import argparse import collections import json import pathlib import re import sys from datetime import UTC, datetime def parse_iso(s: str | None) -> datetime | None: if not s: return None # Handle both '2026-04-17T15:59:06.123Z' and POSIX floats. if isinstance(s, int | float): return datetime.fromtimestamp(float(s), tz=UTC) try: # Python needs +00:00 instead of Z, and fractional seconds have # varying widths. return datetime.fromisoformat(s.replace("Z", "+00:00")) except (ValueError, TypeError): return None def short_arn(arn: str | None) -> str: if not arn: return "?" # task-definition/claude-ecs-task-protection-experiment:5 -> :5 m = re.search(r":(\d+)$", arn) if m: return f"rev{m.group(1)}" # task/cluster/ -> short id parts = arn.rsplit("/", 1) return parts[-1][:8] if parts else arn[:8] def load_jsonl(path: pathlib.Path) -> list[dict]: if not path.exists(): return [] records = [] for line in path.read_text().splitlines(): line = line.strip() if not line: continue try: records.append(json.loads(line)) except json.JSONDecodeError: continue return records def ts_of(record: dict) -> datetime | None: for key in ("ts", "snapshot_ts", "first_seen_ts", "createdAt"): if key in record: return parse_iso(record[key]) return None def analyze(results_dir: pathlib.Path) -> None: worker_logs = load_jsonl(results_dir / "worker-logs.jsonl") service_state = load_jsonl(results_dir / "service-state.jsonl") service_events = load_jsonl(results_dir / "service-events.jsonl") task_state = load_jsonl(results_dir / "task-state.jsonl") deployments = load_jsonl(results_dir / "service-deployments.jsonl") task_defs = load_jsonl(results_dir / "task-definitions.jsonl") timeline: list[tuple[datetime, str, str]] = [] for rec in worker_logs: ts = ts_of(rec) if not ts: continue event = rec.get("event", "?") rev = rec.get("revision") or "?" task = short_arn(rec.get("task_arn")) extras = [] if event == "enable_failed": extras.append(f"reason={rec.get('reason')}") if event == "enable_success": extras.append(f"iter={rec.get('iteration')}") if "elapsed_ms" in rec: extras.append(f"{rec['elapsed_ms']}ms") detail = f"task={task} rev={rev} {event}" + (f" [{', '.join(extras)}]" if extras else "") timeline.append((ts, "worker", detail)) for rec in service_events: ts = parse_iso(rec.get("createdAt")) if not ts: continue message = rec.get("message", "") timeline.append((ts, "svc-event", message[:180])) for rec in deployments: ts = ts_of(rec) if not ts: continue status = rec.get("status", "?") target = short_arn(rec.get("targetServiceRevisionArn")) rollout = rec.get("rollout") or {} timeline.append( (ts, "deployment", f"{status} target={target} rollout={json.dumps(rollout)[:120]}") ) # Dedup consecutive identical deployment records so we don't spam. timeline.sort(key=lambda r: r[0]) dedup_timeline: list[tuple[datetime, str, str]] = [] last_deploy: tuple[str, str] | None = None for ts, stream, detail in timeline: if stream == "deployment": key = (stream, detail) if key == last_deploy: continue last_deploy = key else: last_deploy = None dedup_timeline.append((ts, stream, detail)) # Write timeline.txt timeline_path = results_dir / "timeline.txt" with timeline_path.open("w") as f: for ts, stream, detail in dedup_timeline: f.write(f"{ts.isoformat()} [{stream:11}] {detail}\n") # Summary stats worker_events_by_revision: dict[str, collections.Counter] = collections.defaultdict( collections.Counter ) for rec in worker_logs: rev = rec.get("revision") or "?" event = rec.get("event") or "?" worker_events_by_revision[str(rev)][event] += 1 deployment_blocked_events = [ rec for rec in worker_logs if rec.get("event") == "enable_failed" and rec.get("reason") == "DEPLOYMENT_BLOCKED" ] sigterm_events = [rec for rec in worker_logs if rec.get("event") == "sigterm_received"] # Per-task lifecycle (from task_state, one record per task, latest status) task_latest: dict[str, dict] = {} for rec in task_state: arn = rec.get("taskArn") if not arn: continue existing = task_latest.get(arn) if existing is None: task_latest[arn] = rec continue new_ts = ts_of(rec) existing_ts = ts_of(existing) if new_ts and existing_ts and new_ts > existing_ts: task_latest[arn] = rec # Service state first vs last service_first = service_state[0] if service_state else None service_last = service_state[-1] if service_state else None # Write summary.md summary_path = results_dir / "summary.md" lines: list[str] = [] lines.append(f"# Summary — {results_dir.name}") lines.append("") lines.append("## Service state") if service_first and service_last: lines.append( f"- First snapshot: `{service_first.get('ts')}` " f"running={service_first.get('runningCount')} " f"desired={service_first.get('desiredCount')} " f"deps={len(service_first.get('deployments') or [])}" ) lines.append( f"- Last snapshot: `{service_last.get('ts')}` " f"running={service_last.get('runningCount')} " f"desired={service_last.get('desiredCount')} " f"deps={len(service_last.get('deployments') or [])}" ) lines.append("") lines.append("## Task definitions seen") for rec in task_defs: rev = rec.get("revision") env = {} try: env = {e["name"]: e["value"] for e in rec["containerDefinitions"][0].get("environment", [])} except (KeyError, IndexError, TypeError): pass env_str = ", ".join(f"{k}={v}" for k, v in env.items() if k in {"WORK_DURATION_MS", "IDLE_MS", "EXPIRES_MINUTES"}) lines.append(f"- rev{rev} first seen `{rec.get('first_seen_ts')}`: {env_str}") lines.append("") lines.append("## Worker events by revision") for rev, counts in sorted(worker_events_by_revision.items()): lines.append(f"- rev{rev}: " + ", ".join(f"{k}={v}" for k, v in sorted(counts.items()))) lines.append("") lines.append("## DEPLOYMENT_BLOCKED events") lines.append(f"- Total: {len(deployment_blocked_events)}") by_rev: collections.Counter = collections.Counter() for rec in deployment_blocked_events: by_rev[str(rec.get("revision"))] += 1 for rev, count in sorted(by_rev.items()): lines.append(f" - rev{rev}: {count}") lines.append("") lines.append("## SIGTERM events") lines.append(f"- Total: {len(sigterm_events)}") by_rev = collections.Counter() for rec in sigterm_events: by_rev[str(rec.get("revision"))] += 1 for rev, count in sorted(by_rev.items()): lines.append(f" - rev{rev}: {count}") lines.append("") lines.append("## Tasks observed") for arn, rec in sorted(task_latest.items(), key=lambda kv: kv[1].get("createdAt") or ""): stop_code = rec.get("stopCode") stop_reason = rec.get("stoppedReason") stop_info = "" if stop_code or stop_reason: stop_info = f" | stopped: {stop_code} — {stop_reason}" lines.append( f"- {short_arn(arn)} ({short_arn(rec.get('taskDefinitionArn'))}): " f"{rec.get('lastStatus')}/{rec.get('desiredStatus')}{stop_info}" ) lines.append("") lines.append("## First 20 ECS service events") sorted_events = sorted( (rec for rec in service_events if rec.get("createdAt")), key=lambda r: r["createdAt"], ) for rec in sorted_events[:20]: lines.append(f"- `{rec['createdAt']}` {rec.get('message', '')[:140]}") summary_path.write_text("\n".join(lines) + "\n") print(f"wrote {timeline_path}") print(f"wrote {summary_path}") print( f"timeline: {len(dedup_timeline)} events, " f"workers: {len(worker_logs)}, tasks seen: {len(task_latest)}, " f"DEPLOYMENT_BLOCKED: {len(deployment_blocked_events)}, " f"SIGTERM: {len(sigterm_events)}" ) def main() -> int: parser = argparse.ArgumentParser(description=__doc__) parser.add_argument("results_dir", type=pathlib.Path) args = parser.parse_args() if not args.results_dir.is_dir(): sys.stderr.write(f"not a directory: {args.results_dir}\n") return 1 analyze(args.results_dir) return 0 if __name__ == "__main__": sys.exit(main())