File size: 11,322 Bytes
51defdc
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
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
224
225
226
227
228
#!/usr/bin/env python3
# SPDX-License-Identifier: Apache-2.0
"""Contained stress runs of the ``frame`` trace (the METEOR hang investigation, RESET_INVESTIGATION.md section 6).

    METEOR_DEVICE_OK=1 bin/devrun -t 900 -- python -u code/scripts/stress_frames.py --mode replay --frames 200
    METEOR_DEVICE_OK=1 bin/devrun -t 1500 -- python -u code/scripts/stress_frames.py --mode rigs --frames 100

``--mode replay`` (R1): the model is built with the sample rig's tables and cameras as its warm-up data, captured,
then the frame trace is replayed ``--frames`` times with NO upload at all, reading back every ``--read-every``-th
replay. ``--mode rigs`` (R2): the model is built like the API (zero tables), then every frame cycles to the next of
the golden rigs: table write (when the rig changes) + camera upload + replay + segmented readback.

Every read is checked: bit-identical to the first read of the same rig (corruption / non-determinism shows up as a
mismatch) and agreement with the CPU golden (``lane`` argmax, ``hm`` PCC).

``--rigs shipped[:N]`` needs no golden (for the container, which has none): N rigs (default 4) made from the shipped
synthetic sample, rig k with every camera raised by 5k cm (a different calibration, so different lift tables, and the
same cameras); only the bit-identity per rig is checked. Inside the container image (``/opt/tt-metal``)::

    python -u /opt/tt-metal/scripts/stress_frames.py --mode rigs --rigs shipped --frames 200 One heartbeat line per frame goes to
stdout (run with ``python -u``) and one JSON line to ``--jsonl`` (fsync'd); ``faulthandler`` dumps the Python stack
when no frame finishes for ``--stall-s`` seconds (a hang leaves evidence of where it stopped, without touching the
device). Nothing here arms tt-triage or an operation timeout.
"""
from __future__ import annotations

import argparse
import faulthandler
import hashlib
import json
import os
import sys
import time
from pathlib import Path

import numpy as np

GOLDENS = Path(os.environ.get("METEOR_GOLDENS", "/home/ubuntu/experiments/tt-models/research/meteor/goldens"))
RIGS = ["pandaset_019_f40", "pandaset_090_f40", "nuscenes_0103_kf09", "meteor_valday_f040"]
INPUT_KEYS = ("input.imgs", "input.K", "input.T_cam_ego", "input.v0", "input.present")
CHECK_KEYS = ("lane", "hm", "ego")


def log(fh, **rec) -> None:
    rec["t"] = round(time.time(), 3)
    line = json.dumps(rec, default=str)
    print(line, flush=True)
    if fh is not None:
        fh.write(line + "\n")
        fh.flush()
        os.fsync(fh.fileno())


def digest(out) -> str:
    h = hashlib.sha256()
    for k in sorted(out):
        a = np.ascontiguousarray(out[k])
        h.update(k.encode())
        h.update(a.tobytes())
    return h.hexdigest()[:16]


def checks(out, gold) -> dict:
    from tt_meteor.ttaw.metrics import pcc

    lane = float((np.asarray(out["lane"]).reshape(-1) == np.asarray(gold["lane"]).reshape(-1)).mean())
    hm = float(pcc(np.asarray(out["hm"], np.float64), np.asarray(gold["hm"], np.float64)))
    return {"lane_agree": round(lane, 5), "hm_pcc": round(hm, 6)}


def main() -> int:
    ap = argparse.ArgumentParser()
    ap.add_argument("--mode", choices=["replay", "rigs", "segrigs"], required=True)
    ap.add_argument("--frames", type=int, default=100)
    ap.add_argument("--read-every", type=int, default=1)
    ap.add_argument("--rigs", default=",".join(RIGS))
    ap.add_argument("--dispatch", default="eth", choices=["eth", "worker"])
    ap.add_argument("--open-cqs", type=int, default=None, help="HW command queues to open (default: DEVICE_DEFAULTS)")
    ap.add_argument("--runner-cqs", type=int, default=None, help="CQs the TraceRunner uses (default: as opened)")
    ap.add_argument("--lift-reps", type=int, default=1,
                    help="segrigs: replay the lift segment(s) this many times per frame (sync + heartbeat after each)")
    ap.add_argument("--split-lift", nargs="?", const=True, default=False,
                    help="segrigs: lift as lift_gs (grid_samples) + lift_rest; 'fine': lift_rest as 4 op groups")
    ap.add_argument("--no-tables", action="store_true", help="segrigs: one rig only, tables written once")
    ap.add_argument("--stall-s", type=float, default=120.0)
    ap.add_argument("--jsonl")
    a = ap.parse_args()

    faulthandler.enable()
    fh = open(a.jsonl, "a") if a.jsonl else None
    rigs = a.rigs.split(",")
    faulthandler.dump_traceback_later(900, repeat=True)          # build / warm-up / capture
    log(fh, ev="start", mode=a.mode, frames=a.frames, read_every=a.read_every, rigs=rigs, pid=os.getpid())

    import ttnn

    from tt_meteor.device import close_device, describe_device, open_device
    from tt_meteor.host.calib import lift_geometry
    from tt_meteor.host.preprocess import MeteorFrame
    from tt_meteor.tt.lift import lift_tables
    from tt_meteor.tt.model import TtMETEOR
    from tt_meteor.tt.params import MeteorParams
    from tt_meteor.tt.unpack import outputs_from_device

    t0 = time.perf_counter()
    params = MeteorParams.load()
    data = {}
    if len(rigs) == 1 and rigs[0].startswith("shipped"):
        from tt_meteor.api import load_sample
        from tt_meteor.host.inputs import SAMPLES_DIR, prepare_request

        n = int(rigs[0].partition(":")[2] or 4)
        kw = load_sample(SAMPLES_DIR / "synthetic_8cam.json")
        base = prepare_request(kw["images"], kw["calibration"], kw["ego_speed"], None)
        rigs = [f"shipped{k}" for k in range(n)]
        for k, name in enumerate(rigs):
            shift = np.eye(4, dtype=np.float32)
            shift[2, 3] = -0.05 * k                 # ego -> ego lowered by 5k cm = every camera raised by 5k cm
            T = (base.T_cam_ego.astype(np.float64) @ shift).astype(np.float32)
            fr = MeteorFrame(base.imgs, base.K, T, base.v0, base.present)
            data[name] = (fr, lift_geometry(fr.K, fr.T_cam_ego, points=params.ground_points()), None)
    for name in ([] if data else rigs):
        with np.load(GOLDENS / name / "taps.npz") as z:
            g = {k: z[k] for k in INPUT_KEYS + CHECK_KEYS}
        fr = MeteorFrame(g["input.imgs"], g["input.K"], g["input.T_cam_ego"], g["input.v0"], g["input.present"])
        geom = lift_geometry(fr.K, fr.T_cam_ego, points=params.ground_points())
        data[name] = (fr, geom, g)
    log(fh, ev="host_ready", s=round(time.perf_counter() - t0, 1), keys={n: d[1].key for n, d in data.items()})

    dev = open_device(dispatch=a.dispatch, allow_fallback=False, num_command_queues=a.open_cqs)
    rc = 0
    try:
        log(fh, ev="device_open", info=describe_device(dev))
        first = rigs[0]
        fr0, geom0, _ = data[first]
        if a.mode == "replay":
            warm = {"imgs": TtMETEOR.image_rows(fr0.imgs), "v0": float(np.asarray(fr0.v0).reshape(-1)[0]),
                    **lift_tables(geom0)}
            tt = TtMETEOR(dev, params, warmup=warm, tables_key=geom0.key, num_command_queues=a.runner_cqs)
        elif a.mode == "segrigs":
            tt = TtMETEOR(dev, params, segmented=True, split_lift=a.split_lift, num_command_queues=a.runner_cqs)
        else:
            tt = TtMETEOR(dev, params, num_command_queues=a.runner_cqs)
        t1 = time.perf_counter()
        tt.capture()
        ttnn.synchronize_device(dev)
        log(fh, ev="captured", s=round(time.perf_counter() - t1, 1), timings=tt.runner.timings_ms,
            program_cache_entries=tt.describe().get("program_cache_entries"))
        faulthandler.dump_traceback_later(a.stall_s, repeat=True)
        ref_digest = {}
        mismatches = 0
        for i in range(a.frames):
            faulthandler.dump_traceback_later(a.stall_s, repeat=True)   # re-armed per frame: fires only on a stall
            ts = time.perf_counter()
            if a.mode == "replay":
                name = first
                tt.runner.replay("frame")
                ttnn.synchronize_device(dev)
                rec = {"ev": "frame", "i": i, "rig": name, "replay_ms": round((time.perf_counter() - ts) * 1e3, 1)}
                if i % a.read_every == 0 or i == a.frames - 1:
                    tr = time.perf_counter()
                    out = outputs_from_device(tt.unpack(tt.runner.read("frame")))
                    rec["read_ms"] = round((time.perf_counter() - tr) * 1e3, 1)
                else:
                    out = None
            elif a.mode == "segrigs":
                name = rigs[0] if a.no_tables else rigs[i % len(rigs)]
                fr, geom, _ = data[name]
                tw = time.perf_counter()
                wrote = tt.set_tables(geom)
                rec = {"ev": "frame", "i": i, "rig": name, "tables": wrote,
                       "tables_ms": round((time.perf_counter() - tw) * 1e3, 1)}
                r = tt.runner
                for k, stage in enumerate(tt.segments):
                    ta = time.perf_counter()
                    print(f"[seg] i={i} {stage} start", flush=True)
                    if k == 0:
                        r.run("seg_image", inputs={"imgs": tt.image_rows(fr.imgs)}, params=tt.frame_params(fr))
                    else:
                        r.run(f"seg_{stage}")
                    ttnn.synchronize_device(dev)
                    if stage.startswith(("lift", "lr_")):
                        for rep in range(1, a.lift_reps):
                            print(f"[seg] i={i} {stage} rep {rep} start", flush=True)
                            r.run(f"seg_{stage}")
                            ttnn.synchronize_device(dev)
                    rec[f"{stage}_ms"] = round((time.perf_counter() - ta) * 1e3, 1)
                ego = np.array(r.read("seg_head"), copy=True)
                out = {"ego": ego}
            else:
                name = rigs[i % len(rigs)]
                fr, geom, _ = data[name]
                tw = time.perf_counter()
                wrote = tt.set_tables(geom)
                tables_ms = (time.perf_counter() - tw) * 1e3
                tr = time.perf_counter()
                out = tt.run_frame(fr, geom)
                rec = {"ev": "frame", "i": i, "rig": name, "tables": wrote, "tables_ms": round(tables_ms, 1),
                       "run_ms": round((time.perf_counter() - tr) * 1e3, 1)}
            if out is not None:
                d = digest(out)
                rec["digest"] = d
                if name not in ref_digest:
                    ref_digest[name] = d
                    if "lane" in out and data[name][2] is not None:
                        rec.update(checks(out, data[name][2]))
                elif d != ref_digest[name]:
                    mismatches += 1
                    rec["MISMATCH_vs"] = ref_digest[name]
                    if "lane" in out and data[name][2] is not None:
                        rec.update(checks(out, data[name][2]))
            log(fh, **rec)
        faulthandler.cancel_dump_traceback_later()
        log(fh, ev="done", frames=a.frames, mismatches=mismatches, digests=ref_digest,
            s=round(time.perf_counter() - t0, 1))
        rc = 1 if mismatches else 0
        tt.release()
    finally:
        faulthandler.dump_traceback_later(300, repeat=True)
        close_device(dev)
        faulthandler.cancel_dump_traceback_later()
        log(fh, ev="closed", rc=rc)
    return rc


if __name__ == "__main__":
    sys.exit(main())