diff --git a/inferencex-e2e/benchmarks/multi_node/llm-d/README.md b/inferencex-e2e/benchmarks/multi_node/llm-d/README.md index d9cdb687df..1282374954 100644 --- a/inferencex-e2e/benchmarks/multi_node/llm-d/README.md +++ b/inferencex-e2e/benchmarks/multi_node/llm-d/README.md @@ -9,7 +9,7 @@ srt-slurm launch path. |---|---| | `submit.sh` | sbatch wrapper. Validates env, exports tuning vars, returns `JOB_ID`. May read `slurm.time_limit` from the recipe to override `TIME_LIMIT`. | | `job.slurm` | sbatch entrypoint. Allocates `PREFILL_NODES + DECODE_NODES` nodes, derives per-node IPs, runs one Docker container per node via `srun`, threads role assignment env into each. | -| `server.sh` | Per-node entry. Reads `NODE_RANK = SLURM_PROCID`, picks role, starts vLLM (with the wide-EP / DeepEP / NIXL flag set from the llm-d wide-EP-lws guide), starts the pd-sidecar on each leader, and on the decode leader additionally writes `endpoints.yaml`, starts EPP + Envoy, runs `benchmark_serving.py`, and `scancel`s the job. | +| `server.sh` | Per-node entry. Reads `NODE_RANK = SLURM_PROCID`, picks role, starts vLLM (with the wide-EP / DeepEP / NIXL flag set from the llm-d wide-EP-lws guide), starts the pd-sidecar on each leader, and on the decode leader additionally writes `endpoints.yaml`, starts EPP + Envoy, runs `benchmark_serving.py`, and signals `job.slurm` to end the job. | ## Topology diff --git a/inferencex-e2e/benchmarks/multi_node/llm-d/job.slurm b/inferencex-e2e/benchmarks/multi_node/llm-d/job.slurm index fe36301b6a..433dc33ee5 100644 --- a/inferencex-e2e/benchmarks/multi_node/llm-d/job.slurm +++ b/inferencex-e2e/benchmarks/multi_node/llm-d/job.slurm @@ -63,25 +63,26 @@ export DOCKER_CONT_NAME : "${BENCHMARK_LOGS_DIR:?BENCHMARK_LOGS_DIR not set}" DOCKER_MOUNT_PATH="/workspace" -cleanup() { - echo "[${SLURM_JOB_ID}] cleanup on $(hostname)" - [[ -n "${WATCHER_PID:-}" ]] && kill "$WATCHER_PID" 2>/dev/null || true -} -trap cleanup INT TERM HUP EXIT - -# Coordinator-done watcher. server.sh on the decode coordinator writes -# this marker after the bench finishes; we then scancel the allocation -# from outside the container (the image has no SLURM client tools). -# Without this, workers `wait` on local vLLM forever and the job runs -# to TIME_LIMIT. +# Stop the step and exit 0 once the coordinator is done, so the job ends COMPLETED, not CANCELLED. BENCH_DONE_MARKER="$BENCHMARK_LOGS_DIR/.bench_done.$SLURM_JOB_ID" rm -f "$BENCH_DONE_MARKER" -( - while [[ ! -f "$BENCH_DONE_MARKER" ]]; do sleep 5; done - echo "[${SLURM_JOB_ID}] coordinator finished; scancel'ing job" - scancel "$SLURM_JOB_ID" 2>/dev/null || true -) & -WATCHER_PID=$! + +run_until_bench_done() { + "$@" & + local step=$! rc=0 + while kill -0 "$step" 2>/dev/null; do + if [[ -f "$BENCH_DONE_MARKER" ]]; then + echo "[${SLURM_JOB_ID}] coordinator finished; stopping the srun step" + kill -TERM "$step" 2>/dev/null || true + break + fi + sleep 5 + done + wait "$step" || rc=$? + [[ -f "$BENCH_DONE_MARKER" ]] && return 0 + echo "[${SLURM_JOB_ID}] srun step exited rc=$rc before the coordinator finished" >&2 + return "$rc" +} # Container engine: 'docker' (default) for clusters where the SLURM # user can talk to /var/run/docker.sock (e.g. h200-dgxc-slurm); 'pyxis' @@ -102,7 +103,7 @@ done if [[ "$LLMD_CONTAINER_ENGINE" == "docker" ]]; then # One docker run per node, one task per node. server.sh dispatches by NODE_RANK. - srun \ + run_until_bench_done srun \ --kill-on-bad-exit=1 \ --signal=TERM@30 \ --unbuffered \ @@ -261,7 +262,7 @@ elif [[ "$LLMD_CONTAINER_ENGINE" == "pyxis" ]]; then # MODEL_DIR / BENCHMARK_LOGS_DIR / NODE_RANK are translated to their # in-container values inside bash -lc (host MODEL_DIR is the source # path of the bind mount, but server.sh expects /models inside). - srun \ + run_until_bench_done srun \ --kill-on-bad-exit=1 \ --signal=TERM@30 \ --unbuffered \ diff --git a/inferencex-e2e/benchmarks/multi_node/llm-d/server.sh b/inferencex-e2e/benchmarks/multi_node/llm-d/server.sh index a4e8a94047..0933a380ee 100755 --- a/inferencex-e2e/benchmarks/multi_node/llm-d/server.sh +++ b/inferencex-e2e/benchmarks/multi_node/llm-d/server.sh @@ -569,8 +569,7 @@ PY ) fi - # Signal job.slurm (outside the container, where scancel exists) to release - # the allocation; without it workers wait until TIME_LIMIT. + # job.slurm stops the srun step once this marker exists. touch "$BENCHMARK_LOGS_DIR/.bench_done.$SLURM_JOB_ID" else # Workers (prefill leader, prefill/decode workers): keep vLLM alive. diff --git a/inferencex-e2e/infx/launch/drivers/__init__.py b/inferencex-e2e/infx/launch/drivers/__init__.py index f97953d3c2..8d2797b041 100644 --- a/inferencex-e2e/infx/launch/drivers/__init__.py +++ b/inferencex-e2e/infx/launch/drivers/__init__.py @@ -13,7 +13,7 @@ from infx.launch import policy from infx.launch.backends import backend_class from infx.launch.context import Launch, LaunchError -from infx.launch.drivers import script, srt +from infx.launch.drivers import llmd, script, srt from infx.launch.policy import LaunchPath, launch_path if TYPE_CHECKING: @@ -36,6 +36,7 @@ class Route: LaunchPath.SRT_MULTI: Route("slurm", srt.run_multinode), LaunchPath.SRT_NATIVE: Route("slurm", srt.run_multinode), LaunchPath.SCRIPT: Route(None, script.run), + LaunchPath.LLMD: Route("slurm", llmd.run), } diff --git a/inferencex-e2e/infx/launch/drivers/llmd.py b/inferencex-e2e/infx/launch/drivers/llmd.py new file mode 100644 index 0000000000..92587dd90b --- /dev/null +++ b/inferencex-e2e/infx/launch/drivers/llmd.py @@ -0,0 +1,158 @@ +"""llm-d vLLM multinode jobs submitted through benchmarks/multi_node/llm-d/submit.sh.""" + +from __future__ import annotations + +import os +import shutil +import subprocess +import sys +from pathlib import Path + +from infx.launch import artifacts, policy, proc +from infx.launch.backends.base import BackendError +from infx.launch.backends.slurm import cli +from infx.launch.context import Launch, LaunchError +from infx.launch.drivers.srt import models +from infx.launch.drivers.srt.run import slurm_backend +from infx.launch.request import LlmdRequest, RequestError + +CANCEL_TIMEOUT_S = 600.0 +DEFAULT_TIME_LIMIT = "08:00:00" +LLMD_DIR = "benchmarks/multi_node/llm-d" + + +def _find_eval_dir(logs_dir: Path) -> Path | None: + for root, dirs, _files in os.walk(logs_dir): + if "eval_results" in dirs: + return Path(root) / "eval_results" + return None + + +def _stage_agentic(logs_dir: Path, workspace: Path) -> None: + agentic = logs_dir / "agentic" + if not agentic.is_dir(): + return + staged = workspace / "LOGS" / "agentic" + staged.mkdir(parents=True, exist_ok=True) + for entry in agentic.iterdir(): + destination = staged / entry.name + if entry.is_dir(): + shutil.copytree(entry, destination, dirs_exist_ok=True) + elif entry.is_file(): + shutil.copy2(entry, destination) + + +def run(launch: Launch) -> int: + """Submit the llm-d Slurm job, follow its log, and stage benchmark artifacts.""" + backend = slurm_backend(launch) + request = LlmdRequest.from_env(launch.request.env) + if backend.settings.squash is None: + raise LaunchError(f"llmd-vllm: cluster {launch.cluster.id!r} has no slurm.squash") + if not request.disagg: + raise LaunchError("llmd-vllm supports only P/D disaggregated points") + + checkpoint = models.checkpoint(launch.cluster, request) + if checkpoint is None: + raise LaunchError( + f"cluster {launch.cluster.id!r} stages no checkpoint for MODEL={request.model}" + ) + model_path = models.host_path(launch.cluster, checkpoint) + if not checkpoint.node_local and not (model_path / "config.json").is_file(): + raise LaunchError(f"model checkpoint is unavailable: {model_path / 'config.json'}") + + squash = backend.prepare_image(request.image) + logs_dir = request.workspace / "benchmark_logs" + logs_dir.mkdir(parents=True, exist_ok=True) + + account = backend.settings.account or cli.default_account() + if not account: + raise RequestError.missing("SLURM_ACCOUNT") + + env = policy.runtime_env( + launch.cluster, + request, + models.job_env(launch.cluster, request, str(model_path)), + { + "SLURM_PARTITION": backend.settings.partition, + "SLURM_ACCOUNT": account, + "MODEL_PATH": str(model_path), + "MODEL_NAME": request.model, + "CONTAINER_IMAGE": request.image, + "GPUS_PER_NODE": str(launch.cluster.gpus_per_node), + "TIME_LIMIT": request.env.get("TIME_LIMIT") or DEFAULT_TIME_LIMIT, + "PREFILL_WORKERS": str(request.prefill_num_workers), + "DECODE_WORKERS": str(request.decode_num_workers), + "LLMD_CONTAINER_ENGINE": "pyxis", + "LLMD_SQUASH_FILE": squash.reference, + "BENCHMARK_LOGS_DIR": str(logs_dir), + }, + ) + + argv = [ + "bash", + "submit.sh", + str(request.prefill_nodes), + str(request.decode_nodes), + str(request.isl), + str(request.osl), + "x".join(map(str, request.conc_list)), + "inf", + request.random_range_ratio, + ] + proc.echo(argv, env) + submitted = subprocess.run( + argv, + stdout=subprocess.PIPE, + stderr=sys.stderr, + text=True, + env=env, + cwd=request.workspace / LLMD_DIR, + check=False, + ) + job_id = submitted.stdout.strip() + if submitted.returncode != 0 or not job_id: + print("ERROR: llm-d submit.sh failed before returning a Slurm job id", file=sys.stderr) + return 1 + if not (job_id.isascii() and job_id.isdigit()): + print( + f"ERROR: llm-d submit.sh printed {job_id!r} instead of a Slurm job id", + file=sys.stderr, + ) + return 1 + + log_file = logs_dir / f"slurm_job-{job_id}.out" + job = backend.attach(job_id, log=log_file, outputs=logs_dir) + print(f"Submitted llm-d job: {job_id}", flush=True) + + launch.life.callback( + artifacts.bundle_server_logs, logs_dir, request.workspace / "multinode_server_logs.tar.gz" + ) + launch.life.callback(backend.cancel, job, wait_s=CANCEL_TIMEOUT_S) + + try: + backend.stream_logs(job) + except BackendError: + return 1 + + status = backend.state(job) + rc = 0 if status.succeeded else 1 + + for result_file in sorted(logs_dir.glob(f"{request.result_filename}*.json")): + try: + artifacts.copy_to_workspace(result_file, request.workspace / result_file.name) + except artifacts.ArtifactError as error: + print(f"ERROR: {error}", file=sys.stderr) + rc = 1 + + if request.is_agentic and not request.eval_only: + _stage_agentic(logs_dir, request.workspace) + + if request.run_eval: + eval_dir = _find_eval_dir(logs_dir) or logs_dir / "eval_results" + try: + artifacts.copy_eval_artifacts(eval_dir, request.workspace) + except artifacts.ArtifactError as error: + print(f"ERROR: {error}", file=sys.stderr) + rc = 1 + + return rc diff --git a/inferencex-e2e/infx/launch/drivers/srt/models.py b/inferencex-e2e/infx/launch/drivers/srt/models.py index 39799fdb0f..28ac33998c 100644 --- a/inferencex-e2e/infx/launch/drivers/srt/models.py +++ b/inferencex-e2e/infx/launch/drivers/srt/models.py @@ -46,6 +46,11 @@ class Override: Override(Match(model_glob="*/DeepSeek-V4-Pro-0813"), entry="DeepSeek-V4-Pro-0813"), ), "gb200-nv": ( + Override( + Match(any_of("dsv4"), any_of("fp4"), any_of("llmd-vllm")), + entry="DeepSeek-V4-Pro@numa1", + served_name="deepseek-ai/DeepSeek-V4-Pro", + ), Override( Match(any_of("dsr1"), any_of("fp4"), any_of("dynamo-sglang")), entry="deepseek-r1-0528-fp4-v2", diff --git a/inferencex-e2e/infx/launch/policy.py b/inferencex-e2e/infx/launch/policy.py index d6bf0edba4..3a6ed913e8 100644 --- a/inferencex-e2e/infx/launch/policy.py +++ b/inferencex-e2e/infx/launch/policy.py @@ -57,6 +57,7 @@ class LaunchPath(StrEnum): SRT_NATIVE = "srt-native" SRT_BATCH = "srt-batch" SCRIPT = "script" + LLMD = "llmd" NATIVE_SRT_LANES: dict[str, tuple[Match, ...]] = { @@ -76,6 +77,8 @@ class LaunchPath(StrEnum): def launch_path(cluster_id: str, request: LaunchRequest) -> LaunchPath: if request.is_multinode: + if request.framework == "llmd-vllm": + return LaunchPath.LLMD if any(lane(request) for lane in NATIVE_SRT_LANES.get(cluster_id, ())): return LaunchPath.SRT_NATIVE return LaunchPath.SRT_MULTI diff --git a/inferencex-e2e/infx/launch/request.py b/inferencex-e2e/infx/launch/request.py index f293892211..77cc4b1012 100644 --- a/inferencex-e2e/infx/launch/request.py +++ b/inferencex-e2e/infx/launch/request.py @@ -158,3 +158,17 @@ class ScriptRequest(LaunchRequest): model: str = Field(alias="MODEL") gpu_count: int = Field(alias="GPU_COUNT") salloc_time_limit: int = Field(alias="SALLOC_TIME_LIMIT") + + +class LlmdRequest(SrtRequest): + """An llm-d vLLM multinode job submitted through benchmarks/multi_node/llm-d.""" + + model: str = Field(alias="MODEL") + disagg: TrueFlag = Field(alias="DISAGG") + prefill_nodes: int = Field(alias="PREFILL_NODES") + decode_nodes: int = Field(alias="DECODE_NODES") + prefill_num_workers: int = Field(1, alias="PREFILL_NUM_WORKERS") + decode_num_workers: int = Field(1, alias="DECODE_NUM_WORKERS") + isl: int = Field(alias="ISL") + osl: int = Field(alias="OSL") + random_range_ratio: str = Field(alias="RANDOM_RANGE_RATIO") diff --git a/inferencex-e2e/infx/tests/launch/test_llmd_driver.py b/inferencex-e2e/infx/tests/launch/test_llmd_driver.py new file mode 100644 index 0000000000..179fa083b4 --- /dev/null +++ b/inferencex-e2e/infx/tests/launch/test_llmd_driver.py @@ -0,0 +1,136 @@ +"""The llm-d driver: run submit.sh, attach, stage artifacts.""" + +import json +import subprocess +import tarfile +from pathlib import Path + +import pytest + +from infx.tests.launch.fake_slurm import ( + base_env, + install_fakes, + launch, + make_workspace, + runner_for, + sandbox_runner_config, +) + +LLMD_SUBMIT = """#!/usr/bin/env bash +set -e +env > "$GITHUB_WORKSPACE/submitted.env" +echo "$PWD $*" > "$GITHUB_WORKSPACE/submitted.args" +logs="$BENCHMARK_LOGS_DIR" +job="$logs/slurm_job-4299" +mkdir -p "$logs/agentic/conc_128" "$job/eval_results" +echo '{"conc": 128}' > "$logs/point-identity_conc128.json" +echo trace > "$logs/agentic/conc_128/profile.json" +echo '{"score": 1}' > "$job/eval_results/results_gsm8k.json" +echo 'server log' > "$logs/server.log" +echo 'benchmark done' > "$logs/slurm_job-4299.out" +echo 'worker warning' > "$logs/slurm_job-4299.err" +echo 'submitting' >&2 +[[ "${NO_JOB_ID:-}" == 1 ]] && exit 1 +echo 4299 +""" + + +@pytest.fixture +def harness(tmp_path): + """Sandboxed gb200-nv cluster, fake Slurm binaries, and a workspace with a fake submit.sh.""" + sandbox = tmp_path / "sandbox" + sandbox.mkdir() + config = sandbox_runner_config(sandbox) + workspace = make_workspace(tmp_path / "workspace") + wrapper = workspace / "benchmarks/multi_node/llm-d/submit.sh" + wrapper.parent.mkdir(parents=True, exist_ok=True) + wrapper.write_text(LLMD_SUBMIT) + wrapper.chmod(0o755) + + model_root = sandbox / "mnt/numa1/models/DeepSeek-V4-Pro" + model_root.mkdir(parents=True) + (model_root / "config.json").write_text("{}\n") + + logs = tmp_path / "logs" + env = base_env( + fakes=install_fakes(tmp_path / "bin"), logs=logs, workspace=workspace, sandbox=sandbox + ) + env.update( + RUNNER_NAME=runner_for("gb200-nv"), + IS_MULTINODE="true", + IS_AGENTIC="1", + RUN_EVAL="true", + EVAL_ONLY="false", + FRAMEWORK="llmd-vllm", + MODEL="deepseek-ai/DeepSeek-V4-Pro", + MODEL_PREFIX="dsv4", + PRECISION="fp4", + SPEC_DECODING="mtp", + THINKING_MODE="thinking_on", + PREFILL_NODES="2", + DECODE_NODES="2", + ISL="8192", + OSL="1024", + RANDOM_RANGE_RATIO="0.8", + CONC_LIST="256 512", + DISAGG="true", + IMAGE="vllm/vllm-openai:v0.21.0", + RESULT_FILENAME="point-identity", + ENROOT_IMPORT_TIME_LIMIT="10", + ) + return config, workspace, env + + +def run_launch(harness) -> subprocess.CompletedProcess[str]: + config, workspace, env = harness + return launch(env, config, workspace) + + +def test_llmd_driver_submits_the_wrapper_and_stages_artifacts(harness): + result = run_launch(harness) + config, workspace, env = harness + + assert result.returncode == 0, result.stdout + result.stderr + submitted = dict( + line.split("=", 1) for line in (workspace / "submitted.env").read_text().splitlines() if "=" in line + ) + model_path = f"{config.parent}/mnt/numa1/models/DeepSeek-V4-Pro" + assert submitted["MODEL_PATH"] == model_path + assert submitted["MODEL_NAME"] == env["MODEL"] + assert submitted["LLMD_CONTAINER_ENGINE"] == "pyxis" + assert submitted["LLMD_SQUASH_FILE"] + assert submitted["BENCHMARK_LOGS_DIR"] == f"{workspace}/benchmark_logs" + assert submitted["SLURM_PARTITION"] == "batch" + assert submitted["SLURM_ACCOUNT"] == "benchmark" + assert submitted["GPUS_PER_NODE"] == "4" + assert submitted["CONTAINER_IMAGE"] == env["IMAGE"] + assert (workspace / "submitted.args").read_text().split() == [ + str(workspace / "benchmarks/multi_node/llm-d"), + *"2 2 8192 1024 256x512 inf 0.8".split(), + ] + + assert json.loads((workspace / "point-identity_conc128.json").read_text()) == {"conc": 128} + assert (workspace / "LOGS/agentic/conc_128/profile.json").read_text() == "trace\n" + assert json.loads((workspace / "results_gsm8k.json").read_text()) == {"score": 1} + with tarfile.open(workspace / "multinode_server_logs.tar.gz") as bundle: + assert "./server.log" in bundle.getnames() + assert "submitting" in result.stderr + + +def test_llmd_driver_fails_when_the_wrapper_prints_no_job_id(harness): + _, workspace, env = harness + env["NO_JOB_ID"] = "1" + + result = launch(env, harness[0], workspace) + + assert result.returncode == 1 + assert "submit.sh failed before returning a Slurm job id" in result.stderr + + +def test_llmd_driver_leaves_node_local_checkpoints_to_the_job(harness): + config, workspace, env = harness + (config.parent / "mnt/numa1/models/DeepSeek-V4-Pro/config.json").unlink() + + result = launch(env, config, workspace) + + assert result.returncode == 0, result.stdout + result.stderr diff --git a/inferencex-e2e/infx/tests/launch/test_srt_policy.py b/inferencex-e2e/infx/tests/launch/test_srt_policy.py index 6da6766573..dc5a883fed 100644 --- a/inferencex-e2e/infx/tests/launch/test_srt_policy.py +++ b/inferencex-e2e/infx/tests/launch/test_srt_policy.py @@ -55,6 +55,7 @@ def cluster(tmp_path, single_node_models: str = "staged") -> Cluster: ("b200-nscale", dict(MULTI, MODEL_PREFIX="glm5.3", PRECISION="fp8", FRAMEWORK="tilert", SPEC_DECODING="mtp"), LaunchPath.SRT_MULTI), ("mi355x-amds", dict(IS_MULTINODE="true", FRAMEWORK="atom-disagg"), LaunchPath.SRT_MULTI), ("mi355x-amds", dict(MULTI, FRAMEWORK="sglang-disagg"), LaunchPath.SRT_MULTI), + ("gb200-nv", dict(MULTI, FRAMEWORK="llmd-vllm", MODEL_PREFIX="dsv4", PRECISION="fp4", SPEC_DECODING="mtp"), LaunchPath.LLMD), ("gb200-nv", dict(MULTI, FRAMEWORK="tilert"), LaunchPath.SRT_MULTI), ("b300-dsxe", dict(SINGLE, MODEL_PREFIX="dsv41flash", FRAMEWORK="sglang", IS_AGENTIC="1"), LaunchPath.SRT_BATCH), ("b300-dsxe", dict(SINGLE, MODEL_PREFIX="dsv41flash", FRAMEWORK="sglang", IS_AGENTIC="1", INFX_BATCH_REENTRY="1"), LaunchPath.SRT_SINGLE),