Skip to content
2 changes: 1 addition & 1 deletion inferencex-e2e/benchmarks/multi_node/llm-d/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
39 changes: 20 additions & 19 deletions inferencex-e2e/benchmarks/multi_node/llm-d/job.slurm
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Comment on lines +70 to +84

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🔴 Any startup failure in the decode coordinator (EPP/Envoy not ready, prefill health timeout) now makes the whole llm-d job report success instead of failure, with no benchmark results. server.sh:303 has trap 'touch "$BENCH_DONE_MARKER"' EXIT, which fires on every exit of that branch, including the exit 1 failure paths, writing the same marker file as the real success touch at server.sh:573. job.slurm:82 [[ -f "$BENCH_DONE_MARKER" ]] && return 0 treats marker-exists as unconditional success, ignoring rc. Fix: use a distinct marker/exit path for genuine completion so job.slurm only succeeds when the benchmark actually ran, not merely when the coordinator exited for any reason.

Why this was flagged

Trigger: the decode-leader branch in server.sh (ROLE==decode && LWS_WORKER_INDEX==0) hits any of its exit 1 checks - EPP not binding within 60s (server.sh:384/388), Envoy /ready failing within 120s (server.sh:407/412), or prefill /health timing out within 300s (server.sh:509) - reached via job.slurm's run_until_bench_done -> srun -> server.sh. The EXIT trap at server.sh:303 touches $BENCH_DONE_MARKER on that exit, same file as the genuine-success touch at server.sh:573. job.slurm's run_until_bench_done (job.slurm:70-85) only checks file existence at line 82, not the step's rc, so it returns 0. set -eo pipefail (job.slurm:10) then lets the script finish normally, Slurm records COMPLETED, and drivers/llmd.py:137-138 computes rc=0 from that state; the result-JSON glob (llmd.py:140) finds nothing but the loop simply doesn't run, so no error is raised. On the base branch this same failure ends CANCELLED, which the launcher already treats as failure.

Verification: server.sh:303's trap 'touch "$BENCH_DONE_MARKER"' EXIT fires on every exit of the decode-coordinator branch, including the startup exit 1 failures at server.sh:384, 388, 407, 412, and 509, writing the same marker as the genuine-success touch at server.sh:573. job.slurm's [[ -f "$BENCH_DONE_MARKER" ]] && return 0 then discards the failing rc whenever that marker exists.

}

# Container engine: 'docker' (default) for clusters where the SLURM
# user can talk to /var/run/docker.sock (e.g. h200-dgxc-slurm); 'pyxis'
Expand All @@ -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 \
Expand Down Expand Up @@ -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 \
Expand Down
3 changes: 1 addition & 2 deletions inferencex-e2e/benchmarks/multi_node/llm-d/server.sh
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
3 changes: 2 additions & 1 deletion inferencex-e2e/infx/launch/drivers/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand All @@ -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),
}


Expand Down
158 changes: 158 additions & 0 deletions inferencex-e2e/infx/launch/drivers/llmd.py
Original file line number Diff line number Diff line change
@@ -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
5 changes: 5 additions & 0 deletions inferencex-e2e/infx/launch/drivers/srt/models.py
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
3 changes: 3 additions & 0 deletions inferencex-e2e/infx/launch/policy.py
Original file line number Diff line number Diff line change
Expand Up @@ -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, ...]] = {
Expand All @@ -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
Expand Down
14 changes: 14 additions & 0 deletions inferencex-e2e/infx/launch/request.py
Original file line number Diff line number Diff line change
Expand Up @@ -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")
Loading
Loading