diff --git a/miles/dashboard/collector.py b/miles/dashboard/collector.py index 245ba30e..989accb8 100644 --- a/miles/dashboard/collector.py +++ b/miles/dashboard/collector.py @@ -7,8 +7,8 @@ from dataclasses import dataclass, field from typing import Any, ClassVar -from miles.dashboard.events import PhaseEvent, TrajectoryEvent -from miles.dashboard.store import DashboardStore, Record, Stream +from miles.dashboard.events import PhaseEvent, RequestEvent, TrajectoryEvent +from miles.dashboard.store import DashboardStore, Record, stream_for logger = logging.getLogger(__name__) @@ -59,8 +59,12 @@ def push_trajectories(self, batch: list[TrajectoryEvent]) -> None: for event in batch: self._append(event) + def push_requests(self, batch: list[RequestEvent]) -> None: + for event in batch: + self._append(event) + def _append(self, record: Record) -> None: - stream = Stream.PHASES if isinstance(record, PhaseEvent) else Stream.TRAJECTORIES + stream = stream_for(record) with self._lock: if self._store.buffered_count(stream) >= self.MAX_BUFFERED_PER_STREAM: self._dropped_since_flush += self._store.drop_oldest_buffered(stream) diff --git a/miles/dashboard/events.py b/miles/dashboard/events.py index f5559f61..49ef202f 100644 --- a/miles/dashboard/events.py +++ b/miles/dashboard/events.py @@ -53,6 +53,28 @@ class TrajectoryEventKind(StrEnum): } +@dataclass +class RequestEvent: + """One rollout request's wall-clock marks, keyed by miles.utils.request_timing. + + The marks, not the leg durations derived from them: the reader places every + leg at the time it happened, and durations are a subtraction away. + ``engine_stages`` is the engine's own breakdown of the one leg it owns. + """ + + rollout_id: int + request_id: str + group_index: int + sample_indices: list[int] + marks: dict[str, float] + engine_stages: dict[str, float] + worker: str + resp_bytes: int + + def to_dict(self) -> dict: + return asdict(self) + + @dataclass class TrajectoryEvent: ts: float diff --git a/miles/dashboard/hooks.py b/miles/dashboard/hooks.py index e4566d65..62c1db84 100644 --- a/miles/dashboard/hooks.py +++ b/miles/dashboard/hooks.py @@ -10,7 +10,18 @@ from contextlib import contextmanager from pathlib import Path -from miles.dashboard.events import STAGE_KINDS, PhaseEvent, TrajectoryEvent +from miles.dashboard.events import STAGE_KINDS, PhaseEvent, RequestEvent, TrajectoryEvent +from miles.utils.request_timing import ( + ROUTER_TIMING_HEADER, + ROUTER_WIRE_KEYS, + ROUTER_WORKER_HEADER, + SGLD_STAGES_HEADER, + SGLD_TIMING_HEADER, + SGLD_WIRE_KEYS, + Marks, + parse_header, + parse_stages, +) logger = logging.getLogger(__name__) @@ -119,8 +130,82 @@ def flush(self) -> None: logger.warning("dashboard trajectory sink flush failed; dropping events", exc_info=True) +class RequestSink: + def __init__(self, handle) -> None: + self.handle = handle + self._buffer: list[RequestEvent] = [] + self._lock = threading.Lock() + self._last_flush = time.monotonic() + + def record(self, tracer: RequestTracer, samples) -> None: + try: + sample = samples[0] + event = RequestEvent( + rollout_id=tracer.rollout_id, + request_id=sample.request_id or "", + group_index=sample.group_index if sample.group_index is not None else -1, + sample_indices=[s.index if s.index is not None else -1 for s in samples], + marks={name: round(ts, 6) for name, ts in tracer.marks.marks.items()}, + engine_stages=tracer.engine_stages, + worker=tracer.worker, + resp_bytes=tracer.resp_bytes, + ) + with self._lock: + self._buffer.append(event) + batch = self._take_batch_if_due() + if batch: + self.handle.push_requests.remote(batch) + except Exception: # noqa: BLE001 + logger.warning("dashboard request sink failed; dropping events", exc_info=True) + + def _take_batch_if_due(self) -> list[RequestEvent] | None: + if len(self._buffer) < BATCH_MAX_EVENTS and time.monotonic() - self._last_flush < BATCH_MAX_SECONDS: + return None + batch, self._buffer = self._buffer, [] + self._last_flush = time.monotonic() + return batch + + def flush(self) -> None: + try: + with self._lock: + batch, self._buffer = self._buffer, [] + if batch: + _ray_get(self.handle.push_requests.remote(batch)) + except Exception: # noqa: BLE001 + logger.warning("dashboard request sink flush failed; dropping events", exc_info=True) + + +class RequestTracer: + """Collects one rollout request's marks; ``done`` is a no-op with the + dashboard off, so the rollout path can carry a tracer unconditionally.""" + + def __init__(self, rollout_id: int) -> None: + self.rollout_id = rollout_id + self.marks = Marks() + self.engine_stages: dict[str, float] = {} + self.worker = "" + self.resp_bytes = 0 + + def mark(self, name: str) -> None: + self.marks.mark(name) + + def absorb_response(self, result) -> None: + """Take the sender's own marks plus the ones the reply carries back.""" + self.marks.absorb(result.marks) + headers = {key.lower(): value for key, value in result.headers.items()} + self.marks.absorb(parse_header(headers.get(SGLD_TIMING_HEADER), SGLD_WIRE_KEYS)) + self.marks.absorb(parse_header(headers.get(ROUTER_TIMING_HEADER), ROUTER_WIRE_KEYS)) + self.engine_stages = parse_stages(headers.get(SGLD_STAGES_HEADER)) + self.worker = headers.get(ROUTER_WORKER_HEADER, "") + + def done(self, samples) -> None: + if _request_sink is not None: + _request_sink.record(self, samples) + + _phase_sink: PhaseSink | None = None _trajectory_sink: TrajectorySink | None = None +_request_sink: RequestSink | None = None _rollout_id = -1 _GPU_SAMPLER: GpuUtilSampler | None = None @@ -152,6 +237,7 @@ def register_rollout_manager(args) -> None: return attach_phase_sink(handle, "rollout") attach_trajectory_sink(handle) + attach_request_sink(handle) def set_rollout_id(rollout_id: int) -> None: @@ -159,6 +245,10 @@ def set_rollout_id(rollout_id: int) -> None: _rollout_id = rollout_id +def current_rollout_id() -> int: + return _rollout_id + + def record_trajectory(sample) -> None: if _trajectory_sink is not None: _trajectory_sink.record(sample, _rollout_id) @@ -201,8 +291,14 @@ def attach_trajectory_sink(handle) -> None: _trajectory_sink = TrajectorySink(handle) +def attach_request_sink(handle) -> None: + global _request_sink + if _request_sink is None: + _request_sink = RequestSink(handle) + + def detach_and_flush() -> None: - global _phase_sink, _trajectory_sink, _GPU_SAMPLER + global _phase_sink, _trajectory_sink, _request_sink, _GPU_SAMPLER from miles.utils.timer import Timer if _phase_sink is not None: @@ -212,6 +308,9 @@ def detach_and_flush() -> None: if _trajectory_sink is not None: _trajectory_sink.flush() _trajectory_sink = None + if _request_sink is not None: + _request_sink.flush() + _request_sink = None if _GPU_SAMPLER is not None: _GPU_SAMPLER.stop() _GPU_SAMPLER = None diff --git a/miles/dashboard/store.py b/miles/dashboard/store.py index c08dc869..e8462cf9 100644 --- a/miles/dashboard/store.py +++ b/miles/dashboard/store.py @@ -8,7 +8,7 @@ from enum import StrEnum from pathlib import Path -from miles.dashboard.events import PhaseEvent, TrajectoryEvent +from miles.dashboard.events import PhaseEvent, RequestEvent, TrajectoryEvent RUN_DIR_PREFIX = "run_" @@ -17,20 +17,25 @@ class Stream(StrEnum): PHASES = "phases" TRAJECTORIES = "trajectories" + REQUESTS = "requests" -Record = PhaseEvent | TrajectoryEvent +Record = PhaseEvent | TrajectoryEvent | RequestEvent -def _stream(record: Record) -> Stream: +def stream_for(record: Record) -> Stream: if isinstance(record, PhaseEvent): return Stream.PHASES + if isinstance(record, RequestEvent): + return Stream.REQUESTS return Stream.TRAJECTORIES def _timestamp(record: Record) -> float: if isinstance(record, PhaseEvent): return record.t1 + if isinstance(record, RequestEvent): + return max(record.marks.values(), default=0.0) return record.ts @@ -63,7 +68,7 @@ def write_meta(self, *, run_name: str, start_ts: float, args: dict) -> None: ) def append(self, record: Record) -> None: - self._buffers[_stream(record)].append(record) + self._buffers[stream_for(record)].append(record) def buffered_count(self, stream: Stream) -> int: return len(self._buffers[stream]) diff --git a/miles/dashboard/viewer.py b/miles/dashboard/viewer.py index ca9384b1..6908023f 100644 --- a/miles/dashboard/viewer.py +++ b/miles/dashboard/viewer.py @@ -13,6 +13,23 @@ from miles.dashboard.events import SPAN_KINDS from miles.dashboard.store import resolve_run_dir +from miles.utils.request_timing import CROSS, LEG_NAMES, LEGS, derive, leg_sources + +# one hue per process, so a waterfall reads first as "where was the request"; +# the cross-process hops stay near-neutral, distinguished by lightness alone +_SOURCE_HUE = {"client": (212, 60), "router": (22, 62), "sgld": (158, 55), CROSS: (35, 10)} + + +def _leg_colors() -> dict[str, str]: + sources = leg_sources() + seen: dict[str, int] = {} + colors = {} + for name in LEG_NAMES: + hue, saturation = _SOURCE_HUE[sources[name]] + index = seen.get(sources[name], 0) + seen[sources[name]] = index + 1 + colors[name] = f"hsl({hue},{saturation - index * 3}%,{38 + min(index, 6) * 6}%)" + return colors def _read_jsonl(paths): @@ -54,17 +71,26 @@ def _fold_trajectory(events): return segments -def load_streams(workspace: str): +def load_streams(workspace: str, max_rollouts: int = 0): phases = _read_jsonl(sorted(glob.glob(os.path.join(workspace, "phases", "*.jsonl")))) gpu = _read_jsonl(sorted(glob.glob(os.path.join(workspace, "gpu_util", "*.jsonl")))) traj = _read_jsonl(sorted(glob.glob(os.path.join(workspace, "trajectories", "*.jsonl")))) - return phases, gpu, _fold_trajectory(traj) + reqs = _read_jsonl(sorted(glob.glob(os.path.join(workspace, "requests", "*.jsonl")))) + if max_rollouts: + # the page embeds every record, so a long run needs a window + keep = sorted({r.get("rollout_id", -1) for r in reqs})[-max_rollouts:] + reqs = [r for r in reqs if r.get("rollout_id", -1) in keep] + return phases, gpu, _fold_trajectory(traj), reqs -def _compute_data(phases, gpu, life=None) -> dict: +def _compute_data(phases, gpu, life=None, reqs=None) -> dict: life = life or [] + timings = [(r, derive(r.get("marks") or {})) for r in (reqs or [])] + timings = [(r, t) for r, t in timings if t.t_end > t.t_start] # normalize all timestamps to the earliest event across streams - ts_candidates = [p["t0"] for p in phases] + [g["ts"] for g in gpu] + [x["t0"] for x in life] + ts_candidates = ( + [p["t0"] for p in phases] + [g["ts"] for g in gpu] + [x["t0"] for x in life] + [t.t_start for _, t in timings] + ) t0 = min(ts_candidates) if ts_candidates else 0.0 ph = [ { @@ -98,12 +124,36 @@ def _compute_data(phases, gpu, life=None) -> dict: for x in life if x.get("t1", 0) >= x.get("t0", 0) ] - span = max([r["e"] for r in ph] + [r["t"] for r in gp] + [r["e"] for r in lf] + [1.0]) - return {"phases": ph, "gpu": gp, "life": lf, "span": span} + rq = [ + { + "k": (r.get("request_id") or "?")[:12], + "r": r.get("rollout_id", -1), + "s": round(timing.t_start - t0, 3), + "e": round(timing.t_end - t0, 3), + "x": round(timing.cross_total, 4), + "w": r.get("worker", ""), + "n": len(r.get("sample_indices") or []), + "m": {name: round(ts - t0, 4) for name, ts in (r.get("marks") or {}).items()}, + "g": {name: round(sec, 4) for name, sec in (r.get("engine_stages") or {}).items()}, + } + for r, timing in timings + ] + rq.sort(key=lambda r: r["s"]) + span = max([r["e"] for r in ph] + [r["t"] for r in gp] + [r["e"] for r in lf] + [r["e"] for r in rq] + [1.0]) + return { + "phases": ph, + "gpu": gp, + "life": lf, + "reqs": rq, + "legs": [list(leg) for leg in LEGS], + "legSource": leg_sources(), + "legColors": _leg_colors(), + "span": span, + } -def build_html(phases, gpu, life=None, title="miles-D rollout dashboard") -> str: - data = json.dumps(_compute_data(phases, gpu, life), separators=(",", ":")) +def build_html(phases, gpu, life=None, reqs=None, title="miles-D rollout dashboard") -> str: + data = json.dumps(_compute_data(phases, gpu, life, reqs), separators=(",", ":")) return _TEMPLATE.replace("__TITLE__", title).replace("__SERVE__", "false").replace("__DATA__", data) @@ -141,6 +191,13 @@ def build_html(phases, gpu, life=None, title="miles-D rollout dashboard") -> str .grid{stroke:var(--border);stroke-width:1;} .rowlab{fill:var(--text);font:11px ui-monospace,monospace;} .roundline{stroke:var(--accent);stroke-width:1;stroke-dasharray:3 3;opacity:.5;} +table.pct{border-collapse:collapse;width:100%;font-size:12px;font-family:ui-monospace,monospace;} +table.pct th,table.pct td{text-align:right;padding:3px 8px;border-bottom:1px solid var(--border);} +table.pct th{color:var(--muted);font-weight:600;} +table.pct th.leg,table.pct td.leg{text-align:left;} +.sw{width:10px;height:10px;border-radius:2px;display:inline-block;margin-right:6px;vertical-align:-1px;} +select{background:var(--bg);color:var(--text);border:1px solid var(--border);border-radius:4px;padding:3px 6px;font:inherit;font-size:12px;} +.warn{color:var(--accent);} #tooltip{position:fixed;pointer-events:none;background:var(--topbar-bg);color:var(--topbar-text);border:1px solid var(--topbar-border);border-radius:4px;padding:6px 9px;font-size:12px;z-index:100;opacity:0;font-family:ui-monospace,monospace;white-space:pre;line-height:1.5;}
| ${label} | n | p50 | p90 | ` + +`p99 | max | mean | share |
|---|---|---|---|---|---|---|---|
| ${r.k} | ` + +`${r.n} | ${ms(r.p50)} | ${ms(r.p90)} | ${ms(r.p99)} | ${ms(r.max)} | ` + +`${ms(r.mean)} | ${(100*r.share).toFixed(1)}% |