diff --git a/.github/workflows/run-sweep.yml b/.github/workflows/run-sweep.yml index b983982af8..fdb34e1e4d 100644 --- a/.github/workflows/run-sweep.yml +++ b/.github/workflows/run-sweep.yml @@ -556,7 +556,7 @@ jobs: precision: ${{ matrix.config.precision }} router: ${{ matrix.config.router && toJson(matrix.config.router) || '' }} kv-p2p-transfer: ${{ matrix.config['kv-p2p-transfer'] || '' }} - conc-list: '[${{ matrix.config.conc }}]' + conc-list: ${{ toJson(matrix.config.conc) }} spec-decoding: ${{ matrix.config.spec-decoding }} disagg: ${{ matrix.config.disagg }} prefill-hardware: ${{ matrix.config.prefill.hardware }} @@ -577,7 +577,7 @@ jobs: decode-ep: ${{ matrix.config.decode.ep }} decode-dp-attn: ${{ matrix.config.decode.dp-attn }} decode-additional-settings: ${{ toJson(matrix.config.decode.additional-settings) }} - conc: ${{ matrix.config.conc }} + conc: ${{ matrix.config.conc[0] }} kv-offloading: ${{ matrix.config.kv-offloading }} kv-offload-backend: ${{ matrix.config['kv-offload-backend'].name }} kv-offload-backend-metadata: ${{ matrix.config['kv-offload-backend'] && toJson(matrix.config['kv-offload-backend']) || '' }} diff --git a/benchmarks/benchmark_lib.sh b/benchmarks/benchmark_lib.sh index 2fc68c013a..cb1933254f 100644 --- a/benchmarks/benchmark_lib.sh +++ b/benchmarks/benchmark_lib.sh @@ -1346,6 +1346,18 @@ ensure_hf_cli() { } resolve_trace_source() { + if [[ -n "${AIPERF_TRACE_SOURCE_FLAG:-}" ]]; then + TRACE_SOURCE_FLAG="$AIPERF_TRACE_SOURCE_FLAG" + echo "Loading traces via AIPERF_TRACE_SOURCE_FLAG: $TRACE_SOURCE_FLAG" + return 0 + fi + + if [[ -n "${AIPERF_WEKA_TRACE_DIR:-}" ]]; then + TRACE_SOURCE_FLAG="--custom-dataset-type weka_trace --input-file $AIPERF_WEKA_TRACE_DIR" + echo "Loading traces via local Weka trace dir: $AIPERF_WEKA_TRACE_DIR" + return 0 + fi + # Per-recipe override: set WEKA_LOADER_OVERRIDE to one of the aiperf # public-dataset loader names allowed by the inferencex-agentx-mvp # scenario. Used by recipes whose servers have non-default context @@ -1527,8 +1539,9 @@ build_replay_cmd() { # Default --num-dataset-entries is 100; the with-subagents Weka corpus # has 393. Cap at 393 so all unique traces are loaded (the loader treats # this as a ``min(cap, available)`` ceiling, not a target — see - # semianalysis_cc_traces_weka.py). - REPLAY_CMD+=" --num-dataset-entries 393" + # semianalysis_cc_traces_weka.py). Smoke tests can override this for a + # small local fixture without changing production recipes. + REPLAY_CMD+=" --num-dataset-entries ${AIPERF_NUM_DATASET_ENTRIES:-393}" # 1-second timeslices on the server-metrics scrape so the post-run # plotter has per-window time series (KV usage, cache hit rate, # throughput, etc.). Matches kv-cache-tester's poll_interval=1.0 diff --git a/benchmarks/multi_node/amd_utils/job.slurm b/benchmarks/multi_node/amd_utils/job.slurm index 4638e5cadb..495d26308a 100755 --- a/benchmarks/multi_node/amd_utils/job.slurm +++ b/benchmarks/multi_node/amd_utils/job.slurm @@ -83,6 +83,9 @@ DECODE_MTP_SIZE=${DECODE_MTP_SIZE:-0} ROUTER_TYPE="${ROUTER_TYPE:-vllm-router}" ROUTER_PORT="${ROUTER_PORT:-30000}" PROXY_PING_PORT="${PROXY_PING_PORT:-36367}" +VLLM_ROUTER_POLICY="${VLLM_ROUTER_POLICY:-consistent_hash}" +VLLM_ROUTER_PREFILL_POLICY="${VLLM_ROUTER_PREFILL_POLICY:-consistent_hash}" +VLLM_ROUTER_DECODE_POLICY="${VLLM_ROUTER_DECODE_POLICY:-round_robin}" # ============================================================================= # Model Path Resolution @@ -113,6 +116,20 @@ if [[ "$ENGINE" == "vllm-disagg" ]]; then resolve_hf_cache_path() { local base_path=$1 if [[ -d "${base_path}/snapshots" ]]; then + if [[ -f "${base_path}/refs/main" ]]; then + local ref_snapshot + ref_snapshot="$(cat "${base_path}/refs/main" 2>/dev/null | tr -d '[:space:]')" + if [[ -n "$ref_snapshot" && -d "${base_path}/snapshots/${ref_snapshot}" ]]; then + echo "${base_path}/snapshots/${ref_snapshot}" + return 0 + fi + fi + local processor_snapshot + processor_snapshot=$(find "${base_path}/snapshots" -mindepth 2 -maxdepth 2 -name preprocessor_config.json -print -quit 2>/dev/null | xargs -r dirname) + if [[ -n "$processor_snapshot" ]]; then + echo "$processor_snapshot" + return 0 + fi local snapshot=$(ls -1 "${base_path}/snapshots" 2>/dev/null | head -1) if [[ -n "$snapshot" ]]; then echo "${base_path}/snapshots/${snapshot}" @@ -124,20 +141,26 @@ if [[ "$ENGINE" == "vllm-disagg" ]]; then } MODEL_PATH="" - SEARCH_PATHS=( - "${MODEL_DIR}/${DISK_DIR_NAME}" - "${MODEL_DIR}/${MODEL_NAME}" - "/nfsdata/hf_hub_cache-0/${DISK_DIR_NAME}" - "/nfsdata/hf_hub_cache-0/${MODEL_NAME}" - ) - - for search_path in "${SEARCH_PATHS[@]}"; do - if [[ -d "$search_path" ]]; then - RESOLVED=$(resolve_hf_cache_path "$search_path") - MODEL_PATH="$RESOLVED" - echo "Found MODEL_PATH: $MODEL_PATH" - break - fi + MODEL_MOUNT_DIR="" + SEARCH_ROOTS=("$MODEL_DIR") + if [[ "$(hostname -s)" == mia1* ]]; then + SEARCH_ROOTS+=("/it-share/models" "/it-share/hf_cache" "/it-share/data") + fi + SEARCH_ROOTS+=("/nfsdata/hf_hub_cache-0" "/nfsdata/hf_hub_cache" "/nfsdata") + + SEARCH_PATHS=() + for search_root in "${SEARCH_ROOTS[@]}"; do + SEARCH_PATHS+=("${search_root}/${DISK_DIR_NAME}" "${search_root}/${MODEL_NAME}") + for search_path in "${search_root}/${DISK_DIR_NAME}" "${search_root}/${MODEL_NAME}"; do + if [[ -d "$search_path" ]]; then + RESOLVED=$(resolve_hf_cache_path "$search_path") + MODEL_PATH="$RESOLVED" + MODEL_MOUNT_DIR="$search_root" + echo "Found MODEL_PATH: $MODEL_PATH" + echo "Found MODEL_MOUNT_DIR: $MODEL_MOUNT_DIR" + break 2 + fi + done done if [[ -z "$MODEL_PATH" ]]; then @@ -247,11 +270,13 @@ echo "NNODES: ${NNODES}" echo "REPO DIR: ${DI_REPO_DIR}" echo "USER: ${USER_NAME}" + # Reduce log spam export TQDM_MININTERVAL=20 # Translate the host-resolved MODEL_PATH to the Docker mount namespace -DOCKER_MODEL_PATH="${MODEL_PATH/#$MODEL_DIR//models}" +MODEL_MOUNT_DIR="${MODEL_MOUNT_DIR:-$MODEL_DIR}" +DOCKER_MODEL_PATH="${MODEL_PATH/#$MODEL_MOUNT_DIR//models}" export DI_REPO_DIR=$DI_REPO_DIR export WS_PATH=$WS_PATH @@ -259,6 +284,7 @@ export NNODES=$NNODES export NODE0_ADDR=$NODE0_ADDR export MODEL_PATH=$MODEL_PATH export MODEL_DIR=$MODEL_DIR +export MODEL_MOUNT_DIR=$MODEL_MOUNT_DIR export xP=$xP export yD=$yD export MODEL_NAME=$MODEL_NAME @@ -288,13 +314,39 @@ export RESULT_FILENAME="${RESULT_FILENAME:-}" export SPEC_DECODING="${SPEC_DECODING:-}" export IS_MULTINODE="${IS_MULTINODE:-false}" +# Agentic / custom vLLM-disagg connector knobs (threaded from submit.sh) +export IS_AGENTIC="${IS_AGENTIC:-0}" +export DURATION="${DURATION:-1800}" +export MODEL="${MODEL:-}" +export ROUTER_TYPE="${ROUTER_TYPE:-vllm-router}" +export ROUTER_PORT="${ROUTER_PORT:-30000}" +export VLLM_ROUTER_POLICY="${VLLM_ROUTER_POLICY:-consistent_hash}" +export VLLM_ROUTER_PREFILL_POLICY="${VLLM_ROUTER_PREFILL_POLICY:-consistent_hash}" +export VLLM_ROUTER_DECODE_POLICY="${VLLM_ROUTER_DECODE_POLICY:-round_robin}" +export ENABLE_PREFIX_CACHING="${ENABLE_PREFIX_CACHING:-}" +export MAX_MODEL_LEN="${MAX_MODEL_LEN:-}" +export WEKA_LOADER_OVERRIDE="${WEKA_LOADER_OVERRIDE:-}" +export AIPERF_TRACE_SOURCE_FLAG="${AIPERF_TRACE_SOURCE_FLAG:-}" +export AIPERF_WEKA_TRACE_DIR="${AIPERF_WEKA_TRACE_DIR:-}" +export AIPERF_NUM_DATASET_ENTRIES="${AIPERF_NUM_DATASET_ENTRIES:-}" +export VLLM_BIND_IP="${VLLM_BIND_IP:-}" +export SERVED_MODEL_NAME="${SERVED_MODEL_NAME:-}" +export PREFILL_KV_CONNECTOR="${PREFILL_KV_CONNECTOR:-moriio}" +export DECODE_KV_CONNECTOR="${DECODE_KV_CONNECTOR:-moriio}" +export MORI_RDMA_DEVICES="${MORI_RDMA_DEVICES:-}" +export MORI_IO_TC="${MORI_IO_TC:-}" +export MORI_RDMA_SL="${MORI_RDMA_SL:-}" +export MORI_IO_SL="${MORI_IO_SL:-}" +export UCX_IB_GID_INDEX="${UCX_IB_GID_INDEX:-}" +export DOCKER_PULL_POLICY="${DOCKER_PULL_POLICY:-missing}" + SANITIZED_USER=$(echo "$USER_NAME" | tr -c 'a-zA-Z0-9_.-' '_') export DOCKER_CONT_NAME="container_${ENGINE}_${SANITIZED_USER}_${MODEL_NAME}_${SLURM_JOB_ID}" # vLLM external router container. # NOTE: vllm/vllm-router only retains ~16 recent nightlies on Docker Hub; older # dated tags are garbage-collected (manifest unknown) -VLLM_ROUTER_IMAGE="${VLLM_ROUTER_IMAGE:-vllm/vllm-router:nightly-20260629-e667ebb}" +VLLM_ROUTER_IMAGE="${VLLM_ROUTER_IMAGE:-vllm/vllm-router:nightly}" ROUTER_CONT_NAME="router_vllm_${SANITIZED_USER}_${SLURM_JOB_ID}" export RUN_FILE_FULL="$WS_PATH/${RUN_FILE}" @@ -371,6 +423,37 @@ DOCKER_ENV_COMMON=( -e DECODE_ENABLE_DP=\$DECODE_ENABLE_DP -e DECODE_MTP_SIZE=\$DECODE_MTP_SIZE -e IS_MULTINODE=\$IS_MULTINODE + -e MODEL=\$MODEL + -e ROUTER_TYPE=\$ROUTER_TYPE + -e ROUTER_PORT=\$ROUTER_PORT + -e VLLM_ROUTER_POLICY=\$VLLM_ROUTER_POLICY + -e VLLM_ROUTER_PREFILL_POLICY=\$VLLM_ROUTER_PREFILL_POLICY + -e VLLM_ROUTER_DECODE_POLICY=\$VLLM_ROUTER_DECODE_POLICY + -e ENABLE_PREFIX_CACHING=\$ENABLE_PREFIX_CACHING + -e MAX_MODEL_LEN=\$MAX_MODEL_LEN + -e MAX_NUM_SEQS=\$MAX_NUM_SEQS + -e WEKA_LOADER_OVERRIDE=\$WEKA_LOADER_OVERRIDE + -e AIPERF_TRACE_SOURCE_FLAG=\$AIPERF_TRACE_SOURCE_FLAG + -e AIPERF_WEKA_TRACE_DIR=\$AIPERF_WEKA_TRACE_DIR + -e AIPERF_NUM_DATASET_ENTRIES=\$AIPERF_NUM_DATASET_ENTRIES + -e VLLM_BIND_IP=\$VLLM_BIND_IP + -e SERVED_MODEL_NAME=\$SERVED_MODEL_NAME + -e PREFILL_KV_CONNECTOR=\$PREFILL_KV_CONNECTOR + -e DECODE_KV_CONNECTOR=\$DECODE_KV_CONNECTOR + -e MORI_RDMA_DEVICES=\$MORI_RDMA_DEVICES + -e MORI_IO_TC=\$MORI_IO_TC + -e MORI_RDMA_SL=\$MORI_RDMA_SL + -e MORI_IO_SL=\$MORI_IO_SL + -e UCX_IB_GID_INDEX=\$UCX_IB_GID_INDEX + -e LMCACHE_HOST=\${LMCACHE_HOST:-} + -e LMCACHE_PORT=\${LMCACHE_PORT:-} + -e LMCACHE_HTTP_PORT=\${LMCACHE_HTTP_PORT:-} + -e LMCACHE_L1_SIZE_GB=\${LMCACHE_L1_SIZE_GB:-} + -e LMCACHE_L1_INIT_SIZE_GB=\${LMCACHE_L1_INIT_SIZE_GB:-} + -e LMCACHE_L1_READ_TTL_SECONDS=\${LMCACHE_L1_READ_TTL_SECONDS:-} + -e LMCACHE_CHUNK_SIZE=\${LMCACHE_CHUNK_SIZE:-} + -e LMCACHE_MAX_WORKERS=\${LMCACHE_MAX_WORKERS:-} + -e LMCACHE_MP_MQ_TIMEOUT=\${LMCACHE_MP_MQ_TIMEOUT:-} ) # Engine-specific env vars @@ -387,6 +470,13 @@ if [[ "$ENGINE" == "vllm-disagg" ]]; then -e UCX_LOG_LEVEL=warn -e HSA_ENABLE_SDMA=1 -e PROXY_STREAM_IDLE_TIMEOUT=\${PROXY_STREAM_IDLE_TIMEOUT:-300} + -e VLLM_MORIIO_CONNECTOR_READ_MODE=\${VLLM_MORIIO_CONNECTOR_READ_MODE:-1} + -e IBDEVICES=\${IBDEVICES:-} + -e MORI_RDMA_TC=\${MORI_RDMA_TC:-} + -e MORI_IO_TC=\${MORI_IO_TC:-} + -e MORI_RDMA_SL=\${MORI_RDMA_SL:-} + -e MORI_IO_SL=\${MORI_IO_SL:-} + -e NCCL_IB_DISABLE=\${NCCL_IB_DISABLE:-} -e PYTHONPYCACHEPREFIX=/tmp/pycache ) elif [[ "$ENGINE" == "atom-disagg" ]]; then @@ -428,6 +518,36 @@ echo \"Rank \$SLURM_PROCID on \$(hostname)\" eval \"\$DOCKER_CMD_DETECT\" echo \"[docker-detect] rank \$SLURM_PROCID: DOCKER_CMD=\$DOCKER_CMD\" +pull_docker_image_if_needed() { + local image=\"\$1\" + case \"${DOCKER_PULL_POLICY}\" in + always) + echo \"[docker-pull] pulling \$image on \$(hostname) (policy=always)\" + \$DOCKER_CMD pull \"\$image\" + ;; + missing|if-not-present) + if ! \$DOCKER_CMD image inspect \"\$image\" >/dev/null 2>&1; then + echo \"[docker-pull] pulling \$image on \$(hostname) (policy=missing)\" + \$DOCKER_CMD pull \"\$image\" + else + echo \"[docker-pull] using existing \$image on \$(hostname) (policy=missing)\" + fi + ;; + never|false|0) + echo \"[docker-pull] skipping pull for \$image on \$(hostname) (policy=${DOCKER_PULL_POLICY})\" + ;; + *) + echo \"[docker-pull] invalid DOCKER_PULL_POLICY=${DOCKER_PULL_POLICY}; expected always, missing, or never\" >&2 + exit 1 + ;; + esac +} + +pull_docker_image_if_needed \"$DOCKER_IMAGE_NAME\" +if [[ \"$ENGINE\" == \"vllm-disagg\" && \"$ROUTER_TYPE\" == \"vllm-router\" && \"\$SLURM_PROCID\" == \"0\" ]]; then + pull_docker_image_if_needed \"$VLLM_ROUTER_IMAGE\" +fi + # Enable out-of-tree RDMA library mounts for atom-disagg (mooncake requires host RDMA stack) RDMA_MOUNTS=() if [[ "$ENGINE" == "atom-disagg" ]]; then @@ -507,27 +627,47 @@ fi fi # end: if ENGINE == atom-disagg # Pre-clean (idempotent) -\$DOCKER_CMD ps -aq --filter \"$CONT_FILTER\" | xargs -r \$DOCKER_CMD rm -f || true -\$DOCKER_CMD ps -aq | xargs -r \$DOCKER_CMD stop || true +mkdir -p \"${BENCHMARK_LOGS_DIR}/runtime_logs\" +\$DOCKER_CMD rm -f \"$DOCKER_CONT_NAME\" 2>/dev/null || true +if [[ \"$ENGINE\" == \"vllm-disagg\" && \"$ROUTER_TYPE\" == \"vllm-router\" && \"\$SLURM_PROCID\" == \"0\" ]]; then + \$DOCKER_CMD rm -f \"$ROUTER_CONT_NAME\" 2>/dev/null || true +fi # Start vLLM external router container on node 0 if [[ \"$ENGINE\" == \"vllm-disagg\" && \"$ROUTER_TYPE\" == \"vllm-router\" && \"\$SLURM_PROCID\" == \"0\" ]]; then \$DOCKER_CMD rm -f \"$ROUTER_CONT_NAME\" 2>/dev/null || true + ROUTER_KV_CONNECTOR=\"\${VLLM_ROUTER_KV_CONNECTOR:-moriio}\" + if [[ \"\${PREFILL_KV_CONNECTOR:-}\" == \"lmcache-nixl\" ]]; then + ROUTER_KV_CONNECTOR=\"\${VLLM_ROUTER_KV_CONNECTOR:-nixl}\" + fi + ROUTER_ENDPOINT_ARGS=\"--vllm-discovery-address 0.0.0.0:${PROXY_PING_PORT}\" + if [[ \"\${PREFILL_KV_CONNECTOR:-}\" == \"lmcache-nixl\" ]]; then + ROUTER_ENDPOINT_ARGS=\"\" + IFS=',' read -r -a _router_ips <<< \"${IPADDRS}\" + for ((i=0; i<${xP} && i<\${#_router_ips[@]}; i++)); do + ROUTER_ENDPOINT_ARGS+=\" --prefill http://\${_router_ips[\$i]}:${SERVER_PORT:-2584}\" + done + for ((i=${xP}; i<\${#_router_ips[@]}; i++)); do + ROUTER_ENDPOINT_ARGS+=\" --decode http://\${_router_ips[\$i]}:${SERVER_PORT:-2584}\" + done + fi + echo \"[ROUTER_CMD] image=$VLLM_ROUTER_IMAGE vllm-router --vllm-pd-disaggregation --kv-connector \${ROUTER_KV_CONNECTOR} \${ROUTER_ENDPOINT_ARGS} --worker-startup-timeout-secs \${VLLM_ROUTER_WORKER_STARTUP_TIMEOUT_SECS:-1800} --port ${ROUTER_PORT} --host 0.0.0.0 --policy ${VLLM_ROUTER_POLICY} --prefill-policy ${VLLM_ROUTER_PREFILL_POLICY} --decode-policy ${VLLM_ROUTER_DECODE_POLICY} --log-level info\" \$DOCKER_CMD run -d \ --name \"$ROUTER_CONT_NAME\" \ --network host \ --ulimit nofile=1048576:1048576 \ - -v /tmp:/run_logs \ + -v ${BENCHMARK_LOGS_DIR}/runtime_logs:/run_logs \ \"$VLLM_ROUTER_IMAGE\" \ bash -lc \"mkdir -p /run_logs/slurm_job-${SLURM_JOB_ID} && exec vllm-router \ --vllm-pd-disaggregation \ - --kv-connector moriio \ - --vllm-discovery-address 0.0.0.0:${PROXY_PING_PORT} \ + --kv-connector \${ROUTER_KV_CONNECTOR} \ + \${ROUTER_ENDPOINT_ARGS} \ + --worker-startup-timeout-secs \${VLLM_ROUTER_WORKER_STARTUP_TIMEOUT_SECS:-1800} \ --port ${ROUTER_PORT} \ --host 0.0.0.0 \ - --policy consistent_hash \ - --prefill-policy consistent_hash \ - --decode-policy consistent_hash \ + --policy ${VLLM_ROUTER_POLICY} \ + --prefill-policy ${VLLM_ROUTER_PREFILL_POLICY} \ + --decode-policy ${VLLM_ROUTER_DECODE_POLICY} \ --log-level info 2>&1 | tee /run_logs/slurm_job-${SLURM_JOB_ID}/vllm_router_\$(hostname).log \" fi @@ -566,10 +706,10 @@ fi --privileged \ -v /sys:/sys \ $(command -v nicctl >/dev/null 2>&1 && echo "-v $(which nicctl):/usr/sbin/nicctl") \ - -v ${MODEL_DIR}:/models \ + -v ${MODEL_MOUNT_DIR}:/models \ -v \$HOME/.ssh:/root/.ssh \ --shm-size 128G \ - -v /tmp:/run_logs \ + -v ${BENCHMARK_LOGS_DIR}/runtime_logs:/run_logs \ -v ${BENCHMARK_LOGS_DIR}:/benchmark_logs \ -v ${DI_REPO_DIR}:${DOCKER_MOUNT_PATH} \ ${EXTRA_DOCKER_MOUNTS:-} \ diff --git a/benchmarks/multi_node/amd_utils/models_vllm.yaml b/benchmarks/multi_node/amd_utils/models_vllm.yaml index 9c046b4cf8..9b099f3615 100644 --- a/benchmarks/multi_node/amd_utils/models_vllm.yaml +++ b/benchmarks/multi_node/amd_utils/models_vllm.yaml @@ -25,22 +25,58 @@ amd-Llama-3.3-70B-Instruct-FP8-KV: env: "VLLM_USE_V1=1 VLLM_V1_USE_PREFILL_DECODE_ATTENTION=1 AMDGCN_USE_BUFFER_OPS=1 VLLM_ROCM_USE_AITER=1 VLLM_ROCM_USE_AITER_RMSNORM=1 VLLM_USE_AITER_TRITON_ROPE=1 TRITON_HIP_ASYNC_COPY_BYPASS_PERMUTE=1 TRITON_HIP_USE_ASYNC_COPY=1 TRITON_HIP_USE_BLOCK_PINGPONG=1 TRITON_HIP_ASYNC_FAST_SWIZZLE=1" Kimi-K2.5-MXFP4: - prefill_flags: "--tensor-parallel-size 8 --compilation-config '{\"cudagraph_mode\":\"PIECEWISE\"}' --no-enable-prefix-caching --block-size 1 --gpu-memory-utilization 0.90 --mm-encoder-tp-mode data" - decode_flags: "--tensor-parallel-size 8 --enable-expert-parallel --all2all-backend mori_low_latency --compilation-config '{\"cudagraph_mode\":\"PIECEWISE\"}' --no-enable-prefix-caching --block-size 1 --gpu-memory-utilization 0.90 --mm-encoder-tp-mode data" + prefill_flags: "--tensor-parallel-size 8 --compilation-config '{\"cudagraph_mode\":\"PIECEWISE\"}' --max-model-len 16384 --no-enable-prefix-caching --block-size 1 --gpu-memory-utilization 0.80 --mm-encoder-tp-mode data" + decode_flags: "--tensor-parallel-size 8 --enable-expert-parallel --all2all-backend mori_low_latency --compilation-config '{\"cudagraph_mode\":\"PIECEWISE\"}' --max-model-len 16384 --no-enable-prefix-caching --block-size 1 --gpu-memory-utilization 0.80 --mm-encoder-tp-mode data" env: "VLLM_USE_V1=1 VLLM_ROCM_USE_AITER=1 VLLM_ROCM_USE_AITER_PAGED_ATTN=0 VLLM_ROCM_USE_AITER_RMSNORM=1 VLLM_USE_AITER_TRITON_SILU_MUL=0 VLLM_ENGINE_READY_TIMEOUT_S=3600" hf_dir: "models--amd--Kimi-K2.5-MXFP4" +Kimi-K2.5-MXFP4-MoRI-LMCache-Agentic: + prefill_flags: "--tensor-parallel-size 8 --compilation-config '{\"cudagraph_mode\":\"PIECEWISE\"}' --max-model-len 16384 --no-enable-prefix-caching --kv-cache-dtype fp8 --block-size 1 --gpu-memory-utilization 0.80 --mm-encoder-tp-mode data" + decode_flags: "--tensor-parallel-size 8 --compilation-config '{\"cudagraph_mode\":\"PIECEWISE\"}' --max-model-len 16384 --no-enable-prefix-caching --kv-cache-dtype fp8 --block-size 1 --gpu-memory-utilization 0.80 --mm-encoder-tp-mode data" + env: "VLLM_USE_V1=1 VLLM_ROCM_USE_AITER=1 VLLM_ROCM_USE_AITER_PAGED_ATTN=0 VLLM_ROCM_USE_AITER_RMSNORM=1 VLLM_USE_AITER_TRITON_SILU_MUL=0 VLLM_ROCM_USE_AITER_FUSION_SHARED_EXPERTS=0 VLLM_ROCM_QUICK_REDUCE_QUANTIZATION=INT4 VLLM_ENGINE_READY_TIMEOUT_S=3600 VLLM_EXECUTE_MODEL_TIMEOUT_SECONDS=1200" + hf_dir: "Kimi-K2.5-MXFP4" + +Kimi-K2.7-Code-MXFP4: + prefill_flags: "--tensor-parallel-size 8 --compilation-config '{\"cudagraph_mode\":\"PIECEWISE\"}' --max-model-len 16384 --no-enable-prefix-caching --kv-cache-dtype fp8 --block-size 1 --gpu-memory-utilization 0.80 --mm-encoder-tp-mode data" + decode_flags: "--tensor-parallel-size 8 --compilation-config '{\"cudagraph_mode\":\"PIECEWISE\"}' --max-model-len 16384 --no-enable-prefix-caching --kv-cache-dtype fp8 --block-size 1 --gpu-memory-utilization 0.80 --mm-encoder-tp-mode data" + env: "VLLM_USE_V1=1 VLLM_ROCM_USE_AITER=1 VLLM_ROCM_USE_AITER_PAGED_ATTN=0 VLLM_ROCM_USE_AITER_RMSNORM=1 VLLM_USE_AITER_TRITON_SILU_MUL=0 VLLM_ROCM_USE_AITER_FUSION_SHARED_EXPERTS=0 VLLM_ROCM_QUICK_REDUCE_QUANTIZATION=INT4 VLLM_ENGINE_READY_TIMEOUT_S=3600 VLLM_EXECUTE_MODEL_TIMEOUT_SECONDS=1200" + hf_dir: "models--amd--Kimi-K2.7-Code-MXFP4" + +# DEP (DP-attention + Expert-Parallel) variant of Kimi-K2.7-Code-MXFP4. +# - tensor-parallel-size is a placeholder; server_vllm.sh sed-rewrites it to +# PREFILL_TP_SIZE / DECODE_TP_SIZE (DEP target = TP1 per node). +# - --enable-dp-attention / --enable-expert-parallel are appended by +# server_vllm.sh from the submit booleans (PREFILL/DECODE_ENABLE_DP/EP), so +# they are NOT hard-coded here. +# - AITER MLA is DISABLED: K2.7-Code has 64 KV heads, which the AITER MLA +# kernel (16/128 only) cannot handle on ROCm. +# - Default KV dtype = BF16 (no --kv-cache-dtype): the FP8-KV path crashes in +# AITER a8w8 MLA decode (roadmap/k27_dp_ep_fp8kv_debug). To try the FP8-KV + +# BF16-q workaround, set decode_env FP8KV_REAL_Q_DTYPE=bf16 and append +# --kv-cache-dtype fp8 (see decode_env below / config additional-settings). +Kimi-K2.7-Code-MXFP4-DEP: + prefill_flags: "--tensor-parallel-size 8 --compilation-config '{\"cudagraph_mode\":\"PIECEWISE\"}' --max-model-len 16384 --no-enable-prefix-caching --block-size 16 --gpu-memory-utilization 0.80 --mm-encoder-tp-mode data --trust-remote-code --tool-call-parser kimi_k2 --reasoning-parser kimi_k2 --dcp-comm-backend a2a --disable-custom-all-reduce" + decode_flags: "--tensor-parallel-size 8 --compilation-config '{\"cudagraph_mode\":\"PIECEWISE\"}' --max-model-len 16384 --no-enable-prefix-caching --block-size 16 --gpu-memory-utilization 0.80 --mm-encoder-tp-mode data --trust-remote-code --tool-call-parser kimi_k2 --reasoning-parser kimi_k2 --dcp-comm-backend a2a --disable-custom-all-reduce" + env: "VLLM_USE_V1=1 VLLM_ROCM_USE_AITER=1 VLLM_ROCM_USE_AITER_MLA=0 VLLM_ROCM_USE_AITER_FUSION_SHARED_EXPERTS=0 VLLM_ROCM_USE_AITER_FP4BMM=0 VLLM_ROCM_USE_AITER_PAGED_ATTN=0 VLLM_ROCM_USE_AITER_RMSNORM=1 VLLM_USE_AITER_TRITON_SILU_MUL=0 VLLM_ROCM_QUICK_REDUCE_QUANTIZATION=INT4 VLLM_ENGINE_READY_TIMEOUT_S=3600 VLLM_EXECUTE_MODEL_TIMEOUT_SECONDS=1200" + hf_dir: "models--amd--Kimi-K2.7-Code-MXFP4" + +facebook-opt-125m: + prefill_flags: "--tensor-parallel-size 1 --max-model-len 1024 --max-num-seqs 4 --gpu-memory-utilization 0.25 --chat-template /workspace/benchmarks/multi_node/amd_utils/simple_chat_template.jinja" + decode_flags: "--tensor-parallel-size 1 --max-model-len 1024 --max-num-seqs 4 --gpu-memory-utilization 0.25 --chat-template /workspace/benchmarks/multi_node/amd_utils/simple_chat_template.jinja" + env: "VLLM_USE_V1=1 VLLM_ENGINE_READY_TIMEOUT_S=600" + hf_dir: "facebook-opt-125m" + MiniMax-M2.5: # AITER fused-MoE kernel fmoe_bf16_blockscaleFp8_g1u1_vs_silu_32x384 for gfx950 writes OOB when run with MiniMax's shapes at M=8K(=num batched tokens), crashing vllm during AITER warmup. # Set token budget to 4k to avoid using that shape, instead of disabling AITER_MOE. - prefill_flags: "--max-num-batched-tokens 4K --tensor-parallel-size 8 --enable-expert-parallel --all2all-backend mori_low_latency --no-enable-prefix-caching --gpu-memory-utilization 0.95 --block-size 32" - decode_flags: "--max-num-batched-tokens 4K --tensor-parallel-size 8 --enable-expert-parallel --all2all-backend mori_low_latency --no-enable-prefix-caching --gpu-memory-utilization 0.95 --block-size 32" + prefill_flags: "--max-num-batched-tokens 4K --tensor-parallel-size 8 --enable-expert-parallel --all2all-backend mori_low_latency --max-model-len 16384 --no-enable-prefix-caching --gpu-memory-utilization 0.95 --block-size 32" + decode_flags: "--max-num-batched-tokens 4K --tensor-parallel-size 8 --enable-expert-parallel --all2all-backend mori_low_latency --max-model-len 16384 --no-enable-prefix-caching --gpu-memory-utilization 0.95 --block-size 32" env: "VLLM_USE_V1=1 VLLM_ROCM_USE_AITER=1 VLLM_ROCM_QUICK_REDUCE_QUANTIZATION=INT4 VLLM_ENGINE_READY_TIMEOUT_S=3600 VLLM_ROCM_SHUFFLE_KV_CACHE_LAYOUT=1" hf_dir: "models--MiniMaxAI--MiniMax-M2.5" MiniMax-M3-MXFP4: - prefill_flags: "--tensor-parallel-size 8 --max-num-batched-tokens 32768 --max-num-seqs 512 --block-size 128 --language-model-only --attention-backend TRITON_ATTN --moe-backend aiter --no-enable-prefix-caching --gpu-memory-utilization 0.90 --tool-call-parser minimax_m3 --reasoning-parser minimax_m3 --enable-auto-tool-choice" - decode_flags: "--tensor-parallel-size 8 --max-num-batched-tokens 32768 --max-num-seqs 512 --block-size 128 --language-model-only --attention-backend TRITON_ATTN --moe-backend aiter --no-enable-prefix-caching --gpu-memory-utilization 0.90 --tool-call-parser minimax_m3 --reasoning-parser minimax_m3 --enable-auto-tool-choice" + prefill_flags: "--tensor-parallel-size 8 --max-num-batched-tokens 32768 --max-num-seqs 512 --block-size 128 --language-model-only --attention-backend TRITON_ATTN --moe-backend aiter --max-model-len 16384 --no-enable-prefix-caching --gpu-memory-utilization 0.80 --tool-call-parser minimax_m3 --reasoning-parser minimax_m3 --enable-auto-tool-choice" + decode_flags: "--tensor-parallel-size 8 --max-num-batched-tokens 32768 --max-num-seqs 512 --block-size 128 --language-model-only --attention-backend TRITON_ATTN --moe-backend aiter --max-model-len 16384 --no-enable-prefix-caching --gpu-memory-utilization 0.80 --tool-call-parser minimax_m3 --reasoning-parser minimax_m3 --enable-auto-tool-choice" env: "VLLM_USE_V1=1 VLLM_ROCM_USE_AITER=1 VLLM_ROCM_USE_AITER_MOE=1 VLLM_USE_BREAKABLE_CUDAGRAPH=0 VLLM_ENGINE_READY_TIMEOUT_S=3600" hf_dir: "models--amd--MiniMax-M3-MXFP4" diff --git a/benchmarks/multi_node/amd_utils/server_vllm.sh b/benchmarks/multi_node/amd_utils/server_vllm.sh index f19ce8560b..9bcb54809f 100755 --- a/benchmarks/multi_node/amd_utils/server_vllm.sh +++ b/benchmarks/multi_node/amd_utils/server_vllm.sh @@ -55,6 +55,9 @@ MODEL_PATH="${MODEL_PATH:-${MODEL_DIR}/${MODEL_NAME}}" source $WS_PATH/env.sh host_ip=$(ip route get 1.1.1.1 2>/dev/null | awk '/src/ {print $7}') +IFS=',' read -ra NODE_IP_ARRAY <<< "$IPADDRS" +node_ip="${NODE_IP_ARRAY[$NODE_RANK]:-${NODE0_ADDR}}" +host_ip="${host_ip:-$node_ip}" # RDMA IP for Nixl KV transfer (prefer 192.168.x.x subnet if available) rdma_ip=$(hostname -I | tr ' ' '\n' | grep '^192\.168\.' | head -1) rdma_ip="${rdma_ip:-$host_ip}" @@ -171,37 +174,76 @@ fi if [[ "${PREFILL_ENABLE_EP:-false}" == "true" ]] && ! echo "$PREFILL_SERVER_CONFIG" | grep -q -- '--enable-expert-parallel'; then PREFILL_SERVER_CONFIG+=" --enable-expert-parallel" fi -if [[ "${PREFILL_ENABLE_DP:-false}" == "true" ]] && ! echo "$PREFILL_SERVER_CONFIG" | grep -q -- '--enable-dp-attention'; then - PREFILL_SERVER_CONFIG+=" --enable-dp-attention" +if [[ "${PREFILL_ENABLE_DP:-false}" == "true" ]] && ! echo "$PREFILL_SERVER_CONFIG" | grep -q -- '--decode-context-parallel-size'; then + PREFILL_SERVER_CONFIG+=" --decode-context-parallel-size ${PREFILL_TP_SIZE:-1}" fi if [[ "${DECODE_ENABLE_EP:-false}" == "true" ]] && ! echo "$DECODE_SERVER_CONFIG" | grep -q -- '--enable-expert-parallel'; then DECODE_SERVER_CONFIG+=" --enable-expert-parallel" fi -if [[ "${DECODE_ENABLE_DP:-false}" == "true" ]] && ! echo "$DECODE_SERVER_CONFIG" | grep -q -- '--enable-dp-attention'; then - DECODE_SERVER_CONFIG+=" --enable-dp-attention" +if [[ "${DECODE_ENABLE_DP:-false}" == "true" ]] && ! echo "$DECODE_SERVER_CONFIG" | grep -q -- '--decode-context-parallel-size'; then + DECODE_SERVER_CONFIG+=" --decode-context-parallel-size ${DECODE_TP_SIZE:-1}" fi echo "PREFILL_SERVER_CONFIG (after TP/EP/DP): $PREFILL_SERVER_CONFIG" echo "DECODE_SERVER_CONFIG (after TP/EP/DP): $DECODE_SERVER_CONFIG" +if [[ "${ENABLE_PREFIX_CACHING:-0}" == "1" ]]; then + for _cfg in PREFILL_SERVER_CONFIG DECODE_SERVER_CONFIG; do + _val="${!_cfg}" + _val="${_val//--no-enable-prefix-caching/}" + if ! echo "$_val" | grep -q -- '--enable-prefix-caching'; then + _val+=" --enable-prefix-caching" + fi + printf -v "$_cfg" '%s' "$_val" + done + echo "[vLLM] ENABLE_PREFIX_CACHING=1 -> prefix cache enabled on prefill + decode" +fi + +if [[ -n "${MAX_MODEL_LEN:-}" && "${MAX_MODEL_LEN}" != "0" ]]; then + for _cfg in PREFILL_SERVER_CONFIG DECODE_SERVER_CONFIG; do + _val="${!_cfg}" + if echo "$_val" | grep -q -- '--max-model-len'; then + _val=$(echo "$_val" | sed -E "s/--max-model-len[=[:space:]]+[0-9]+/--max-model-len ${MAX_MODEL_LEN}/g") + else + _val+=" --max-model-len ${MAX_MODEL_LEN}" + fi + printf -v "$_cfg" '%s' "$_val" + done + echo "[vLLM] MAX_MODEL_LEN=${MAX_MODEL_LEN}" +fi + +if [[ -n "${MAX_NUM_SEQS:-}" && "${MAX_NUM_SEQS}" != "0" ]]; then + for _cfg in PREFILL_SERVER_CONFIG DECODE_SERVER_CONFIG; do + _val="${!_cfg}" + if echo "$_val" | grep -q -- '--max-num-seqs'; then + _val=$(echo "$_val" | sed -E "s/--max-num-seqs[=[:space:]]+[0-9]+/--max-num-seqs ${MAX_NUM_SEQS}/g") + else + _val+=" --max-num-seqs ${MAX_NUM_SEQS}" + fi + printf -v "$_cfg" '%s' "$_val" + done + echo "[vLLM] MAX_NUM_SEQS=${MAX_NUM_SEQS}" +fi + # ============================================================================= # Container Synchronization # ============================================================================= -echo "Waiting at the container creation barrier on $host_name" +CONTAINER_BARRIER_PORT="${CONTAINER_BARRIER_PORT:-5000}" +echo "Waiting at the container creation barrier on $host_name (port ${CONTAINER_BARRIER_PORT})" python3 $WS_PATH/sync.py barrier \ --local-ip ${host_ip} \ - --local-port 5000 \ + --local-port ${CONTAINER_BARRIER_PORT} \ --enable-port \ --node-ips ${IPADDRS} \ - --node-ports 5000 \ + --node-ports ${CONTAINER_BARRIER_PORT} \ --wait-for-all-ports \ --timeout 600 # ============================================================================= # Cluster Topology Configuration # ============================================================================= -IFS=',' read -ra IP_ARRAY <<< "$IPADDRS" +IP_ARRAY=("${NODE_IP_ARRAY[@]}") PREFILL_ARGS="" DECODE_ARGS="" @@ -220,13 +262,340 @@ echo "Decode node IPs: ${DECODE_ARGS}" # MoRI-IO proxy ZMQ registration port (must match vllm-router --vllm-discovery-address) PROXY_PING_PORT="${PROXY_PING_PORT:-36367}" +# ============================================================================= +# KV connector selection +# ============================================================================= +_MORIIO_EXTRA="\"kv_connector_extra_config\": {\"proxy_ip\": \"${NODE0_ADDR}\", \"proxy_ping_port\": \"${PROXY_PING_PORT}\", \"http_port\": \"${SERVER_PORT}\"}" +MORIIO_PREFILL_CONN="{\"kv_connector\": \"MoRIIOConnector\", \"kv_role\": \"kv_producer\", ${_MORIIO_EXTRA}}" +MORIIO_DECODE_CONN="{\"kv_connector\": \"MoRIIOConnector\", \"kv_role\": \"kv_consumer\", ${_MORIIO_EXTRA}}" +KVT_PREFILL="$MORIIO_PREFILL_CONN" +KVT_DECODE="$MORIIO_DECODE_CONN" + +LMCACHE_CONNECT_HOST="${LMCACHE_CONNECT_HOST:-tcp://${LMCACHE_CONNECT_HOST_IP:-127.0.0.1}}" +LMC_CONN="{\"kv_connector\":\"LMCacheMPConnector\",\"kv_connector_module_path\":\"lmcache.integration.vllm.lmcache_mp_connector\",\"kv_role\":\"kv_both\",\"kv_connector_extra_config\":{\"lmcache.mp.host\":\"${LMCACHE_CONNECT_HOST}\",\"lmcache.mp.port\":${LMCACHE_PORT:-5555},\"lmcache.mp.mq_timeout\":${LMCACHE_MP_MQ_TIMEOUT:-1200}}}" + +# ----------------------------------------------------------------------------- +# LMCache MP + native NIXL/RIXL P/D transfer. +# ----------------------------------------------------------------------------- +# The current LMCache MP docs compose vLLM P/D with MultiConnector: +# NixlConnector : per-request prefill -> decode KV handoff +# LMCacheMPConnector : local LMCache server offload/load for reuse +# P2P sharing is implemented by lmcache server + coordinator. +LMC_NIXL_PREFILL_XFER="{\"kv_connector\":\"NixlConnector\",\"kv_role\":\"kv_producer\",\"kv_load_failure_policy\":\"fail\"}" +LMC_NIXL_DECODE_XFER="{\"kv_connector\":\"NixlConnector\",\"kv_role\":\"kv_consumer\",\"kv_load_failure_policy\":\"fail\"}" +LMC_NIXL_PREFILL_CONN="{\"kv_connector\":\"MultiConnector\",\"kv_role\":\"kv_producer\",\"kv_connector_extra_config\":{\"connectors\":[${LMC_NIXL_PREFILL_XFER},${LMC_CONN}]}}" +LMC_NIXL_DECODE_CONN="{\"kv_connector\":\"MultiConnector\",\"kv_role\":\"kv_consumer\",\"kv_connector_extra_config\":{\"connectors\":[${LMC_NIXL_DECODE_XFER},${LMC_CONN}]}}" + +case "${PREFILL_KV_CONNECTOR:-moriio}" in + moriio-lmcachemp) + KVT_PREFILL="{\"kv_connector\":\"MultiConnector\",\"kv_role\":\"kv_producer\",\"kv_connector_extra_config\":{\"connectors\":[${MORIIO_PREFILL_CONN},${LMC_CONN}]}}" + ;; + lmcache-nixl) + KVT_PREFILL="$LMC_NIXL_PREFILL_CONN" + ;; + moriio|"") + ;; + *) + echo "ERROR: unsupported PREFILL_KV_CONNECTOR=${PREFILL_KV_CONNECTOR}" >&2 + exit 1 + ;; +esac + +case "${DECODE_KV_CONNECTOR:-moriio}" in + moriio-lmcachemp) + KVT_DECODE="{\"kv_connector\":\"MultiConnector\",\"kv_role\":\"kv_consumer\",\"kv_connector_extra_config\":{\"connectors\":[${MORIIO_DECODE_CONN},${LMC_CONN}]}}" + ;; + lmcache-nixl) + KVT_DECODE="$LMC_NIXL_DECODE_CONN" + ;; + moriio|"") + ;; + *) + echo "ERROR: unsupported DECODE_KV_CONNECTOR=${DECODE_KV_CONNECTOR}" >&2 + exit 1 + ;; +esac + +if [[ "$NODE_RANK" -lt "$xP" ]]; then + ROLE_KV_CONNECTOR="${PREFILL_KV_CONNECTOR:-moriio}" +else + ROLE_KV_CONNECTOR="${DECODE_KV_CONNECTOR:-moriio}" +fi +echo "[KV] PREFILL_KV_CONNECTOR=${PREFILL_KV_CONNECTOR:-moriio}; DECODE_KV_CONNECTOR=${DECODE_KV_CONNECTOR:-moriio}; rank connector=${ROLE_KV_CONNECTOR}" + + # vLLM runtime environment (static vars moved to env.sh; these depend on per-node state) setup_vllm_env() { - export VLLM_NIXL_SIDE_CHANNEL_HOST=${rdma_ip} - export VLLM_NIXL_SIDE_CHANNEL_PORT=5600 + local bind_ip="${VLLM_BIND_IP:-${rdma_ip}}" + export VLLM_HOST_IP="${bind_ip}" + export VLLM_NIXL_SIDE_CHANNEL_HOST="${bind_ip}" + export VLLM_NIXL_SIDE_CHANNEL_PORT="${VLLM_NIXL_SIDE_CHANNEL_PORT:-$((5600 + NODE_RANK))}" + if [[ "$ROLE_KV_CONNECTOR" == "lmcache-nixl" ]]; then + # LMCache MP P2P docs recommend UCX_NET_DEVICES=all. The generic + # MoRI/UCX auto-detect path may set ionic_0:1, which is not visible to + # LMCache's NIXL backend in the current nightly container. + export UCX_NET_DEVICES="${LMCACHE_NIXL_UCX_NET_DEVICES:-all}" + export NCCL_CUMEM_ENABLE="${NCCL_CUMEM_ENABLE:-1}" + fi for env_pair in ${MODEL_ENVS}; do export "$env_pair" done + echo "[vLLM] VLLM_HOST_IP=${VLLM_HOST_IP} VLLM_NIXL_SIDE_CHANNEL_HOST=${VLLM_NIXL_SIDE_CHANNEL_HOST} VLLM_NIXL_SIDE_CHANNEL_PORT=${VLLM_NIXL_SIDE_CHANNEL_PORT}" +} + +# LMCache MP P2P uses the coordinator for peer discovery. It is a control-plane +# service only; KV data still moves peer-to-peer via the server transfer channel. +start_lmcache_coordinator_if_needed() { + if [[ "${LMCACHE_ENABLE_P2P:-false}" != "true" ]]; then + return 0 + fi + if [[ "$ROLE_KV_CONNECTOR" != *"-lmcachemp" && "$ROLE_KV_CONNECTOR" != "lmcache-nixl" ]]; then + return 0 + fi + if [[ "$NODE_RANK" -ne 0 ]]; then + return 0 + fi + + local coord_host="${LMCACHE_COORDINATOR_HOST:-0.0.0.0}" + local coord_port="${LMCACHE_COORDINATOR_PORT:-9300}" + local coord_log="/run_logs/slurm_job-${SLURM_JOB_ID}/lmcache_coordinator_${host_name}.log" + echo "[LMCacheMP] starting coordinator on ${coord_host}:${coord_port}" + lmcache coordinator --host "$coord_host" --port "$coord_port" > "$coord_log" 2>&1 & + LMCACHE_COORDINATOR_PID=$! + + for _i in $(seq 1 60); do + if curl -sf --max-time 3 "http://${NODE0_ADDR}:${coord_port}/healthz" >/dev/null 2>&1; then + echo "[LMCacheMP] coordinator healthy" + return 0 + fi + sleep 1 + done + echo "ERROR: LMCache coordinator failed to become healthy; tailing $coord_log" >&2 + tail -n 80 "$coord_log" >&2 || true + return 1 +} + +build_vllm_metrics_urls() { + local urls=() + local ip + for ip in "${IP_ARRAY[@]}"; do + [[ -n "$ip" ]] || continue + urls+=("http://${ip}:${SERVER_PORT}/metrics") + done + (IFS=,; echo "${urls[*]}") +} + +capture_vllm_cache_metrics() { + local label="${1:-snapshot}" + local out="/run_logs/slurm_job-${SLURM_JOB_ID}/vllm_cache_metrics_${label}.txt" + local urls="${AIPERF_SERVER_METRICS_URLS:-$(build_vllm_metrics_urls)}" + local url + mkdir -p "$(dirname "$out")" + { + echo "=== vLLM cache metrics snapshot: ${label} $(date --iso-8601=seconds) ===" + echo "AIPERF_SERVER_METRICS_URLS=${urls}" + IFS=',' read -r -a _metric_urls <<< "$urls" + for url in "${_metric_urls[@]}"; do + [[ -n "$url" ]] || continue + echo "--- ${url} ---" + curl -fsS --max-time 5 "$url" 2>/dev/null \ + | grep -E '^(vllm:num_requests_running|vllm:num_requests_waiting|vllm:num_requests_waiting_by_reason|vllm:prompt_tokens$|vllm:generation_tokens|vllm:prompt_tokens_cached|vllm:prompt_tokens_by_source|vllm:external_prefix_cache_hits|vllm:external_prefix_cache_queries|vllm:kv_cache_usage_perc|vllm:gpu_cache_usage_perc|vllm:prefix_cache_hits|vllm:prefix_cache_queries|vllm:num_preemptions|vllm:request_success|vllm:kv_offload|vllm:cpu_kv|vllm:.*lmcache)' \ + || true + done + } | tee -a "$out" +} + +capture_lmcache_metrics() { + local label="${1:-snapshot}" + if [[ "$ROLE_KV_CONNECTOR" != *"-lmcachemp" && "$ROLE_KV_CONNECTOR" != "lmcache-nixl" ]]; then + return 0 + fi + + local http_host="${LMCACHE_HOST:-127.0.0.1}" + local http_port="${LMCACHE_HTTP_PORT:-8080}" + local base_url="http://${http_host}:${http_port}" + local out="/run_logs/slurm_job-${SLURM_JOB_ID}/lmcache_metrics_${host_name}_${label}.txt" + mkdir -p "$(dirname "$out")" + + { + echo "=== LMCacheMP metrics snapshot: ${label} $(date --iso-8601=seconds) ===" + echo "ROLE_KV_CONNECTOR=${ROLE_KV_CONNECTOR}" + echo "LMCACHE_URL=${base_url}" + echo "--- /healthcheck ---" + curl -fsS --max-time 5 "${base_url}/healthcheck" 2>&1 || true + echo + echo "--- /metrics ---" + if ! curl -fsS --max-time 10 "${base_url}/metrics" 2>&1; then + echo "[LMCacheMP] /metrics unavailable" + fi + echo + echo "--- filtered cache keywords from /metrics ---" + curl -fsS --max-time 10 "${base_url}/metrics" 2>/dev/null \ + | grep -Ei 'hit|miss|evict|eviction|lookup|retrieve|store|read|write|cache|l1|chunk|token|byte|latency|duration' \ + || true + echo + echo "--- lmcache log tail ---" + tail -n 120 "/run_logs/slurm_job-${SLURM_JOB_ID}/lmcache_${host_name}.log" 2>/dev/null || true + } | tee -a "$out" +} + +start_lmcache_mp_if_needed() { + if [[ "$ROLE_KV_CONNECTOR" != *"-lmcachemp" && "$ROLE_KV_CONNECTOR" != "lmcache-nixl" ]]; then + return 0 + fi + + LMCACHE_PATCH_DIR="/run_logs/slurm_job-${SLURM_JOB_ID}/lmcache_mp_patch" + mkdir -p "$LMCACHE_PATCH_DIR" + cat > "$LMCACHE_PATCH_DIR/sitecustomize.py" <<'PY' +"""Keep LMCacheMP from producing proxy-visible PD transfer params. + +MultiConnector permits only one child connector to return kv_transfer_params. +The PD transfer connector owns the prefill->decode protocol; LMCacheMP should +only provide local L2 lookup/retrieve/store for prefix reuse. +""" +import builtins +import sys + +_orig_import = builtins.__import__ + + +def _patch_module(mod): + cls = getattr(mod, "LMCacheMPConnector", None) + if cls is None or getattr(cls, "_inferencex_pd_params_patch", False): + return + orig = cls.request_finished + + def request_finished(self, request, block_ids): + async_save, params = orig(self, request, block_ids) + req_params = getattr(request, "kv_transfer_params", None) + if req_params and ( + req_params.get("do_remote_decode") or req_params.get("do_remote_prefill") + ): + return async_save, None + return async_save, params + + cls.request_finished = request_finished + orig_process_new = cls._process_new_requests + + def _process_new_requests(self, scheduler_output, metadata): + for new_request in scheduler_output.scheduled_new_reqs: + block_ids = getattr(new_request, "block_ids", None) + if not block_ids: + continue + request_tracker = self._get_request_tracker(new_request.req_id) + try: + allocated = request_tracker.num_allocated_blocks() + except Exception: + allocated = {} + if not allocated: + request_tracker.append_block_ids(block_ids) + return orig_process_new(self, scheduler_output, metadata) + + cls._process_new_requests = _process_new_requests + cls._inferencex_pd_params_patch = True + + +def _import(name, globals=None, locals=None, fromlist=(), level=0): + mod = _orig_import(name, globals, locals, fromlist, level) + target = "lmcache.integration.vllm.lmcache_mp_connector" + if name == target or target in sys.modules: + _patch_module(sys.modules[target]) + return mod + + +builtins.__import__ = _import +if "lmcache.integration.vllm.lmcache_mp_connector" in sys.modules: + _patch_module(sys.modules["lmcache.integration.vllm.lmcache_mp_connector"]) + +try: + from vllm.distributed.kv_transfer.kv_connector.v1 import multi_connector + _orig_observe = multi_connector.MultiConnectorPromMetrics.observe + + def _safe_observe(self, transfer_stats_data, engine_idx): + if not isinstance(transfer_stats_data, dict): + return _orig_observe(self, transfer_stats_data, engine_idx) + filtered = { + connector_id: stats_data + for connector_id, stats_data in transfer_stats_data.items() + if connector_id in getattr(self, "_prom_metrics", {}) + } + if filtered: + return _orig_observe(self, filtered, engine_idx) + + multi_connector.MultiConnectorPromMetrics.observe = _safe_observe +except Exception: + pass +PY + export PYTHONPATH="$LMCACHE_PATCH_DIR${PYTHONPATH:+:$PYTHONPATH}" + + python3 -c 'import lmcache.integration.vllm.lmcache_mp_connector; print("LMCacheMPConnector import OK")' + + LMCACHE_LOG="/run_logs/slurm_job-${SLURM_JOB_ID}/lmcache_${host_name}.log" + local enable_p2p="${LMCACHE_ENABLE_P2P:-false}" + local server_host="${LMCACHE_SERVER_HOST:-${LMCACHE_HOST:-127.0.0.1}}" + if [[ "$enable_p2p" == "true" ]]; then + server_host="${LMCACHE_SERVER_HOST:-0.0.0.0}" + fi + + local server_args=( + --host "$server_host" --port "${LMCACHE_PORT:-5555}" + --http-port "${LMCACHE_HTTP_PORT:-8080}" + --l1-size-gb "${LMCACHE_L1_SIZE_GB:-1200}" + --l1-init-size-gb "${LMCACHE_L1_INIT_SIZE_GB:-20}" + --l1-read-ttl-seconds "${LMCACHE_L1_READ_TTL_SECONDS:-7200}" + --chunk-size "${LMCACHE_CHUNK_SIZE:-256}" + --max-workers "${LMCACHE_MAX_WORKERS:-8}" + --eviction-policy LRU + --instance-id "${LMCACHE_INSTANCE_ID:-${host_name}-${SLURM_JOB_ID}}" + ) + if [[ "$enable_p2p" == "true" ]]; then + local coord_port="${LMCACHE_COORDINATOR_PORT:-9300}" + server_args+=( + --l1-align-bytes "${LMCACHE_L1_ALIGN_BYTES:-65536}" + --coordinator-url "${LMCACHE_COORDINATOR_URL:-http://${NODE0_ADDR}:${coord_port}}" + --coordinator-advertise-ip "${LMCACHE_COORDINATOR_ADVERTISE_IP:-${host_ip}}" + --p2p-advertise-url "${LMCACHE_P2P_ADVERTISE_URL:-${rdma_ip}:${LMCACHE_P2P_PORT:-8500}}" + ) + if [[ -n "${LMCACHE_P2P_LISTEN_URL:-}" ]]; then + server_args+=(--p2p-listen-url "${LMCACHE_P2P_LISTEN_URL}") + fi + fi + + echo "[LMCacheMP] starting server on ${server_host}:${LMCACHE_PORT:-5555} (http-port ${LMCACHE_HTTP_PORT:-8080}) p2p=${enable_p2p}" + lmcache server "${server_args[@]}" \ + > "$LMCACHE_LOG" 2>&1 & + + for _i in $(seq 1 120); do + if curl -sf --max-time 3 "http://${LMCACHE_HOST:-127.0.0.1}:${LMCACHE_HTTP_PORT:-8080}/healthcheck" >/dev/null 2>&1; then + echo "[LMCacheMP] server healthy" + capture_lmcache_metrics "after_lmcache_start" + return 0 + fi + sleep 2 + done + echo "ERROR: LMCache MP server failed to become healthy; tailing $LMCACHE_LOG" >&2 + tail -n 80 "$LMCACHE_LOG" >&2 || true + return 1 +} + +dump_runtime_logs() { + local label="${1:-runtime failure}" + echo "ERROR: ${label}" >&2 + if [[ "$DRY_RUN" -eq 0 ]]; then + local logs_output="${BENCHMARK_LOGS_DIR:-/run_logs}/logs" + mkdir -p "$logs_output" 2>/dev/null || true + cp -r "/run_logs/slurm_job-${SLURM_JOB_ID}" "$logs_output/" 2>/dev/null || \ + sudo cp -r "/run_logs/slurm_job-${SLURM_JOB_ID}" "$logs_output/" 2>/dev/null || true + echo "==== staged runtime logs to ${logs_output}/slurm_job-${SLURM_JOB_ID} ====" >&2 + fi + echo "==== /run_logs/slurm_job-${SLURM_JOB_ID} files ====" >&2 + find "/run_logs/slurm_job-${SLURM_JOB_ID}" -maxdepth 2 -type f -printf '%p %s bytes\n' 2>/dev/null | sort >&2 || true + echo "==== recent server logs ====" >&2 + for _log in /run_logs/slurm_job-"${SLURM_JOB_ID}"/*.log; do + [[ -f "$_log" ]] || continue + echo "----- $_log -----" >&2 + tail -n 200 "$_log" >&2 || true + done } # ============================================================================= @@ -250,16 +619,18 @@ if [ "$NODE_RANK" -eq 0 ]; then echo "================================================" setup_vllm_env + start_lmcache_coordinator_if_needed || exit 1 + start_lmcache_mp_if_needed || exit 1 # Router is started as an external container by job.slurm (VLLM_ROUTER_IMAGE) echo "Using external vllm-router container (started by job.slurm on this node)" - SERVED_MODEL="${MODEL_NAME}" + SERVED_MODEL="${SERVED_MODEL_NAME:-${MODEL_NAME}}" PREFILL_CMD="vllm serve ${MODEL_PATH} \ --served-model-name ${SERVED_MODEL} \ --port $SERVER_PORT \ --trust-remote-code \ - --kv-transfer-config '{\"kv_connector\": \"MoRIIOConnector\", \"kv_role\": \"kv_producer\", \"kv_connector_extra_config\": {\"proxy_ip\": \"${NODE0_ADDR}\", \"proxy_ping_port\": \"${PROXY_PING_PORT}\", \"http_port\": \"${SERVER_PORT}\", \"read_mode\": true}}' \ + --kv-transfer-config '${KVT_PREFILL}' \ ${PREFILL_SERVER_CONFIG}" if [[ "$DRY_RUN" -eq 1 ]]; then @@ -276,11 +647,15 @@ if [ "$NODE_RANK" -eq 0 ]; then if [[ "$DRY_RUN" -eq 1 ]]; then echo "DRY RUN: skipping barrier (wait-for-all-ports)" else - python3 $WS_PATH/sync.py barrier \ + if ! python3 $WS_PATH/sync.py barrier \ --node-ips ${IPADDRS} \ --node-ports $SERVER_PORT \ --wait-for-all-ports \ - --timeout 1800 + --timeout 1800; then + dump_runtime_logs "prefill/decode server readiness barrier failed" + [[ -n "${prefill_pid:-}" ]] && kill $prefill_pid 2>/dev/null || true + exit 1 + fi fi echo "Congratulations!!! All prefill and decode servers are up . . ." @@ -296,7 +671,11 @@ if [ "$NODE_RANK" -eq 0 ]; then if [[ "$DRY_RUN" -eq 1 ]]; then echo "DRY RUN: $HEALTH_BARRIER_CMD" else - eval "$HEALTH_BARRIER_CMD" + if ! eval "$HEALTH_BARRIER_CMD"; then + dump_runtime_logs "proxy health barrier failed" + [[ -n "${prefill_pid:-}" ]] && kill $prefill_pid 2>/dev/null || true + exit 1 + fi echo "MoRI-IO proxy is ready for benchmarking" fi @@ -305,21 +684,79 @@ if [ "$NODE_RANK" -eq 0 ]; then cd $WS_PATH export ROUTER_PORT=$ROUTER_PORT - BENCH_CMD="bash $WS_PATH/bench.sh ${xP} ${yD} $((PREFILL_TP_SIZE*xP)) $((DECODE_TP_SIZE*yD)) \ - $MODEL_DIR $MODEL_NAME /run_logs/slurm_job-${SLURM_JOB_ID} ${BENCH_INPUT_LEN} \ - ${BENCH_OUTPUT_LEN} \"${BENCH_MAX_CONCURRENCY}\" ${BENCH_REQUEST_RATE} \ - ${BENCH_RANDOM_RANGE_RATIO} ${BENCH_NUM_PROMPTS_MULTIPLIER}" + export AIPERF_SERVER_METRICS_URLS="${AIPERF_SERVER_METRICS_URLS:-$(build_vllm_metrics_urls)}" + echo "AIPERF_SERVER_METRICS_URLS=${AIPERF_SERVER_METRICS_URLS}" + capture_vllm_cache_metrics "before_benchmark" + capture_lmcache_metrics "before_benchmark" if [[ "${EVAL_ONLY:-false}" == "true" ]]; then echo "EVAL_ONLY mode: skipping throughput benchmark" + elif [[ "${IS_AGENTIC:-0}" == "1" || "${SCENARIO_TYPE:-}" == "agentic-coding" ]]; then + AGENTIC_CMD=' + set -euo pipefail + cd /workspace + export AIPERF_RUNTIME_DIR="${AIPERF_RUNTIME_DIR:-/run_logs/slurm_job-'"${SLURM_JOB_ID}"'/aiperf_runtime}" + source /workspace/benchmarks/benchmark_lib.sh + + export PORT='"${ROUTER_PORT}"' + export MODEL="${SERVED_MODEL_NAME:-${MODEL:-${MODEL_NAME}}}" + export MODEL_PREFIX="${MODEL_PREFIX:-'"${MODEL_NAME}"'}" + export FRAMEWORK="${FRAMEWORK:-vllm-disagg}" + export PRECISION="${PRECISION:-}" + export KV_OFFLOADING="${KV_OFFLOADING:-none}" + export DURATION="${DURATION:-1800}" + export RESULT_FILENAME="${RESULT_FILENAME:-agentic_${MODEL_NAME}_${SLURM_JOB_ID}}" + export AGENTIC_OUTPUT_DIR="/run_logs/slurm_job-'"${SLURM_JOB_ID}"'/agentic" + export AIPERF_SERVER_METRICS_URLS="${AIPERF_SERVER_METRICS_URLS:-'"${AIPERF_SERVER_METRICS_URLS}"'}" + + IFS="x" read -r -a CONCURRENCIES <<< "${BENCH_MAX_CONCURRENCY:-1}" + if [ "${#CONCURRENCIES[@]}" -eq 0 ]; then + echo "ERROR: no agentic concurrency configured" >&2 + exit 1 + fi + + resolve_trace_source + install_agentic_deps + + for concurrency in "${CONCURRENCIES[@]}"; do + if ! [[ "$concurrency" =~ ^[1-9][0-9]*$ ]]; then + echo "ERROR: invalid agentic concurrency: $concurrency" >&2 + exit 1 + fi + export CONC="$concurrency" + result_dir="/run_logs/slurm_job-'"${SLURM_JOB_ID}"'/agentic/conc_${concurrency}" + mkdir -p "$result_dir" + echo "Running vLLM disagg agentic replay at concurrency $concurrency" + build_replay_cmd "$result_dir" + run_agentic_replay_and_write_outputs "$result_dir" + done + ' + if [[ "$DRY_RUN" -eq 1 ]]; then + echo "DRY RUN: $AGENTIC_CMD" + else + set -x + bash -lc "$AGENTIC_CMD" + set +x + fi elif [[ "$DRY_RUN" -eq 1 ]]; then + BENCH_CMD="bash $WS_PATH/bench.sh ${xP} ${yD} $((PREFILL_TP_SIZE*xP)) $((DECODE_TP_SIZE*yD)) \ + $MODEL_DIR $MODEL_NAME /run_logs/slurm_job-${SLURM_JOB_ID} ${BENCH_INPUT_LEN} \ + ${BENCH_OUTPUT_LEN} \"${BENCH_MAX_CONCURRENCY}\" ${BENCH_REQUEST_RATE} \ + ${BENCH_RANDOM_RANGE_RATIO} ${BENCH_NUM_PROMPTS_MULTIPLIER}" echo "DRY RUN: $BENCH_CMD" else + BENCH_CMD="bash $WS_PATH/bench.sh ${xP} ${yD} $((PREFILL_TP_SIZE*xP)) $((DECODE_TP_SIZE*yD)) \ + $MODEL_DIR $MODEL_NAME /run_logs/slurm_job-${SLURM_JOB_ID} ${BENCH_INPUT_LEN} \ + ${BENCH_OUTPUT_LEN} \"${BENCH_MAX_CONCURRENCY}\" ${BENCH_REQUEST_RATE} \ + ${BENCH_RANDOM_RANGE_RATIO} ${BENCH_NUM_PROMPTS_MULTIPLIER}" set -x eval "$BENCH_CMD" set +x fi + capture_vllm_cache_metrics "after_benchmark" + capture_lmcache_metrics "after_benchmark" + # Run evaluation if requested (before killing router) if [[ "${RUN_EVAL:-false}" == "true" ]]; then echo "Running lm-eval evaluation on Node 0..." @@ -391,6 +828,17 @@ if [ "$NODE_RANK" -eq 0 ]; then popd fi + capture_vllm_cache_metrics "after_eval" + capture_lmcache_metrics "after_eval" + fi + + # Stage workspace-level artifacts such as aggregated benchmark JSON and + # AIPerf exports before the container exits. + if [[ "$DRY_RUN" -eq 0 ]]; then + WORKSPACE_ARTIFACT_DIR="/run_logs/slurm_job-${SLURM_JOB_ID}/workspace_artifacts" + mkdir -p "$WORKSPACE_ARTIFACT_DIR" + find /workspace -maxdepth 1 -type f \( -name '*.json' -o -name '*.jsonl' -o -name '*.csv' -o -name '*.png' \) \ + -exec cp -f {} "$WORKSPACE_ARTIFACT_DIR/" \; 2>/dev/null || true fi # Copy benchmark/eval results to BENCHMARK_LOGS_DIR (mounted from host) @@ -419,13 +867,15 @@ elif [ "$NODE_RANK" -gt 0 ] && [ "$NODE_RANK" -lt "$xP" ]; then echo "Using prefill config: $PREFILL_SERVER_CONFIG" setup_vllm_env + start_lmcache_coordinator_if_needed || exit 1 + start_lmcache_mp_if_needed || exit 1 - SERVED_MODEL="${MODEL_NAME}" + SERVED_MODEL="${SERVED_MODEL_NAME:-${MODEL_NAME}}" PREFILL_CMD="vllm serve ${MODEL_PATH} \ --served-model-name ${SERVED_MODEL} \ --port $SERVER_PORT \ --trust-remote-code \ - --kv-transfer-config '{\"kv_connector\": \"MoRIIOConnector\", \"kv_role\": \"kv_producer\", \"kv_connector_extra_config\": {\"proxy_ip\": \"${NODE0_ADDR}\", \"proxy_ping_port\": \"${PROXY_PING_PORT}\", \"http_port\": \"${SERVER_PORT}\", \"read_mode\": true}}' \ + --kv-transfer-config '${KVT_PREFILL}' \ ${PREFILL_SERVER_CONFIG}" if [[ "$DRY_RUN" -eq 1 ]]; then @@ -448,7 +898,11 @@ elif [ "$NODE_RANK" -gt 0 ] && [ "$NODE_RANK" -lt "$xP" ]; then if [[ "$DRY_RUN" -eq 1 ]]; then echo "DRY RUN: $BARRIER_CMD" else - eval "$BARRIER_CMD" + if ! eval "$BARRIER_CMD"; then + dump_runtime_logs "proxy port barrier failed on additional prefill node" + [[ -n "${prefill_pid:-}" ]] && kill $prefill_pid 2>/dev/null || true + exit 1 + fi fi echo "Waiting until proxy server closes..." @@ -461,6 +915,7 @@ elif [ "$NODE_RANK" -gt 0 ] && [ "$NODE_RANK" -lt "$xP" ]; then else eval "$WAIT_CMD" fi + capture_lmcache_metrics "after_router_wait" echo "Killing the prefill server" [[ "$DRY_RUN" -eq 0 ]] && kill $prefill_pid 2>/dev/null || true @@ -470,18 +925,20 @@ else echo "Using decode config: $DECODE_SERVER_CONFIG" setup_vllm_env + start_lmcache_coordinator_if_needed || exit 1 + start_lmcache_mp_if_needed || exit 1 for env_pair in ${DECODE_MODEL_ENVS}; do export "$env_pair" echo "[DECODE_ENV] $env_pair" done - SERVED_MODEL="${MODEL_NAME}" + SERVED_MODEL="${SERVED_MODEL_NAME:-${MODEL_NAME}}" DECODE_CMD="vllm serve ${MODEL_PATH} \ --served-model-name ${SERVED_MODEL} \ --port $SERVER_PORT \ --trust-remote-code \ - --kv-transfer-config '{\"kv_connector\": \"MoRIIOConnector\", \"kv_role\": \"kv_consumer\", \"kv_connector_extra_config\": {\"proxy_ip\": \"${NODE0_ADDR}\", \"proxy_ping_port\": \"${PROXY_PING_PORT}\", \"http_port\": \"${SERVER_PORT}\", \"read_mode\": true}}' \ + --kv-transfer-config '${KVT_DECODE}' \ ${DECODE_SERVER_CONFIG}" if [[ "$DRY_RUN" -eq 1 ]]; then @@ -504,7 +961,11 @@ else if [[ "$DRY_RUN" -eq 1 ]]; then echo "DRY RUN: $BARRIER_CMD" else - eval "$BARRIER_CMD" + if ! eval "$BARRIER_CMD"; then + dump_runtime_logs "proxy port barrier failed on decode node" + [[ -n "${decode_pid:-}" ]] && kill $decode_pid 2>/dev/null || true + exit 1 + fi fi echo "Waiting until proxy server closes..." diff --git a/benchmarks/multi_node/amd_utils/setup_deps.sh b/benchmarks/multi_node/amd_utils/setup_deps.sh index 35eaf17dc0..3cd5b1337f 100644 --- a/benchmarks/multi_node/amd_utils/setup_deps.sh +++ b/benchmarks/multi_node/amd_utils/setup_deps.sh @@ -79,6 +79,602 @@ install_amd_quark() { _SETUP_INSTALLED+=("amd-quark") } +# --------------------------------------------------------------------------- +# 8. Patch vLLM MoRI-IO save_kv_layer busy-spin (C128 tail-batch deadlock) +# In WRITE mode, save_kv_layer spins forever waiting for the handshake +# callback to set write_ready_flags. This blocks the model worker thread, +# preventing it from responding to EngineCore shm_broadcast, causing a +# TimeoutError cascade and crash. +# Patch: add time.sleep(0.001) and a 30s timeout to yield CPU and prevent +# the model worker from deadlocking. +# --------------------------------------------------------------------------- +patch_moriio_save_kv_timeout() { + python3 -c ' +import os, sys + +try: + import vllm.distributed.kv_transfer.kv_connector.v1.moriio.moriio_connector as mc + f = mc.__file__ + src = open(f).read() + + # Already patched? + if "[PATCHED] save_kv_layer timeout" in src: + print("[SETUP] save_kv_layer timeout patch already applied") + sys.exit(0) + + old = """ while True: + if ( + self._ready_requests.empty() + and remote_engine_id not in self.write_ready_flags + ): + continue""" + + if old not in src: + print("[SETUP] WARN: save_kv_layer busy-spin pattern not found, skipping patch") + sys.exit(0) + + new = """ # [PATCHED] save_kv_layer — null guard + timeout + sleep + if remote_engine_id is None: + return + import time as _time, os as _os + _wait_start = _time.monotonic() + _SAVE_KV_TIMEOUT = float(_os.environ.get("VLLM_MORIIO_HANDSHAKE_TIMEOUT", "30")) + while True: + if ( + self._ready_requests.empty() + and remote_engine_id not in self.write_ready_flags + ): + _elapsed = _time.monotonic() - _wait_start + if _elapsed > _SAVE_KV_TIMEOUT: + import logging as _logging + _logging.getLogger("vllm.moriio").warning( + "[HANGFIX] save_kv_layer: timeout (%.1fs) waiting for " + "write_ready_flags[%s], breaking to unblock model " + "worker", _elapsed, remote_engine_id) + break + _time.sleep(0.001) + continue""" + + new_src = src.replace(old, new) + if new_src == src: + print("[SETUP] WARN: replacement had no effect") + sys.exit(0) + + open(f, "w").write(new_src) + print("[SETUP] Patched save_kv_layer: null guard + timeout + sleep") +except Exception as e: + print(f"[SETUP] WARN patch save_kv_layer: {e}", file=sys.stderr) +' + _SETUP_INSTALLED+=("MoRIIO-save-kv-timeout-patch") +} + +# --------------------------------------------------------------------------- +# 9. Patch MoRIIO waiting_for_transfer_complete with bounded timeout +# The original status.Wait() blocks forever if an RDMA completion never +# arrives (e.g., NIC queue saturation at C256). This replaces the unbounded +# wait with a polling loop using status.Succeeded() + configurable timeout. +# Also adds error handling to the write worker loop so a single failed +# transfer doesn't kill the background thread. +# --------------------------------------------------------------------------- +patch_moriio_transfer_timeout() { + python3 -c ' +import os, sys, textwrap + +try: + import vllm.distributed.kv_transfer.kv_connector.v1.moriio.moriio_engine as me + f = me.__file__ + src = open(f).read() + + if "[PATCHED] transfer completion timeout" in src: + print("[SETUP] transfer completion timeout patch already applied") + sys.exit(0) + + # --- Patch 1: Replace waiting_for_transfer_complete with polling + timeout --- + old_wait = """ def waiting_for_transfer_complete(self): + if not self.transfer_status: + return + + transfers_to_wait = [] + with self.lock: + transfers_to_wait = self.transfer_status[:] + self.transfer_status.clear() + + for status in transfers_to_wait: + try: + status.Wait() + if not status.Succeeded(): + logger.error( + "Transfer failed: %s, Code: %s", status.Message(), status.Code() + ) + raise TransferError("MoRIIO transfer failed!") + except Exception as e: + logger.error("Transfer %s failed: %s", status, e) + raise""" + + new_wait = """ def waiting_for_transfer_complete(self): + # [PATCHED] transfer completion timeout — bounded polling loop + import time as _time, os as _os + if not self.transfer_status: + return + + _timeout = float(_os.environ.get("VLLM_MORIIO_TRANSFER_TIMEOUT", "120")) + + transfers_to_wait = [] + with self.lock: + transfers_to_wait = self.transfer_status[:] + self.transfer_status.clear() + + _start = _time.monotonic() + remaining = list(transfers_to_wait) + _polls = 0 + _completed = 0 + + while remaining: + _elapsed = _time.monotonic() - _start + if _elapsed > _timeout: + logger.error( + "[HANGFIX] transfer_timeout elapsed=%.1fs " + "pending=%d/%d completed=%d polls=%d " + "action=raise_transfer_error", + _elapsed, len(remaining), len(transfers_to_wait), + _completed, _polls, + ) + raise TransferError( + f"RDMA transfer timeout after {_elapsed:.1f}s, " + f"{len(remaining)}/{len(transfers_to_wait)} pending" + ) + + still_waiting = [] + for status in remaining: + try: + if status.Succeeded(): + _completed += 1 + continue + still_waiting.append(status) + except Exception as e: + logger.error( + "[HANGFIX] transfer_poll_error error=%s", e) + raise TransferError( + f"Transfer failed during poll: {e}" + ) from e + + remaining = still_waiting + if remaining: + _time.sleep(0.005) + _polls += 1 + if _polls % 2000 == 0: + logger.warning( + "[HANGFIX] transfer_wait pending=%d " + "completed=%d elapsed=%.1fs timeout=%.0fs", + len(remaining), _completed, + _time.monotonic() - _start, _timeout, + )""" + + if old_wait not in src: + print("[SETUP] WARN: waiting_for_transfer_complete pattern not found") + sys.exit(0) + + new_src = src.replace(old_wait, new_wait) + + # --- Patch 2: Add error handling + cleanup to _write_worker_loop --- + old_loop = """ self._execute_write_task(task)""" + + new_loop = """ try: + self._execute_write_task(task) + except Exception as _e: + logger.error( + "[HANGFIX] req=%s write_task_failed error=%s " + "action=cleanup_and_mark_done", + task.request_id, _e, + ) + try: + _wr = self.worker.moriio_wrapper + with _wr.lock: + _wr.done_req_ids.append(task.request_id) + _wr.done_remote_allocate_req_dict.pop( + task.request_id, None + ) + except Exception: + pass""" + + if old_loop in new_src: + new_src = new_src.replace(old_loop, new_loop, 1) + else: + print("[SETUP] WARN: _write_worker_loop pattern not found for error handling") + + # --- Patch 3: Add deferred task timeout to _process_deferred_tasks --- + old_deferred = """ def _process_deferred_tasks(self) -> None: + \"\"\"Process tasks that were previously deferred.\"\"\" + if not self._deferred_tasks: + return + + still_deferred: list[WriteTask] = [] + for task in self._deferred_tasks: + if self._is_remote_ready(task): + self._execute_write_task(task) + else: + still_deferred.append(task) + + self._deferred_tasks = still_deferred""" + + new_deferred = """ def _process_deferred_tasks(self) -> None: + \"\"\"Process tasks that were previously deferred.\"\"\" + # [PATCHED] deferred task timeout — prune stale tasks + import time as _time, os as _os + if not self._deferred_tasks: + return + + _DEFER_TIMEOUT = float( + _os.environ.get("VLLM_MORIIO_DEFER_TIMEOUT", "60")) + + still_deferred: list[WriteTask] = [] + for task in self._deferred_tasks: + _age = _time.monotonic() - getattr(task, "_defer_ts", _time.monotonic()) + if _age > _DEFER_TIMEOUT: + logger.error( + "[HANGFIX] req=%s deferred_task_expired age=%.1fs " + "action=drop_and_mark_done", + task.request_id, _age, + ) + try: + _wr = self.worker.moriio_wrapper + with _wr.lock: + _wr.done_req_ids.append(task.request_id) + _wr.done_remote_allocate_req_dict.pop( + task.request_id, None) + except Exception: + pass + continue + if self._is_remote_ready(task): + try: + self._execute_write_task(task) + except Exception as _e: + logger.error( + "[HANGFIX] req=%s deferred_write_failed error=%s", + task.request_id, _e, + ) + try: + _wr = self.worker.moriio_wrapper + with _wr.lock: + _wr.done_req_ids.append(task.request_id) + _wr.done_remote_allocate_req_dict.pop( + task.request_id, None) + except Exception: + pass + else: + still_deferred.append(task) + + self._deferred_tasks = still_deferred""" + + if old_deferred in new_src: + new_src = new_src.replace(old_deferred, new_deferred, 1) + else: + print("[SETUP] WARN: _process_deferred_tasks pattern not found") + + # --- Patch 4: Stamp defer time when task is deferred --- + old_defer_add = """ self._deferred_tasks.append(task)""" + new_defer_add = """ import time as _time2 + if not hasattr(task, "_defer_ts"): + task._defer_ts = _time2.monotonic() + self._deferred_tasks.append(task)""" + if old_defer_add in new_src: + new_src = new_src.replace(old_defer_add, new_defer_add, 1) + else: + print("[SETUP] WARN: deferred task timestamp patch target not found") + + open(f, "w").write(new_src) + print("[SETUP] Patched: transfer timeout + writer error handling") + +except Exception as e: + print(f"[SETUP] WARN patch transfer_timeout: {e}", file=sys.stderr) +' + _SETUP_INSTALLED+=("MoRIIO-transfer-timeout-patch") +} + +# --------------------------------------------------------------------------- +# 10. Patch MoRIIO notify parser for vllm-router service-discovery request IDs +# vllm-router's MoRIIO PD path sends OpenAI request IDs such as +# chatcmpl-___prefill_addr_host:...___decode_addr_host:... to the notify +# socket. Upstream MoRIIO only accepts msgpack allocation messages or +# tx-prefixed completion messages, so the notify listener thread exits with +# HandshakeError on the first routed request. Treat router request IDs as +# completion notifications. +# --------------------------------------------------------------------------- +patch_moriio_router_request_id_notify() { + python3 -c ' +import os, sys + +try: + import vllm.distributed.kv_transfer.kv_connector.v1.moriio.moriio_engine as me + f = me.__file__ + src = open(f).read() + + if "[PATCHED] vllm-router request id notify" in src: + print("[SETUP] vllm-router request-id notify patch already applied") + sys.exit(0) + + old = """ if msg_str.startswith(MoRIIOConstants.TRANSFER_PREFIX): + self._handle_completion_message(msg_str) + handled = True""" + + new = """ if msg_str.startswith(MoRIIOConstants.TRANSFER_PREFIX): + self._handle_completion_message(msg_str) + handled = True + elif "___prefill_addr_" in msg_str and "___decode_addr_" in msg_str: + # [PATCHED] vllm-router request id notify + self._handle_completion_message(msg_str) + handled = True""" + + if old not in src: + print("[SETUP] WARN: MoRIIO notify parser pattern not found") + sys.exit(0) + + open(f, "w").write(src.replace(old, new, 1)) + print("[SETUP] Patched MoRIIO notify parser for vllm-router request IDs") +except Exception as e: + print(f"[SETUP] WARN patch vllm-router request-id notify: {e}", file=sys.stderr) +' + _SETUP_INSTALLED+=("MoRIIO-router-request-id-notify-patch") +} + +# --------------------------------------------------------------------------- +# 11. Patch MoRIIO start_load_kv busy-spin (same pattern as save_kv_layer) +# The READ-mode spin loop in start_load_kv has the same unbounded-spin +# issue as save_kv_layer. Add timeout + sleep + null guard. +# --------------------------------------------------------------------------- +patch_moriio_load_kv_timeout() { + python3 -c ' +import os, sys + +try: + import vllm.distributed.kv_transfer.kv_connector.v1.moriio.moriio_connector as mc + f = mc.__file__ + src = open(f).read() + + if "[PATCHED] start_load_kv timeout" in src: + print("[SETUP] start_load_kv timeout patch already applied") + sys.exit(0) + + old = """ while True: + if ( + self._ready_requests.empty() + and remote_engine_id not in self.load_ready_flag + and wait_handshake_readd_req + ): + continue""" + + if old not in src: + print("[SETUP] WARN: start_load_kv busy-spin pattern not found, skipping") + sys.exit(0) + + new = """ # [PATCHED] start_load_kv timeout — prevent model worker deadlock + if remote_engine_id is None and not wait_handshake_readd_req: + self._reqs_to_send.update(metadata.reqs_to_send) + return + import time as _time, os as _os + _wait_start = _time.monotonic() + _LOAD_KV_TIMEOUT = float(_os.environ.get("VLLM_MORIIO_HANDSHAKE_TIMEOUT", "30")) + while True: + if ( + self._ready_requests.empty() + and remote_engine_id not in self.load_ready_flag + and wait_handshake_readd_req + ): + if _time.monotonic() - _wait_start > _LOAD_KV_TIMEOUT: + import logging as _logging + _logging.getLogger("vllm.moriio").warning( + "[HANGFIX] start_load_kv: timeout (%.1fs) waiting for " + "load_ready_flag[%s]", _time.monotonic() - _wait_start, + remote_engine_id) + break + _time.sleep(0.001) + continue""" + + new_src = src.replace(old, new) + if new_src == src: + print("[SETUP] WARN: start_load_kv replacement had no effect") + sys.exit(0) + + open(f, "w").write(new_src) + print("[SETUP] Patched start_load_kv busy-spin with timeout + sleep") +except Exception as e: + print(f"[SETUP] WARN patch start_load_kv: {e}", file=sys.stderr) +' + _SETUP_INSTALLED+=("MoRIIO-load-kv-timeout-patch") +} + +# --------------------------------------------------------------------------- +# 11. Fix READ-mode scheduler assertion in _update_from_kv_xfer_finished +# vLLM asserts that a request in finished_recving must be either +# WAITING_FOR_REMOTE_KVS or finished. In READ mode the request can +# transition to RUNNING before the aggregated recv notification arrives, +# crashing the engine with AssertionError. +# (present in v0.17.1 & v0.18.0) +# --------------------------------------------------------------------------- +patch_scheduler_read_mode_fix() { + python3 -c ' +import os, sys + +try: + import vllm.v1.core.sched.scheduler as smod + f = smod.__file__ + src = open(f).read() + + if "[PATCHED] read-mode recv assertion" in src: + print("[SETUP] scheduler read-mode assertion fix already applied") + sys.exit(0) + + old_recv = """ for req_id in kv_connector_output.finished_recving or (): + logger.debug("Finished recving KV transfer for request %s", req_id) + assert req_id in self.requests + req = self.requests[req_id] + if req.status == RequestStatus.WAITING_FOR_REMOTE_KVS: + self.finished_recving_kv_req_ids.add(req_id) + else: + assert RequestStatus.is_finished(req.status) + self._free_blocks(self.requests[req_id])""" + + new_recv = """ # [PATCHED] read-mode recv assertion — handle intermediate states + for req_id in kv_connector_output.finished_recving or (): + logger.debug("Finished recving KV transfer for request %s", req_id) + if req_id not in self.requests: + logger.debug("Request %s already removed, skipping recv", req_id) + continue + req = self.requests[req_id] + if req.status == RequestStatus.WAITING_FOR_REMOTE_KVS: + self.finished_recving_kv_req_ids.add(req_id) + elif RequestStatus.is_finished(req.status): + self._free_blocks(self.requests[req_id]) + else: + logger.debug( + "Request %s recv finished but status=%s (not " + "WAITING_FOR_REMOTE_KVS or finished), skipping " + "block free — will be freed on request completion", + req_id, req.status.name)""" + + if old_recv not in src: + print("[SETUP] WARN: scheduler finished_recving pattern not found, skipping") + sys.exit(0) + + new_src = src.replace(old_recv, new_recv, 1) + + old_send = """ for req_id in kv_connector_output.finished_sending or (): + logger.debug("Finished sending KV transfer for request %s", req_id) + assert req_id in self.requests + self._free_blocks(self.requests[req_id])""" + + new_send = """ for req_id in kv_connector_output.finished_sending or (): + logger.debug("Finished sending KV transfer for request %s", req_id) + if req_id not in self.requests: + logger.debug("Request %s already removed, skipping send", req_id) + continue + self._free_blocks(self.requests[req_id])""" + + if old_send in new_src: + new_src = new_src.replace(old_send, new_send, 1) + else: + print("[SETUP] WARN: scheduler finished_sending pattern not found") + + open(f, "w").write(new_src) + print("[SETUP] Patched: scheduler _update_from_kv_xfer_finished read-mode fix") + +except Exception as e: + print(f"[SETUP] WARN patch scheduler read-mode: {e}", file=sys.stderr) +' + _SETUP_INSTALLED+=("scheduler-read-mode-fix") +} + +# --------------------------------------------------------------------------- +# 12. Idle KV block reaper for disaggregated prefill (READ mode) +# The RIXL notification path can lose `finished_sending` signals under +# high concurrency with ibv_post_send failures. This leaves KV blocks +# permanently allocated on the prefill engine even after the decode has +# finished reading. Over multiple benchmark rounds, leaked blocks +# accumulate and eventually saturate the prefill KV cache. +# +# Fix: instrument the scheduler's `schedule()` method to detect idle +# periods (0 running, 0 waiting for >5s) and force-free blocks for +# any remaining requests whose status is finished. +# --------------------------------------------------------------------------- +patch_prefill_idle_kv_reaper() { + python3 -c ' +import os, sys + +try: + import vllm.v1.core.sched.scheduler as smod + f = smod.__file__ + src = open(f).read() + + if "[PATCHED] idle-kv-reaper" in src: + print("[SETUP] idle KV block reaper already applied") + sys.exit(0) + + # Find the _update_from_kv_xfer_finished method end and add reaper logic + # We inject into the method that processes KV transfer completions. + marker = "[PATCHED] read-mode recv assertion" + if marker not in src: + print("[SETUP] WARN: scheduler read-mode patch not found, skipping reaper") + sys.exit(0) + + # Add reaper state initialization to __init__ + old_init_marker = "self.finished_recving_kv_req_ids" + if old_init_marker not in src: + print("[SETUP] WARN: finished_recving_kv_req_ids not found in scheduler") + sys.exit(0) + + # Find the first occurrence to insert reaper state + init_pos = src.find(old_init_marker) + # Find the line containing it + line_end = src.find("\n", init_pos) + init_line = src[init_pos:line_end] + + # Add reaper state after this line + reaper_init = init_line + """ + # [PATCHED] idle-kv-reaper state + self._idle_kv_reaper_ts = 0.0 + self._idle_kv_reaper_active = False""" + + src = src.replace(init_line, reaper_init, 1) + + # Now add the reaper logic at the end of _update_from_kv_xfer_finished + # Find the finished_sending handler we patched + send_handler = """ for req_id in kv_connector_output.finished_sending or (): + logger.debug("Finished sending KV transfer for request %s", req_id) + if req_id not in self.requests: + logger.debug("Request %s already removed, skipping send", req_id) + continue + self._free_blocks(self.requests[req_id])""" + + reaper_logic = send_handler + """ + + # [PATCHED] idle-kv-reaper — force-free leaked prefill KV blocks + import time as _time + _REAPER_IDLE_SECS = 5.0 + _num_running = sum(1 for r in self.requests.values() + if r.status == RequestStatus.RUNNING) + _should_reap = (_num_running == 0) + + if _should_reap: + if not self._idle_kv_reaper_active: + self._idle_kv_reaper_active = True + self._idle_kv_reaper_ts = _time.monotonic() + elif _time.monotonic() - self._idle_kv_reaper_ts > _REAPER_IDLE_SECS: + _reaped = 0 + _reap_ids = [] + for _rid, _req in list(self.requests.items()): + if RequestStatus.is_finished(_req.status): + _reap_ids.append(_rid) + for _rid in _reap_ids: + try: + _req = self.requests[_rid] + self._free_blocks(_req) + _reaped += 1 + except Exception as _e: + logger.debug("[KV-REAPER] free_blocks failed for %s: %s", _rid, _e) + if _reaped > 0: + logger.warning( + "[KV-REAPER] Force-freed blocks for %d finished " + "requests after %.1fs idle", + _reaped, _time.monotonic() - self._idle_kv_reaper_ts) + self._idle_kv_reaper_ts = _time.monotonic() + else: + self._idle_kv_reaper_active = False""" + + if send_handler in src: + src = src.replace(send_handler, reaper_logic, 1) + else: + print("[SETUP] WARN: send handler not found for reaper injection") + sys.exit(0) + + open(f, "w").write(src) + print("[SETUP] Patched: idle KV block reaper for prefill") + +except Exception as e: + print(f"[SETUP] WARN patch idle-kv-reaper: {e}", file=sys.stderr) +' + _SETUP_INSTALLED+=("idle-kv-reaper") +} + # --------------------------------------------------------------------------- # SGLang: Patch aiter gluon pa_mqa_logits — fix 2D → 3D instr_shape for # Triton ≥ 3.5. @@ -189,9 +785,86 @@ install_transformers_glm5() { # Run installers (engine-gated) # ============================================================================= +install_lmcache_nixl() { + # LMCache (P2P/PD) + NIXL/RIXL data plane, only for the lmcache-nixl connector. + # No-op otherwise. Skips install when already importable, so a prebuilt image + # (e.g. one bundling lmcache+rixl) is used as-is. + case "${PREFILL_KV_CONNECTOR:-}:${DECODE_KV_CONNECTOR:-}" in + *lmcache-nixl*) ;; + *) return 0 ;; + esac + + if python3 -c "import lmcache" 2>/dev/null; then + echo "[SETUP] lmcache already present" + else + # docs.lmcache.ai/mp/disaggregated_prefill: the MultiConnector NIXL P/D + # path needs the lmcache[nixl] extra (pulls nixl>=1.3.0). + echo "[SETUP] Installing LMCache (lmcache[nixl]) for the NIXL P/D path..." + pip install --quiet "${LMCACHE_PIP_SPEC:-lmcache[nixl]}" || echo "[SETUP] WARN: lmcache pip install failed" + python3 -c "import lmcache" 2>/dev/null && _SETUP_INSTALLED+=("lmcache") + fi + + # LMCache 0.5.x declares CUDA wheels such as cupy-cuda13x in its base + # dependencies. Those import on ROCm but fail at runtime when LMCacheMP + # registers vLLM KV cache tensors. Replace them with the ROCm CuPy wheel. + if python3 - <<'PY' 2>/dev/null +import torch +raise SystemExit(0 if getattr(torch.version, "hip", None) else 1) +PY + then + local cupy_rocm_spec="${LMCACHE_ROCM_CUPY_SPEC:-cupy-rocm-7-0}" + if python3 - <<'PY' 2>/dev/null +import cupy +runtime = cupy.cuda.runtime +is_hip = bool(getattr(runtime, "is_hip", False)) +raise SystemExit(0 if is_hip else 1) +PY + then + echo "[SETUP] ROCm CuPy already present" + else + echo "[SETUP] Replacing CUDA CuPy with ${cupy_rocm_spec} for ROCm LMCache..." + pip uninstall -y cupy-cuda13x cupy-cuda12x cupy-cuda11x cupy >/dev/null 2>&1 || true + pip install --quiet --no-cache-dir "${cupy_rocm_spec}" || echo "[SETUP] WARN: ${cupy_rocm_spec} install failed" + fi + if python3 - <<'PY' 2>/dev/null +import cupy +runtime = cupy.cuda.runtime +is_hip = bool(getattr(runtime, "is_hip", False)) +print(f"[SETUP] cupy={cupy.__version__} is_hip={is_hip}") +raise SystemExit(0 if is_hip else 1) +PY + then + _SETUP_INSTALLED+=("cupy-rocm") + else + echo "[SETUP] WARN: ROCm CuPy smoke check failed; LMCacheMP may fail during KV cache registration" + fi + fi + + # NIXL python bindings. On ROCm the RIXL fork is expected prebuilt at + # $RIXL_HOME; only fall back to a pip nixl when neither is importable. + if python3 -c "import rixl" 2>/dev/null; then + echo "[SETUP] rixl (ROCm NIXL) already present" + elif python3 -c "import nixl" 2>/dev/null; then + echo "[SETUP] nixl already present" + else + echo "[SETUP] WARN: neither rixl nor nixl importable; the lmcache-nixl data plane needs one." + echo "[SETUP] WARN: For ROCm build RIXL into the image (RIXL_HOME=${RIXL_HOME}), or set LMCACHE_NIXL_PIP_SPEC." + if [[ -n "${LMCACHE_NIXL_PIP_SPEC:-}" ]]; then + pip install --quiet "${LMCACHE_NIXL_PIP_SPEC}" || echo "[SETUP] WARN: nixl pip install failed" + fi + fi +} + if [[ "$ENGINE" == "vllm-disagg" ]]; then install_recipe_deps install_amd_quark + install_lmcache_nixl + patch_moriio_save_kv_timeout + patch_moriio_transfer_timeout + patch_moriio_router_request_id_notify + patch_moriio_load_kv_timeout + patch_scheduler_read_mode_fix + patch_prefill_idle_kv_reaper # ========================================================================= # vLLM: Export UCX/RIXL paths (persists since this file is sourced) diff --git a/benchmarks/multi_node/amd_utils/simple_chat_template.jinja b/benchmarks/multi_node/amd_utils/simple_chat_template.jinja new file mode 100644 index 0000000000..5b81ef5a4b --- /dev/null +++ b/benchmarks/multi_node/amd_utils/simple_chat_template.jinja @@ -0,0 +1,2 @@ +{% for message in messages %}{{ message['role'] }}: {{ message['content'] }} +{% endfor %}{% if add_generation_prompt %}assistant: {% endif %} diff --git a/benchmarks/multi_node/amd_utils/submit.sh b/benchmarks/multi_node/amd_utils/submit.sh index 262fd25aa0..43f33f0f2c 100755 --- a/benchmarks/multi_node/amd_utils/submit.sh +++ b/benchmarks/multi_node/amd_utils/submit.sh @@ -146,6 +146,21 @@ export RESULT_FILENAME="${RESULT_FILENAME:-}" export SPEC_DECODING="${SPEC_DECODING:-}" export IS_MULTINODE="${IS_MULTINODE:-false}" +# Agentic / custom vLLM-disagg connector knobs (threaded to job.slurm -> Docker). +export IS_AGENTIC="${IS_AGENTIC:-0}" +export DURATION="${DURATION:-1800}" +export MODEL="${MODEL:-}" +export ROUTER_TYPE="${ROUTER_TYPE:-vllm-router}" +export ROUTER_PORT="${ROUTER_PORT:-30000}" +export ENABLE_PREFIX_CACHING="${ENABLE_PREFIX_CACHING:-}" +export MAX_MODEL_LEN="${MAX_MODEL_LEN:-}" +export MAX_NUM_SEQS="${MAX_NUM_SEQS:-}" +export WEKA_LOADER_OVERRIDE="${WEKA_LOADER_OVERRIDE:-}" +export VLLM_BIND_IP="${VLLM_BIND_IP:-}" +export SERVED_MODEL_NAME="${SERVED_MODEL_NAME:-}" +export PREFILL_KV_CONNECTOR="${PREFILL_KV_CONNECTOR:-moriio}" +export DECODE_KV_CONNECTOR="${DECODE_KV_CONNECTOR:-moriio}" + # Log directory: must be on NFS (shared filesystem) so the submit host can read SLURM output. export BENCHMARK_LOGS_DIR="${BENCHMARK_LOGS_DIR:-$(pwd)/benchmark_logs}" mkdir -p "$BENCHMARK_LOGS_DIR" diff --git a/benchmarks/multi_node/amd_utils/sync.py b/benchmarks/multi_node/amd_utils/sync.py index 96e94c1b08..049a484cf1 100755 --- a/benchmarks/multi_node/amd_utils/sync.py +++ b/benchmarks/multi_node/amd_utils/sync.py @@ -160,8 +160,9 @@ def close_port(): if args.enable_port: # Keep the port open long enough for slow nodes to pass their barrier. - # The previous 30s was too short when setup times vary by minutes. - grace = max(60, args.timeout // 2) if args.timeout > 0 else 300 + # Once all ports are observed as open, only a short grace period is + # needed for peers that are in the same polling cycle. + grace = 15 time.sleep(grace) close_port() diff --git a/configs/amd-master.yaml b/configs/amd-master.yaml index 4b1ea4aed3..d9533b301d 100644 --- a/configs/amd-master.yaml +++ b/configs/amd-master.yaml @@ -1238,6 +1238,59 @@ kimik2.5-fp4-mi355x-vllm-disagg: additional-settings: - "DECODE_NODES=2" + +# DCP + Expert-Parallel Kimi-K2.7-Code-MXFP4 disaggregated serving using +# vLLM MultiConnector[NixlConnector + LMCacheMPConnector]. P2P is off for this +# 1P1D CI point because there is only one prefill instance; LMCache still +# provides the local L2 cache with 1.2T capacity per instance. +kimik2.7-code-fp4-mi355x-vllm-disagg-dep-lmcache-nixl-p2poff: + image: vllm/vllm-openai-rocm:nightly + model: amd/Kimi-K2.7-Code-MXFP4 + model-prefix: kimik2.7-code + runner: cluster:mi355x-amds + precision: fp4 + framework: vllm-disagg + router: { name: vllm-router, version: "nightly" } + kv-p2p-transfer: nixl + multinode: true + disagg: true + scenarios: + agentic-coding: + - dram-utilization: 0.80 + search-space: + - spec-decoding: "none" + conc-list: [ 16 ] + kv-offloading: dram + kv-offload-backend: { name: lmcache } + prefill: + num-worker: 1 + tp: 8 + ep: 8 + dp-attn: true + additional-settings: + - "PREFILL_NODES=1" + - "PREFILL_KV_CONNECTOR=lmcache-nixl" + - "LMCACHE_ENABLE_P2P=false" + - "LMCACHE_L1_SIZE_GB=1200" + - "LMCACHE_L1_INIT_SIZE_GB=20" + - "LMCACHE_L1_READ_TTL_SECONDS=7200" + - "DOCKER_PULL_POLICY=always" + - "MODEL_NAME=Kimi-K2.7-Code-MXFP4-DEP" + decode: + num-worker: 1 + tp: 8 + ep: 8 + dp-attn: true + additional-settings: + - "DECODE_NODES=1" + - "DECODE_KV_CONNECTOR=lmcache-nixl" + - "LMCACHE_ENABLE_P2P=false" + - "LMCACHE_L1_SIZE_GB=1200" + - "LMCACHE_L1_INIT_SIZE_GB=20" + - "LMCACHE_L1_READ_TTL_SECONDS=7200" + - "DOCKER_PULL_POLICY=always" + - "MODEL_NAME=Kimi-K2.7-Code-MXFP4-DEP" + dsr1-fp4-mi355x-sglang-disagg: image: lmsysorg/sglang-rocm:v0.5.12-rocm720-mi35x-20260519 model: amd/DeepSeek-R1-0528-MXFP4-v2 diff --git a/perf-changelog.yaml b/perf-changelog.yaml index 57d510dd15..fd33c2fb12 100644 --- a/perf-changelog.yaml +++ b/perf-changelog.yaml @@ -4750,3 +4750,12 @@ - "Image: lmsysorg/sglang:nightly-dev-cu13-20260709-074bb928" - "6 topologies across 1k/1k and 8k/1k: 1P1D TP4 STP + wide-EP (DEP4 prefill / DEP16 decode) from 1P1D up to 8P1D, recipes under benchmarks/multi_node/srt-slurm-recipes/sglang/qwen3.5/gb300-fp8/" pr-link: https://github.com/SemiAnalysisAI/InferenceX/pull/2137 + +- config-keys: + - kimik2.7-code-fp4-mi355x-vllm-disagg-dep-lmcache-nixl-p2poff + scenario-type: + - agentic-coding + description: + - "Add one MI355X 1P1D Kimi-K2.7-Code-MXFP4 DEP vLLM disagg CI point using MultiConnector[NixlConnector + LMCacheMPConnector]." + - "Configure LMCache 1.2T local L2 on prefill and decode with LMCACHE_ENABLE_P2P=false for the single-prefill p2p-off bring-up path." + pr-link: https://github.com/SemiAnalysisAI/InferenceX/pull/2171