mirror of
https://github.com/agent-substrate/substrate.git
synced 2026-10-02 03:24:42 +08:00
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
595 lines
22 KiB
Python
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()
|