From caca9d2bd087644b4c9ca3e39af2f8954fbc6f15 Mon Sep 17 00:00:00 2001 From: Lucky Abolorunke <62015433+Oneimu@users.noreply.github.com> Date: Wed, 30 Sep 2026 09:05:34 -0700 Subject: [PATCH 1/4] benchmarking/analysis: aggregate suspend/resume phase breakdown log records --- benchmarking/README.md | 17 + benchmarking/analysis/README.md | 81 +++++ benchmarking/analysis/collect_logs.sh | 64 ++++ benchmarking/analysis/phase_report.py | 349 +++++++++++++++++++++ benchmarking/analysis/test_phase_report.py | 172 ++++++++++ 5 files changed, 683 insertions(+) create mode 100644 benchmarking/analysis/README.md create mode 100755 benchmarking/analysis/collect_logs.sh create mode 100755 benchmarking/analysis/phase_report.py create mode 100644 benchmarking/analysis/test_phase_report.py diff --git a/benchmarking/README.md b/benchmarking/README.md index 4092f5ea2a..f5df44603a 100644 --- a/benchmarking/README.md +++ b/benchmarking/README.md @@ -207,6 +207,16 @@ The Kubernetes API is not required. If it is unreachable, or discovery was skipped, the affected fields are written as `null` and the run still succeeds. A `null` means the value was not measured. It never means zero. +## Suspend/resume phase breakdown + +`SuspendActor` and `ResumeActor` latency can be attributed to the phases +inside them (object-storage transfer, hypervisor snapshot/restore, rootfs +assembly, ...) from the structured log records atelet and ateom-microvm +emit. `analysis/collect_logs.sh` dumps the node logs of a run and +`analysis/phase_report.py` aggregates them into per-phase percentiles and +waterfalls of the slowest operations. See +[analysis/README.md](analysis/README.md). + ## Optional: Prometheus + Grafana Locust provides graphs, statistics, etc. via the UI. However, you @@ -241,3 +251,10 @@ repository root: ```bash python3 -m unittest discover -s benchmarking/locust/unit_tests ``` + +`analysis/test_phase_report.py` covers the phase-log report's parser and +aggregation: + +```bash +python3 -m unittest discover -s benchmarking/analysis +``` diff --git a/benchmarking/analysis/README.md b/benchmarking/analysis/README.md new file mode 100644 index 0000000000..9772da6353 --- /dev/null +++ b/benchmarking/analysis/README.md @@ -0,0 +1,81 @@ +# Suspend/resume phase analysis + +This directory turns the node side's developer-facing timing logs into the +percentiles a benchmark run needs, without making the phases metric API. The +`ate.actor.{restore,checkpoint}.duration` histograms stop at +`ateom_restore` / `ateom_checkpoint` / `persist` / `download` on purpose: +phases inside those are implementation details of a runtime, and metrics are +a contract. The finer breakdown is emitted as structured log records instead, +and this is the consumer that aggregates them. + +Two joinable JSON log records feed it, each written by both layers: + +| record (`msg`) | emitter | keys | +|---|---|---| +| `Restore timing breakdown` | atelet and ateom-microvm | `ate.actor.restore.duration.` / `ateom.actor.restore.duration.` | +| `Checkpoint timing breakdown` | atelet and ateom-microvm | `ate.actor.checkpoint.duration.` / `ateom.actor.checkpoint.duration.` | + +Every record carries the full actor identity (`ate.actor.uid`, name, +atespace, template) and the snapshot scope, which the histograms are barred +from, so records join per actor and per operation. Durations are float +seconds, the histograms' unit. A phase that never ran is absent, not zero. +A failed atelet operation still writes its record, marked with `error.type` +(the gRPC code); the report excludes those from the percentiles. + +## Making a measurement + +1. Run a suspend/resume-heavy load (e.g. the `glutton` or `sweperf` user + class; see `../README.md`). +2. Dump the node logs for the run's window while the pods still exist. + `automation/orchestrator.py` deletes the workload and ate-system pods right + after each test, and does not call this script yet, so run it before that + teardown (or against a cluster you manage yourself): + + ```bash + ./collect_logs.sh --dest /tmp/run1 --since 30m + ``` + + `--namespace` (default `ate-system`) selects the atelet pods and + `--worker-namespace` (default `benchmark-workloads`) the worker pods; the + ateom records land in the worker pod's stdout. + +3. Aggregate: + + ```bash + python3 phase_report.py /tmp/run1/*.log --csv /tmp/run1/report + ``` + +## Reading the report + +**Phase percentiles.** Per layer, operation, scope and phase: count, p50, +p90, p95 and max. The atelet rows split a checkpoint between +`sandbox_assets`, `ateom_checkpoint` and `persist`, and a restore between +`volume_mount`, `manifest_fetch`, `sandbox_assets`, `download`, `oci_unpack` +and `ateom_restore`. The ateom rows split the `ateom_*` phase further: +`prep` / `pause` / `snapshot` / `durable_dir` / `rootfs_upper` / `teardown` +for a checkpoint, `prep` / `bundles` / `upper_join` / `lowers` / `tap` / +`vmm_launch` / `vm_restore` / `resume` / `wakeup_probe` for a restore. + +Concurrency matters when reading them: the atelet restore phases overlap +(the download runs alongside the asset fetch and OCI unpack), the three +ateom checkpoint captures run concurrently on the paused guest (the paused +window costs their max), and the ateom restore phases are sequential. For +the two checkpoint layers the report derives an `unattributed` row: the total +minus what the logged phases account for, counting the concurrent captures +once. It is the time the instrumentation does not yet name. + +**Waterfalls.** The slowest operations, with the ateom record nested under +the atelet `ateom_*` phase and that phase's gap to the ateom total (RPC and +queueing between the layers). Tail outliers that blow up in one phase every +time are systematic; different phases each time are environmental. + +`--csv` also writes `phase_percentiles.csv` for run-over-run comparison. + +The parser is deliberately tolerant: it scans any line for a JSON object and +matches on `msg`, so raw `kubectl logs` dumps (even with prefixes) work. + +## Tests + +```bash +python3 -m unittest discover -s benchmarking/analysis +``` diff --git a/benchmarking/analysis/collect_logs.sh b/benchmarking/analysis/collect_logs.sh new file mode 100755 index 0000000000..feedf7490f --- /dev/null +++ b/benchmarking/analysis/collect_logs.sh @@ -0,0 +1,64 @@ +#!/usr/bin/env bash +# Copyright 2026 Google LLC +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +# Dump the node-side logs a benchmark run needs for phase_report.py: every +# atelet pod (the atelet-side timing breakdowns) and every worker pod +# (ateom-microvm writes its records to the worker pod's stdout). Point +# --since at the run's start so the report covers exactly one run. +# +# Usage: collect_logs.sh --dest DIR [--since 30m] [--namespace ate-system] +# [--worker-namespace benchmark-workloads] + +set -euo pipefail + +DEST="" +SINCE="1h" +NAMESPACE="ate-system" +WORKER_NAMESPACE="benchmark-workloads" + +while [[ $# -gt 0 ]]; do + case "$1" in + --dest) DEST="$2"; shift 2 ;; + --since) SINCE="$2"; shift 2 ;; + --namespace) NAMESPACE="$2"; shift 2 ;; + --worker-namespace) WORKER_NAMESPACE="$2"; shift 2 ;; + *) echo "unknown flag: $1" >&2; exit 2 ;; + esac +done + +if [[ -z "${DEST}" ]]; then + echo "usage: $0 --dest DIR [--since 30m] [--namespace ate-system] [--worker-namespace benchmark-workloads]" >&2 + exit 2 +fi +mkdir -p "${DEST}" + +collect() { + local ns="$1" selector="$2" + local pods + pods=$(kubectl get pods -n "${ns}" -l "${selector}" -o name) + for pod in ${pods}; do + local name="${pod#pod/}" + echo "collecting ${ns}/${name} (since ${SINCE})" + kubectl logs -n "${ns}" "${name}" --since="${SINCE}" --timestamps=false \ + > "${DEST}/${ns}-${name}.log" + done +} + +collect "${NAMESPACE}" "app=atelet" +# Worker pods carry the pool label whatever the pool's name is. +collect "${WORKER_NAMESPACE}" "ate.dev/worker-pool" + +echo "logs in ${DEST}; next:" +echo " python3 benchmarking/analysis/phase_report.py ${DEST}/*.log" diff --git a/benchmarking/analysis/phase_report.py b/benchmarking/analysis/phase_report.py new file mode 100755 index 0000000000..8da99494d6 --- /dev/null +++ b/benchmarking/analysis/phase_report.py @@ -0,0 +1,349 @@ +#!/usr/bin/env python3 +# Copyright 2026 Google LLC +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +"""Aggregate the suspend/resume phase-breakdown log records into a report. + +Reads JSON-lines logs (kubectl logs dumps of the atelet and worker pods; any +non-JSON or unrelated lines are skipped) and aggregates the two joinable +record kinds the node side emits, each written by both layers: + + - "Restore timing breakdown" atelet (ate.actor.restore.duration.*) + and ateom (ateom.actor.restore.duration.*) + - "Checkpoint timing breakdown" atelet (ate.actor.checkpoint.duration.*) + and ateom (ateom.actor.checkpoint.duration.*) + +The report answers "where does the SuspendActor / ResumeActor time go": +per-phase percentiles at each layer, and the slowest operations as nested +waterfalls (the ateom record joined under the atelet record of the same +actor). + +Usage: + phase_report.py run-logs/*.log [--csv DEST_DIR] [--slowest N] + +The records are developer-facing logs, not metric API; this reader is the +consumer that makes them percentiles. See benchmarking/analysis/README.md. +""" + +from __future__ import annotations + +import argparse +import csv +import json +import re +import statistics +import sys +from collections import defaultdict +from dataclasses import dataclass, field +from datetime import datetime +from pathlib import Path + +# The duration-key prefixes. Source of truth: cmd/atelet/metrics.go +# (restoreDurationMetric, checkpointDurationMetric) and +# cmd/ateom-microvm/phaselog.go. +BREAKDOWN_PREFIXES = { + ("atelet", "restore"): "ate.actor.restore.duration.", + ("atelet", "checkpoint"): "ate.actor.checkpoint.duration.", + ("ateom", "restore"): "ateom.actor.restore.duration.", + ("ateom", "checkpoint"): "ateom.actor.checkpoint.duration.", +} +BREAKDOWN_MSGS = {"Restore timing breakdown", "Checkpoint timing breakdown"} + +ACTOR_UID_KEY = "ate.actor.uid" +SCOPE_KEY = "ate.snapshot.scope" +PHASE_KEY = "ate.snapshot.phase" +KIND_KEY = "ate.snapshot.kind" +TEMPLATE_KEY = "ate.template.name" +ERROR_TYPE_KEY = "error.type" + +# The phase every record reports, and the one the report derives: the part of +# the total no logged phase accounts for. +TOTAL = "total" +UNATTRIBUTED = "unattributed" + +# Sequential order of the known phases, for display. Phases absent from a +# record simply don't print; unknown phases print after, in input order. +PHASE_ORDER = [ + # atelet restore + "volume_mount", "manifest_fetch", "sandbox_assets", "download", + "oci_unpack", "ateom_restore", + # ateom restore + "prep", "bundles", "upper_join", "lowers", "tap", "vmm_launch", + "vm_restore", "resume", "wakeup_probe", + # atelet + ateom checkpoint + "pause", "snapshot", "durable_dir", "rootfs_upper", + "ateom_checkpoint", "teardown", "persist", + UNATTRIBUTED, TOTAL, +] +_PHASE_RANK = {name: i for i, name in enumerate(PHASE_ORDER)} + +# Which atelet phase wraps the ateom record. +INNER_PHASE = {"restore": "ateom_restore", "checkpoint": "ateom_checkpoint"} + +# The ateom checkpoint captures that run concurrently on the paused guest: the +# paused window costs their max, not their sum. +CONCURRENT_CAPTURES = ("snapshot", "durable_dir", "rootfs_upper") + +# Slack when matching the ateom record to the atelet record of the same +# operation: the ateom one is emitted a hair before the atelet one. +JOIN_SLACK_S = 1.0 + + +@dataclass +class Breakdown: + """One parsed timing-breakdown record.""" + source: str # atelet | ateom + op: str # restore | checkpoint + time: str + ts: float | None # epoch seconds parsed from time, None when unparseable + actor_uid: str + template: str + scope: str + kind: str + phases: dict[str, float] # phase name -> seconds + failed: bool + + +@dataclass +class Parsed: + breakdowns: list[Breakdown] = field(default_factory=list) + lines_seen: int = 0 + lines_matched: int = 0 + + +_TIME_RE = re.compile(r"^(.*T\d\d:\d\d:\d\d)(\.\d+)?(Z|[+-]\d\d:?\d\d)?$") + + +def parse_time(s: str) -> float | None: + """RFC 3339 -> epoch seconds. slog writes nanoseconds, which fromisoformat + rejects, so the fraction is trimmed to microseconds. None when unparseable + (a record then still counts, it just cannot be joined by time).""" + m = _TIME_RE.match(s.strip()) if s else None + if not m: + return None + base, frac, tz = m.groups() + frac = (frac or "")[:7] + tz = "+00:00" if tz in (None, "Z") else tz + try: + return datetime.fromisoformat(f"{base}{frac}{tz}").timestamp() + except ValueError: + return None + + +def unattributed(source: str, op: str, phases: dict[str, float]) -> float | None: + """The total minus what the logged phases account for, where that is + well-defined: the atelet checkpoint phases are sequential; the ateom + checkpoint is prep, pause, the concurrent captures (max), then teardown. + The atelet restore phases overlap by design (download runs alongside the + asset fetch and OCI unpack) and the ateom restore phases partition their + total by construction, so neither gets a residual.""" + total = phases.get(TOTAL) + if total is None: + return None + if (source, op) == ("atelet", "checkpoint"): + spent = sum(phases.get(p, 0.0) for p in ("sandbox_assets", "ateom_checkpoint", "persist")) + elif (source, op) == ("ateom", "checkpoint"): + spent = (phases.get("prep", 0.0) + phases.get("pause", 0.0) + + max(phases.get(p, 0.0) for p in CONCURRENT_CAPTURES) + + phases.get("teardown", 0.0)) + else: + return None + return max(total - spent, 0.0) + + +def parse_line(obj: dict, out: Parsed) -> None: + msg = obj.get("msg", "") + if msg in BREAKDOWN_MSGS: + for (source, op), prefix in BREAKDOWN_PREFIXES.items(): + phases = { + k[len(prefix):]: float(v) + for k, v in obj.items() + if k.startswith(prefix) + } + if not phases: + continue + residual = unattributed(source, op, phases) + if residual is not None: + phases[UNATTRIBUTED] = residual + out.breakdowns.append(Breakdown( + source=source, + op=op, + time=obj.get("time", ""), + ts=parse_time(obj.get("time", "")), + actor_uid=obj.get(ACTOR_UID_KEY, ""), + template=obj.get(TEMPLATE_KEY, ""), + scope=obj.get(SCOPE_KEY, ""), + kind=obj.get(KIND_KEY, ""), + phases=phases, + failed=ERROR_TYPE_KEY in obj, + )) + out.lines_matched += 1 + return + + +def parse_files(paths: list[str]) -> Parsed: + out = Parsed() + for path in paths: + f = sys.stdin if path == "-" else open(path, encoding="utf-8", errors="replace") + with f: + for line in f: + # kubectl log dumps may prefix each line (pod name, timestamp); + # recover the JSON object from the first brace. + brace = line.find("{") + if brace < 0: + continue + out.lines_seen += 1 + try: + obj = json.loads(line[brace:]) + except json.JSONDecodeError: + continue + if isinstance(obj, dict): + parse_line(obj, out) + return out + + +def percentile(values: list[float], q: float) -> float: + if not values: + return 0.0 + if len(values) == 1: + return values[0] + return statistics.quantiles(values, n=100, method="inclusive")[int(q) - 1] + + +def phase_sort_key(name: str) -> tuple[int, str]: + return (_PHASE_RANK.get(name, len(PHASE_ORDER)), name) + + +def fmt_s(seconds: float) -> str: + return f"{seconds * 1000:8.1f}" + + +def report_phases(breakdowns: list[Breakdown], writer) -> list[dict]: + """Per (source, op, scope, phase) percentiles. Returns the rows for CSV.""" + groups: dict[tuple, list[float]] = defaultdict(list) + for b in breakdowns: + if b.failed: + continue + for name, seconds in b.phases.items(): + groups[(b.source, b.op, b.scope or "-", name)].append(seconds) + + rows = [] + writer("== Phase percentiles (ms) ==") + writer(f"{'layer':7} {'op':11} {'scope':15} {'phase':16} {'n':>5} " + f"{'p50':>8} {'p90':>8} {'p95':>8} {'max':>8}") + for key in sorted(groups, key=lambda k: (k[0], k[1], k[2], phase_sort_key(k[3]))): + vals = sorted(groups[key]) + source, op, scope, name = key + row = { + "layer": source, "op": op, "scope": scope, "phase": name, + "count": len(vals), + "p50_ms": percentile(vals, 50) * 1000, + "p90_ms": percentile(vals, 90) * 1000, + "p95_ms": percentile(vals, 95) * 1000, + "max_ms": max(vals) * 1000, + } + rows.append(row) + writer(f"{source:7} {op:11} {scope:15} {name:16} {len(vals):5d} " + f"{fmt_s(percentile(vals, 50))} {fmt_s(percentile(vals, 90))} " + f"{fmt_s(percentile(vals, 95))} {fmt_s(max(vals))}") + failed = sum(1 for b in breakdowns if b.failed) + if failed: + writer(f"(excluded {failed} failed operation records)") + return rows + + +def inner_record(op: Breakdown, by_actor: dict[tuple, list[Breakdown]]) -> Breakdown | None: + """The ateom record of the same actor and op nearest before op's own: one + cycle emits exactly one of each, and the ateom one lands first, inside + the atelet ateom_* phase.""" + candidates = [a for a in by_actor.get((op.actor_uid, op.op), []) + if a.ts is not None and op.ts is not None and a.ts <= op.ts + JOIN_SLACK_S] + return candidates[-1] if candidates else None + + +def report_waterfalls(breakdowns: list[Breakdown], writer, slowest: int) -> None: + """The slowest operations: the ateom record nested under the atelet + ateom_* phase, and the gap between that phase and the ateom total (RPC + and queueing between the layers).""" + by_actor: dict[tuple, list[Breakdown]] = defaultdict(list) + for b in breakdowns: + if b.source == "ateom": + by_actor[(b.actor_uid, b.op)].append(b) + for v in by_actor.values(): + v.sort(key=lambda b: (b.ts or 0.0, b.time)) + + atelet = [b for b in breakdowns if b.source == "atelet" and not b.failed] + atelet.sort(key=lambda b: b.phases.get(TOTAL, 0), reverse=True) + + writer("") + writer(f"== Slowest {slowest} operations (waterfall, ms) ==") + for b in atelet[:slowest]: + inner = inner_record(b, by_actor) + total = b.phases.get(TOTAL, 0) + writer(f"{b.op} actor={b.actor_uid} template={b.template} " + f"scope={b.scope or '-'} kind={b.kind or '-'} " + f"total={total * 1000:.1f}") + for name in sorted(b.phases, key=phase_sort_key): + if name == TOTAL: + continue + writer(f" atelet {name:15} {fmt_s(b.phases[name])}") + if name == INNER_PHASE.get(b.op) and inner: + for iname in sorted(inner.phases, key=phase_sort_key): + if iname == TOTAL: + continue + writer(f" ateom {iname:13} {fmt_s(inner.phases[iname])}") + if TOTAL in inner.phases: + gap = b.phases[name] - inner.phases[TOTAL] + writer(f" (gap) {'rpc/queueing':13} {fmt_s(gap)}") + + +def write_csv(dest: Path, name: str, rows: list[dict]) -> None: + if not rows: + return + dest.mkdir(parents=True, exist_ok=True) + path = dest / name + with open(path, "w", newline="", encoding="utf-8") as f: + w = csv.DictWriter(f, fieldnames=list(rows[0].keys())) + w.writeheader() + w.writerows(rows) + + +def main() -> int: + ap = argparse.ArgumentParser(description=__doc__.splitlines()[0]) + ap.add_argument("logs", nargs="+", help="JSON-lines log files ('-' for stdin)") + ap.add_argument("--csv", type=Path, default=None, + help="also write phase_percentiles.csv here") + ap.add_argument("--slowest", type=int, default=3, + help="number of slowest operations to print as waterfalls") + args = ap.parse_args() + + parsed = parse_files(args.logs) + print(f"parsed {parsed.lines_matched} timing breakdown records " + f"out of {parsed.lines_seen} JSON log lines") + if not parsed.breakdowns: + print("no matching records; are these the atelet and worker pod logs?", file=sys.stderr) + return 1 + + print() + phase_rows = report_phases(parsed.breakdowns, print) + report_waterfalls(parsed.breakdowns, print, args.slowest) + + if args.csv: + write_csv(args.csv, "phase_percentiles.csv", phase_rows) + print(f"\nCSV written to {args.csv}/") + return 0 + + +if __name__ == "__main__": + sys.exit(main()) diff --git a/benchmarking/analysis/test_phase_report.py b/benchmarking/analysis/test_phase_report.py new file mode 100644 index 0000000000..39b6fc89db --- /dev/null +++ b/benchmarking/analysis/test_phase_report.py @@ -0,0 +1,172 @@ +# Copyright 2026 Google LLC +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +"""Parser and aggregation tests for phase_report. No cluster needed: + + python3 -m unittest discover -s benchmarking/analysis +""" + +import io +import json +import os +import sys +import tempfile +import unittest + +sys.path.insert(0, os.path.dirname(os.path.abspath(__file__))) + +import phase_report # noqa: E402 + + +def parse(*records): + out = phase_report.Parsed() + for rec in records: + phase_report.parse_line(json.loads(json.dumps(rec)), out) + return out + + +def render(fn, *args, **kwargs): + buf = io.StringIO() + fn(*args, writer=lambda line: buf.write(line + "\n"), **kwargs) + return buf.getvalue() + + +ATELET_RESTORE = { + "time": "2026-09-23T10:00:05.500000000Z", "level": "INFO", + "msg": "Restore timing breakdown", + "ate.actor.uid": "uid-1", "ate.template.name": "swebench-astropy-7336", + "ate.snapshot.scope": "full", "ate.snapshot.kind": "latest", + "ate.actor.restore.duration.download": 2.4, + "ate.actor.restore.duration.ateom_restore": 1.1, + "ate.actor.restore.duration.total": 3.9, +} + +ATEOM_RESTORE = { + "time": "2026-09-23T10:00:05.400000000Z", "level": "INFO", + "msg": "Restore timing breakdown", + "ate.actor.uid": "uid-1", "ate.template.name": "swebench-astropy-7336", + "ate.snapshot.scope": "full", + "ateom.actor.restore.duration.vm_restore": 0.8, + "ateom.actor.restore.duration.wakeup_probe": 0.2, + "ateom.actor.restore.duration.total": 1.05, +} + +ATELET_CHECKPOINT = { + "time": "2026-09-23T10:01:00.000000000Z", "level": "INFO", + "msg": "Checkpoint timing breakdown", + "ate.actor.uid": "uid-1", "ate.template.name": "swebench-astropy-7336", + "ate.snapshot.scope": "full", "ate.snapshot.kind": "latest", + "ate.actor.checkpoint.duration.sandbox_assets": 0.01, + "ate.actor.checkpoint.duration.ateom_checkpoint": 1.14, + "ate.actor.checkpoint.duration.persist": 4.48, + "ate.actor.checkpoint.duration.total": 6.0, +} + +ATEOM_CHECKPOINT = { + "time": "2026-09-23T10:00:55.000000000Z", "level": "INFO", + "msg": "Checkpoint timing breakdown", + "ate.actor.uid": "uid-1", "ate.template.name": "swebench-astropy-7336", + "ate.snapshot.scope": "full", + "ateom.actor.checkpoint.duration.prep": 0.04, + "ateom.actor.checkpoint.duration.pause": 0.01, + "ateom.actor.checkpoint.duration.snapshot": 0.5, + "ateom.actor.checkpoint.duration.rootfs_upper": 0.9, + "ateom.actor.checkpoint.duration.teardown": 0.2, + "ateom.actor.checkpoint.duration.total": 1.2, +} + +FAILED = { + "time": "2026-09-23T10:02:00Z", "level": "INFO", + "msg": "Restore timing breakdown", + "ate.actor.uid": "uid-2", "ate.snapshot.scope": "full", + "error.type": "DeadlineExceeded", + "ate.actor.restore.duration.download": 30.0, + "ate.actor.restore.duration.total": 30.0, +} + + +class ParseTest(unittest.TestCase): + def test_parses_both_layers_of_the_same_msg(self): + out = parse(ATELET_RESTORE, ATEOM_RESTORE) + self.assertEqual([(b.source, b.op) for b in out.breakdowns], + [("atelet", "restore"), ("ateom", "restore")]) + self.assertEqual(out.breakdowns[0].phases["download"], 2.4) + self.assertEqual(out.breakdowns[1].phases["vm_restore"], 0.8) + self.assertEqual(out.breakdowns[1].actor_uid, "uid-1") + + def test_nanosecond_timestamps_parse(self): + out = parse(ATELET_RESTORE, FAILED) + self.assertIsNotNone(out.breakdowns[0].ts) + self.assertIsNotNone(out.breakdowns[1].ts) + self.assertAlmostEqual(out.breakdowns[0].ts % 1, 0.5, places=6) + self.assertIsNone(phase_report.parse_time("yesterday")) + + def test_ateom_checkpoint_parses(self): + b = parse(ATEOM_CHECKPOINT).breakdowns[0] + self.assertEqual((b.source, b.op), ("ateom", "checkpoint")) + self.assertEqual(b.phases["rootfs_upper"], 0.9) + + def test_prefixed_kubectl_lines_still_parse(self): + with tempfile.TemporaryDirectory() as d: + path = os.path.join(d, "pod.log") + with open(path, "w") as f: + f.write("pod/atelet-abc " + json.dumps(ATELET_RESTORE) + "\nnot json\n{\"msg\": \"other\"}\n") + out = phase_report.parse_files([path]) + self.assertEqual(len(out.breakdowns), 1) + self.assertEqual(out.lines_seen, 2) + + +class UnattributedTest(unittest.TestCase): + def test_atelet_checkpoint_residual_is_total_minus_sequential_phases(self): + b = parse(ATELET_CHECKPOINT).breakdowns[0] + self.assertAlmostEqual(b.phases["unattributed"], 6.0 - (0.01 + 1.14 + 4.48)) + + def test_ateom_checkpoint_counts_the_concurrent_captures_once(self): + b = parse(ATEOM_CHECKPOINT).breakdowns[0] + # prep + pause + max(snapshot, rootfs_upper) + teardown; the two + # captures overlapped, so only the slower one is spent wall time. + self.assertAlmostEqual(b.phases["unattributed"], 1.2 - (0.04 + 0.01 + 0.9 + 0.2)) + + def test_restore_layers_get_no_residual(self): + out = parse(ATELET_RESTORE, ATEOM_RESTORE) + for b in out.breakdowns: + self.assertNotIn("unattributed", b.phases) + + def test_residual_never_negative(self): + rec = dict(ATELET_CHECKPOINT, **{"ate.actor.checkpoint.duration.total": 1.0}) + self.assertEqual(parse(rec).breakdowns[0].phases["unattributed"], 0.0) + + +class ReportTest(unittest.TestCase): + def test_failed_records_are_flagged_and_excluded_from_percentiles(self): + out = parse(ATELET_RESTORE, FAILED) + self.assertEqual([b.failed for b in out.breakdowns], [False, True]) + rows = phase_report.report_phases(out.breakdowns, lambda _: None) + downloads = [r for r in rows if r["phase"] == "download"] + self.assertEqual(len(downloads), 1) + self.assertEqual(downloads[0]["count"], 1) # the failed 30s never entered + + def test_waterfall_nests_ateom_and_gap_under_atelet(self): + out = parse(ATEOM_RESTORE, ATELET_RESTORE) + text = render(phase_report.report_waterfalls, out.breakdowns, slowest=1) + lines = [line.strip() for line in text.splitlines()] + self.assertTrue(any(line.startswith("atelet ateom_restore") for line in lines), text) + self.assertTrue(any(line.startswith("ateom vm_restore") for line in lines), text) + gap = [line for line in lines if line.startswith("(gap)")] + self.assertEqual(len(gap), 1, text) + self.assertIn("50.0", gap[0]) # 1.1s atelet ateom_restore - 1.05s ateom total + + +if __name__ == "__main__": + unittest.main() From c90ea8da67e950ea937a8394ee471b9f8527d5d4 Mon Sep 17 00:00:00 2001 From: Lucky Abolorunke <62015433+Oneimu@users.noreply.github.com> Date: Wed, 30 Sep 2026 14:22:41 -0700 Subject: [PATCH 2/4] benchmarking: pair records by operation window, split by class and kind, read logging exports --- CONTRIBUTING.md | 4 +- benchmarking/analysis/README.md | 36 +++-- benchmarking/analysis/collect_logs.sh | 50 +++++-- benchmarking/analysis/phase_report.py | 116 ++++++++++----- benchmarking/analysis/test_phase_report.py | 72 +++++++++ docs/dev/suspend-resume-phase-breakdown.md | 164 +++++++++++++++++++++ 6 files changed, 389 insertions(+), 53 deletions(-) create mode 100644 docs/dev/suspend-resume-phase-breakdown.md diff --git a/CONTRIBUTING.md b/CONTRIBUTING.md index c51f08ed70..c3a5c7a7dd 100644 --- a/CONTRIBUTING.md +++ b/CONTRIBUTING.md @@ -31,7 +31,9 @@ The [Quickstart (Development)](README.md#quickstart-development) in the README covers bringing up a local cluster with the default (gVisor) runtime. To run the microVM runtime locally — which needs `/dev/kvm`, or Lima nested virtualization on Apple Silicon — see -[docs/dev/microvm-local.md](docs/dev/microvm-local.md). +[docs/dev/microvm-local.md](docs/dev/microvm-local.md). To measure where a +suspend or resume spends its time on such a cluster, follow +[docs/dev/suspend-resume-phase-breakdown.md](docs/dev/suspend-resume-phase-breakdown.md). ## Contribution process diff --git a/benchmarking/analysis/README.md b/benchmarking/analysis/README.md index 9772da6353..6d1160b742 100644 --- a/benchmarking/analysis/README.md +++ b/benchmarking/analysis/README.md @@ -32,14 +32,26 @@ A failed atelet operation still writes its record, marked with `error.type` teardown (or against a cluster you manage yourself): ```bash - ./collect_logs.sh --dest /tmp/run1 --since 30m + ./collect_logs.sh --dest /tmp/run1 --since-time 2026-09-24T18:00:00Z ``` - `--namespace` (default `ate-system`) selects the atelet pods and - `--worker-namespace` (default `benchmark-workloads`) the worker pods; the - ateom records land in the worker pod's stdout. + `--since-time` (or a relative `--since 30m`) should cover the run and + nothing before it. `--namespace` (default `ate-system`) selects the atelet + pods and `--worker-namespace` (default `benchmark-workloads`) the worker + pods; the ateom records land in the worker pod's stdout. `kubectl logs` + returns only a container's current log file, so on a long or busy run + kubelet's rotation can drop earlier records — collect soon after the run. + A container that restarted mid-run is dumped twice (`.previous.log`). -3. Aggregate: + On GKE the same records are in Cloud Logging, which keeps them past + rotation and teardown; an export is accepted as input directly: + + ```bash + gcloud logging read '(jsonPayload.msg="Checkpoint timing breakdown" OR jsonPayload.msg="Restore timing breakdown") AND resource.labels.cluster_name="" AND timestamp>="2026-09-24T18:00:00Z"' \ + --project --format json > /tmp/run1/export.json + ``` + +3. Aggregate (kubectl dumps and Cloud Logging exports can be mixed): ```bash python3 phase_report.py /tmp/run1/*.log --csv /tmp/run1/report @@ -47,8 +59,10 @@ A failed atelet operation still writes its record, marked with `error.type` ## Reading the report -**Phase percentiles.** Per layer, operation, scope and phase: count, p50, -p90, p95 and max. The atelet rows split a checkpoint between +**Phase percentiles.** Per layer, operation, sandbox class, snapshot kind, +scope and phase: count, p50, p90, p95 and max. A `golden` restore downloads +the golden image and a `latest` one the actor's own, so they are separate +rows, as are gVisor and micro-VM checkpoints. The atelet rows split a checkpoint between `sandbox_assets`, `ateom_checkpoint` and `persist`, and a restore between `volume_mount`, `manifest_fetch`, `sandbox_assets`, `download`, `oci_unpack` and `ateom_restore`. The ateom rows split the `ateom_*` phase further: @@ -64,9 +78,11 @@ the two checkpoint layers the report derives an `unattributed` row: the total minus what the logged phases account for, counting the concurrent captures once. It is the time the instrumentation does not yet name. -**Waterfalls.** The slowest operations, with the ateom record nested under -the atelet `ateom_*` phase and that phase's gap to the ateom total (RPC and -queueing between the layers). Tail outliers that blow up in one phase every +**Waterfalls.** The slowest operations, with the ateom record of the same +actor and time window nested under the atelet `ateom_*` phase and that +phase's gap to the ateom total (RPC and queueing between the layers). An +operation whose ateom record is missing prints without one rather than with +another cycle's. Tail outliers that blow up in one phase every time are systematic; different phases each time are environmental. `--csv` also writes `phase_percentiles.csv` for run-over-run comparison. diff --git a/benchmarking/analysis/collect_logs.sh b/benchmarking/analysis/collect_logs.sh index feedf7490f..641a0b16ae 100755 --- a/benchmarking/analysis/collect_logs.sh +++ b/benchmarking/analysis/collect_logs.sh @@ -16,15 +16,24 @@ # Dump the node-side logs a benchmark run needs for phase_report.py: every # atelet pod (the atelet-side timing breakdowns) and every worker pod # (ateom-microvm writes its records to the worker pod's stdout). Point -# --since at the run's start so the report covers exactly one run. +# --since-time (RFC 3339) or --since at the run's start so the report covers +# exactly one run. # -# Usage: collect_logs.sh --dest DIR [--since 30m] [--namespace ate-system] -# [--worker-namespace benchmark-workloads] +# kubectl logs returns only a container's current log file: once kubelet +# rotates it under load, earlier records are gone, so collect soon after the +# run. A container that restarted mid-run is dumped twice, its previous log +# under .previous.log. On GKE, Cloud Logging keeps every record past +# rotation and teardown; phase_report.py reads `gcloud logging read +# --format json` output directly. +# +# Usage: collect_logs.sh --dest DIR [--since 30m | --since-time 2026-09-24T18:00:00Z] +# [--namespace ate-system] [--worker-namespace benchmark-workloads] -set -euo pipefail +set -uo pipefail DEST="" SINCE="1h" +SINCE_TIME="" NAMESPACE="ate-system" WORKER_NAMESPACE="benchmark-workloads" @@ -32,6 +41,7 @@ while [[ $# -gt 0 ]]; do case "$1" in --dest) DEST="$2"; shift 2 ;; --since) SINCE="$2"; shift 2 ;; + --since-time) SINCE_TIME="$2"; shift 2 ;; --namespace) NAMESPACE="$2"; shift 2 ;; --worker-namespace) WORKER_NAMESPACE="$2"; shift 2 ;; *) echo "unknown flag: $1" >&2; exit 2 ;; @@ -39,20 +49,42 @@ while [[ $# -gt 0 ]]; do done if [[ -z "${DEST}" ]]; then - echo "usage: $0 --dest DIR [--since 30m] [--namespace ate-system] [--worker-namespace benchmark-workloads]" >&2 + echo "usage: $0 --dest DIR [--since 30m | --since-time RFC3339] [--namespace ate-system] [--worker-namespace benchmark-workloads]" >&2 exit 2 fi mkdir -p "${DEST}" +if [[ -n "${SINCE_TIME}" ]]; then + WINDOW=(--since-time="${SINCE_TIME}") +else + WINDOW=(--since="${SINCE}") +fi + +# One pod that cannot be read (still starting, evicted) must not cost the +# others, so every kubectl call warns and moves on rather than aborting. collect() { local ns="$1" selector="$2" local pods - pods=$(kubectl get pods -n "${ns}" -l "${selector}" -o name) + if ! pods=$(kubectl get pods -n "${ns}" -l "${selector}" -o name); then + echo "warn: could not list pods in ${ns} (${selector}); skipping" >&2 + return 0 + fi for pod in ${pods}; do local name="${pod#pod/}" - echo "collecting ${ns}/${name} (since ${SINCE})" - kubectl logs -n "${ns}" "${name}" --since="${SINCE}" --timestamps=false \ - > "${DEST}/${ns}-${name}.log" + echo "collecting ${ns}/${name} (${WINDOW[*]})" + kubectl logs -n "${ns}" "${name}" "${WINDOW[@]}" --timestamps=false \ + > "${DEST}/${ns}-${name}.log" \ + || { echo "warn: kubectl logs ${ns}/${name} failed; skipping" >&2; rm -f "${DEST}/${ns}-${name}.log"; } + # --previous is only valid after a restart; kubectl rejects it otherwise. + local restarts + restarts=$(kubectl get pod -n "${ns}" "${name}" \ + -o jsonpath='{.status.containerStatuses[0].restartCount}' 2>/dev/null || echo 0) + if [[ "${restarts:-0}" -gt 0 ]]; then + echo "collecting ${ns}/${name} previous container (${restarts} restarts)" + kubectl logs -n "${ns}" "${name}" --previous "${WINDOW[@]}" --timestamps=false \ + > "${DEST}/${ns}-${name}.previous.log" \ + || { echo "warn: kubectl logs --previous ${ns}/${name} failed; skipping" >&2; rm -f "${DEST}/${ns}-${name}.previous.log"; } + fi done } diff --git a/benchmarking/analysis/phase_report.py b/benchmarking/analysis/phase_report.py index 8da99494d6..7424011d6e 100755 --- a/benchmarking/analysis/phase_report.py +++ b/benchmarking/analysis/phase_report.py @@ -15,9 +15,11 @@ """Aggregate the suspend/resume phase-breakdown log records into a report. -Reads JSON-lines logs (kubectl logs dumps of the atelet and worker pods; any -non-JSON or unrelated lines are skipped) and aggregates the two joinable -record kinds the node side emits, each written by both layers: +Reads the atelet and worker pod logs, either as kubectl logs dumps (JSON +lines; unrelated lines are skipped) or as a Cloud Logging export +(`gcloud logging read --format json`: a JSON array of entries with the record +under jsonPayload), and aggregates the two joinable record kinds the node +side emits, each written by both layers: - "Restore timing breakdown" atelet (ate.actor.restore.duration.*) and ateom (ateom.actor.restore.duration.*) @@ -64,9 +66,14 @@ SCOPE_KEY = "ate.snapshot.scope" PHASE_KEY = "ate.snapshot.phase" KIND_KEY = "ate.snapshot.kind" +SANDBOX_CLASS_KEY = "ate.sandbox.class" TEMPLATE_KEY = "ate.template.name" ERROR_TYPE_KEY = "error.type" +# Only ateom-microvm writes the ateom records, and they carry no sandbox +# class of their own; the atelet record of the same operation does. +ATEOM_SANDBOX_CLASS = "microvm" + # The phase every record reports, and the one the report derives: the part of # the total no logged phase accounts for. TOTAL = "total" @@ -109,6 +116,7 @@ class Breakdown: ts: float | None # epoch seconds parsed from time, None when unparseable actor_uid: str template: str + sandbox_class: str scope: str kind: str phases: dict[str, float] # phase name -> seconds @@ -120,6 +128,7 @@ class Parsed: breakdowns: list[Breakdown] = field(default_factory=list) lines_seen: int = 0 lines_matched: int = 0 + bad_values: int = 0 # duration keys whose value was not a number _TIME_RE = re.compile(r"^(.*T\d\d:\d\d:\d\d)(\.\d+)?(Z|[+-]\d\d:?\d\d)?$") @@ -163,14 +172,24 @@ def unattributed(source: str, op: str, phases: dict[str, float]) -> float | None def parse_line(obj: dict, out: Parsed) -> None: + # A Cloud Logging entry nests the record under jsonPayload and moves + # slog's time onto the entry's timestamp; a kubectl dump is the record. + payload = obj.get("jsonPayload") + if isinstance(payload, dict): + entry, obj = obj, dict(payload) + obj.setdefault("time", entry.get("timestamp", "")) msg = obj.get("msg", "") if msg in BREAKDOWN_MSGS: for (source, op), prefix in BREAKDOWN_PREFIXES.items(): - phases = { - k[len(prefix):]: float(v) - for k, v in obj.items() - if k.startswith(prefix) - } + phases: dict[str, float] = {} + for k, v in obj.items(): + if not k.startswith(prefix): + continue + try: + phases[k[len(prefix):]] = float(v) + except (TypeError, ValueError): + # One bad value costs that phase, not the run's report. + out.bad_values += 1 if not phases: continue residual = unattributed(source, op, phases) @@ -183,6 +202,8 @@ def parse_line(obj: dict, out: Parsed) -> None: ts=parse_time(obj.get("time", "")), actor_uid=obj.get(ACTOR_UID_KEY, ""), template=obj.get(TEMPLATE_KEY, ""), + sandbox_class=obj.get(SANDBOX_CLASS_KEY, "") + or (ATEOM_SANDBOX_CLASS if source == "ateom" else ""), scope=obj.get(SCOPE_KEY, ""), kind=obj.get(KIND_KEY, ""), phases=phases, @@ -197,19 +218,32 @@ def parse_files(paths: list[str]) -> Parsed: for path in paths: f = sys.stdin if path == "-" else open(path, encoding="utf-8", errors="replace") with f: - for line in f: - # kubectl log dumps may prefix each line (pod name, timestamp); - # recover the JSON object from the first brace. - brace = line.find("{") - if brace < 0: - continue - out.lines_seen += 1 - try: - obj = json.loads(line[brace:]) - except json.JSONDecodeError: - continue + text = f.read() + if text.lstrip().startswith("["): + # A Cloud Logging export: one JSON array of entries, pretty-printed + # across lines, so it cannot be read line by line. + try: + entries = json.loads(text) + except json.JSONDecodeError: + entries = [] + for obj in entries: if isinstance(obj, dict): + out.lines_seen += 1 parse_line(obj, out) + continue + for line in text.splitlines(): + # kubectl log dumps may prefix each line (pod name, timestamp); + # recover the JSON object from the first brace. + brace = line.find("{") + if brace < 0: + continue + out.lines_seen += 1 + try: + obj = json.loads(line[brace:]) + except json.JSONDecodeError: + continue + if isinstance(obj, dict): + parse_line(obj, out) return out @@ -230,23 +264,27 @@ def fmt_s(seconds: float) -> str: def report_phases(breakdowns: list[Breakdown], writer) -> list[dict]: - """Per (source, op, scope, phase) percentiles. Returns the rows for CSV.""" + """Per (source, op, sandbox class, kind, scope, phase) percentiles. A + golden restore (which downloads the golden image) and a latest restore, + or a gVisor and a micro-VM checkpoint, are different distributions and + must not be pooled. Returns the rows for CSV.""" groups: dict[tuple, list[float]] = defaultdict(list) for b in breakdowns: if b.failed: continue for name, seconds in b.phases.items(): - groups[(b.source, b.op, b.scope or "-", name)].append(seconds) + groups[(b.source, b.op, b.sandbox_class or "-", b.kind or "-", b.scope or "-", name)].append(seconds) rows = [] writer("== Phase percentiles (ms) ==") - writer(f"{'layer':7} {'op':11} {'scope':15} {'phase':16} {'n':>5} " + writer(f"{'layer':7} {'op':11} {'class':8} {'kind':7} {'scope':15} {'phase':16} {'n':>5} " f"{'p50':>8} {'p90':>8} {'p95':>8} {'max':>8}") - for key in sorted(groups, key=lambda k: (k[0], k[1], k[2], phase_sort_key(k[3]))): + for key in sorted(groups, key=lambda k: (k[:5], phase_sort_key(k[5]))): vals = sorted(groups[key]) - source, op, scope, name = key + source, op, sandbox_class, kind, scope, name = key row = { - "layer": source, "op": op, "scope": scope, "phase": name, + "layer": source, "op": op, "class": sandbox_class, "kind": kind, + "scope": scope, "phase": name, "count": len(vals), "p50_ms": percentile(vals, 50) * 1000, "p90_ms": percentile(vals, 90) * 1000, @@ -254,7 +292,7 @@ def report_phases(breakdowns: list[Breakdown], writer) -> list[dict]: "max_ms": max(vals) * 1000, } rows.append(row) - writer(f"{source:7} {op:11} {scope:15} {name:16} {len(vals):5d} " + writer(f"{source:7} {op:11} {sandbox_class:8} {kind:7} {scope:15} {name:16} {len(vals):5d} " f"{fmt_s(percentile(vals, 50))} {fmt_s(percentile(vals, 90))} " f"{fmt_s(percentile(vals, 95))} {fmt_s(max(vals))}") failed = sum(1 for b in breakdowns if b.failed) @@ -264,11 +302,18 @@ def report_phases(breakdowns: list[Breakdown], writer) -> list[dict]: def inner_record(op: Breakdown, by_actor: dict[tuple, list[Breakdown]]) -> Breakdown | None: - """The ateom record of the same actor and op nearest before op's own: one - cycle emits exactly one of each, and the ateom one lands first, inside - the atelet ateom_* phase.""" + """The successful ateom record of the same actor and op inside op's own + window (its record time minus its total): one cycle emits exactly one of + each, and the ateom one lands first, inside the atelet ateom_* phase. An + older record belongs to an earlier cycle whose partner is missing (pod + gone, log rotated, window cut), and pairing it would print a gap that + never happened, so nothing is better than the wrong one.""" + total = op.phases.get(TOTAL) + if op.ts is None or total is None: + return None + lo, hi = op.ts - total - JOIN_SLACK_S, op.ts + JOIN_SLACK_S candidates = [a for a in by_actor.get((op.actor_uid, op.op), []) - if a.ts is not None and op.ts is not None and a.ts <= op.ts + JOIN_SLACK_S] + if not a.failed and a.ts is not None and lo <= a.ts <= hi] return candidates[-1] if candidates else None @@ -292,7 +337,7 @@ def report_waterfalls(breakdowns: list[Breakdown], writer, slowest: int) -> None inner = inner_record(b, by_actor) total = b.phases.get(TOTAL, 0) writer(f"{b.op} actor={b.actor_uid} template={b.template} " - f"scope={b.scope or '-'} kind={b.kind or '-'} " + f"class={b.sandbox_class or '-'} scope={b.scope or '-'} kind={b.kind or '-'} " f"total={total * 1000:.1f}") for name in sorted(b.phases, key=phase_sort_key): if name == TOTAL: @@ -321,7 +366,10 @@ def write_csv(dest: Path, name: str, rows: list[dict]) -> None: def main() -> int: ap = argparse.ArgumentParser(description=__doc__.splitlines()[0]) - ap.add_argument("logs", nargs="+", help="JSON-lines log files ('-' for stdin)") + ap.add_argument("logs", nargs="+", + help="kubectl logs dumps (JSON lines) or Cloud Logging exports " + "(JSON array, as gcloud logging read --format json writes); " + "'-' for stdin") ap.add_argument("--csv", type=Path, default=None, help="also write phase_percentiles.csv here") ap.add_argument("--slowest", type=int, default=3, @@ -330,7 +378,9 @@ def main() -> int: parsed = parse_files(args.logs) print(f"parsed {parsed.lines_matched} timing breakdown records " - f"out of {parsed.lines_seen} JSON log lines") + f"out of {parsed.lines_seen} JSON log entries") + if parsed.bad_values: + print(f"(skipped {parsed.bad_values} non-numeric duration values)") if not parsed.breakdowns: print("no matching records; are these the atelet and worker pod logs?", file=sys.stderr) return 1 diff --git a/benchmarking/analysis/test_phase_report.py b/benchmarking/analysis/test_phase_report.py index 39b6fc89db..f8efaf3d4f 100644 --- a/benchmarking/analysis/test_phase_report.py +++ b/benchmarking/analysis/test_phase_report.py @@ -86,6 +86,33 @@ def render(fn, *args, **kwargs): "ateom.actor.checkpoint.duration.total": 1.2, } +# The same actor's ateom restore record from an hour earlier: a different +# cycle, whose own atelet partner is missing from the dump. +ATEOM_RESTORE_STALE = dict(ATEOM_RESTORE, time="2026-09-23T09:00:05.400000000Z") + +# One entry as `gcloud logging read --format json` writes it, trimmed to the +# fields that matter: the record under jsonPayload, slog's time promoted to +# the entry's timestamp, dotted keys kept as they are. +CLOUD_LOGGING_ENTRY = { + "insertId": "m9xljnzmzppup1vj", + "jsonPayload": { + "ate.actor.checkpoint.duration.ateom_checkpoint": 1.024788329, + "ate.actor.checkpoint.duration.persist": 4.18451256, + "ate.actor.checkpoint.duration.sandbox_assets": 2.4358e-05, + "ate.actor.checkpoint.duration.total": 5.385861508, + "ate.actor.name": "sb-d790f7ed", "ate.actor.uid": "cc7a2ac7", + "ate.atespace": "benchmark", "ate.sandbox.class": "microvm", + "ate.snapshot.kind": "latest", "ate.snapshot.scope": "full", + "ate.template.atespace": "benchmark-workloads", "ate.template.name": "glutton", + "level": "INFO", "msg": "Checkpoint timing breakdown", + }, + "logName": "projects/p/logs/stdout", + "resource": {"labels": {"container_name": "atelet", "namespace_name": "ate-system"}, + "type": "k8s_container"}, + "severity": "INFO", + "timestamp": "2026-09-28T17:49:47.665836237Z", +} + FAILED = { "time": "2026-09-23T10:02:00Z", "level": "INFO", "msg": "Restore timing breakdown", @@ -117,6 +144,35 @@ def test_ateom_checkpoint_parses(self): self.assertEqual((b.source, b.op), ("ateom", "checkpoint")) self.assertEqual(b.phases["rootfs_upper"], 0.9) + def test_cloud_logging_entry_parses_with_the_entry_timestamp(self): + b = parse(CLOUD_LOGGING_ENTRY).breakdowns[0] + self.assertEqual((b.source, b.op), ("atelet", "checkpoint")) + self.assertEqual(b.time, "2026-09-28T17:49:47.665836237Z") + self.assertIsNotNone(b.ts) + self.assertEqual((b.sandbox_class, b.kind, b.scope), ("microvm", "latest", "full")) + self.assertAlmostEqual(b.phases["persist"], 4.18451256) + + def test_cloud_logging_export_file_is_a_json_array(self): + with tempfile.TemporaryDirectory() as d: + path = os.path.join(d, "export.json") + with open(path, "w") as f: + json.dump([CLOUD_LOGGING_ENTRY, CLOUD_LOGGING_ENTRY], f, indent=2) + out = phase_report.parse_files([path]) + self.assertEqual((out.lines_seen, len(out.breakdowns)), (2, 2)) + + def test_ateom_records_default_to_the_microvm_class(self): + out = parse(ATEOM_RESTORE, ATELET_RESTORE) + self.assertEqual(out.breakdowns[0].sandbox_class, "microvm") + self.assertEqual(out.breakdowns[1].sandbox_class, "") # not on this fixture + + def test_non_numeric_duration_is_skipped_not_fatal(self): + rec = dict(ATELET_RESTORE, **{"ate.actor.restore.duration.download": None, + "ate.actor.restore.duration.oci_unpack": "fast"}) + out = parse(rec) + self.assertEqual(out.bad_values, 2) + self.assertEqual(out.breakdowns[0].phases["total"], 3.9) + self.assertNotIn("download", out.breakdowns[0].phases) + def test_prefixed_kubectl_lines_still_parse(self): with tempfile.TemporaryDirectory() as d: path = os.path.join(d, "pod.log") @@ -157,6 +213,22 @@ def test_failed_records_are_flagged_and_excluded_from_percentiles(self): self.assertEqual(len(downloads), 1) self.assertEqual(downloads[0]["count"], 1) # the failed 30s never entered + def test_percentiles_split_by_sandbox_class_and_kind(self): + golden = dict(ATELET_RESTORE, **{"ate.snapshot.kind": "golden", "ate.sandbox.class": "gvisor", + "ate.actor.restore.duration.download": 9.0, + "ate.actor.restore.duration.total": 10.0}) + latest = dict(ATELET_RESTORE, **{"ate.sandbox.class": "gvisor"}) + rows = phase_report.report_phases(parse(golden, latest).breakdowns, lambda _: None) + downloads = {(r["class"], r["kind"]): r["p50_ms"] for r in rows if r["phase"] == "download"} + self.assertEqual(downloads, {("gvisor", "golden"): 9000.0, ("gvisor", "latest"): 2400.0}) + + def test_waterfall_ignores_an_ateom_record_from_another_cycle(self): + out = parse(ATEOM_RESTORE_STALE, ATELET_RESTORE) + text = render(phase_report.report_waterfalls, out.breakdowns, slowest=1) + self.assertIn("atelet ateom_restore", text) + self.assertNotIn("ateom vm_restore", text) + self.assertNotIn("(gap)", text) + def test_waterfall_nests_ateom_and_gap_under_atelet(self): out = parse(ATEOM_RESTORE, ATELET_RESTORE) text = render(phase_report.report_waterfalls, out.breakdowns, slowest=1) diff --git a/docs/dev/suspend-resume-phase-breakdown.md b/docs/dev/suspend-resume-phase-breakdown.md new file mode 100644 index 0000000000..34bf1e042e --- /dev/null +++ b/docs/dev/suspend-resume-phase-breakdown.md @@ -0,0 +1,164 @@ +# Suspend/resume phase breakdown: a runbook + +This guide walks through getting a per-phase breakdown of `SuspendActor` and +`ResumeActor` latency from a cluster you run yourself, from deploying the +workloads to reading the report. It strings together tools documented +elsewhere; follow the links when a step needs more than the command shown. + +The breakdown comes from two developer-facing log records, `Checkpoint timing +breakdown` and `Restore timing breakdown`, written once per operation by +atelet (the outer phases: `download`, `persist`, `ateom_*`, …) and by +ateom-microvm (the phases inside `ateom_*`: `pause`, `snapshot`, `vm_restore`, +`wakeup_probe`, …). `benchmarking/analysis/` collects the records from the pod +logs and aggregates them. See +[docs/observability.md](../observability.md#actor-attributed-component-logs) +for the record shape and +[benchmarking/analysis/README.md](../../benchmarking/analysis/README.md) for +the report. + +## What you need + +* A cluster with Agent Substrate installed from a tree that has the records + and the analysis tools: the + [Quickstart (Development)](../../README.md#quickstart-development) for kind, + or the [GKE Quickstart](../../README.md#gke-quickstart-development). For GKE, + `source .ate-dev-env.sh` first; every script below reads it. +* The **microVM** sandbox class, if you want the inner (ateom) phases: see + [microvm-local.md](microvm-local.md). ateom-gvisor writes no breakdown + record, so on a gVisor pool the report shows the atelet layer only. +* `kubectl`, `kubectl-ate` (`go install ./cmd/kubectl-ate`), and `python3` + (standard library only). + +## 1. Deploy the benchmark workloads + +```sh +./benchmarking/deploy_locust.sh --deploy \ + --sandbox-class microvm --worker-count 3 --actor-memory 1536Mi +``` + +This deploys the benchmark `WorkerPool` and `ActorTemplate`s into the +`benchmark-workloads` atespace, builds and pushes the locust image, and deploys +the locust master and workers (`benchmarking/README.md`, +[Deploy benchmarks](../../benchmarking/README.md#deploy-benchmarks)). On a kind +cluster with limited memory use smaller values, e.g. `--worker-count 1 +--actor-memory 512Mi`, and scale the load flags in step 2 down to match. + +Wait until every template has a golden snapshot — the column is a UUID once +the golden build (a full boot plus checkpoint) has finished: + +```sh +kubectl ate get actor-templates -a benchmark-workloads +``` + +Those golden builds are themselves checkpoints, so the worker pods already +carry `Checkpoint timing breakdown` records at this point: + +```sh +kubectl logs -n benchmark-workloads -l ate.dev/worker-pool --since=15m \ + | grep -c "timing breakdown" +``` + +A zero here means the running image predates the records; redeploy ate-system +from the current tree before going on. + +## 2. Run suspend/resume load + +The `glutton` user class creates an actor per user and cycles it through +suspend and resume. Give it a realistically sized working set so the memory +phases have something to measure: + +```sh +kubectl port-forward svc/locust -n benchmarking 8089:8089 +``` + +Open , pick `GluttonUser`, and set: + +| field | value | why | +|---|---|---| +| users / spawn rate | 2 / 1 | two concurrent cycles; keep it small so transfers do not contend for one NIC | +| `--mem-target` | `1Gi` | resident working set before the first suspend (below `--actor-memory`) | +| `--mem-churn` | `64Mi` | re-dirty part of it each cycle, so successive snapshots differ | +| `--mem-read` | `all` | walk the whole set after each resume, so a demand-paged restore pays for its pages | +| `--min-wait-time` / `--max-wait-time` | `1.0` / `1.0` | one cycle per second per user | + +The form fields are the flags documented in +`benchmarking/locust/common/memload_config.py`; the boomer workers fetch the +values from the master on each spawn. Let it run for about 90 seconds (≈10 +cycles per user), then stop the test. For a single operation without locust, +`kubectl ate suspend actor -a benchmark-workloads` and +`kubectl ate resume actor -a benchmark-workloads` on any actor in the +atespace produce one record pair each. + +## 3. Collect the node logs + +Collect **while the pods still exist**. `benchmarking/automation/ +orchestrator.py` deletes the workload and ate-system pods right after each +test and does not run this step yet, so on a cluster you manage yourself this +is the moment to do it: + +```sh +./benchmarking/analysis/collect_logs.sh --dest /tmp/run1 --since-time 2026-09-24T18:00:00Z +``` + +`--since-time` (or a relative `--since 30m`) should cover the run and nothing +before it. The script writes one file per atelet pod (`ate-system`) and per +worker pod (`benchmark-workloads`) under `--dest`; a pod it cannot read is +skipped with a warning. `kubectl logs` returns only a container's current log +file, so collect soon after the run, before kubelet rotates it. + +On GKE the same records are also in Cloud Logging, where they outlive both +rotation and teardown; an export is accepted by the report as-is: + +```sh +gcloud logging read '(jsonPayload.msg="Checkpoint timing breakdown" OR jsonPayload.msg="Restore timing breakdown") AND resource.labels.cluster_name="" AND timestamp>="2026-09-24T18:00:00Z"' \ + --project --format json > /tmp/run1/export.json +``` + +## 4. Produce the report + +```sh +python3 benchmarking/analysis/phase_report.py /tmp/run1/*.log --csv /tmp/run1/report +``` + +The first line says how many records it found; expect two per suspend and two +per resume (one per layer). If it is zero, the window in step 3 missed the run +or the images predate the records. + +## 5. Read it + +* **Phase percentiles** rank the phases. For a suspend, compare atelet's + `persist` against `ateom_checkpoint`, then the ateom rows (`snapshot`, + `teardown`, `pause`, …) against each other; for a resume, `download` against + `ateom_restore`, then `vm_restore` and `wakeup_probe`. The `unattributed` + row is the part of a checkpoint no logged phase names. +* **Waterfalls** show the slowest operations with the ateom record nested + under the atelet phase it decomposes; a tail that is always the same phase + is systematic, one that moves between phases is environmental. +* Two sanity checks on a healthy run: the atelet and ateom rows for one + operation have the same `n`, and atelet's `ateom_checkpoint` p50 is a few + tens of milliseconds above ateom's `total` p50 (the `(gap)` line in a + waterfall is that RPC overhead). + +`/tmp/run1/report/phase_percentiles.csv` holds the same table for comparing +runs, for example one at `--mem-target 1Gi` against one at `2Gi` to see which +phases scale with dirty memory. + +## 6. Tear down + +```sh +./benchmarking/deploy_locust.sh --delete +``` + +## Troubleshooting + +* **Only atelet rows, no ateom rows.** The pool is gVisor (ateom-gvisor emits + no record), or the worker image predates the records. Check + `kubectl get workerpools -A` and step 1's grep. +* **Counts differ between layers.** The collection window cut through a + cycle, or a worker pod restarted and took its log with it (the script dumps + a restarted container's previous log as `.previous.log`). On GKE the + Cloud Logging export has everything. +* **A record has `error.type`.** That operation failed; the report leaves it + out of the percentiles but the waterfall input still lists it. A + `DeadlineExceeded` on `download` or `persist` with a normal-looking + `ateom_*` phase points at object storage, not the runtime. From 69cb37729f1daa31066d044c3c0ab890d4085fc3 Mon Sep 17 00:00:00 2001 From: Lucky Abolorunke <62015433+Oneimu@users.noreply.github.com> Date: Fri, 2 Oct 2026 08:38:11 -0700 Subject: [PATCH 3/4] benchmarking: tighten the atelet/ateom pairing and surface parse problems --- benchmarking/analysis/README.md | 41 +++++--- benchmarking/analysis/collect_logs.sh | 19 +++- benchmarking/analysis/phase_report.py | 108 +++++++++++++++------ benchmarking/analysis/test_phase_report.py | 50 +++++++++- docs/dev/suspend-resume-phase-breakdown.md | 15 +-- 5 files changed, 177 insertions(+), 56 deletions(-) diff --git a/benchmarking/analysis/README.md b/benchmarking/analysis/README.md index 6d1160b742..957b0f9894 100644 --- a/benchmarking/analysis/README.md +++ b/benchmarking/analysis/README.md @@ -35,28 +35,40 @@ A failed atelet operation still writes its record, marked with `error.type` ./collect_logs.sh --dest /tmp/run1 --since-time 2026-09-24T18:00:00Z ``` - `--since-time` (or a relative `--since 30m`) should cover the run and - nothing before it. `--namespace` (default `ate-system`) selects the atelet - pods and `--worker-namespace` (default `benchmark-workloads`) the worker - pods; the ateom records land in the worker pod's stdout. `kubectl logs` + The script sources `.ate-dev-env.sh` from the repository root when present + (set `NO_DEV_ENV` to skip) and honors `KUBECTL_CONTEXT`, like the `hack/` + scripts. `--since-time` (or a relative `--since 30m`) should cover the run + and nothing before it. `--namespace` (default `ate-system`) selects the + atelet pods and `--worker-namespace` (default `benchmark-workloads`) the + worker pods; the ateom records land in the worker pod's stdout. `kubectl logs` returns only a container's current log file, so on a long or busy run kubelet's rotation can drop earlier records — collect soon after the run. A container that restarted mid-run is dumped twice (`.previous.log`). On GKE the same records are in Cloud Logging, which keeps them past - rotation and teardown; an export is accepted as input directly: + rotation and teardown; an export is accepted as input directly, instead of + the kubectl dumps: ```bash gcloud logging read '(jsonPayload.msg="Checkpoint timing breakdown" OR jsonPayload.msg="Restore timing breakdown") AND resource.labels.cluster_name="" AND timestamp>="2026-09-24T18:00:00Z"' \ --project --format json > /tmp/run1/export.json ``` -3. Aggregate (kubectl dumps and Cloud Logging exports can be mixed): +3. Aggregate, from the kubectl dumps: ```bash python3 phase_report.py /tmp/run1/*.log --csv /tmp/run1/report ``` + or from the Cloud Logging export: + + ```bash + python3 phase_report.py /tmp/run1/export.json --csv /tmp/run1/report + ``` + + Use one source per run: the dumps and an export of the same pods hold the + same records, and passing both would count each twice. + ## Reading the report **Phase percentiles.** Per layer, operation, sandbox class, snapshot kind, @@ -76,14 +88,19 @@ ateom checkpoint captures run concurrently on the paused guest (the paused window costs their max), and the ateom restore phases are sequential. For the two checkpoint layers the report derives an `unattributed` row: the total minus what the logged phases account for, counting the concurrent captures -once. It is the time the instrumentation does not yet name. +once. It is the time the instrumentation does not yet name. The ateom +records carry no snapshot kind of their own; a paired one takes its kind from +the atelet record, so both layers split into the same rows. **Waterfalls.** The slowest operations, with the ateom record of the same -actor and time window nested under the atelet `ateom_*` phase and that -phase's gap to the ateom total (RPC and queueing between the layers). An -operation whose ateom record is missing prints without one rather than with -another cycle's. Tail outliers that blow up in one phase every -time are systematic; different phases each time are environmental. +actor nested under the atelet `ateom_*` phase and that phase's gap to the +ateom total (RPC and queueing between the layers). The ateom record is the +one written inside that phase's window (for a checkpoint, before `persist` +began), and each ateom record pairs with at most one operation, so an +operation whose own record is missing prints without one rather than with +another cycle's. Failed operations (`error.type` set) are left out of both +the percentiles and the waterfalls. Tail outliers that blow up in one phase +every time are systematic; different phases each time are environmental. `--csv` also writes `phase_percentiles.csv` for run-over-run comparison. diff --git a/benchmarking/analysis/collect_logs.sh b/benchmarking/analysis/collect_logs.sh index 641a0b16ae..e0555ca15b 100755 --- a/benchmarking/analysis/collect_logs.sh +++ b/benchmarking/analysis/collect_logs.sh @@ -26,11 +26,22 @@ # rotation and teardown; phase_report.py reads `gcloud logging read # --format json` output directly. # +# Like the other hack scripts, this sources .ate-dev-env.sh from the repo root +# for the cluster settings unless NO_DEV_ENV is set, and respects +# KUBECTL_CONTEXT. +# # Usage: collect_logs.sh --dest DIR [--since 30m | --since-time 2026-09-24T18:00:00Z] # [--namespace ate-system] [--worker-namespace benchmark-workloads] set -uo pipefail +ROOT="$(git rev-parse --show-toplevel 2>/dev/null || pwd)" +if [[ -r "${ROOT}/.ate-dev-env.sh" ]] && [[ -z "${NO_DEV_ENV:-}" ]]; then + # shellcheck source=/dev/null + source "${ROOT}/.ate-dev-env.sh" +fi +run_kubectl() { kubectl ${KUBECTL_CONTEXT:+--context=${KUBECTL_CONTEXT}} "$@"; } + DEST="" SINCE="1h" SINCE_TIME="" @@ -65,23 +76,23 @@ fi collect() { local ns="$1" selector="$2" local pods - if ! pods=$(kubectl get pods -n "${ns}" -l "${selector}" -o name); then + if ! pods=$(run_kubectl get pods -n "${ns}" -l "${selector}" -o name); then echo "warn: could not list pods in ${ns} (${selector}); skipping" >&2 return 0 fi for pod in ${pods}; do local name="${pod#pod/}" echo "collecting ${ns}/${name} (${WINDOW[*]})" - kubectl logs -n "${ns}" "${name}" "${WINDOW[@]}" --timestamps=false \ + run_kubectl logs -n "${ns}" "${name}" "${WINDOW[@]}" --timestamps=false \ > "${DEST}/${ns}-${name}.log" \ || { echo "warn: kubectl logs ${ns}/${name} failed; skipping" >&2; rm -f "${DEST}/${ns}-${name}.log"; } # --previous is only valid after a restart; kubectl rejects it otherwise. local restarts - restarts=$(kubectl get pod -n "${ns}" "${name}" \ + restarts=$(run_kubectl get pod -n "${ns}" "${name}" \ -o jsonpath='{.status.containerStatuses[0].restartCount}' 2>/dev/null || echo 0) if [[ "${restarts:-0}" -gt 0 ]]; then echo "collecting ${ns}/${name} previous container (${restarts} restarts)" - kubectl logs -n "${ns}" "${name}" --previous "${WINDOW[@]}" --timestamps=false \ + run_kubectl logs -n "${ns}" "${name}" --previous "${WINDOW[@]}" --timestamps=false \ > "${DEST}/${ns}-${name}.previous.log" \ || { echo "warn: kubectl logs --previous ${ns}/${name} failed; skipping" >&2; rm -f "${DEST}/${ns}-${name}.previous.log"; } fi diff --git a/benchmarking/analysis/phase_report.py b/benchmarking/analysis/phase_report.py index 7424011d6e..11a5ac8959 100755 --- a/benchmarking/analysis/phase_report.py +++ b/benchmarking/analysis/phase_report.py @@ -34,6 +34,9 @@ Usage: phase_report.py run-logs/*.log [--csv DEST_DIR] [--slowest N] +Feed it one source per run: the kubectl dumps, or the Cloud Logging export, +not both, or every record counts twice. + The records are developer-facing logs, not metric API; this reader is the consumer that makes them percentiles. See benchmarking/analysis/README.md. """ @@ -102,8 +105,10 @@ # paused window costs their max, not their sum. CONCURRENT_CAPTURES = ("snapshot", "durable_dir", "rootfs_upper") -# Slack when matching the ateom record to the atelet record of the same -# operation: the ateom one is emitted a hair before the atelet one. +# Slack on the window in which the ateom record of an operation must land: +# atelet does a little work after the ateom call returns before it writes its +# own record (unmounting, resetting the actor dirs), and the two clocks are +# the same node's but not the same goroutine's. JOIN_SLACK_S = 1.0 @@ -156,7 +161,9 @@ def unattributed(source: str, op: str, phases: dict[str, float]) -> float | None checkpoint is prep, pause, the concurrent captures (max), then teardown. The atelet restore phases overlap by design (download runs alongside the asset fetch and OCI unpack) and the ateom restore phases partition their - total by construction, so neither gets a residual.""" + total by construction, so neither gets a residual. The phases are + sub-intervals of the total, so a negative result is a bug in the emitter + or in this formula; it is not clamped, so that it shows.""" total = phases.get(TOTAL) if total is None: return None @@ -168,7 +175,7 @@ def unattributed(source: str, op: str, phases: dict[str, float]) -> float | None + phases.get("teardown", 0.0)) else: return None - return max(total - spent, 0.0) + return total - spent def parse_line(obj: dict, out: Parsed) -> None: @@ -192,14 +199,15 @@ def parse_line(obj: dict, out: Parsed) -> None: out.bad_values += 1 if not phases: continue + time_s = obj.get("time", "") residual = unattributed(source, op, phases) if residual is not None: phases[UNATTRIBUTED] = residual out.breakdowns.append(Breakdown( source=source, op=op, - time=obj.get("time", ""), - ts=parse_time(obj.get("time", "")), + time=time_s, + ts=parse_time(time_s), actor_uid=obj.get(ACTOR_UID_KEY, ""), template=obj.get(TEMPLATE_KEY, ""), sandbox_class=obj.get(SANDBOX_CLASS_KEY, "") @@ -224,8 +232,10 @@ def parse_files(paths: list[str]) -> Parsed: # across lines, so it cannot be read line by line. try: entries = json.loads(text) - except json.JSONDecodeError: - entries = [] + except json.JSONDecodeError as e: + print(f"warn: {path}: not a valid JSON array (truncated export?): {e}; skipped", + file=sys.stderr) + continue for obj in entries: if isinstance(obj, dict): out.lines_seen += 1 @@ -301,40 +311,75 @@ def report_phases(breakdowns: list[Breakdown], writer) -> list[dict]: return rows -def inner_record(op: Breakdown, by_actor: dict[tuple, list[Breakdown]]) -> Breakdown | None: - """The successful ateom record of the same actor and op inside op's own - window (its record time minus its total): one cycle emits exactly one of - each, and the ateom one lands first, inside the atelet ateom_* phase. An - older record belongs to an earlier cycle whose partner is missing (pod - gone, log rotated, window cut), and pairing it would print a gap that - never happened, so nothing is better than the wrong one.""" +def ateom_window(op: Breakdown) -> tuple[float, float] | None: + """When, relative to the atelet record's time, the ateom record of the + same operation was written. ateom writes it as its RPC returns, i.e. at + the end of the atelet ateom_* phase: for a checkpoint that is before + persist runs, for a restore it is the last phase. An operation that never + reached the ateom call has no window and nothing to pair with.""" total = op.phases.get(TOTAL) - if op.ts is None or total is None: + if op.ts is None or total is None or INNER_PHASE[op.op] not in op.phases: return None - lo, hi = op.ts - total - JOIN_SLACK_S, op.ts + JOIN_SLACK_S - candidates = [a for a in by_actor.get((op.actor_uid, op.op), []) - if not a.failed and a.ts is not None and lo <= a.ts <= hi] - return candidates[-1] if candidates else None - - -def report_waterfalls(breakdowns: list[Breakdown], writer, slowest: int) -> None: - """The slowest operations: the ateom record nested under the atelet - ateom_* phase, and the gap between that phase and the ateom total (RPC - and queueing between the layers).""" + if op.op == "checkpoint": + return (op.ts - total - JOIN_SLACK_S, + op.ts - op.phases.get("persist", 0.0) + JOIN_SLACK_S) + return (op.ts - op.phases[INNER_PHASE[op.op]] - JOIN_SLACK_S, op.ts + JOIN_SLACK_S) + + +def pair_records(breakdowns: list[Breakdown]) -> dict[int, Breakdown]: + """Match each atelet record to the ateom record of the same operation: + same actor and op, written inside the atelet ateom_* phase's window, and + not already claimed by another operation. Returns id(atelet) -> ateom. + + One cycle emits exactly one of each. An ateom record outside the window + belongs to another cycle whose partner is missing (pod gone, log rotated, + window cut); claiming it would print a gap that never happened, so an + operation without a match prints without one. Matching in time order + with a consumed set keeps rapid cycles of one actor from sharing or + swapping records. A paired ateom record also inherits the snapshot kind + its layer does not log, so the two layers' percentiles split alike.""" by_actor: dict[tuple, list[Breakdown]] = defaultdict(list) for b in breakdowns: - if b.source == "ateom": + if b.source == "ateom" and b.ts is not None: by_actor[(b.actor_uid, b.op)].append(b) for v in by_actor.values(): - v.sort(key=lambda b: (b.ts or 0.0, b.time)) + v.sort(key=lambda b: b.ts) + + pairs: dict[int, Breakdown] = {} + consumed: set[int] = set() + atelet = sorted((b for b in breakdowns if b.source == "atelet" and b.ts is not None), + key=lambda b: b.ts) + for op in atelet: + window = ateom_window(op) + if window is None: + continue + lo, hi = window + candidates = [a for a in by_actor.get((op.actor_uid, op.op), []) + if id(a) not in consumed and lo <= a.ts <= hi] + if not candidates: + continue + inner = candidates[-1] + consumed.add(id(inner)) + pairs[id(op)] = inner + if not inner.kind: + inner.kind = op.kind + return pairs + +def report_waterfalls(breakdowns: list[Breakdown], writer, slowest: int, + pairs: dict[int, Breakdown] | None = None) -> None: + """The slowest operations: the ateom record nested under the atelet + ateom_* phase, and the gap between that phase and the ateom total (RPC + and queueing between the layers).""" + if pairs is None: + pairs = pair_records(breakdowns) atelet = [b for b in breakdowns if b.source == "atelet" and not b.failed] atelet.sort(key=lambda b: b.phases.get(TOTAL, 0), reverse=True) writer("") writer(f"== Slowest {slowest} operations (waterfall, ms) ==") for b in atelet[:slowest]: - inner = inner_record(b, by_actor) + inner = pairs.get(id(b)) total = b.phases.get(TOTAL, 0) writer(f"{b.op} actor={b.actor_uid} template={b.template} " f"class={b.sandbox_class or '-'} scope={b.scope or '-'} kind={b.kind or '-'} " @@ -385,9 +430,12 @@ def main() -> int: print("no matching records; are these the atelet and worker pod logs?", file=sys.stderr) return 1 + # Pair before the percentiles: a paired ateom record takes its kind from + # the atelet record, so both layers' rows split the same way. + pairs = pair_records(parsed.breakdowns) print() phase_rows = report_phases(parsed.breakdowns, print) - report_waterfalls(parsed.breakdowns, print, args.slowest) + report_waterfalls(parsed.breakdowns, print, args.slowest, pairs) if args.csv: write_csv(args.csv, "phase_percentiles.csv", phase_rows) diff --git a/benchmarking/analysis/test_phase_report.py b/benchmarking/analysis/test_phase_report.py index f8efaf3d4f..6e3ff7840f 100644 --- a/benchmarking/analysis/test_phase_report.py +++ b/benchmarking/analysis/test_phase_report.py @@ -17,6 +17,7 @@ python3 -m unittest discover -s benchmarking/analysis """ +import contextlib import io import json import os @@ -155,8 +156,9 @@ def test_cloud_logging_entry_parses_with_the_entry_timestamp(self): def test_cloud_logging_export_file_is_a_json_array(self): with tempfile.TemporaryDirectory() as d: path = os.path.join(d, "export.json") + second = dict(CLOUD_LOGGING_ENTRY, timestamp="2026-09-28T17:50:47.000000000Z") with open(path, "w") as f: - json.dump([CLOUD_LOGGING_ENTRY, CLOUD_LOGGING_ENTRY], f, indent=2) + json.dump([CLOUD_LOGGING_ENTRY, second], f, indent=2) out = phase_report.parse_files([path]) self.assertEqual((out.lines_seen, len(out.breakdowns)), (2, 2)) @@ -173,6 +175,18 @@ def test_non_numeric_duration_is_skipped_not_fatal(self): self.assertEqual(out.breakdowns[0].phases["total"], 3.9) self.assertNotIn("download", out.breakdowns[0].phases) + def test_malformed_export_warns_instead_of_vanishing(self): + with tempfile.TemporaryDirectory() as d: + path = os.path.join(d, "export.json") + with open(path, "w") as f: + f.write('[{"jsonPayload": {"msg": "Checkpoint timing breakdown"') # truncated + err = io.StringIO() + with contextlib.redirect_stderr(err): + out = phase_report.parse_files([path]) + self.assertEqual(out.breakdowns, []) + self.assertIn("export.json", err.getvalue()) + self.assertIn("truncated export", err.getvalue()) + def test_prefixed_kubectl_lines_still_parse(self): with tempfile.TemporaryDirectory() as d: path = os.path.join(d, "pod.log") @@ -199,9 +213,10 @@ def test_restore_layers_get_no_residual(self): for b in out.breakdowns: self.assertNotIn("unattributed", b.phases) - def test_residual_never_negative(self): + def test_negative_residual_is_not_clamped(self): + # Cannot happen on a well-formed record; if it does, it must show. rec = dict(ATELET_CHECKPOINT, **{"ate.actor.checkpoint.duration.total": 1.0}) - self.assertEqual(parse(rec).breakdowns[0].phases["unattributed"], 0.0) + self.assertAlmostEqual(parse(rec).breakdowns[0].phases["unattributed"], 1.0 - (0.01 + 1.14 + 4.48)) class ReportTest(unittest.TestCase): @@ -214,7 +229,8 @@ def test_failed_records_are_flagged_and_excluded_from_percentiles(self): self.assertEqual(downloads[0]["count"], 1) # the failed 30s never entered def test_percentiles_split_by_sandbox_class_and_kind(self): - golden = dict(ATELET_RESTORE, **{"ate.snapshot.kind": "golden", "ate.sandbox.class": "gvisor", + golden = dict(ATELET_RESTORE, **{"time": "2026-09-23T09:59:00.000000000Z", + "ate.snapshot.kind": "golden", "ate.sandbox.class": "gvisor", "ate.actor.restore.duration.download": 9.0, "ate.actor.restore.duration.total": 10.0}) latest = dict(ATELET_RESTORE, **{"ate.sandbox.class": "gvisor"}) @@ -222,6 +238,32 @@ def test_percentiles_split_by_sandbox_class_and_kind(self): downloads = {(r["class"], r["kind"]): r["p50_ms"] for r in rows if r["phase"] == "download"} self.assertEqual(downloads, {("gvisor", "golden"): 9000.0, ("gvisor", "latest"): 2400.0}) + def test_paired_ateom_record_inherits_the_atelet_kind(self): + out = parse(ATEOM_RESTORE, ATELET_RESTORE) + pairs = phase_report.pair_records(out.breakdowns) + self.assertEqual(len(pairs), 1) + self.assertEqual(out.breakdowns[0].kind, "latest") + rows = phase_report.report_phases(out.breakdowns, lambda _: None) + self.assertEqual({r["kind"] for r in rows if r["layer"] == "ateom"}, {"latest"}) + + def test_checkpoint_pairing_window_ends_where_persist_starts(self): + # The ateom record is written when ateom_checkpoint ends, before the + # 4.48 s persist; a record from inside the persist window belongs to + # a later cycle of a rapidly cycling actor. + during_persist = dict(ATEOM_CHECKPOINT, time="2026-09-23T10:00:59.000000000Z") + out = parse(ATEOM_CHECKPOINT, during_persist, ATELET_CHECKPOINT) + pairs = phase_report.pair_records(out.breakdowns) + atelet = next(b for b in out.breakdowns if b.source == "atelet") + self.assertIs(pairs[id(atelet)], out.breakdowns[0]) + + def test_an_ateom_record_pairs_with_at_most_one_operation(self): + second = dict(ATELET_RESTORE, time="2026-09-23T10:00:05.900000000Z") + out = parse(ATEOM_RESTORE, ATELET_RESTORE, second) + pairs = phase_report.pair_records(out.breakdowns) + self.assertEqual(len(pairs), 1) + first = next(b for b in out.breakdowns if b.source == "atelet") + self.assertIn(id(first), pairs) + def test_waterfall_ignores_an_ateom_record_from_another_cycle(self): out = parse(ATEOM_RESTORE_STALE, ATELET_RESTORE) text = render(phase_report.report_waterfalls, out.breakdowns, slowest=1) diff --git a/docs/dev/suspend-resume-phase-breakdown.md b/docs/dev/suspend-resume-phase-breakdown.md index 34bf1e042e..8cc72e3b6a 100644 --- a/docs/dev/suspend-resume-phase-breakdown.md +++ b/docs/dev/suspend-resume-phase-breakdown.md @@ -21,8 +21,9 @@ the report. * A cluster with Agent Substrate installed from a tree that has the records and the analysis tools: the [Quickstart (Development)](../../README.md#quickstart-development) for kind, - or the [GKE Quickstart](../../README.md#gke-quickstart-development). For GKE, - `source .ate-dev-env.sh` first; every script below reads it. + or the [GKE Quickstart](../../README.md#gke-quickstart-development). For GKE + the scripts below read `.ate-dev-env.sh` from the repository root + (`collect_logs.sh` sources it itself; set `NO_DEV_ENV` to opt out). * The **microVM** sandbox class, if you want the inner (ateom) phases: see [microvm-local.md](microvm-local.md). ateom-gvisor writes no breakdown record, so on a gVisor pool the report shows the atelet layer only. @@ -103,7 +104,8 @@ is the moment to do it: `--since-time` (or a relative `--since 30m`) should cover the run and nothing before it. The script writes one file per atelet pod (`ate-system`) and per worker pod (`benchmark-workloads`) under `--dest`; a pod it cannot read is -skipped with a warning. `kubectl logs` returns only a container's current log +skipped with a warning. Use either these dumps or the Cloud Logging export +below for a run, not both. `kubectl logs` returns only a container's current log file, so collect soon after the run, before kubelet rotates it. On GKE the same records are also in Cloud Logging, where they outlive both @@ -159,6 +161,7 @@ phases scale with dirty memory. a restarted container's previous log as `.previous.log`). On GKE the Cloud Logging export has everything. * **A record has `error.type`.** That operation failed; the report leaves it - out of the percentiles but the waterfall input still lists it. A - `DeadlineExceeded` on `download` or `persist` with a normal-looking - `ateom_*` phase points at object storage, not the runtime. + out of the percentiles and the waterfalls. To see where it died, grep the + record itself: a `DeadlineExceeded` whose time sits in `download` or + `persist` with a normal-looking `ateom_*` phase points at object storage, + not the runtime. From 5dd747d5832b1c0cb8a97704852a07c3413a7dab Mon Sep 17 00:00:00 2001 From: Lucky Abolorunke <62015433+Oneimu@users.noreply.github.com> Date: Mon, 5 Oct 2026 09:18:27 -0700 Subject: [PATCH 4/4] benchmarking: read a bracket-prefixed dump line by line instead of skipping it --- benchmarking/analysis/phase_report.py | 20 ++++++++++++-------- benchmarking/analysis/test_phase_report.py | 16 +++++++++++++++- 2 files changed, 27 insertions(+), 9 deletions(-) diff --git a/benchmarking/analysis/phase_report.py b/benchmarking/analysis/phase_report.py index 11a5ac8959..962046f793 100755 --- a/benchmarking/analysis/phase_report.py +++ b/benchmarking/analysis/phase_report.py @@ -233,14 +233,18 @@ def parse_files(paths: list[str]) -> Parsed: try: entries = json.loads(text) except json.JSONDecodeError as e: - print(f"warn: {path}: not a valid JSON array (truncated export?): {e}; skipped", - file=sys.stderr) - continue - for obj in entries: - if isinstance(obj, dict): - out.lines_seen += 1 - parse_line(obj, out) - continue + # Not an export after all: a dump whose lines carry a bracket + # prefix, or a truncated export. Say so, then read it line by + # line like any other dump rather than drop the whole file. + print(f"warn: {path}: starts with '[' but is not a JSON array ({e}); " + f"reading it line by line", file=sys.stderr) + else: + if isinstance(entries, list): + for obj in entries: + if isinstance(obj, dict): + out.lines_seen += 1 + parse_line(obj, out) + continue for line in text.splitlines(): # kubectl log dumps may prefix each line (pod name, timestamp); # recover the JSON object from the first brace. diff --git a/benchmarking/analysis/test_phase_report.py b/benchmarking/analysis/test_phase_report.py index 6e3ff7840f..45e61169d6 100644 --- a/benchmarking/analysis/test_phase_report.py +++ b/benchmarking/analysis/test_phase_report.py @@ -185,7 +185,21 @@ def test_malformed_export_warns_instead_of_vanishing(self): out = phase_report.parse_files([path]) self.assertEqual(out.breakdowns, []) self.assertIn("export.json", err.getvalue()) - self.assertIn("truncated export", err.getvalue()) + self.assertIn("not a JSON array", err.getvalue()) + + def test_bracket_prefixed_dump_is_read_line_by_line(self): + # A dump whose lines start with "[INFO]" also starts with "[", but + # it is not an export; its records must not be lost. + with tempfile.TemporaryDirectory() as d: + path = os.path.join(d, "pod.log") + with open(path, "w") as f: + f.write("[INFO] " + json.dumps(ATELET_RESTORE) + "\n" + "[2026-09-23 10:00:05] " + json.dumps(ATEOM_RESTORE) + "\n") + err = io.StringIO() + with contextlib.redirect_stderr(err): + out = phase_report.parse_files([path]) + self.assertEqual(len(out.breakdowns), 2) + self.assertIn("reading it line by line", err.getvalue()) def test_prefixed_kubectl_lines_still_parse(self): with tempfile.TemporaryDirectory() as d: