Team Ai
Apppublic

lablab-ai-amd-developer-hackathon/gpu-goblin

sourceHugging Facemitupdated 5mo agoView on Hugging Face
0likes
protocol.py370 linesDownload Raw Back to runner
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