lablab-ai-amd-developer-hackathon/gpu-goblin
0
1"""RunnerProtocol — the seam between profile_run/benchmark and the actual GPU.2 3This is the testability fix from the brooks-audit (Warning #2): without this4abstraction, every change to a tool that touches profiling required an5MI300X cloud session. With it, Backend Lead develops on a laptop using6FakeRunner; Day-3 swaps in the real runner for the canonical demo.7 8Real implementations subclass `Runner` and call into goblin_runner.sh.9Tests and laptop dev use FakeRunner, which loads canned RunMetrics from10workloads/synthetic/.11 12`LiveRunner` is the production path: it shells out to `goblin_runner.sh`13(which itself wraps rocprofv3 + torch.profiler), parses the resulting14artefacts via `runner.profile_parser.parse`, and on ANY failure15(missing tools, no GPU, subprocess error) falls back to FakeRunner so the16demo still works on a laptop.17"""18 19from __future__ import annotations20 21import json22import logging23import os24import shutil25import subprocess26import tempfile27from pathlib import Path28from typing import Protocol29 30from agent.schemas import RunMetrics, WorkloadConfig31 32_LOG = logging.getLogger(__name__)33 34 35class Runner(Protocol):36 """Anything that can take a WorkloadConfig and produce RunMetrics."""37 38 def run(self, config: WorkloadConfig, steps: int) -> RunMetrics: # pragma: no cover39 ...40 41 42# ---------------------------------------------------------------------------43# GPU detection — used by LiveRunner to decide live vs. fallback.44# ---------------------------------------------------------------------------45 46 47def _has_render_device() -> bool:48 """At least one /dev/dri/renderD* device is present."""49 dri = Path("/dev/dri")50 if not dri.exists():51 return False52 try:53 return any(child.name.startswith("renderD") for child in dri.iterdir())54 except OSError:55 return False56 57 58def gpu_available() -> tuple[bool, str | None]:59 """Return `(ok, reason_if_missing)`.60 61 LiveRunner is only safe to invoke when ALL of the following hold:62 1. `rocprofv3` is on PATH (the kernel-trace driver).63 2. `amd-smi` is on PATH (HBM/power telemetry sampler).64 3. /dev/dri has at least one `renderD*` node (a real AMD GPU).65 66 If any check fails we fall back to FakeRunner with a clear warning.67 """68 if shutil.which("rocprofv3") is None:69 return False, "rocprofv3 not found on PATH"70 if shutil.which("amd-smi") is None:71 return False, "amd-smi not found on PATH"72 if not _has_render_device():73 return False, "no /dev/dri/renderD* device present"74 return True, None75 76 77# ---------------------------------------------------------------------------78# FakeRunner — loads canned RunMetrics from workloads/synthetic/.79# ---------------------------------------------------------------------------80 81 82class FakeRunner:83 """Loads pre-recorded RunMetrics from workloads/synthetic/<scenario>/cached_metrics.json.84 85 The scenario is selected by matching `WorkloadConfig` fields against each86 synthetic scenario's `match` block in its manifest. If multiple scenarios87 match, the most specific one wins (highest number of matched keys).88 If none match, returns a generic baseline.89 90 This lets us:91 1. Develop the agent loop without an MI300X.92 2. Demo when MI300X cloud is unreachable (offline-replay lane).93 3. Run integration tests deterministically.94 """95 96 def __init__(self, corpus_dir: Path | str = "workloads/synthetic") -> None:97 self.corpus_dir = Path(corpus_dir)98 99 def run(self, config: WorkloadConfig, steps: int) -> RunMetrics:100 scenario = self._match_scenario(config)101 if scenario is None:102 return self._default_metrics(steps)103 104 metrics_path = scenario / "cached_metrics.json"105 if not metrics_path.exists():106 return self._default_metrics(steps)107 108 data = json.loads(metrics_path.read_text())109 # The cached file may not have steps populated; let the caller's110 # request take precedence so profile_run vs benchmark return as expected.111 data["steps"] = steps112 data["runner_kind"] = "fake"113 return RunMetrics.model_validate(data)114 115 # ------------------------------------------------------------------116 # Scenario matching117 # ------------------------------------------------------------------118 119 def _match_scenario(self, config: WorkloadConfig) -> Path | None:120 if not self.corpus_dir.exists():121 return None122 best: tuple[int, Path] | None = None123 for scenario_dir in sorted(self.corpus_dir.iterdir()):124 if not scenario_dir.is_dir():125 continue126 manifest = scenario_dir / "manifest.json"127 if not manifest.exists():128 continue129 try:130 spec = json.loads(manifest.read_text())131 except json.JSONDecodeError:132 continue133 match_block = spec.get("match", {})134 score = self._score(config, match_block)135 if score < 0:136 continue137 if best is None or score > best[0]:138 best = (score, scenario_dir)139 return best[1] if best else None140 141 @staticmethod142 def _score(config: WorkloadConfig, match: dict) -> int:143 """Return number of keys that match, or -1 if any key conflicts."""144 cfg = config.model_dump()145 score = 0146 for key, expected in match.items():147 if cfg.get(key) != expected:148 return -1149 score += 1150 return score151 152 @staticmethod153 def _default_metrics(steps: int) -> RunMetrics:154 from agent.schemas import KernelEntry, WasteBudget155 156 return RunMetrics(157 steps=steps,158 tokens_per_sec=120.0,159 mfu_pct=22.0,160 hbm_peak_gb=72.0,161 hbm_avg_gb=58.0,162 gpu_util_pct=45.0,163 top_kernels=[164 KernelEntry(name="aten::matmul", pct_time=42.0),165 KernelEntry(name="aten::scaled_dot_product_attention", pct_time=18.0),166 KernelEntry(name="aten::layer_norm", pct_time=7.0),167 ],168 attention_kernel_loaded="sdpa",169 waste_budget=WasteBudget(170 useful_gpu=0.55,171 data_wait=0.18,172 host_gap=0.07,173 comm_excess=0.0,174 memory_headroom=0.10,175 precision_path=0.06,176 kernel_shape=0.04,177 ),178 warnings=["FakeRunner: no matching scenario, returning generic baseline."],179 runner_kind="fake",180 )181 182 183# ---------------------------------------------------------------------------184# LiveRunner — production path. Spawns goblin_runner.sh under rocprofv3.185# ---------------------------------------------------------------------------186 187 188# Defaults are pinned to the repo layout. Override via env vars in tests / CI.189_REPO_ROOT = Path(__file__).resolve().parent.parent190_DEFAULT_RUNNER_SCRIPT = _REPO_ROOT / "runner" / "goblin_runner.sh"191_DEFAULT_USER_SCRIPT = _REPO_ROOT / "workloads" / "train_qwen_lora.py"192_FAILURE_ARCHIVE_ROOT = _REPO_ROOT / "bench_cache"193 194 195def _archive_failure(out_dir: Path, proc: subprocess.CompletedProcess) -> Path:196 """Copy a failed runner's out_dir into bench_cache/last_runner_failure_<ts>/197 along with the subprocess's captured stdout/stderr. The directory survives198 after the tempdir cleanup so the user can `tail -n 100 stderr.log` etc.199 """200 import shutil201 import time202 203 ts = time.strftime("%Y%m%dT%H%M%S")204 dest = _FAILURE_ARCHIVE_ROOT / f"last_runner_failure_{ts}"205 try:206 dest.mkdir(parents=True, exist_ok=True)207 if out_dir.exists():208 for child in out_dir.iterdir():209 target = dest / child.name210 if child.is_dir():211 shutil.copytree(child, target, dirs_exist_ok=True)212 else:213 shutil.copy2(child, target)214 # Also persist the subprocess's own captured output — these are what215 # goblin_runner.sh's failure trap dumped.216 (dest / "subprocess_stdout.log").write_text(proc.stdout or "")217 (dest / "subprocess_stderr.log").write_text(proc.stderr or "")218 (dest / "subprocess_returncode").write_text(str(proc.returncode))219 except OSError as exc:220 _LOG.warning("LiveRunner: could not archive failure logs (%s)", exc)221 return dest222 return dest223 224 225class LiveRunner:226 """Real-MI300X path: shells out to goblin_runner.sh and parses artefacts.227 228 Auto-falls-back to FakeRunner whenever the host can't actually run a live229 profile (missing rocprofv3/amd-smi, no AMD GPU, subprocess error, or230 parser failure). The fallback path is the demo safety net.231 232 Public API matches the `Runner` Protocol: `run(config, steps) -> RunMetrics`.233 """234 235 def __init__(236 self,237 runner_script: Path | str = _DEFAULT_RUNNER_SCRIPT,238 user_script: Path | str = _DEFAULT_USER_SCRIPT,239 timeout_seconds: int = 600,240 fake_fallback: FakeRunner | None = None,241 ) -> None:242 # Default 600s (10 min). Profile runs (10 steps) finish in seconds243 # on a healthy MI300X; benchmarks (50 steps) in a couple of minutes.244 # 30 minutes was a leftover from a workload that wasn't honoring245 # --max_steps and silently trained for hours. With max_steps wired246 # correctly, 600s is generous.247 self.runner_script = Path(runner_script)248 self.user_script = Path(user_script)249 self.timeout_seconds = timeout_seconds250 self._fake = fake_fallback or FakeRunner()251 252 # ------------------------------------------------------------------253 254 def run(self, config: WorkloadConfig, steps: int) -> RunMetrics:255 ok, reason = gpu_available()256 if not ok:257 return self._fallback(258 config,259 steps,260 f"LiveRunner: GPU/profiler unavailable ({reason}); using FakeRunner.",261 )262 263 # Sanity-check the runner script before spawning anything.264 if not self.runner_script.exists():265 return self._fallback(266 config,267 steps,268 f"LiveRunner: runner script not found at {self.runner_script}; using FakeRunner.",269 )270 if not os.access(self.runner_script, os.X_OK):271 return self._fallback(272 config,273 steps,274 f"LiveRunner: runner script {self.runner_script} not executable; using FakeRunner.",275 )276 277 # Late import — only needed on the live path. Keeps laptop-only test278 # runs from importing parser dependencies (csv stdlib is fine, but279 # this also keeps the dependency direction explicit).280 from runner import profile_parser281 282 with tempfile.TemporaryDirectory(prefix="goblin_run_") as out_dir_str:283 out_dir = Path(out_dir_str)284 285 env = os.environ.copy()286 env["USER_SCRIPT"] = str(self.user_script)287 env["OUT_DIR"] = str(out_dir)288 env["STEPS"] = str(steps)289 290 cmd = [str(self.runner_script)]291 try:292 proc = subprocess.run(293 cmd,294 env=env,295 capture_output=True,296 text=True,297 timeout=self.timeout_seconds,298 check=False,299 )300 except subprocess.TimeoutExpired:301 return self._fallback(302 config,303 steps,304 f"LiveRunner: goblin_runner.sh timed out after "305 f"{self.timeout_seconds}s; using FakeRunner.",306 )307 except OSError as exc:308 return self._fallback(309 config,310 steps,311 f"LiveRunner: failed to spawn goblin_runner.sh ({exc}); using FakeRunner.",312 )313 314 if proc.returncode != 0:315 # Archive the full out_dir to bench_cache/last_runner_failure_<ts>/316 # so the user can inspect stdout.log / stderr.log / amd_smi.err317 # after the tempdir cleanup. The path goes into the warning318 # message so it's surfaced through ToolResult.warnings.319 archive_path = _archive_failure(out_dir, proc)320 stderr_tail = (proc.stderr or "").strip().splitlines()[-15:]321 stdout_tail = (proc.stdout or "").strip().splitlines()[-5:]322 return self._fallback(323 config,324 steps,325 "LiveRunner: goblin_runner.sh exited with "326 f"code {proc.returncode}; using FakeRunner. "327 f"Failure logs archived at {archive_path}. "328 f"stderr tail: {stderr_tail}. "329 f"stdout tail: {stdout_tail}.",330 )331 332 try:333 metrics = profile_parser.parse(out_dir, config=config, steps=steps)334 except Exception as exc: # pragma: no cover — defensive335 return self._fallback(336 config,337 steps,338 f"LiveRunner: profile_parser.parse failed ({type(exc).__name__}: {exc}); "339 "using FakeRunner.",340 )341 342 metrics.runner_kind = "live"343 return metrics344 345 # ------------------------------------------------------------------346 347 def _fallback(self, config: WorkloadConfig, steps: int, warning: str) -> RunMetrics:348 _LOG.warning(warning)349 metrics = self._fake.run(config, steps)350 # Make the fallback observable to upstream tools — they surface351 # warnings into the final report.352 metrics.warnings = [warning, *metrics.warnings]353 metrics.runner_kind = "fake"354 return metrics355 356 357# ---------------------------------------------------------------------------358# Module-level factory used by agent/tools/{profile_run,benchmark}.py.359# ---------------------------------------------------------------------------360 361 362def _default_runner() -> Runner:363 """Return the runner profile_run / benchmark should use by default.364 365 Always returns a `LiveRunner` — `LiveRunner.run` itself decides whether to366 actually invoke the GPU pipeline or fall back to FakeRunner. Centralising367 this here means the live-vs-fake decision lives in exactly one place.368 """369 return LiveRunner()370 