meteor-p150 / code /scripts /bench.py
changh95's picture
tt-model push meteor-p150 (container)
51defdc verified
Raw History Blame Contribute Delete
14.4 kB
#!/usr/bin/env python3
# SPDX-License-Identifier: Apache-2.0
"""Stage bench of warm forwards: the numbers OPT_BASELINE.md, OPT_REPORT.md and the card quote.
bin/devrun -t 1800 -- python code/scripts/bench.py --iters 60 --inputs sample,synthetic --json out.json
bin/devrun -t 1800 -- python code/scripts/bench.py --dispatch worker --num-cqs 2 --json ... # the D14 matrix
Inputs (``--inputs``, comma separated; each one is a request of ``model(...)``, batch 1):
- ``sample``: the PandaSet 019 frame 40 sample ``pandaset_019_f40.json`` (the baseline's sample: the git-ignored
``staging_samples_pandaset/`` copy until the user approves shipping it, else ``samples/``);
- ``shipped``: the shipped synthetic sample ``samples/synthetic_8cam.json`` (``make_synthetic_sample.py``);
- ``synthetic``: eight 768x432 uint8 noise images made here (seeded; data generated by us), the ``pandaset_019``
calibration (inline from the staged preset; the ``synthetic_8cam`` preset when it is absent), ego speed 10 m/s;
- ``<name>=<manifest.json>`` or a manifest path: any ``load_sample`` manifest (e.g. a public-dataset frame).
Images are decoded when the request is loaded, before timing (the camera driver hands METEOR decoded images). Per
input, after ``--warm`` calls (``ttaw.profiling.StageBench``: p50 / p99 / mean / min / max in ms):
1. **plain**: ``model(...)`` exactly as a client calls it: ``e2e`` and the API's own ``timing_ms``
(``api.preprocess`` / ``api.device`` / ``api.postprocess``), with ``run_frame`` split by wrapping its parts:
``dev.tables`` (rig check, and the lift-table write when the rig changed), ``dev.image_rows`` (NCHW -> NHWC rows),
``dev.runner`` (upload + replay + segmented read), ``dev.unpack`` (join of the readback segments),
``dev.outputs`` (``outputs_from_device``: device layouts -> the 19 ONNX outputs on the host);
2. **stages**: ``ttaw.profiling.bench_trace_runner`` on the ``frame`` variant: ``host_in`` (numpy -> ttnn host
tensor), ``h2d`` (upload + synchronize), ``trace`` (one replay + synchronize), ``d2h`` (segmented read), ``post``
(``unpack``), ``e2e`` (runner call + unpack) and ``b2b`` (back-to-back replays: the device time per frame).
``--rig-switch N``: N calls alternating between the first two inputs (each call writes its rig's lift tables).
An agreement check of every input with a stored CPU reference (``<stem>.reference.json``, fresh stream, the container
smoke's ``compare_with_reference``) runs first, so every configuration's numbers come with its accuracy. AICLK /
power / temperature are sampled during the loops (``AiclkSampler``).
"""
from __future__ import annotations
import argparse
import json
import os
import platform
import subprocess
import sys
import time
from pathlib import Path
from typing import Any, Dict, List, Optional, Tuple
import numpy as np
HERE = Path(__file__).resolve()
CODE = HERE.parents[1]
sys.path.insert(0, str(CODE))
SAMPLES = CODE / "tt_meteor" / "samples"
STAGING = CODE.parent / "staging_samples_pandaset"
ROOT = Path(os.environ.get("TT_MODELS_ROOT", "/home/ubuntu/experiments/tt-models"))
SAMPLE_NAME = "pandaset_019_f40.json"
# ------------------------------------------------------------------------------------------------- inputs
def sample_path() -> Path:
for d in (SAMPLES, STAGING / "samples", STAGING):
if (d / SAMPLE_NAME).is_file():
return d / SAMPLE_NAME
return SAMPLES / SAMPLE_NAME
def synthetic_request(seed: int = 1) -> Dict[str, Any]:
from tt_meteor.reference import config as C
rng = np.random.default_rng(seed)
images = {cam: rng.integers(0, 256, size=(C.IMG_H, C.IMG_W, 3), dtype=np.uint8) for cam in C.CAMERAS}
staged = [d / "pandaset_019.json" for d in (CODE / "tt_meteor" / "calib", STAGING / "calib") if
(d / "pandaset_019.json").is_file()]
calib = json.loads(staged[0].read_text()) if staged else {"preset": "synthetic_8cam"}
return {"images": images, "calibration": calib, "ego_speed": 10.0,
"stream": {"id": "synthetic", "pose": [0.0, 0.0, 0.0]}}
def load_input(spec: str) -> Tuple[str, Dict[str, Any], Optional[Path]]:
from tt_meteor import load_sample
if spec == "synthetic":
return "synthetic", synthetic_request(), None
if spec == "sample":
path = sample_path()
return "sample", load_sample(path), path
if spec == "shipped":
path = SAMPLES / "synthetic_8cam.json"
return "shipped", load_sample(path), path
name, _, p = spec.partition("=")
path = Path(p or name)
return (name if p else path.stem), load_sample(path), path
def find_reference(manifest: Optional[Path]) -> Optional[Path]:
if manifest is None:
return None
p = manifest.parent / f"{manifest.stem}.reference.json"
return p if p.is_file() else None
# --------------------------------------------------------------------------------------------- environment
def _git(*args: str) -> str:
try:
return subprocess.run(["git", *args], capture_output=True, text=True,
env=dict(os.environ, GIT_OPTIONAL_LOCKS="0")).stdout.strip()
except OSError:
return ""
def environment() -> Dict[str, Any]:
env: Dict[str, Any] = {"host": platform.node(), "python": platform.python_version(),
"loadavg_start": os.getloadavg(), "cpus": os.cpu_count(),
"omp_num_threads": os.environ.get("OMP_NUM_THREADS"),
"env": {k: v for k, v in os.environ.items() if k.startswith("METEOR_")}}
try:
env["cpu_model"] = next(line.split(":", 1)[1].strip() for line in open("/proc/cpuinfo")
if line.startswith("model name"))
except (OSError, StopIteration):
pass
tm = ROOT / "tt-metal"
env["tt_metal_commit"] = _git("-C", str(tm), "rev-parse", "HEAD")
try:
vend = json.loads((CODE / "tt_meteor/ttaw/VENDORED.json").read_text())
env["ttaw"] = vend.get("version")
env["ttaw_vendored"] = vend.get("source_commit")
except (OSError, ValueError):
pass
env["bundle_commit"] = _git("-C", str(CODE.parent), "rev-parse", "HEAD")
env["bundle_dirty"] = bool(_git("-C", str(CODE.parent), "status", "--porcelain", "code"))
return env
def cache_entries(device) -> Optional[int]:
fn = getattr(device, "num_program_cache_entries", None)
try:
return int(fn()) if fn else None
except Exception: # noqa: BLE001
return None
# ------------------------------------------------------------------------------------------------- measure
class Split:
"""Times the parts of ``TtMETEOR.run_frame`` (module docstring) while ``bench`` is set."""
def __init__(self, model):
import tt_meteor.tt.model as tmodel
self.bench = None
tt = model.tt
self._undo = []
def timed(name, fn):
def wrapper(*a, **k):
if self.bench is None:
return fn(*a, **k)
t0 = time.perf_counter()
try:
return fn(*a, **k)
finally:
self.bench.add(name, (time.perf_counter() - t0) * 1e3)
return wrapper
for attr, name in (("set_tables", "dev.tables"), ("image_rows", "dev.image_rows"), ("unpack", "dev.unpack"),
("runner", "dev.runner")):
orig = getattr(tt, attr)
setattr(tt, attr, timed(name, orig) if attr != "runner" else _CallTimer(orig, self, name))
self._undo.append(lambda tt=tt, attr=attr, orig=orig: setattr(tt, attr, orig))
orig_out = tmodel.outputs_from_device
tmodel.outputs_from_device = timed("dev.outputs", orig_out)
self._undo.append(lambda: setattr(tmodel, "outputs_from_device", orig_out))
def remove(self) -> None:
for fn in reversed(self._undo):
fn()
class _CallTimer:
"""Stands in for the TraceRunner inside ``run_frame``: times ``runner(...)``, forwards every attribute."""
def __init__(self, inner, split: Split, name: str):
self._inner, self._split, self._name = inner, split, name
def __call__(self, *a, **k):
if self._split.bench is None:
return self._inner(*a, **k)
t0 = time.perf_counter()
try:
return self._inner(*a, **k)
finally:
self._split.bench.add(self._name, (time.perf_counter() - t0) * 1e3)
def __getattr__(self, item):
return getattr(self._inner, item)
def agreement(model, name: str, req: Dict[str, Any], ref_path: Path) -> Dict[str, Any]:
sys.path.insert(0, str(CODE / "tt_meteor" / "ttaw" / "server"))
import smoke as ttaw_smoke # noqa: E402
fresh = dict(req)
fresh["stream"] = dict(req.get("stream") or {}, id=f"agreement-{name}-{time.time_ns()}")
body = model(**fresh).to_dict()
ref = json.loads(ref_path.read_text())
metrics, failures = ttaw_smoke.compare_with_reference(body, ref)
path_dev = float(np.linalg.norm(np.asarray(body["trajectory"]) - np.asarray(ref["trajectory"]), axis=-1).max())
return {"reference": str(ref_path), "metrics": metrics, "failures": failures, "ego.path_dev_m": path_dev,
"pass": not failures}
def bench_input(model, split: Split, name: str, req: Dict[str, Any], iters: int, warm: int,
b2b: int) -> Dict[str, Any]:
from tt_meteor.host.inputs import prepare_request
from tt_meteor.ttaw.profiling import AiclkSampler, StageBench, bench_trace_runner
res: Dict[str, Any] = {"loadavg_start": os.getloadavg()}
for _ in range(warm):
model(**req)
plain = StageBench(f"meteor {name} plain")
with AiclkSampler(interval_s=0.05) as clk_p:
split.bench = plain
try:
for _ in range(iters):
t0 = time.perf_counter()
out = model(**req)
plain.add("e2e", (time.perf_counter() - t0) * 1e3)
for k, v in out.timing_ms.items():
plain.add(f"api.{k}", v)
finally:
split.bench = None
res["plain"] = plain.summary()
res["aiclk_plain"] = clk_p.summary()
res["boxes3d"] = len(out.boxes3d)
print(plain.table(), flush=True)
frame = prepare_request(req["images"], req["calibration"], req["ego_speed"], req.get("stream"))
model.tt.set_tables(model.calib_cache.get(frame.K, frame.T_cam_ego))
tt = model.tt
with AiclkSampler(interval_s=0.05) as clk_s:
stages = bench_trace_runner(model.runner, "frame", {"imgs": tt.image_rows(frame.imgs)},
params=tt.frame_params(frame), iters=iters, warmup=warm,
name=f"meteor {name} stages", post=tt.unpack, b2b_iters=b2b)
res["stages"] = stages.summary()
res["aiclk_stages"] = clk_s.summary()
print(stages.table(), flush=True)
res["loadavg_end"] = os.getloadavg()
return res
def rig_switch(model, split: Split, reqs: List[Tuple[str, Dict[str, Any]]], n: int) -> Dict[str, Any]:
from tt_meteor.ttaw.profiling import StageBench
bench = StageBench("meteor rig switch")
split.bench = bench
try:
for i in range(n):
name, req = reqs[i % 2]
t0 = time.perf_counter()
out = model(**req)
bench.add("e2e", (time.perf_counter() - t0) * 1e3)
for k, v in out.timing_ms.items():
bench.add(f"api.{k}", v)
finally:
split.bench = None
print(bench.table(), flush=True)
return {"inputs": [r[0] for r in reqs[:2]], "calls": n, **{"stages": bench.summary()}}
def main() -> int:
ap = argparse.ArgumentParser(description=__doc__.split("\n\n")[0])
ap.add_argument("--inputs", default="sample", help="comma list: sample | synthetic | <name>=<manifest> | path")
ap.add_argument("--iters", type=int, default=60)
ap.add_argument("--warm", type=int, default=3)
ap.add_argument("--b2b", type=int, default=20, help="back-to-back replays of the device-only row")
ap.add_argument("--rig-switch", type=int, default=0, help="calls alternating between the first two inputs")
ap.add_argument("--dispatch", default=None, choices=["eth", "worker"])
ap.add_argument("--num-cqs", type=int, default=None, choices=[1, 2])
ap.add_argument("--no-agreement", dest="agreement", action="store_false")
ap.add_argument("--json")
a = ap.parse_args()
from tt_meteor import METEOR
inputs = [load_input(s.strip()) for s in a.inputs.split(",") if s.strip()]
res: Dict[str, Any] = {"args": vars(a), "environment": environment(),
"inputs": {n: str(p) if p else "synthetic" for n, _, p in inputs}}
t_load = time.perf_counter()
with METEOR.from_pretrained(dispatch=a.dispatch, num_command_queues=a.num_cqs) as model:
res["load_s"] = time.perf_counter() - t_load
res["warmup_ms"] = model.warmup_ms
res["config"] = model.device_info
res["program_cache_entries"] = cache_entries(model.device)
res["trace"] = model.runner.describe()
print("config:", json.dumps(res["config"]), "| load %.1f s" % res["load_s"], "| programs:",
res["program_cache_entries"], flush=True)
if a.agreement:
res["agreement"] = {}
for name, req, path in inputs:
ref = find_reference(path)
if ref is not None:
res["agreement"][name] = agreement(model, name, req, ref)
print("agreement", name, json.dumps(res["agreement"][name], default=str), flush=True)
split = Split(model)
res["per_input"] = {}
for name, req, _ in inputs:
res["per_input"][name] = bench_input(model, split, name, req, a.iters, a.warm, a.b2b)
if a.rig_switch and len(inputs) >= 2:
res["rig_switch"] = rig_switch(model, split, [(n, r) for n, r, _ in inputs], a.rig_switch)
split.remove()
res["program_cache_entries_end"] = cache_entries(model.device)
res["loadavg_end"] = os.getloadavg()
text = json.dumps(res, indent=1, default=str)
if a.json:
Path(a.json).parent.mkdir(parents=True, exist_ok=True)
Path(a.json).write_text(text + "\n")
print(text)
return 0
if __name__ == "__main__":
sys.exit(main())