Download gpu-sft/scripts/gpu_eval/gpu_eval_monitor.py from fzzhang/svd-code: direct link, hf CLI and curl.
- Browser
- Download file 19.6 kB
-
https://huggingface.co/fzzhang/svd-code/resolve/main/gpu-sft/scripts/gpu_eval/gpu_eval_monitor.py
- Command line
-
hf download hf://fzzhang/svd-code/gpu-sft/scripts/gpu_eval/gpu_eval_monitor.py
-
curl -L -o gpu_eval_monitor.py https://huggingface.co/fzzhang/svd-code/resolve/main/gpu-sft/scripts/gpu_eval/gpu_eval_monitor.py
19.6 kB
| #!/usr/bin/env python3 | |
| # Copyright The Marin Authors | |
| # SPDX-License-Identifier: Apache-2.0 | |
| """Eval-during-training monitor for a single local GPU box. | |
| Local replacement for ``claude/eval_monitor.py``: GCS listing becomes a | |
| filesystem walk, ``iris job run`` becomes a ``gpu_eval_driver.py`` subprocess, | |
| and the "is another job already running this?" question becomes an OS file lock. | |
| python gpu_eval_monitor.py \\ | |
| --run-dir /data/runs/exp_sft_qwen3_8b_x \\ | |
| --experiment exp_sft_qwen3_8b_x \\ | |
| --suite math \\ | |
| --results-root /data/marin/evaluation/evalchemy \\ | |
| --gpu-lock /var/tmp/marin-gpu.lock | |
| HOW GPU ACCESS IS SERIALISED AGAINST TRAINING | |
| One advisory ``flock(2)`` on --gpu-lock. The lock is taken and released | |
| around EACH TASK, not around a whole suite, so a training restart never | |
| waits hours behind a 10-seed AIME sweep. | |
| The training launcher MUST take the same lock. Wrap it: | |
| flock /var/tmp/marin-gpu.lock -c 'python -m levanter.main.train_lm ...' | |
| (util-linux ``flock`` and Python's ``fcntl.flock`` both call flock(2) on the | |
| file, so they interlock correctly. ``flock`` blocks by default; add ``-w N`` | |
| for a timeout.) | |
| The lock alone is not sufficient, because JAX preallocates ~75% of every | |
| visible GPU at import and holds it until the process exits, and because | |
| nothing forces a stray notebook or embedding server to cooperate. So the | |
| monitor ALSO refuses to launch until the busiest visible GPU reports at | |
| least --min-free-vram-mib free. Both gates must pass. | |
| If you have spare cards, the better answer is physical separation: give | |
| training ``CUDA_VISIBLE_DEVICES=0,1,2,3`` and run this monitor with | |
| ``--cuda-visible-devices 4,5,6,7 --no-gpu-lock``. Then nothing contends and | |
| evals never stall the run. | |
| DEDUPE | |
| Three layers, matching the TPU monitor: | |
| 1. Artifact probe -- gpu_eval_driver.is_task_complete() checks that every | |
| expected per-seed results file parses and the compile file has the full | |
| seed count. This is the source of truth; markers are only a fast path. | |
| 2. In-process set of (step, task) pairs launched this session. | |
| 3. The lock itself, which prevents two monitors from racing one GPU. | |
| CHECKPOINT READINESS | |
| A step is only eligible once its HF export is structurally complete: every | |
| shard named in model.safetensors.index.json exists, config.json exists, and | |
| the tokenizer is present. Levanter writes hf/step-N while training | |
| continues, so a naive directory-mtime check will load a truncated export. | |
| A --settle-seconds quiet period on the newest file guards the tail. | |
| """ | |
| from __future__ import annotations | |
| import argparse | |
| import contextlib | |
| import errno | |
| import fcntl | |
| import json | |
| import logging | |
| import os | |
| import re | |
| import signal | |
| import subprocess | |
| import sys | |
| import time | |
| from collections.abc import Iterator | |
| from datetime import UTC, datetime | |
| from pathlib import Path | |
| from typing import Any | |
| sys.path.insert(0, str(Path(__file__).resolve().parent)) | |
| import gpu_eval_driver as driver # noqa: E402 | |
| logger = logging.getLogger("gpu_eval_monitor") | |
| STEP_DIR_RE = re.compile(r"^(?:step-|checkpoint-)(\d+)$") | |
| DEFAULT_POLL_SECONDS = 300 | |
| DEFAULT_MAX_RUNTIME_DAYS = 14 | |
| _SHUTDOWN = False | |
| def _handle_signal(signum: int, _frame: Any) -> None: | |
| global _SHUTDOWN | |
| logger.warning("received signal %d; finishing the current task then exiting", signum) | |
| _SHUTDOWN = True | |
| # -------------------------------------------------------------------------- | |
| # GPU lock | |
| # -------------------------------------------------------------------------- | |
| def gpu_lock(path: Path | None, *, poll_seconds: int, wait_seconds: int | None) -> Iterator[bool]: | |
| """Exclusive flock(2) on ``path``. Yields True when held, False when disabled.""" | |
| if path is None: | |
| yield False | |
| return | |
| path.parent.mkdir(parents=True, exist_ok=True) | |
| deadline = None if wait_seconds is None else time.time() + wait_seconds | |
| handle = path.open("a+") | |
| try: | |
| announced = False | |
| while True: | |
| try: | |
| fcntl.flock(handle.fileno(), fcntl.LOCK_EX | fcntl.LOCK_NB) | |
| break | |
| except OSError as error: | |
| if error.errno not in (errno.EACCES, errno.EAGAIN): | |
| raise | |
| if deadline is not None and time.time() >= deadline: | |
| raise TimeoutError(f"could not acquire {path} within {wait_seconds}s") from None | |
| if not announced: | |
| logger.info("waiting for GPU lock %s (held by the training job?)", path) | |
| announced = True | |
| if _SHUTDOWN: | |
| raise TimeoutError("shutdown requested while waiting for the GPU lock") from None | |
| time.sleep(min(poll_seconds, 15)) | |
| handle.seek(0) | |
| handle.truncate() | |
| handle.write(f"{os.getpid()} gpu_eval_monitor {datetime.now(UTC).isoformat()}\n") | |
| handle.flush() | |
| logger.info("acquired GPU lock %s", path) | |
| yield True | |
| finally: | |
| with contextlib.suppress(OSError): | |
| fcntl.flock(handle.fileno(), fcntl.LOCK_UN) | |
| handle.close() | |
| def free_vram_mib() -> int: | |
| """Free MiB on the busiest visible GPU.""" | |
| return min(gpu.free_mib for gpu in driver.query_gpus()) | |
| # -------------------------------------------------------------------------- | |
| # Checkpoint discovery | |
| # -------------------------------------------------------------------------- | |
| def discover_steps(run_dir: Path, subdir: str) -> list[tuple[int, Path]]: | |
| """Return sorted (step, path) for every checkpoint directory under run_dir/subdir.""" | |
| root = run_dir / subdir if subdir else run_dir | |
| if not root.is_dir(): | |
| return [] | |
| found: list[tuple[int, Path]] = [] | |
| for child in root.iterdir(): | |
| if not child.is_dir(): | |
| continue | |
| match = STEP_DIR_RE.match(child.name) | |
| if match: | |
| found.append((int(match.group(1)), child)) | |
| return sorted(found) | |
| def checkpoint_is_settled(checkpoint: Path, settle_seconds: int) -> bool: | |
| """True when nothing under the checkpoint has been written for settle_seconds.""" | |
| if settle_seconds <= 0: | |
| return True | |
| newest = 0.0 | |
| for path in checkpoint.rglob("*"): | |
| if path.is_file(): | |
| newest = max(newest, path.stat().st_mtime) | |
| return newest > 0 and (time.time() - newest) >= settle_seconds | |
| def checkpoint_is_ready(checkpoint: Path, settle_seconds: int) -> tuple[bool, str]: | |
| try: | |
| driver.validate_checkpoint(checkpoint) | |
| except RuntimeError as error: | |
| return False, str(error) | |
| if not checkpoint_is_settled(checkpoint, settle_seconds): | |
| return False, f"still being written (quiet period {settle_seconds}s not met)" | |
| return True, "ready" | |
| # -------------------------------------------------------------------------- | |
| # Work items | |
| # -------------------------------------------------------------------------- | |
| def pending_tasks( | |
| *, | |
| experiment: str, | |
| checkpoint: Path, | |
| tasks: dict[str, int], | |
| results_root: Path, | |
| ) -> list[str]: | |
| """Tasks for this checkpoint whose artifacts are not already complete.""" | |
| step = driver.step_of(checkpoint) | |
| model_dir_name = driver.sanitize_model_name(str(checkpoint.resolve())) | |
| pending = [] | |
| for task, n_seeds in tasks.items(): | |
| effective = 1 if task in driver.SINGLE_PASS_TASKS else n_seeds | |
| layout = driver.Layout( | |
| results_root=results_root.resolve(), | |
| experiment=experiment, | |
| step=step, | |
| task=task, | |
| seeds=tuple(driver.SEED_BASE + i for i in range(effective)), | |
| model_dir_name=model_dir_name, | |
| ) | |
| if not driver.is_task_complete(layout): | |
| pending.append(task) | |
| return pending | |
| def run_driver(args: argparse.Namespace, checkpoint: Path, task: str) -> tuple[int, float]: | |
| """Run gpu_eval_driver.py for exactly one task in a fresh process. | |
| A fresh process per task is deliberate: vLLM does not reliably return all | |
| device memory to the allocator on engine teardown, so a long-lived driver | |
| would slowly starve itself (and the training job) across a 5-task suite. | |
| """ | |
| command = [ | |
| sys.executable, | |
| str(Path(__file__).resolve().parent / "gpu_eval_driver.py"), | |
| "--checkpoint", | |
| str(checkpoint), | |
| "--experiment", | |
| args.experiment, | |
| "--tasks", | |
| task, | |
| "--results-root", | |
| str(args.results_root), | |
| "--evalchemy-dir", | |
| str(args.evalchemy_dir), | |
| "--num-proc", | |
| str(args.num_proc), | |
| ] | |
| if args.gpu_profile: | |
| command += ["--gpu-profile", args.gpu_profile] | |
| if args.tensor_parallel_size: | |
| command += ["--tensor-parallel-size", str(args.tensor_parallel_size)] | |
| if args.hf_cache: | |
| command += ["--hf-cache", str(args.hf_cache)] | |
| if args.task_timeout: | |
| command += ["--task-timeout", str(args.task_timeout)] | |
| if args.prune_raw: | |
| command.append("--prune-raw") | |
| if args.debug: | |
| command.append("--debug") | |
| env = os.environ.copy() | |
| if args.cuda_visible_devices is not None: | |
| env["CUDA_VISIBLE_DEVICES"] = args.cuda_visible_devices | |
| logger.info("launching driver: %s %s", task, checkpoint) | |
| started = time.time() | |
| process = subprocess.run(command, env=env, check=False) | |
| return process.returncode, time.time() - started | |
| # -------------------------------------------------------------------------- | |
| # Event log | |
| # -------------------------------------------------------------------------- | |
| class EventLog: | |
| def __init__(self, path: Path | None) -> None: | |
| self.path = path | |
| if path is not None: | |
| path.parent.mkdir(parents=True, exist_ok=True) | |
| def emit(self, event: str, **fields: Any) -> None: | |
| record = {"ts": datetime.now(UTC).isoformat(), "event": event, **fields} | |
| logger.info("%s %s", event, {k: v for k, v in fields.items() if k != "ts"}) | |
| if self.path is not None: | |
| with self.path.open("a") as handle: | |
| handle.write(json.dumps(record, default=str) + "\n") | |
| # -------------------------------------------------------------------------- | |
| # Main loop | |
| # -------------------------------------------------------------------------- | |
| def select_steps(steps: list[tuple[int, Path]], args: argparse.Namespace) -> list[tuple[int, Path]]: | |
| selected = [ | |
| (step, path) | |
| for step, path in steps | |
| if step >= args.min_step | |
| and (args.max_step is None or step <= args.max_step) | |
| and (args.step_stride <= 1 or step % args.step_stride == 0) | |
| ] | |
| if args.newest_first: | |
| selected.reverse() | |
| if args.latest_only and selected: | |
| selected = selected[:1] if args.newest_first else selected[-1:] | |
| return selected | |
| def monitor(args: argparse.Namespace) -> int: | |
| events = EventLog(args.event_log) | |
| tasks = driver.SUITES[args.suite] if args.suite else {t: driver.ALL_TASK_SEEDS[t] for t in args.tasks} | |
| if args.seeds is not None: | |
| tasks = {task: args.seeds for task in tasks} | |
| deadline = time.time() + args.max_runtime_days * 24 * 3600 | |
| launched: set[tuple[int, str]] = set() | |
| events.emit( | |
| "start", | |
| run_dir=str(args.run_dir), | |
| experiment=args.experiment, | |
| tasks=list(tasks), | |
| poll_seconds=args.poll, | |
| gpu_lock=str(args.gpu_lock) if args.gpu_lock else None, | |
| min_free_vram_mib=args.min_free_vram_mib, | |
| ) | |
| while not _SHUTDOWN and time.time() < deadline: | |
| steps = select_steps(discover_steps(args.run_dir, args.checkpoint_subdir), args) | |
| if not steps: | |
| events.emit("no_checkpoints", root=str(args.run_dir / args.checkpoint_subdir)) | |
| for step, checkpoint in steps: | |
| if _SHUTDOWN: | |
| break | |
| ready, reason = checkpoint_is_ready(checkpoint, args.settle_seconds) | |
| if not ready: | |
| events.emit("checkpoint_not_ready", step=step, reason=reason) | |
| continue | |
| todo = pending_tasks( | |
| experiment=args.experiment, | |
| checkpoint=checkpoint, | |
| tasks=tasks, | |
| results_root=args.results_root, | |
| ) | |
| todo = [t for t in todo if (step, t) not in launched or args.retry_failed] | |
| if not todo: | |
| continue | |
| events.emit("step_pending", step=step, tasks=todo) | |
| for task in todo: | |
| if _SHUTDOWN: | |
| break | |
| free = free_vram_mib() | |
| if free < args.min_free_vram_mib: | |
| events.emit( | |
| "vram_blocked", | |
| step=step, | |
| task=task, | |
| free_mib=free, | |
| required_mib=args.min_free_vram_mib, | |
| hint="training or another process still holds the card", | |
| ) | |
| break # nothing on this box will run; go back to sleep | |
| try: | |
| with gpu_lock(args.gpu_lock, poll_seconds=args.poll, wait_seconds=args.lock_wait_seconds): | |
| # Re-check under the lock: free VRAM can change while we | |
| # queued behind the training job, and another monitor may | |
| # have finished this task in the meantime. | |
| free = free_vram_mib() | |
| if free < args.min_free_vram_mib: | |
| events.emit("vram_blocked_after_lock", step=step, task=task, free_mib=free) | |
| break | |
| if task not in pending_tasks( | |
| experiment=args.experiment, | |
| checkpoint=checkpoint, | |
| tasks={task: tasks[task]}, | |
| results_root=args.results_root, | |
| ): | |
| events.emit("raced", step=step, task=task) | |
| continue | |
| launched.add((step, task)) | |
| returncode, elapsed = run_driver(args, checkpoint, task) | |
| except TimeoutError as error: | |
| events.emit("lock_timeout", step=step, task=task, detail=str(error)) | |
| break | |
| if returncode == 0: | |
| events.emit("task_done", step=step, task=task, elapsed_seconds=round(elapsed, 1)) | |
| else: | |
| events.emit( | |
| "task_failed", | |
| step=step, | |
| task=task, | |
| returncode=returncode, | |
| elapsed_seconds=round(elapsed, 1), | |
| hint=f"see {args.results_root}/{args.experiment}-step{step}/_raw/{task}.evalchemy.log", | |
| ) | |
| if args.stop_on_failure: | |
| events.emit("stop_on_failure", step=step, task=task) | |
| return 1 | |
| if args.once: | |
| break | |
| events.emit("sleep", seconds=args.poll) | |
| slept = 0 | |
| while slept < args.poll and not _SHUTDOWN: | |
| time.sleep(min(5, args.poll - slept)) | |
| slept += 5 | |
| events.emit("exit", shutdown=_SHUTDOWN, expired=time.time() >= deadline) | |
| return 0 | |
| def build_parser() -> argparse.ArgumentParser: | |
| parser = argparse.ArgumentParser( | |
| description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter | |
| ) | |
| parser.add_argument("--run-dir", type=Path, required=True, help="Training output dir containing hf/step-N") | |
| parser.add_argument("--checkpoint-subdir", default="hf", help="Subdir holding checkpoints (default: hf)") | |
| parser.add_argument("--experiment", required=True, help="Experiment name WITHOUT a -stepN suffix") | |
| parser.add_argument("--results-root", type=Path, required=True) | |
| parser.add_argument("--suite", choices=sorted(driver.SUITES)) | |
| parser.add_argument("--tasks", nargs="+", help="Explicit task names (overrides --suite)") | |
| parser.add_argument("--seeds", type=int, help="Override per-task seed count") | |
| parser.add_argument( | |
| "--evalchemy-dir", | |
| type=Path, | |
| default=Path(os.environ.get("EVALCHEMY_DIR", "/opt/marin-gpu-eval/evalchemy")), | |
| ) | |
| parser.add_argument("--gpu-lock", type=Path, default=Path("/var/tmp/marin-gpu.lock")) | |
| parser.add_argument("--no-gpu-lock", action="store_true", help="Skip locking (use with disjoint --cuda-visible-devices)") | |
| parser.add_argument("--lock-wait-seconds", type=int, default=None, help="Give up waiting for the lock (default: wait forever)") | |
| parser.add_argument( | |
| "--min-free-vram-mib", | |
| type=int, | |
| default=70_000, | |
| help="Refuse to launch below this much free VRAM on the busiest visible GPU. " | |
| "70000 suits an 8B bf16 model at 36864 ctx on an 80GB card.", | |
| ) | |
| parser.add_argument("--cuda-visible-devices", help="CUDA_VISIBLE_DEVICES for eval subprocesses") | |
| parser.add_argument("--gpu-profile", choices=sorted(driver.GPU_PROFILES)) | |
| parser.add_argument("--tensor-parallel-size", type=int) | |
| parser.add_argument("--poll", type=int, default=DEFAULT_POLL_SECONDS, help="Seconds between passes") | |
| parser.add_argument("--settle-seconds", type=int, default=120, help="Quiet period a checkpoint must show") | |
| parser.add_argument("--step-stride", type=int, default=1, help="Only eval steps divisible by this") | |
| parser.add_argument("--min-step", type=int, default=0) | |
| parser.add_argument("--max-step", type=int, default=None) | |
| parser.add_argument("--latest-only", action="store_true", help="Only ever evaluate the newest eligible step") | |
| parser.add_argument("--newest-first", action="store_true", help="Walk steps newest-first") | |
| parser.add_argument("--retry-failed", action="store_true", help="Re-attempt tasks that failed earlier this session") | |
| parser.add_argument("--stop-on-failure", action="store_true") | |
| parser.add_argument("--once", action="store_true", help="Single pass, then exit") | |
| parser.add_argument("--max-runtime-days", type=float, default=DEFAULT_MAX_RUNTIME_DAYS) | |
| parser.add_argument("--task-timeout", type=int, default=None, help="Seconds per task before the driver is killed") | |
| parser.add_argument("--num-proc", type=int, default=8) | |
| parser.add_argument("--hf-cache", type=Path) | |
| parser.add_argument("--prune-raw", action="store_true") | |
| parser.add_argument("--debug", action="store_true", help="Smoke mode (10 examples per task)") | |
| parser.add_argument("--event-log", type=Path, help="Append JSONL events here") | |
| 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", | |
| ) | |
| if not args.suite and not args.tasks: | |
| raise SystemExit("pass --suite or --tasks") | |
| if args.tasks: | |
| unknown = [t for t in args.tasks if t not in driver.ALL_TASK_SEEDS] | |
| if unknown: | |
| raise SystemExit(f"unknown tasks: {unknown}") | |
| if args.no_gpu_lock: | |
| args.gpu_lock = None | |
| signal.signal(signal.SIGINT, _handle_signal) | |
| signal.signal(signal.SIGTERM, _handle_signal) | |
| return monitor(args) | |
| if __name__ == "__main__": | |
| sys.exit(main()) | |