File size: 12,338 Bytes
4d9b003 be62f78 4d9b003 be62f78 4d9b003 be62f78 4d9b003 be62f78 4d9b003 be62f78 4d9b003 be62f78 4d9b003 be62f78 4d9b003 be62f78 4d9b003 be62f78 4d9b003 be62f78 4d9b003 be62f78 4d9b003 be62f78 4d9b003 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 | #!/usr/bin/env python3
# SPDX-License-Identifier: Apache-2.0
"""Stage breakdown of warm plans -- the numbers OPT_BASELINE.md / OPT_REPORT.md / the card quote.
bin/devrun -t 900 -- python code/scripts/bench.py --iters 100 --json out.json
bin/devrun -t 900 -- python code/scripts/bench.py --dispatch worker --num-cqs 2 --iters 100 --json out.json
bin/devrun -t 900 -- python code/scripts/bench.py --input <a.npz> --input <b.npz> ...
``--input`` takes the planner tensors as an ``.npz`` with the 15 ``INPUT_SCHEMA`` names (the shipped samples) or with
``raw/<name>`` keys (the research / public-data scene files); default: the shipped ``kashiwanoha_dense.npz``.
Per input, ``ttaw.profiling.StageBench`` collects ``--iters`` warm iterations of each stage (p50 / p99 / mean / min):
- ``load``: decoding the ``.npz`` + the ``INPUT_SCHEMA`` check (``model(inputs=<path>)`` does it outside timing_ms);
- ``host_pre``: the node's pre-processing (``host.prepare``: normalization, speed masks, encoder host features,
decoder masks, the solver's initial state);
- ``pack``: the 19 persistent trace inputs (``tt.inputs.plan_inputs``);
- ``host_in``: their ttnn host tensors (fp32 / bf16 TILE, ``ttaw.tensors.to_host_tensor``);
- ``h2d``: the upload into the persistent device inputs + device sync;
- ``trace``: one replay of the ``plan`` trace + device sync (the device latency of one plan);
- ``d2h``: the one packed readback (``final_x0`` + logits + the ego rows of the 11 iterates);
- ``host_post``: the node's post-processing (``host.make_output``);
- ``e2e``: ``model(inputs=<decoded arrays>)``, the in-process API call (schema check included);
- ``e2e_path``: ``model(inputs=<path>)`` (adds ``load``; schema-named ``.npz`` files only);
- ``b2b``: back-to-back replays with no host work in between (device time per plan), ``--b2b-rounds`` rounds of
``--b2b-iters`` replays; ``plans_per_s`` = 1000 / median b2b.
Also: ``model(...).timing_ms`` (preprocess / device / postprocess / total), the first call after ``from_pretrained``
and the load time (weights, build, warm-up + capture), AICLK / power / temperature sampled from sysfs during the
timed loops (``ttaw.profiling.AiclkSampler``), the staged path checked bit for bit against ``model()``, the device
configuration (dispatch, CQs, grid) and the numerics options in effect (``DIFFUSION_PLANNER_*`` knobs). Always quote
the configuration line with the numbers.
"""
from __future__ import annotations
import argparse
import json
import statistics
import time
from pathlib import Path
from typing import Any, Dict, Optional
import numpy as np
from tt_diffusion_planner import DiffusionPlanner
SAMPLE = Path(__file__).resolve().parents[1] / "tt_diffusion_planner" / "samples" / "kashiwanoha_dense.npz"
def load_scene(path: str) -> Dict[str, Any]:
"""``{"name", "path", "arrays", "schema_npz"}``: the 15 raw tensors of a schema-named or ``raw/``-prefixed npz."""
from tt_diffusion_planner.reference import config as C
p = Path(path)
with np.load(p, allow_pickle=False) as z:
files = set(z.files)
if all(k in files for k in C.INPUT_NAMES):
arrays, schema_npz = {k: np.array(z[k]) for k in C.INPUT_NAMES}, True
elif all(f"raw/{k}" in files for k in C.INPUT_NAMES):
arrays, schema_npz = {k: np.array(z[f"raw/{k}"]) for k in C.INPUT_NAMES}, False
else:
raise SystemExit(f"{p}: neither the INPUT_SCHEMA names nor raw/<name> keys")
name = p.stem[len("golden_"):] if p.stem.startswith("golden_") else p.stem
if p.parent.name not in ("samples", "ort"):
name = f"{p.parent.name}/{name}"
return {"name": name, "path": str(p), "arrays": arrays, "schema_npz": schema_npz}
def counts(arrays: Dict[str, np.ndarray]) -> Dict[str, int]:
"""Valid entities of a scene (non-empty rows), for the report."""
def rows(a, axis):
return int(np.any(np.abs(a) > 0, axis=axis).sum())
return {"neighbors": rows(arrays["neighbor_agents_past"][0], (1, 2)), "lanes": rows(arrays["lanes"][0], (1, 2)),
"route_lanes": rows(arrays["route_lanes"][0], (1, 2)), "polygons": rows(arrays["polygons"][0], (1, 2)),
"line_strings": rows(arrays["line_strings"][0], (1, 2))}
def unpack(out: Dict[str, np.ndarray], tt: Any) -> Dict[str, Any]:
"""The packed readback -> the raw outputs of ``TtDiffusionPlanner.forward`` (same reshapes)."""
from tt_diffusion_planner.reference import config as C
from tt_diffusion_planner.tt import config as T
final = tt.unpack(out)
steps = out["ego_steps"].reshape(-1, T.STATE_COLS)
return {"final_x0": final.reshape(C.MAX_NUM_AGENTS, C.OUTPUT_T + 1, C.POSE_DIM).astype(np.float32),
"logit": out["logit"].reshape(-1)[:C.TURN_INDICATOR_OUTPUT_DIM].astype(np.float32),
"denoising_steps": [s.reshape(1, C.OUTPUT_T + 1, C.POSE_DIM).astype(np.float32) for s in steps]}
def bench_scene(model, scene: Dict[str, Any], a: argparse.Namespace) -> Dict[str, Any]:
import ttnn
from tt_diffusion_planner.host import pipeline as hp
from tt_diffusion_planner.reference import config as C
from tt_diffusion_planner.tt import inputs as I
from tt_diffusion_planner.ttaw.io import load_named_arrays
from tt_diffusion_planner.ttaw.profiling import AiclkSampler, StageBench, time_b2b
from tt_diffusion_planner.ttaw.tensors import to_host_tensor
runner, dev = model.runner, model.device
arrays, params = scene["arrays"], model.validate_params({})
obs = model.normalization.observation
bench = StageBench(f"diffusion-planner {scene['name']}")
timing: Dict[str, list] = {}
variant = model.tt.variant_for(hp.prepare(load_named_arrays(arrays, C.INPUT_SCHEMA), obs)) # COMPACT bucket
sync = lambda: ttnn.synchronize_device(dev) # noqa: E731
for _ in range(a.warmup):
ref = model(inputs=arrays)
with AiclkSampler(chip=a.chip, interval_s=0.05) as clk:
for _ in range(a.iters): # the in-process API call
with bench.stage("e2e"):
ref = model(inputs=arrays)
for k, v in ref.timing_ms.items():
timing.setdefault(k, []).append(v)
slots = {k: runner._slot(k, "input") for k in I.INPUT_SPECS}
for _ in range(a.iters): # the same path, stage by stage
if scene["schema_npz"]:
with bench.stage("load"):
raw = load_named_arrays(scene["path"], C.INPUT_SCHEMA)
else:
raw = load_named_arrays(arrays, C.INPUT_SCHEMA)
with bench.stage("host_pre"):
prep = hp.prepare(raw, obs)
with bench.stage("pack"):
packed = model.tt.filter_inputs(I.plan_inputs(prep)) # INPUT_TRIM (as model() does)
with bench.stage("host_in"):
host = {k: to_host_tensor(v, slots[k].dtype, slots[k].layout, shape=slots[k].shape)
for k, v in packed.items()}
with bench.stage("h2d"):
runner.upload(host)
sync()
with bench.stage("trace"):
runner.replay(variant)
sync()
with bench.stage("d2h"):
out = runner.read(variant)
with bench.stage("host_post"):
res = model._postprocess(unpack(out, model.tt), prep, params)
if scene["schema_npz"]:
for _ in range(a.iters):
with bench.stage("e2e_path"):
model(inputs=scene["path"])
rounds = [time_b2b(lambda: runner.replay(variant), sync, n=a.b2b_iters, warmup=3)
for _ in range(a.b2b_rounds)]
for r in rounds:
bench.add("b2b", r)
if getattr(a, "dump", None):
Path(a.dump).parent.mkdir(parents=True, exist_ok=True)
flat: Dict[str, np.ndarray] = {}
def walk(prefix, v):
if isinstance(v, dict):
for k, w in v.items():
walk(f"{prefix}{k}/", w)
else:
flat[prefix.rstrip("/") or "out"] = np.asarray(v)
walk("", out)
np.savez(f"{a.dump}.{scene['name'].replace('/', '_')}.npz", **flat)
same = bool(np.array_equal(res.poses, ref.poses) and np.array_equal(res.predicted_agents, ref.predicted_agents)
and res.turn_indicator["command"] == ref.turn_indicator["command"])
summary = bench.summary()
b2b = statistics.median(rounds)
return {"name": scene["name"], "path": scene["path"], "valid": counts(arrays), "variant": variant, "stages_ms": summary,
"timing_ms": {k: {"p50": float(np.percentile(v, 50)), "p99": float(np.percentile(v, 99)),
"min": float(min(v))} for k, v in timing.items()},
"b2b_rounds_ms": rounds, "plans_per_s_b2b": 1000.0 / b2b, "aiclk": clk.summary(),
"staged_equals_model": same, "turn_command": int(ref.turn_indicator["command"]),
"table": bench.table()}
def main() -> None:
ap = argparse.ArgumentParser(description=__doc__.split("\n")[0])
ap.add_argument("--iters", type=int, default=100)
ap.add_argument("--warmup", type=int, default=5)
ap.add_argument("--b2b-iters", type=int, default=50)
ap.add_argument("--b2b-rounds", type=int, default=3)
ap.add_argument("--input", action="append", default=None, help="scene .npz (repeatable)")
ap.add_argument("--dispatch", default=None, choices=["eth", "worker"])
ap.add_argument("--num-cqs", type=int, default=None, choices=[1, 2])
ap.add_argument("--chip", type=int, default=0, help="sysfs chip index for the AICLK sampler")
ap.add_argument("--tag", default="")
ap.add_argument("--json")
ap.add_argument("--dump", help="prefix: save each scene's raw packed readback as <prefix>.<scene>.npz "
"(bit-identity checks of structural rewrites)")
a = ap.parse_args()
from tt_diffusion_planner.reference.weights import find_weights_dir
scenes = [load_scene(p) for p in (a.input or [str(SAMPLE)])]
wd = find_weights_dir() # None: from_pretrained resolves the pinned HF snapshot
t0 = time.perf_counter()
with DiffusionPlanner.from_pretrained(dispatch=a.dispatch, num_command_queues=a.num_cqs,
weights_dir=str(wd) if wd else None) as model:
load_s = time.perf_counter() - t0
t1 = time.perf_counter()
model(inputs=scenes[0]["arrays"]) # first call after from_pretrained (traces captured)
first_ms = (time.perf_counter() - t1) * 1e3
info = model.tt.describe()
res: Dict[str, Any] = {
"tag": a.tag, "config": model.device_info, "iters": a.iters, "load_s": round(load_s, 2),
"warmup_ms": {k: round(v, 1) for k, v in model.warmup_ms.items()}, "first_call_ms": round(first_ms, 2),
"options": info["options"], "precision": info["precision"], "uploaded_mb": info["uploaded_mb"],
"trace_buffers_mb": info["trace_buffers_mb"],
"program_cache_entries": info["trace"].get("program_cache_entries"),
"num_command_queues": info["trace"]["num_command_queues"], "scenes": {}}
for scene in scenes:
r = bench_scene(model, scene, a)
res["scenes"][r["name"]] = r
cfg = model.device_info
print(f"\n## {r['name']} [{cfg.get('dispatch')} {res['num_command_queues']}CQ {cfg.get('grid')}] "
f"valid {r['valid']} variant {r['variant']}")
print(r["table"])
print("timing_ms p50:", {k: round(v["p50"], 3) for k, v in r["timing_ms"].items()},
f"| b2b rounds {[round(x, 3) for x in r['b2b_rounds_ms']]} ms -> {r['plans_per_s_b2b']:.2f} plans/s",
f"| aiclk {r['aiclk'].get('aiclk_mhz')}", f"| check: staged == model() {r['staged_equals_model']}",
flush=True)
print(json.dumps({k: v for k, v in res.items() if k != "scenes"}, default=str))
if a.json:
Path(a.json).parent.mkdir(parents=True, exist_ok=True)
Path(a.json).write_text(json.dumps(res, indent=1, default=str) + "\n")
if __name__ == "__main__":
main()
|