Files
Nishanth Kotla 3f789b309a benchmarking/locust: harvest server-side telemetry from Prometheus (#1725)
Part of #1590

## What this PR does

Locust measures the client side only. This adds a post-run harvest of
server-side ground truth from Prometheus, written to a new
`server_summary.json` and summarized as one row in `stats.jsonl`.

## Proposed Changes

### Steady-state scoping

`server_telemetry.py` derives the steady-state window from
`stats_history.csv`: it starts at the first sample at 90% of peak user
count and ends at the last one. Ramp-up is excluded from every metric,
and the snapshot block also stops at the last full-load sample, so
teardown suspends are left out. Under a step-ladder load shape the
window covers only the top step.

### Harvested metrics

**Cluster packing**, from `ate_workerpool_workers`: busy workers
(`partial` + `at_capacity`) over the pool size at each sample, where the
pool size is the sum of all worker states, so it follows scale-ups and
scale-downs mid-run. Only the live ateapi is read: after a redeploy the
collector keeps re-exporting exited ateapi processes' last values for a
few minutes, which would otherwise inflate the counts. Reported as min,
p50, p90, p95, p99, max and mean (`avg`), plus the underlying timeseries
so transient spikes remain visible.

**Kernel pressure**, from cAdvisor PSI: CPU, memory and IO stall
percentages for the node and for the worker pods, each as the same
percentile set over the window.

**Snapshots**: size mean and p50/p90/p95/p99, checkpoint counts both
in-window and cumulative, restore and checkpoint latency mean and
p50/p90/p95/p99 from the AteomHerder RPC histograms, and checkpoint
throughput as `checkpoint_mb_s`. The atelet exports metrics on an
interval, so its data reaches Prometheus late: the harvest waits
`--atelet-lag-s` (default 70s, enough for the OTel SDK's 60s default
export and a 10s scrape) and reads both window edges half that late.

### Constraints and failure behavior

The module uses only the standard library, since the locust image is
distroless, and every request carries a timeout. An unreachable
Prometheus records nulls and does not fail the run. `--prometheus-url`
overrides the in-cluster default, and `--atelet-lag-s` sets the wait for
the atelet's last export.

Unmeasured fields are `null` and a measured zero is `0`, consistent with
the rest of the runner. Prometheus exposes no byte counter on the
restore path, so no restore throughput field is emitted rather than
deriving one indirectly.

### Output

`server_summary.json` holds the full nested artifact. `stats.jsonl`
receives a single `server_summary` row with 45 flat keys for graphing:
packing, node PSI and pod PSI at p50/p90/p95/p99, plus the snapshot
means, percentiles, `checkpoints_in_window` and `checkpoint_mb_s`.

`status.json` is unchanged.

## How this was tested

14 unit tests in `test_server_telemetry.py` covering steady-state
detection, percentile boundaries, the range-query window guard,
malformed Prometheus responses, the packing and checkpoint arithmetic,
the per-sample pool size, ignoring exited ateapi series, a missing
denominator returning null, snapshot fields returning null rather than
zero, and telemetry surviving a missing stats CSV.

Verified on 2 user / 2 worker and 4 user / 2 worker sympy runs. All 45
keys matched an independent recomputation from raw Prometheus, and the
snapshot block was cross-checked against atelet logs. A run with
`--prometheus-url` pointed at an unreachable address completes normally
with the affected fields null.

Re-harvested 5 past runs (Glutton and SWE-perf, 3 × 60 and 30 × 8
workers) from Prometheus with the new code. Runs with a steady pool
match the previous output exactly, and runs that followed a redeploy now
read the real pool size on every sample.

## References

[Agent Substrate: Actor Density Benchmark
Specs](https://docs.google.com/document/d/1sv5aBXvGOQ69iaqFxdPZ-d6tNbh5AMuwKODDyYjzYQw/edit?tab=t.0)

- [x] Tests pass
- [x] Appropriate changes to documentation are included in the PR
2026-10-01 16:59:41 +00:00

595 lines
22 KiB
Python

# 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.
"""Headless locust runner that publishes results as JSONL + CSVs.
Runs locust with the given flags, converts the resulting stats CSV
to JSONL, and uploads everything to either GCS or local disk under
<dest>/runs/<tag>/<timestamp>/.
When the test target is glutton.py, also spawns the boomer-worker Go
worker as a subprocess (locust runs in --master + --expect-workers=1
mode) so the GluttonUser load comes from boomer instead of Python+gevent.
-f accepts a comma-separated list of files, which locust reads as one test.
Give the test file first, then a load shape, as in
`/app/tests/glutton.py,/app/shapes/ladder_shape.py`. A shape holds no user
class, thus this module finds the test in the list; see test_file().
"""
import argparse
import csv
import json
import os
import re
import shutil
import signal
import subprocess
import sys
import threading
import time
from datetime import datetime, timezone
from pathlib import Path
from typing import IO, Any, TextIO
from cluster_facts import (
EMPTY_FACTS,
append_trial_summary,
get_cluster_hardware_facts,
)
from common.boomer_config import build_config_json
from server_telemetry import extract_and_record_server_telemetry
# Path inside the locust image to the boomer-worker binary baked in by
# benchmarking/locust/Dockerfile.
BOOMER_BINARY = "/app/boomer-worker"
# Port for the headless /boomer-config server (common/boomer_config.py), which
# gives boomer the values that change while a run continues. Locust already
# holds 5557 (master) and 8089 (web UI) in this container.
BOOMER_CONFIG_PORT = 5560
# In-cluster Prometheus that benchmarking/monitoring.yaml deploys. Override
# with --prometheus-url, for example when port-forwarding to a local run.
DEFAULT_PROMETHEUS_URL = "http://prometheus.benchmarking.svc.cluster.local:9090"
# Tab-separated columns written to traces.txt. Order matters — readers split
# on \t and index positionally.
TRACE_COLUMNS = ("time", "name", "duration_ms", "latency_source", "trace_id", "err")
# Python locust per-trace log line. Emitted by common/grpc_tracing.py and
# tests/counter_demo.py as:
# Traced {name}[ (failed)]: trace_id={32hex}, duration_ms={float} ({src})
PY_TRACE_RE = re.compile(
r"Traced\s+(?P<name>\S+)(?:\s+\(failed\))?:\s*"
r"trace_id=(?P<trace_id>[0-9a-f]{32}),\s*"
r"duration_ms=(?P<duration>[0-9.]+)\s+"
r"\((?P<source>\w+)\)"
)
def parse_args() -> argparse.Namespace:
p = argparse.ArgumentParser(description=__doc__)
p.add_argument("-f", required=True, dest="file", help="Locust test file (-f)")
p.add_argument("-t", required=True, dest="duration", help="Run duration (-t)")
p.add_argument(
"-u", required=True, type=int, dest="users", help="Number of users (-u)"
)
p.add_argument("--tag", required=True, help="Tag for this run")
p.add_argument(
"--name", required=True, help="Name for this run; used as locust --csv prefix"
)
p.add_argument(
"--dest",
required=True,
help="Root destination (gs://bucket/path or local path)",
)
p.add_argument(
"--allow-empty-stats",
action="store_true",
help=(
"Exit 0 when locust produced no measurement rows. For a run whose "
"result is not in the locust statistics, such as a run that sends "
"no request on purpose"
),
)
# Boomer startup flag; not a locust or dynconfig flag, so runner.py
# parses it here (keeping it out of locust_extra, which locust would
# reject) and appends it to the boomer command line in run_test().
p.add_argument(
"--actors-per-user",
type=int,
default=None,
help=(
"Number of actors each GluttonUser creates and cycles through in "
"round-robin (iteration i targets actor i%%N). Forwarded to "
"boomer-glutton as --actors-per-user. Omit to keep boomer's "
"default of 1."
),
)
p.add_argument(
"--cluster-facts",
action=argparse.BooleanOptionalAction,
default=True,
help=(
"Read node capacity and worker pod count from the Kubernetes API "
"after the run to derive density frontiers. Pass "
"--no-cluster-facts to skip Kubernetes API discovery"
),
)
p.add_argument(
"--prometheus-url",
default=DEFAULT_PROMETHEUS_URL,
help=(
"Prometheus to harvest server-side telemetry from after the run. "
"An unreachable Prometheus is not an error: the affected fields "
"are recorded as null"
),
)
p.add_argument(
"--atelet-lag-s",
type=int,
default=70,
help=(
"Seconds to wait after the run, and to shift the snapshot window "
"by, so the atelet's last export is scraped. Must cover the "
"atelet's OTEL_METRIC_EXPORT_INTERVAL plus one Prometheus scrape; "
"the default fits the OTel SDK's 60s default export and a 10s scrape"
),
)
args, extra = p.parse_known_args()
args.locust_extra = extra
return args
# Tests whose User class is implemented in Python. Closed set: new load
# generators are written in boomer, so anything else runs on boomer. Remove
# entries as these are retired; when empty, drop needs_boomer entirely.
PYTHON_TESTS = frozenset({
"ate_api.py",
"counter_demo.py",
"sleep.py",
"usermem.py",
"kernelmem.py",
})
def test_file(files: str) -> str:
"""Return the one file of `files` that holds the user class.
`files` is the value of -f, which is one path or a comma-separated list
of paths. A load shape is a file in that list with no user class, thus
the name of the run does not come from it. Two values are wrong if this
function takes the last path: needs_boomer below, and the --user-class
of boomer.
Each shape is a file in the shapes directory, thus remove those paths
first. A boomer test has no entry in PYTHON_TESTS, thus the test cannot
be found by that set alone. What stays is the test, and the result does
not change with the order of the paths.
"""
paths = [p for p in files.split(",") if p]
tests = [p for p in paths if os.path.basename(os.path.dirname(p)) != "shapes"]
if tests:
return tests[0]
return paths[0] if paths else files
def needs_boomer(files: str) -> bool:
return os.path.basename(test_file(files)) not in PYTHON_TESTS
def pair_flags(locust_extra: list[str]) -> list[str]:
"""Join each flag of `locust_extra` with its value, for the log.
tests.yaml gives a flag and its value as two entries of a list, thus one
flag becomes two items here. A reader compares the log with tests.yaml,
and a value on a line of its own makes that comparison difficult.
"""
pairs: list[str] = []
i = 0
while i < len(locust_extra):
item = locust_extra[i]
nxt = locust_extra[i + 1] if i + 1 < len(locust_extra) else None
if item.startswith("--") and nxt is not None and not nxt.startswith("--"):
pairs.append(f"{item} {nxt}")
i += 2
continue
pairs.append(item)
i += 1
return pairs
def tee(logs: TextIO, msg: str) -> None:
print(msg, flush=True)
logs.write(msg + "\n")
logs.flush()
def log_run_config(args: argparse.Namespace, dest_prefix: str, work_dir: Path, logs: TextIO) -> None:
"""Emit a structured summary of the test config at the top of every run
so anyone reading logs.txt later can see exactly what was executed
without cross-referencing tests.yaml + the orchestrator's invocation."""
lines = [
"==== Run config ====",
f" name: {args.name}",
f" tag: {args.tag}",
f" files: {args.file}",
f" test_file: {test_file(args.file)}",
f" duration: {args.duration}",
f" users: {args.users}",
f" uses_boomer: {needs_boomer(args.file)}",
f" empty stats ok: {args.allow_empty_stats}",
f" dest_prefix: {dest_prefix}",
f" work_dir: {work_dir}",
# One flag for each line, with its value. A run gives many flags to
# locust, and one line holds them in a form that no reader can
# compare with the entry in tests.yaml.
" extra flags:",
]
lines += [f" {flag}" for flag in pair_flags(args.locust_extra)] or [" (none)"]
lines.append("====================")
for line in lines:
tee(logs, line)
def extract_trace_record(prefix: str, line: str) -> dict[str, str] | None:
"""Parse `line` into a trace record dict (TRACE_COLUMNS keys) if it
describes a sampled span, else return None. Handles boomer slog JSON
lines (msg starts with 'traced span') and Python locust 'Traced ...:'
free-form lines. Failed-span fields land in `err`."""
if prefix == "boomer":
if "trace_id" not in line:
return None
try:
obj = json.loads(line)
except ValueError:
return None
if not isinstance(obj, dict) or not str(obj.get("msg", "")).startswith("traced span"):
return None
if "trace_id" not in obj:
return None
duration = obj.get("duration_ms")
return {
"time": str(obj.get("time", "")),
"name": str(obj.get("name", "")),
"duration_ms": "" if duration is None else f"{float(duration):.3f}",
"latency_source": str(obj.get("source", "")),
"trace_id": str(obj["trace_id"]),
"err": str(obj.get("err", "")),
}
m = PY_TRACE_RE.search(line)
if not m:
return None
return {
"time": "",
"name": m.group("name"),
"duration_ms": m.group("duration"),
"latency_source": m.group("source"),
"trace_id": m.group("trace_id"),
"err": "failed" if "(failed)" in line else "",
}
def pump_stream(prefix: str, stream: IO[str], logs: TextIO, traces: TextIO) -> None:
"""Forward each line of `stream` to stdout + logs (with a per-source
prefix) and append any extracted trace records to `traces` as TSV rows."""
for line in stream:
line = line.rstrip("\n")
tagged = f"[{prefix}] {line}"
sys.stdout.write(tagged + "\n")
sys.stdout.flush()
logs.write(tagged + "\n")
logs.flush()
record = extract_trace_record(prefix, line)
if record is not None:
traces.write("\t".join(record[c] for c in TRACE_COLUMNS) + "\n")
traces.flush()
def run_test(args: argparse.Namespace, csv_prefix: Path, logs: TextIO, traces: TextIO) -> int:
"""Run locust (and boomer, when needed). Returns locust's exit code.
Stdout from each subprocess is forwarded to logs.txt with a `[locust]` /
`[boomer]` prefix so they're distinguishable; trace_id matches are
siphoned into traces.txt as a deduped one-per-line list.
"""
with_boomer = needs_boomer(args.file)
locust_cmd = [
sys.executable, "-m", "locust",
"--headless",
"-f", args.file,
"-t", args.duration,
"-u", str(args.users),
"--csv", str(csv_prefix),
]
if with_boomer:
# Master mode so boomer can connect as a worker on localhost:5557.
# --expect-workers=1 makes locust wait for boomer before starting.
locust_cmd += ["--master", "--expect-workers", "1"]
# Serve /boomer-config for the worker. --config-json gives boomer the
# values one time, at its start, thus a run that changes a value while
# it runs needs the endpoint. A load shape is one such run.
locust_cmd += ["--boomer-config-port", str(BOOMER_CONFIG_PORT)]
locust_cmd += list(args.locust_extra)
tee(logs, f"Running: {' '.join(locust_cmd)}")
locust_proc = subprocess.Popen(
locust_cmd,
stdout=subprocess.PIPE,
stderr=subprocess.STDOUT,
bufsize=1,
text=True,
)
pumps = []
pumps.append(threading.Thread(
target=pump_stream,
args=("locust", locust_proc.stdout, logs, traces),
daemon=True,
))
boomer_proc = None
if with_boomer:
boomer_cmd = [BOOMER_BINARY, "--user-class", Path(test_file(args.file)).stem]
cfg_json = build_config_json(args.locust_extra)
if cfg_json:
boomer_cmd += ["--config-json", cfg_json]
if args.actors_per_user is not None:
boomer_cmd += ["--actors-per-user", str(args.actors_per_user)]
# Read the endpoint again at each spawn message. Thus a value that
# changes while the run continues, such as the sample rate of a load
# shape, reaches boomer at the change. boomer's --master-host default
# is the loopback address, where the server above listens.
boomer_cmd += ["--master-web-port", str(BOOMER_CONFIG_PORT)]
tee(logs, f"Running: {' '.join(boomer_cmd)}")
boomer_proc = subprocess.Popen(
boomer_cmd,
stdout=subprocess.PIPE,
stderr=subprocess.STDOUT,
bufsize=1,
text=True,
)
pumps.append(threading.Thread(
target=pump_stream,
args=("boomer", boomer_proc.stdout, logs, traces),
daemon=True,
))
for t in pumps:
t.start()
locust_exit = locust_proc.wait()
tee(logs, f"Locust exited with code {locust_exit}")
if boomer_proc is not None:
# Locust finishing means the test window is over; let boomer drain
# its actor cleanup before tearing it down hard. The boomer process
# suspends+deletes every actor it created on SIGTERM.
tee(logs, "Stopping boomer...")
boomer_proc.send_signal(signal.SIGTERM)
try:
boomer_proc.wait(timeout=90)
except subprocess.TimeoutExpired:
tee(logs, "Boomer did not exit within 90s; killing")
boomer_proc.kill()
boomer_proc.wait()
tee(logs, f"Boomer exited with code {boomer_proc.returncode}")
for t in pumps:
t.join(timeout=5)
return locust_exit
def stats_to_jsonl(stats_csv: Path, jsonl_path: Path, timestamp: str, tag: str, test_name: str) -> int:
rows_written = 0
with open(stats_csv) as f, open(jsonl_path, "w") as out:
reader = csv.DictReader(f)
for row in reader:
type_val = row.pop("Type", "") or ""
name_val = row.pop("Name", "") or ""
if name_val == "Aggregated":
continue
measurements = {}
for k, v in row.items():
if k is None:
continue
if k.endswith("%"):
# Percentile columns: "50%" -> "p50", "99.99%" -> "p99_99"
key = "p" + k[:-1].replace(".", "_")
else:
# avoid non-alphanumeric characters
key = k.lower().replace("/", "_per_")
key = re.sub(r"[^a-z0-9]+", "_", key).strip("_")
measurements[key] = v
entry = {
"timestamp": timestamp,
"tag": tag,
"test_name": test_name,
"metric": f"{type_val}_{name_val}",
"measurements": measurements,
}
out.write(json.dumps(entry) + "\n")
rows_written += 1
return rows_written
def upload_to_gcs(local_path: Path, gcs_uri: str) -> None:
# Imported here so non-GCS use doesn't require google-cloud-storage.
from google.cloud import storage
bucket_name, _, blob_path = gcs_uri[len("gs://"):].partition("/")
storage.Client().bucket(bucket_name).blob(blob_path).upload_from_filename(
str(local_path)
)
def upload(src: Path, dest: str) -> None:
if dest.startswith("gs://"):
upload_to_gcs(src, dest)
else:
dest_path = Path(dest)
dest_path.parent.mkdir(parents=True, exist_ok=True)
shutil.copy(src, dest_path)
def collect_cluster_facts(
args: argparse.Namespace, logs: TextIO
) -> dict[str, Any]:
"""Returns cluster hardware facts, or empty facts when discovery is off."""
if not args.cluster_facts:
tee(logs, "Skipping cluster hardware discovery (--no-cluster-facts)")
return dict(EMPTY_FACTS)
return get_cluster_hardware_facts(logs)
def main() -> None:
args = parse_args()
now = datetime.now(timezone.utc)
# Path-safe timestamp for the local work dir (no Hive semantics here).
path_ts = now.strftime("%Y%m%dT%H%M%SZ")
# RFC 3339 / ISO 8601 extended for the JSONL data column
data_ts = now.strftime("%Y-%m-%dT%H:%M:%SZ")
# Hive partition values for the GCS layout
run_date = now.strftime("%Y-%m-%d")
run_ts = int(now.timestamp())
work_dir = Path(f"/tmp/{path_ts}-locust-runner")
work_dir.mkdir(parents=True, exist_ok=True)
csv_prefix = work_dir / args.name
stats_csv = work_dir / f"{args.name}_stats.csv"
jsonl_path = work_dir / f"{args.name}.jsonl"
logs_path = work_dir / f"{args.name}_logs.txt"
traces_path = work_dir / f"{args.name}_traces.txt"
status_path = work_dir / f"{args.name}_status.json"
server_summary_json = work_dir / f"{args.name}_server_summary.json"
prefix = (
f"{args.dest.rstrip('/')}/runs/{args.name}"
f"/run_date={run_date}/run_ts={run_ts}/run_tag={args.tag}"
)
with open(logs_path, "w") as logs, open(traces_path, "w") as traces:
traces.write("\t".join(TRACE_COLUMNS) + "\n")
traces.flush()
log_run_config(args, prefix, work_dir, logs)
exit_code = run_test(args, csv_prefix, logs, traces)
run_end_ts = int(datetime.now(timezone.utc).timestamp())
stats_generated = False
if stats_csv.exists():
try:
rows = stats_to_jsonl(
stats_csv, jsonl_path, data_ts, args.tag, args.name
)
if rows == 0:
tee(
logs,
f"Stats CSV {stats_csv} had no measurement rows; "
f"treating as not produced",
)
if jsonl_path.exists():
jsonl_path.unlink()
else:
stats_generated = jsonl_path.exists()
except Exception as e:
tee(logs, f"Failed to generate JSONL from {stats_csv}: {e}")
if jsonl_path.exists():
jsonl_path.unlink()
else:
tee(logs, f"Stats CSV {stats_csv} not produced; skipping JSONL")
# Density frontiers and server-side telemetry are additive. They are
# kept out of the block above so that a failure here cannot discard
# the measurements the trial actually came for.
stats_history_csv = work_dir / f"{args.name}_stats_history.csv"
# Seeded up front so that a discovery failure still leaves a usable
# value for the frontier summary below.
facts = dict(EMPTY_FACTS)
try:
facts = collect_cluster_facts(args, logs)
except Exception as e:
tee(logs, f"Warning: Failed to read cluster facts: {e}")
# The frontiers divide by user counts, so they need the CSV.
if stats_generated:
try:
append_trial_summary(
jsonl_path,
stats_csv,
stats_history_csv,
args,
data_ts,
facts,
logs,
)
except Exception as e:
tee(logs, f"Warning: Failed to record cluster facts: {e}")
# Server telemetry is a Prometheus time-window query, so it runs
# either way. A run too loaded to write a CSV is the one its
# bin-packing, PSI and snapshot numbers matter most for.
try:
extract_and_record_server_telemetry(
prom_url=args.prometheus_url,
start_ts=run_ts,
end_ts=run_end_ts,
stats_history_csv=stats_history_csv,
output_json_path=server_summary_json,
jsonl_path=jsonl_path,
data_ts=data_ts,
tag=args.tag,
test_name=args.name,
logs=logs,
lag_s=args.atelet_lag_s,
)
except Exception as e:
tee(logs, f"Warning: Failed to harvest server telemetry: {e}")
status_path.write_text(
json.dumps(
{"locust_exit_code": exit_code, "stats_generated": stats_generated}
)
)
files: list[tuple[Path, str]] = [
(status_path, "status.json"),
(logs_path, "logs.txt"),
(traces_path, "traces.txt"),
(jsonl_path, "stats.jsonl"),
(stats_csv, "stats.csv"),
(work_dir / f"{args.name}_exceptions.csv", "exceptions.csv"),
(work_dir / f"{args.name}_failures.csv", "failures.csv"),
(work_dir / f"{args.name}_stats_history.csv", "stats_history.csv"),
(server_summary_json, "server_summary.json"),
# TODO: remove after data migration
(jsonl_path, f"{args.name}.jsonl"),
]
for src, basename in files:
if not src.exists():
print(f"Skipping {src}: not produced", flush=True)
continue
dest = f"{prefix}/{basename}"
upload(src, dest)
print(f"Uploaded {src} -> {dest}", flush=True)
if not stats_generated and not args.allow_empty_stats:
sys.exit(1)
if __name__ == "__main__":
main()