Download gpu-sft/scripts/gpu_eval/gpu_eval_driver.py from fzzhang/svd-code: direct link, hf CLI and curl.
- Browser
- Download file 39.4 kB
-
https://huggingface.co/fzzhang/svd-code/resolve/main/gpu-sft/scripts/gpu_eval/gpu_eval_driver.py
- Command line
-
hf download hf://fzzhang/svd-code/gpu-sft/scripts/gpu_eval/gpu_eval_driver.py
-
curl -L -o gpu_eval_driver.py https://huggingface.co/fzzhang/svd-code/resolve/main/gpu-sft/scripts/gpu_eval/gpu_eval_driver.py
39.4 kB
| #!/usr/bin/env python3 | |
| # Copyright The Marin Authors | |
| # SPDX-License-Identifier: Apache-2.0 | |
| """Run Marin's evalchemy benchmark suites against a LOCAL HF checkpoint on one GPU box. | |
| This is the local-disk, CUDA-vLLM replacement for the TPU path | |
| (``experiments/evals/exp_evalchemy_eval.py`` + ``EvalchemyEvaluator``). It keeps | |
| the benchmark set, the sampling parameters and -- critically -- the on-disk | |
| result schema that ``claude/compile_results.py`` reads to build results.md. | |
| python gpu_eval_driver.py \\ | |
| --checkpoint /data/runs/exp_sft_qwen3_8b_x/hf/step-300 \\ | |
| --experiment exp_sft_qwen3_8b_x \\ | |
| --suite math \\ | |
| --results-root /data/marin/evaluation/evalchemy | |
| WHAT IS FAITHFUL TO THE TPU PIPELINE | |
| * Task set and per-task seed counts (math AIME*/AMC23/HMMT x10, MATH500 and | |
| OlympiadBench x1; science x3; code x6), seeds 42..51. | |
| * num_fewshot=0, chat template applied, temperature 0.7, top_p 1.0 | |
| (both are the benchmark/vLLM defaults; Marin's --gen_kwargs never reached | |
| evalchemy chat benchmarks on TPU either), max_gen_toks 32768, | |
| max_model_len = 32768 + 4096 = 36864. | |
| * Real per-benchmark graders (is_equiv, symbolic, execution) -- the naive | |
| string re-grade bug in ``evalchemy_results_compiler.py`` is NOT reproduced. | |
| * Output layout, so the existing results.md compiler needs zero changes: | |
| <root>/<experiment>-step<N>/<TASK>_seed<S>-<h>/<TASK>_0shot/<model>/results_<iso>.json | |
| <root>/<experiment>-step<N>/compile_<TASK>_avg<K>seeds-<h>/compiled_results/averaged_results.json | |
| WHAT DELIBERATELY DIFFERS | |
| * ONE vLLM engine per task instead of one Iris job per (task, seed). | |
| TPU/JAX cannot do per-request seeds, so Marin faked multi-seed by launching | |
| K whole jobs with K engine seeds. CUDA vLLM *can* do per-request seeds, so | |
| we set EVALCHEMY_N_REPEAT=K and pass ``--seed 42,42,42,42``; evalchemy's | |
| benchmarks then use per-request seed ``42 + repetition_index``, which maps | |
| exactly onto Marin seeds 42..42+K-1. One model load instead of K. | |
| Consequence: numbers are NOT bit-comparable to the TPU tables (different | |
| hardware, different kernels, per-request vs engine seeding). Rankings are. | |
| * tensor_parallel_size defaults to the number of visible GPUs, not 32. | |
| * max_num_seqs defaults to 32-64 per the GPU profile, not 256. 256 on a | |
| single 80GB card guarantees KV-cache thrash and preemption at 36864 ctx. | |
| """ | |
| from __future__ import annotations | |
| import argparse | |
| import csv | |
| import dataclasses | |
| import hashlib | |
| import json | |
| import logging | |
| import os | |
| import re | |
| import shutil | |
| import statistics | |
| import subprocess | |
| import sys | |
| import time | |
| from datetime import datetime | |
| from pathlib import Path | |
| from typing import Any | |
| logger = logging.getLogger("gpu_eval_driver") | |
| # -------------------------------------------------------------------------- | |
| # Task tables. Mirrors experiments/evals/exp_evalchemy_eval.py and | |
| # experiments/evals/evalchemy_task_configs.py. | |
| # -------------------------------------------------------------------------- | |
| SEED_BASE = 42 | |
| # task name (as evalchemy knows it) -> number of seeds | |
| SUITES: dict[str, dict[str, int]] = { | |
| "math": { | |
| "MATH500": 1, | |
| "OlympiadBench": 1, | |
| "AIME24": 10, | |
| "AIME25": 10, | |
| "AIME26": 10, | |
| "AMC23": 10, | |
| "HMMT": 10, | |
| }, | |
| "science": { | |
| "GPQADiamond": 3, | |
| "JEEBench": 3, | |
| "HLE": 3, | |
| "OlympiadBench_Physics": 3, | |
| }, | |
| "code": { | |
| "LiveCodeBench": 6, | |
| "LiveCodeBenchv5_official": 6, | |
| "LiveCodeBenchv6_official": 6, | |
| }, | |
| } | |
| ALL_TASK_SEEDS: dict[str, int] = {task: n for suite in SUITES.values() for task, n in suite.items()} | |
| # Benchmarks that are single-pass by construction: they carry no ``n_repeat`` | |
| # attribute, so EVALCHEMY_N_REPEAT is inert and they emit ``accuracy`` rather | |
| # than ``accuracy_avg``/``run_stats``. | |
| SINGLE_PASS_TASKS = frozenset({"MATH500", "OlympiadBench", "OlympiadBench_Physics"}) | |
| # Generation / engine constants, matching the TPU pipeline. | |
| MAX_GEN_TOKS = 32768 | |
| CONTEXT_BUFFER = 4096 | |
| MAX_MODEL_LEN = MAX_GEN_TOKS + CONTEXT_BUFFER # 36864 | |
| TEMPERATURE = 0.7 | |
| TOP_P = 1.0 | |
| NUM_FEWSHOT = 0 | |
| RESULTS_FILE_RE = re.compile(r"^results_.*\.json$") | |
| STEP_RE = re.compile(r"step-(\d+)") | |
| # -------------------------------------------------------------------------- | |
| # GPU engine profiles | |
| # -------------------------------------------------------------------------- | |
| class GpuProfile: | |
| """vLLM engine settings for an 8B bf16 model at 36864 context.""" | |
| name: str | |
| gpu_memory_utilization: float | |
| max_num_seqs: int | |
| max_num_batched_tokens: int | |
| # Minimum bytes vLLM must be allowed to claim on each GPU before the run is | |
| # worth starting: weights + activations + a few full-length sequences of KV. | |
| min_budget_gib_per_gpu: float | |
| min_gpus: int | |
| notes: str | |
| GPU_PROFILES: dict[str, GpuProfile] = { | |
| # 1x 80GB is the reference config: 16.4 GiB bf16 weights, ~4-6 GiB | |
| # activations/CUDA graphs, ~48 GiB KV = ~340k cached tokens = ~9 concurrent | |
| # full-length (36864) sequences, far more at realistic generation lengths. | |
| "a100-80g": GpuProfile( | |
| name="a100-80g", | |
| gpu_memory_utilization=0.90, | |
| max_num_seqs=64, | |
| max_num_batched_tokens=8192, | |
| min_budget_gib_per_gpu=40.0, | |
| min_gpus=1, | |
| notes="A100-SXM/PCIe 80GB. 0.90 leaves ~8 GiB for the driver+NCCL. Do not exceed 0.92.", | |
| ), | |
| "h100-80g": GpuProfile( | |
| name="h100-80g", | |
| gpu_memory_utilization=0.90, | |
| max_num_seqs=64, | |
| max_num_batched_tokens=16384, | |
| min_budget_gib_per_gpu=40.0, | |
| min_gpus=1, | |
| notes="H100 80GB HBM3. Higher prefill budget than A100; FlashAttention-3 kernels.", | |
| ), | |
| "h200-141g": GpuProfile( | |
| name="h200-141g", | |
| gpu_memory_utilization=0.90, | |
| max_num_seqs=128, | |
| max_num_batched_tokens=16384, | |
| min_budget_gib_per_gpu=40.0, | |
| min_gpus=1, | |
| notes="H200 141GB. KV cache stops being the bottleneck; raise max_num_seqs.", | |
| ), | |
| # 40GB: 0.92*39.6 = 36.4 GiB budget, minus 16.4 weights minus ~4 overhead | |
| # leaves ~16 GiB KV = ~113k tokens = 3 full-length sequences. Runs, but slow | |
| # and preemption-heavy. Use two cards. | |
| "a100-40g": GpuProfile( | |
| name="a100-40g", | |
| gpu_memory_utilization=0.92, | |
| max_num_seqs=16, | |
| max_num_batched_tokens=4096, | |
| min_budget_gib_per_gpu=34.0, | |
| min_gpus=1, | |
| notes="A100 40GB. Tight for 8B at 36864 ctx; strongly prefer --tensor-parallel-size 2.", | |
| ), | |
| "l40s-48g": GpuProfile( | |
| name="l40s-48g", | |
| gpu_memory_utilization=0.92, | |
| max_num_seqs=24, | |
| max_num_batched_tokens=4096, | |
| min_budget_gib_per_gpu=34.0, | |
| min_gpus=1, | |
| notes="L40S 48GB, no NVLink. Keep tensor_parallel_size=1; PCIe TP is slower than 1 card.", | |
| ), | |
| } | |
| # nvidia-smi product name substring -> profile | |
| _NAME_TO_PROFILE = ( | |
| ("H200", "h200-141g"), | |
| ("H100", "h100-80g"), | |
| ("A100-SXM4-40GB", "a100-40g"), | |
| ("A100 40GB", "a100-40g"), | |
| ("A100", "a100-80g"), | |
| ("L40S", "l40s-48g"), | |
| ) | |
| class GpuInfo: | |
| index: int | |
| name: str | |
| total_mib: int | |
| used_mib: int | |
| free_mib: int | |
| def query_gpus() -> list[GpuInfo]: | |
| """Read nvidia-smi. Honours CUDA_VISIBLE_DEVICES ordering when set.""" | |
| smi = shutil.which("nvidia-smi") | |
| if smi is None: | |
| raise RuntimeError("nvidia-smi not found; this driver requires an NVIDIA GPU box") | |
| out = subprocess.run( | |
| [smi, "--query-gpu=index,name,memory.total,memory.used,memory.free", "--format=csv,noheader,nounits"], | |
| check=True, | |
| capture_output=True, | |
| text=True, | |
| ).stdout | |
| gpus = [] | |
| for line in out.strip().splitlines(): | |
| index, name, total, used, free = (part.strip() for part in line.split(",")) | |
| gpus.append(GpuInfo(int(index), name, int(total), int(used), int(free))) | |
| visible = os.environ.get("CUDA_VISIBLE_DEVICES") | |
| if visible: | |
| wanted = [int(x) for x in visible.split(",") if x.strip() != ""] | |
| by_index = {g.index: g for g in gpus} | |
| missing = [i for i in wanted if i not in by_index] | |
| if missing: | |
| raise RuntimeError(f"CUDA_VISIBLE_DEVICES names GPUs {missing} that nvidia-smi does not report") | |
| gpus = [by_index[i] for i in wanted] | |
| return gpus | |
| def detect_profile(gpus: list[GpuInfo]) -> GpuProfile: | |
| name = gpus[0].name | |
| for needle, key in _NAME_TO_PROFILE: | |
| if needle.lower() in name.lower(): | |
| profile = GPU_PROFILES[key] | |
| # An "A100" with <60GB is the 40GB SKU regardless of the name string. | |
| if key == "a100-80g" and gpus[0].total_mib < 60_000: | |
| profile = GPU_PROFILES["a100-40g"] | |
| logger.info("detected GPU %r -> profile %s", name, profile.name) | |
| return profile | |
| raise RuntimeError( | |
| f"no built-in profile for GPU {name!r}; pass --gpu-profile explicitly " | |
| f"(one of: {', '.join(sorted(GPU_PROFILES))})" | |
| ) | |
| def resolve_gpu_memory_utilization(profile: GpuProfile, gpus: list[GpuInfo], headroom_mib: int) -> float: | |
| """Clamp gpu_memory_utilization so co-tenants on the same card do not OOM us. | |
| THE PITFALL THIS EXISTS FOR: vLLM interprets ``gpu_memory_utilization`` as a | |
| fraction of TOTAL device memory, not of FREE device memory. If an embedding | |
| server (or a JAX training process, which preallocates 75% by default) already | |
| holds 10 GiB of an 80 GiB card and you ask for 0.90, vLLM will size its KV | |
| pool as if it owned 72 GiB. Total demand becomes 82 GiB and one of the two | |
| processes dies with CUDA OOM -- and it usually dies *mid-run*, after the KV | |
| pool is grown, not at init, so you lose hours of generation. | |
| So: measure what is already resident, and cap the fraction at what is | |
| actually claimable. If that leaves less than the profile's minimum budget, | |
| refuse to start rather than crash later. | |
| """ | |
| worst = min(gpus, key=lambda g: g.free_mib) | |
| claimable_mib = worst.free_mib - headroom_mib | |
| if claimable_mib <= 0: | |
| raise RuntimeError( | |
| f"GPU {worst.index} ({worst.name}) has only {worst.free_mib} MiB free " | |
| f"(headroom {headroom_mib} MiB). Something else owns this card." | |
| ) | |
| allowed_fraction = claimable_mib / worst.total_mib | |
| gmu = min(profile.gpu_memory_utilization, allowed_fraction) | |
| budget_gib = gmu * worst.total_mib / 1024.0 | |
| if budget_gib < profile.min_budget_gib_per_gpu: | |
| raise RuntimeError( | |
| f"only {budget_gib:.1f} GiB claimable on GPU {worst.index} " | |
| f"({worst.used_mib} MiB already in use by another process) but profile " | |
| f"{profile.name} needs {profile.min_budget_gib_per_gpu:.1f} GiB. " | |
| "Free the card, pin the other process to a different GPU with " | |
| "CUDA_VISIBLE_DEVICES, or raise --tensor-parallel-size." | |
| ) | |
| if gmu < profile.gpu_memory_utilization: | |
| logger.warning( | |
| "clamping gpu_memory_utilization %.2f -> %.3f: GPU %d already has %d MiB resident", | |
| profile.gpu_memory_utilization, | |
| gmu, | |
| worst.index, | |
| worst.used_mib, | |
| ) | |
| return round(gmu, 3) | |
| def kv_bytes_per_token(checkpoint: Path) -> int | None: | |
| """Estimate bf16 KV cache cost per token from the checkpoint's config.json.""" | |
| config_path = checkpoint / "config.json" | |
| if not config_path.exists(): | |
| return None | |
| cfg = json.loads(config_path.read_text()) | |
| layers = cfg.get("num_hidden_layers") | |
| kv_heads = cfg.get("num_key_value_heads") or cfg.get("num_attention_heads") | |
| head_dim = cfg.get("head_dim") | |
| if head_dim is None and cfg.get("hidden_size") and cfg.get("num_attention_heads"): | |
| head_dim = cfg["hidden_size"] // cfg["num_attention_heads"] | |
| if not (layers and kv_heads and head_dim): | |
| return None | |
| return layers * kv_heads * head_dim * 2 * 2 # K and V, 2 bytes each | |
| # -------------------------------------------------------------------------- | |
| # Output layout -- must stay byte-compatible with claude/compile_results.py | |
| # -------------------------------------------------------------------------- | |
| def sanitize_model_name(name: str) -> str: | |
| """Reproduce lm-eval's GeneralConfigTracker.model_name_sanitized.""" | |
| return re.sub(r"[\"<>:/\|\\?\*\[\]]+", "__", name) | |
| def _hash6(*parts: str) -> str: | |
| """Stable 6-hex suffix, standing in for the Marin executor's config hash. | |
| The results.md compiler only substring-matches ``<TASK>_seed<S>-`` and | |
| ``compile_<TASK>_avg<K>seeds-``, so the value is opaque -- but it must be | |
| STABLE so re-runs land in the same directory instead of piling up. | |
| """ | |
| return hashlib.md5("\x00".join(parts).encode()).hexdigest()[:6] | |
| def step_of(checkpoint: Path) -> str | None: | |
| match = STEP_RE.search(str(checkpoint)) | |
| return match.group(1) if match else None | |
| class Layout: | |
| """Every path a single (checkpoint, task) evaluation reads or writes.""" | |
| results_root: Path | |
| experiment: str | |
| step: str | None | |
| task: str | |
| seeds: tuple[int, ...] | |
| model_dir_name: str | |
| def eval_dir(self) -> Path: | |
| name = f"{self.experiment}-step{self.step}" if self.step else self.experiment | |
| return self.results_root / name | |
| def seed_dir(self, seed: int) -> Path: | |
| suffix = _hash6(self.experiment, self.step or "", self.task, str(seed)) | |
| return self.eval_dir / f"{self.task}_seed{seed}-{suffix}" | |
| def seed_results_dir(self, seed: int) -> Path: | |
| return self.seed_dir(seed) / f"{self.task}_{NUM_FEWSHOT}shot" / self.model_dir_name | |
| def compile_dir(self) -> Path | None: | |
| if len(self.seeds) < 2: | |
| return None | |
| suffix = _hash6(self.experiment, self.step or "", self.task, "compile") | |
| return self.eval_dir / f"compile_{self.task}_avg{len(self.seeds)}seeds-{suffix}" | |
| def averaged_results(self) -> Path | None: | |
| compile_dir = self.compile_dir | |
| return None if compile_dir is None else compile_dir / "compiled_results" / "averaged_results.json" | |
| def raw_dir(self) -> Path: | |
| return self.eval_dir / "_raw" / self.task | |
| def marker(self) -> Path: | |
| return self.eval_dir / "_markers" / f"{self.task}.complete.json" | |
| def find_results_file(root: Path) -> Path | None: | |
| candidates = [p for p in root.rglob("results_*.json") if RESULTS_FILE_RE.match(p.name)] | |
| if not candidates: | |
| return None | |
| return max(candidates, key=lambda p: p.stat().st_mtime) | |
| def read_task_accuracy(results_file: Path, task: str) -> float | None: | |
| """Read the metric the results.md compiler reads: accuracy_avg, else accuracy.""" | |
| try: | |
| payload = json.loads(results_file.read_text()) | |
| except (OSError, json.JSONDecodeError): | |
| return None | |
| entry = payload.get("results", {}).get(task) | |
| if not isinstance(entry, dict): | |
| return None | |
| for key in ("accuracy_avg", "accuracy"): | |
| value = entry.get(key) | |
| if isinstance(value, (int, float)): | |
| return float(value) | |
| return None | |
| def is_task_complete(layout: Layout) -> bool: | |
| """Same three-part definition the eval monitor and results.md compiler use.""" | |
| for seed in layout.seeds: | |
| results_file = find_results_file(layout.seed_dir(seed)) | |
| if results_file is None or read_task_accuracy(results_file, layout.task) is None: | |
| return False | |
| averaged = layout.averaged_results | |
| if averaged is not None: | |
| if not averaged.exists(): | |
| return False | |
| try: | |
| payload = json.loads(averaged.read_text()) | |
| if payload[0]["num_seeds"] != len(layout.seeds): | |
| return False | |
| except (OSError, json.JSONDecodeError, KeyError, IndexError): | |
| return False | |
| return True | |
| # -------------------------------------------------------------------------- | |
| # Running one task | |
| # -------------------------------------------------------------------------- | |
| def build_model_args(checkpoint: Path, *, profile: GpuProfile, gmu: float, tensor_parallel_size: int) -> str: | |
| parts = [ | |
| f"pretrained={checkpoint}", | |
| "dtype=bfloat16", | |
| f"tensor_parallel_size={tensor_parallel_size}", | |
| f"gpu_memory_utilization={gmu}", | |
| f"max_model_len={MAX_MODEL_LEN}", | |
| f"max_gen_toks={MAX_GEN_TOKS}", | |
| f"max_num_seqs={profile.max_num_seqs}", | |
| f"max_num_batched_tokens={profile.max_num_batched_tokens}", | |
| "trust_remote_code=True", | |
| # Deterministic-ish prefix reuse across the K repetitions of a task: | |
| # every repetition sends the identical prompt, so the prefix cache turns | |
| # K prefills into 1. Big win at 10 seeds. | |
| "enable_prefix_caching=True", | |
| ] | |
| return ",".join(parts) | |
| def subprocess_env( | |
| *, | |
| evalchemy_dir: Path, | |
| n_repeat: int, | |
| num_proc: int, | |
| hf_cache: Path | None, | |
| extra: dict[str, str] | None = None, | |
| ) -> dict[str, str]: | |
| env = os.environ.copy() | |
| env.update( | |
| { | |
| "EVALCHEMY_N_REPEAT": str(n_repeat), | |
| "EVALCHEMY_NUM_PROC": str(num_proc), | |
| # LiveCodeBench loads a dataset SCRIPT; without this it refuses. | |
| "HF_DATASETS_TRUST_REMOTE_CODE": "1", | |
| # lm-eval refuses to execute generated code without this. | |
| "HF_ALLOW_CODE_EVAL": "1", | |
| "TOKENIZERS_PARALLELISM": "false", | |
| "PYTHONUNBUFFERED": "1", | |
| # The LCB grader spawns children; keep them off the GPU entirely. | |
| "VLLM_WORKER_MULTIPROC_METHOD": "spawn", | |
| "PYTHONPATH": os.pathsep.join( | |
| [str(evalchemy_dir), *([os.environ["PYTHONPATH"]] if os.environ.get("PYTHONPATH") else [])] | |
| ), | |
| } | |
| ) | |
| if hf_cache is not None: | |
| env.update( | |
| { | |
| "HF_HOME": str(hf_cache), | |
| "HF_HUB_CACHE": str(hf_cache / "hub"), | |
| "HF_DATASETS_CACHE": str(hf_cache / "datasets"), | |
| } | |
| ) | |
| if extra: | |
| env.update(extra) | |
| return env | |
| def run_evalchemy( | |
| *, | |
| evalchemy_dir: Path, | |
| checkpoint: Path, | |
| task: str, | |
| n_repeat: int, | |
| out_dir: Path, | |
| profile: GpuProfile, | |
| gmu: float, | |
| tensor_parallel_size: int, | |
| batch_size: int, | |
| hf_cache: Path | None, | |
| num_proc: int, | |
| debug: bool, | |
| timeout: int | None, | |
| ) -> Path: | |
| """Invoke evalchemy once for one task, producing one results_*.json.""" | |
| entry = evalchemy_dir / "_marin_gpu_entry.py" | |
| if not entry.exists(): | |
| raise RuntimeError(f"{entry} missing -- run patch_evalchemy_gpu.py against {evalchemy_dir}") | |
| if out_dir.exists(): | |
| # Exactly one results file per run, so downstream globs are unambiguous. | |
| shutil.rmtree(out_dir) | |
| out_dir.mkdir(parents=True, exist_ok=True) | |
| command = [ | |
| sys.executable, | |
| str(entry), | |
| "--model", | |
| "vllm", | |
| "--tasks", | |
| task, | |
| "--model_args", | |
| build_model_args(checkpoint, profile=profile, gmu=gmu, tensor_parallel_size=tensor_parallel_size), | |
| "--batch_size", | |
| str(batch_size), | |
| "--output_path", | |
| str(out_dir), | |
| "--verbosity", | |
| "INFO", | |
| "--apply_chat_template", | |
| "--confirm_run_unsafe_code", | |
| "--max_tokens", | |
| str(MAX_GEN_TOKS), | |
| # evalchemy hands seeds[0] to vLLM as the per-request seed and adds the | |
| # repetition index, so repetition i runs at seed SEED_BASE + i. | |
| "--seed", | |
| ",".join([str(SEED_BASE)] * 4), | |
| "--gen_kwargs", | |
| f"temperature={TEMPERATURE},top_p={TOP_P},max_gen_toks={MAX_GEN_TOKS}", | |
| ] | |
| if debug: | |
| command.append("--debug") | |
| env = subprocess_env(evalchemy_dir=evalchemy_dir, n_repeat=n_repeat, num_proc=num_proc, hf_cache=hf_cache) | |
| log_path = out_dir.parent / f"{task}.evalchemy.log" | |
| log_path.parent.mkdir(parents=True, exist_ok=True) | |
| logger.info("launching %s (n_repeat=%d) -> %s", task, n_repeat, log_path) | |
| logger.debug("command: %s", " ".join(command)) | |
| started = time.time() | |
| with log_path.open("w") as log_file: | |
| log_file.write(f"# {' '.join(command)}\n\n") | |
| log_file.flush() | |
| process = subprocess.Popen( | |
| command, | |
| cwd=str(evalchemy_dir), | |
| env=env, | |
| stdout=log_file, | |
| stderr=subprocess.STDOUT, | |
| ) | |
| try: | |
| returncode = process.wait(timeout=timeout) | |
| except subprocess.TimeoutExpired: | |
| process.kill() | |
| process.wait() | |
| raise RuntimeError(f"{task} exceeded timeout of {timeout}s; see {log_path}") from None | |
| elapsed = time.time() - started | |
| results_file = find_results_file(out_dir) | |
| # Exit 0 with no results file means scoring hung or crashed silently -- the | |
| # same guard the TPU evaluator carries. Treat it as a hard failure. | |
| if returncode != 0: | |
| raise RuntimeError(f"{task} exited {returncode} after {elapsed:.0f}s; see {log_path}") | |
| if results_file is None: | |
| raise RuntimeError( | |
| f"{task} exited 0 after {elapsed:.0f}s but wrote no results_*.json under {out_dir}; " | |
| f"scoring likely hung or crashed silently. See {log_path}" | |
| ) | |
| logger.info("%s finished in %.0fs -> %s", task, elapsed, results_file) | |
| return results_file | |
| # -------------------------------------------------------------------------- | |
| # Fan-out: one combined run -> per-seed result files + a compile file | |
| # -------------------------------------------------------------------------- | |
| def per_repetition_accuracies(task_entry: dict[str, Any], n_repeat: int, task: str) -> list[float]: | |
| run_stats = task_entry.get("run_stats") | |
| if isinstance(run_stats, list) and len(run_stats) == n_repeat: | |
| return [float(stat["accuracy"]) for stat in run_stats] | |
| if n_repeat == 1: | |
| for key in ("accuracy_avg", "accuracy"): | |
| value = task_entry.get(key) | |
| if isinstance(value, (int, float)): | |
| return [float(value)] | |
| raise RuntimeError( | |
| f"{task}: expected {n_repeat} repetitions but the result carries " | |
| f"run_stats={type(run_stats).__name__} of length " | |
| f"{len(run_stats) if isinstance(run_stats, list) else 'n/a'}. " | |
| "The n_repeat patch probably did not apply -- re-run patch_evalchemy_gpu.py." | |
| ) | |
| def slice_examples(examples: list[dict[str, Any]], index: int) -> list[dict[str, Any]] | None: | |
| """Project a multi-repetition example list down to one repetition. | |
| Returns None when the benchmark does not keep per-repetition outputs | |
| (LiveCodeBench* keep only the last run's examples), in which case the | |
| per-seed file simply omits ``examples``. | |
| """ | |
| sliced: list[dict[str, Any]] = [] | |
| for example in examples: | |
| answers = example.get("model_answers") | |
| outputs = example.get("model_outputs") | |
| if not isinstance(answers, list) or index >= len(answers): | |
| return None | |
| projected = dict(example) | |
| projected["model_answers"] = [answers[index]] | |
| if isinstance(outputs, list) and index < len(outputs): | |
| projected["model_outputs"] = [outputs[index]] | |
| sliced.append(projected) | |
| return sliced | |
| def write_per_seed_results( | |
| *, | |
| raw_results_file: Path, | |
| layout: Layout, | |
| task: str, | |
| checkpoint: Path, | |
| ) -> dict[int, float]: | |
| """Split one combined run into K schema-compatible per-seed result files.""" | |
| raw = json.loads(raw_results_file.read_text()) | |
| entry = raw.get("results", {}).get(task) | |
| if not isinstance(entry, dict): | |
| available = sorted(raw.get("results", {})) | |
| raise RuntimeError(f"{raw_results_file} has no results[{task!r}]; found {available}") | |
| seeds = layout.seeds | |
| accuracies = per_repetition_accuracies(entry, len(seeds), task) | |
| run_stats = entry.get("run_stats") if isinstance(entry.get("run_stats"), list) else None | |
| examples = entry.get("examples") if isinstance(entry.get("examples"), list) else None | |
| per_seed: dict[int, float] = {} | |
| timestamp = datetime.now().isoformat().replace(":", "-") | |
| for index, seed in enumerate(seeds): | |
| accuracy = accuracies[index] | |
| stat = run_stats[index] if run_stats and index < len(run_stats) else {} | |
| num_total = stat.get("num_total", entry.get("num_total")) | |
| num_solved = stat.get("num_solved", entry.get("num_solved")) | |
| seed_entry: dict[str, Any] = { | |
| "num_total": num_total, | |
| "num_solved": num_solved, | |
| # Both keys, because the results.md compiler prefers accuracy_avg and | |
| # falls back to accuracy, while some downstream readers expect the | |
| # single-pass shape. | |
| "accuracy": accuracy, | |
| "accuracy_avg": accuracy, | |
| "accuracy_std_err": 0.0, | |
| "num_repeat": 1, | |
| "solved_avg": num_solved, | |
| "run_stats": [ | |
| { | |
| "repetition": 1, | |
| "num_total": num_total, | |
| "num_solved": num_solved, | |
| "accuracy": accuracy, | |
| } | |
| ], | |
| } | |
| if examples is not None: | |
| projected = slice_examples(examples, index) | |
| if projected is not None: | |
| seed_entry["examples"] = projected | |
| # Carry through any extra per-difficulty metrics LCB emits. | |
| for key, value in entry.items(): | |
| if key.startswith("accuracy_") and key.endswith("_avg") and key != "accuracy_avg": | |
| seed_entry[key] = value | |
| payload = {k: v for k, v in raw.items() if k != "results"} | |
| payload["results"] = {task: seed_entry} | |
| payload["_marin_gpu_port"] = { | |
| "source_results_file": str(raw_results_file), | |
| "repetition_index": index, | |
| "seed": seed, | |
| "note": ( | |
| "Synthesised from a single K-repetition evalchemy run. On GPU the " | |
| "seed is per-request (SEED_BASE + repetition_index), not an engine " | |
| "seed as on TPU." | |
| ), | |
| "checkpoint": str(checkpoint), | |
| } | |
| out_dir = layout.seed_results_dir(seed) | |
| out_dir.mkdir(parents=True, exist_ok=True) | |
| for stale in out_dir.glob("results_*.json"): | |
| stale.unlink() | |
| (out_dir / f"results_{timestamp}.json").write_text(json.dumps(payload, indent=2, default=str)) | |
| per_seed[seed] = accuracy | |
| logger.info("%s per-seed accuracies: %s", task, {s: round(a, 4) for s, a in per_seed.items()}) | |
| return per_seed | |
| def write_compile_results(*, layout: Layout, per_seed: dict[int, float], checkpoint: Path) -> None: | |
| """Write the averaged_results.json that claude/compile_results.py reads. | |
| Schema matches ``experiments/evals/evalchemy_results_compiler.py`` exactly: | |
| a one-element JSON array with base_model_name, dataset_name, num_seeds, | |
| seeds[], correct_mean, correct_std, correct_per_seed{}. ``correct_std`` is a | |
| SAMPLE standard deviation (ddof=1), matching pandas Series.std(), and it is | |
| what renders as the bracketed [x.x] in results.md. | |
| Unlike the TPU compile step, the numbers here come from each benchmark's own | |
| grader rather than a naive string re-match, so HMMT / JEEBench / | |
| OlympiadBench_Physics / LiveCodeBench* are correct in this file too. The | |
| results.md compiler still bypasses the compile dir for those six and averages | |
| the per-seed files instead; both paths now agree. | |
| """ | |
| compile_dir = layout.compile_dir | |
| if compile_dir is None: | |
| return | |
| seeds = sorted(per_seed) | |
| values = [per_seed[s] for s in seeds] | |
| mean = statistics.fmean(values) | |
| std = statistics.stdev(values) if len(values) > 1 else 0.0 | |
| record = { | |
| "base_model_name": layout.model_dir_name.lower(), | |
| "dataset_name": layout.task.lower(), | |
| "num_seeds": len(seeds), | |
| "seeds": seeds, | |
| "correct_mean": mean, | |
| "correct_std": std, | |
| "correct_per_seed": {str(s): per_seed[s] for s in seeds}, | |
| } | |
| out_dir = compile_dir / "compiled_results" | |
| out_dir.mkdir(parents=True, exist_ok=True) | |
| (out_dir / "averaged_results.json").write_text(json.dumps([record], indent=2)) | |
| with (out_dir / "averaged_results.csv").open("w", newline="") as handle: | |
| writer = csv.writer(handle) | |
| writer.writerow(["base_model_name", "dataset_name", "num_seeds", "seeds", "correct_mean", "correct_std"]) | |
| writer.writerow( | |
| [record["base_model_name"], record["dataset_name"], record["num_seeds"], seeds, mean, std] | |
| ) | |
| # compiled_results.{json,csv} exist for parity with the TPU compile step. | |
| # NOTE: the TPU version stores one row per graded EXAMPLE; we store one row | |
| # per seed, because we take accuracy from the benchmark's own grader rather | |
| # than re-grading examples with string equality. Nothing downstream reads it. | |
| rows = [ | |
| { | |
| "dataset_name": layout.task.lower(), | |
| "model_name": layout.model_dir_name.lower(), | |
| "seed": seed, | |
| "accuracy": per_seed[seed], | |
| "checkpoint": str(checkpoint), | |
| } | |
| for seed in seeds | |
| ] | |
| (out_dir / "compiled_results.json").write_text(json.dumps(rows, indent=2)) | |
| with (out_dir / "compiled_results.csv").open("w", newline="") as handle: | |
| writer = csv.DictWriter(handle, fieldnames=list(rows[0])) | |
| writer.writeheader() | |
| writer.writerows(rows) | |
| logger.info( | |
| "%s compiled: mean=%.4f std=%.4f over %d seeds -> %s", | |
| layout.task, | |
| mean, | |
| std, | |
| len(seeds), | |
| out_dir / "averaged_results.json", | |
| ) | |
| # -------------------------------------------------------------------------- | |
| # Checkpoint validation | |
| # -------------------------------------------------------------------------- | |
| def validate_checkpoint(checkpoint: Path) -> None: | |
| """Reject half-written HF exports before burning an hour of GPU on them. | |
| Levanter writes hf/step-N incrementally while training continues; a monitor | |
| that races the export will otherwise load a truncated shard set. | |
| """ | |
| if not checkpoint.is_dir(): | |
| raise RuntimeError(f"{checkpoint} is not a directory") | |
| if not (checkpoint / "config.json").exists(): | |
| raise RuntimeError(f"{checkpoint}/config.json missing") | |
| if not any((checkpoint / name).exists() for name in ("tokenizer_config.json", "tokenizer.json")): | |
| raise RuntimeError(f"{checkpoint} has no tokenizer files") | |
| index_path = checkpoint / "model.safetensors.index.json" | |
| if index_path.exists(): | |
| index = json.loads(index_path.read_text()) | |
| shards = sorted(set(index.get("weight_map", {}).values())) | |
| missing = [s for s in shards if not (checkpoint / s).exists()] | |
| if missing: | |
| raise RuntimeError(f"{checkpoint} is incomplete: {len(missing)} of {len(shards)} shards missing: {missing[:3]}") | |
| return | |
| if not (checkpoint / "model.safetensors").exists(): | |
| raise RuntimeError(f"{checkpoint} has neither model.safetensors nor a shard index") | |
| # -------------------------------------------------------------------------- | |
| # Entry point | |
| # -------------------------------------------------------------------------- | |
| def resolve_tasks(args: argparse.Namespace) -> dict[str, int]: | |
| if args.tasks: | |
| selected = {} | |
| for name in args.tasks: | |
| if name not in ALL_TASK_SEEDS: | |
| raise SystemExit(f"unknown task {name!r}; known: {', '.join(sorted(ALL_TASK_SEEDS))}") | |
| selected[name] = ALL_TASK_SEEDS[name] | |
| elif args.suite: | |
| selected = dict(SUITES[args.suite]) | |
| else: | |
| raise SystemExit("pass --suite or --tasks") | |
| if args.seeds is not None: | |
| selected = {task: args.seeds for task in selected} | |
| return selected | |
| def evaluate_task(args: argparse.Namespace, task: str, n_seeds: int, profile: GpuProfile, gmu: float) -> dict[str, Any]: | |
| checkpoint = args.checkpoint.resolve() | |
| effective_seeds = 1 if task in SINGLE_PASS_TASKS else n_seeds | |
| seeds = tuple(SEED_BASE + i for i in range(effective_seeds)) | |
| layout = Layout( | |
| results_root=args.results_root.resolve(), | |
| experiment=args.experiment, | |
| step=step_of(checkpoint), | |
| task=task, | |
| seeds=seeds, | |
| model_dir_name=sanitize_model_name(str(checkpoint)), | |
| ) | |
| if not args.force and is_task_complete(layout): | |
| logger.info("%s already complete for %s; skipping", task, layout.eval_dir.name) | |
| return {"task": task, "status": "cached", "seeds": list(seeds)} | |
| started = time.time() | |
| raw_results_file = run_evalchemy( | |
| evalchemy_dir=args.evalchemy_dir.resolve(), | |
| checkpoint=checkpoint, | |
| task=task, | |
| n_repeat=effective_seeds, | |
| out_dir=layout.raw_dir, | |
| profile=profile, | |
| gmu=gmu, | |
| tensor_parallel_size=args.tensor_parallel_size, | |
| batch_size=args.batch_size, | |
| hf_cache=args.hf_cache, | |
| num_proc=args.num_proc, | |
| debug=args.debug, | |
| timeout=args.task_timeout, | |
| ) | |
| per_seed = write_per_seed_results( | |
| raw_results_file=raw_results_file, layout=layout, task=task, checkpoint=checkpoint | |
| ) | |
| write_compile_results(layout=layout, per_seed=per_seed, checkpoint=checkpoint) | |
| summary = { | |
| "task": task, | |
| "status": "ok", | |
| "seeds": list(seeds), | |
| "mean": statistics.fmean(per_seed.values()), | |
| "std": statistics.stdev(per_seed.values()) if len(per_seed) > 1 else 0.0, | |
| "per_seed": {str(k): v for k, v in per_seed.items()}, | |
| "elapsed_seconds": round(time.time() - started, 1), | |
| } | |
| layout.marker.parent.mkdir(parents=True, exist_ok=True) | |
| layout.marker.write_text(json.dumps({**summary, "checkpoint": str(checkpoint)}, indent=2)) | |
| if args.prune_raw and layout.raw_dir.exists(): | |
| shutil.rmtree(layout.raw_dir) | |
| logger.info("pruned %s", layout.raw_dir) | |
| return summary | |
| def build_parser() -> argparse.ArgumentParser: | |
| parser = argparse.ArgumentParser( | |
| description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter | |
| ) | |
| parser.add_argument("--checkpoint", type=Path, required=True, help="Local HF checkpoint dir (.../hf/step-N)") | |
| parser.add_argument( | |
| "--experiment", | |
| required=True, | |
| help="Experiment name WITHOUT a -stepN suffix; the step is taken from --checkpoint", | |
| ) | |
| parser.add_argument("--suite", choices=sorted(SUITES), help="math | science | code") | |
| parser.add_argument("--tasks", nargs="+", help="Explicit evalchemy task names (overrides --suite)") | |
| parser.add_argument("--seeds", type=int, help="Override the per-task seed count (default: TPU parity)") | |
| parser.add_argument("--results-root", type=Path, required=True, help="Local stand-in for gs://.../evaluation/evalchemy") | |
| parser.add_argument( | |
| "--evalchemy-dir", | |
| type=Path, | |
| default=Path(os.environ.get("EVALCHEMY_DIR", "/opt/marin-gpu-eval/evalchemy")), | |
| help="Patched evalchemy checkout (default: $EVALCHEMY_DIR)", | |
| ) | |
| parser.add_argument("--gpu-profile", choices=sorted(GPU_PROFILES), help="Default: auto-detect from nvidia-smi") | |
| parser.add_argument("--tensor-parallel-size", type=int, help="Default: number of visible GPUs") | |
| parser.add_argument("--batch-size", type=int, default=64, help="lm-eval batch size (default 64)") | |
| parser.add_argument( | |
| "--gpu-memory-utilization", | |
| type=float, | |
| help="Override the profile value. Still clamped by free VRAM unless --no-vram-clamp.", | |
| ) | |
| parser.add_argument( | |
| "--vram-headroom-mib", | |
| type=int, | |
| default=2048, | |
| help="MiB left unclaimed on the busiest GPU when clamping (default 2048)", | |
| ) | |
| parser.add_argument("--no-vram-clamp", action="store_true", help="Trust the profile value verbatim (unsafe)") | |
| parser.add_argument("--hf-cache", type=Path, help="HF_HOME for the eval subprocess") | |
| parser.add_argument("--num-proc", type=int, default=8, help="Dataset map/filter workers (default 8)") | |
| parser.add_argument("--task-timeout", type=int, default=None, help="Seconds per task before SIGKILL") | |
| parser.add_argument("--debug", action="store_true", help="Smoke mode: evalchemy limits each task to 10 examples") | |
| parser.add_argument( | |
| "--prune-raw", | |
| action="store_true", | |
| help="Delete <eval_dir>/_raw/<TASK> after a successful fan-out. The raw file holds all K " | |
| "repetitions at full 32768-token outputs and is the largest artifact on disk.", | |
| ) | |
| parser.add_argument("--force", action="store_true", help="Re-run tasks that already have complete results") | |
| parser.add_argument("--dry-run", action="store_true", help="Print the plan and exit") | |
| parser.add_argument("--log-level", default="INFO") | |
| return parser | |
| def main(argv: list[str] | None = None) -> int: | |
| args = build_parser().parse_args(argv) | |
| logging.basicConfig( | |
| level=getattr(logging, args.log_level.upper()), | |
| format="%(asctime)s %(levelname)-7s %(name)s | %(message)s", | |
| ) | |
| tasks = resolve_tasks(args) | |
| validate_checkpoint(args.checkpoint.resolve()) | |
| gpus = query_gpus() | |
| profile = GPU_PROFILES[args.gpu_profile] if args.gpu_profile else detect_profile(gpus) | |
| if args.gpu_memory_utilization is not None: | |
| profile = dataclasses.replace(profile, gpu_memory_utilization=args.gpu_memory_utilization) | |
| if args.tensor_parallel_size is None: | |
| args.tensor_parallel_size = len(gpus) | |
| gmu = ( | |
| profile.gpu_memory_utilization | |
| if args.no_vram_clamp | |
| else resolve_gpu_memory_utilization(profile, gpus, args.vram_headroom_mib) | |
| ) | |
| logger.info("checkpoint %s (step %s)", args.checkpoint, step_of(args.checkpoint.resolve())) | |
| logger.info("gpus %s", [f"{g.index}:{g.name}:{g.free_mib}MiB free" for g in gpus]) | |
| logger.info("profile %s (%s)", profile.name, profile.notes) | |
| logger.info( | |
| "engine tp=%d gmu=%.3f max_num_seqs=%d max_model_len=%d", | |
| args.tensor_parallel_size, | |
| gmu, | |
| profile.max_num_seqs, | |
| MAX_MODEL_LEN, | |
| ) | |
| per_token = kv_bytes_per_token(args.checkpoint.resolve()) | |
| if per_token: | |
| seq_gib = per_token * MAX_MODEL_LEN / 1024**3 | |
| logger.info("kv cache %d B/token -> %.2f GiB per full-length sequence", per_token, seq_gib) | |
| logger.info("tasks %s", {t: (1 if t in SINGLE_PASS_TASKS else n) for t, n in tasks.items()}) | |
| if args.dry_run: | |
| return 0 | |
| summaries: list[dict[str, Any]] = [] | |
| failures: list[tuple[str, str]] = [] | |
| for task, n_seeds in tasks.items(): | |
| try: | |
| summaries.append(evaluate_task(args, task, n_seeds, profile, gmu)) | |
| except Exception as error: # one bad task must not abandon the rest | |
| logger.exception("task %s failed", task) | |
| failures.append((task, str(error))) | |
| print(json.dumps({"summaries": summaries, "failures": failures}, indent=2)) | |
| return 1 if failures else 0 | |
| if __name__ == "__main__": | |
| sys.exit(main()) | |