From 6a95aa6b104ed4d33792e3a53d2e5ca8df4dfeed Mon Sep 17 00:00:00 2001 From: Duyi-Wang Date: Wed, 16 Sep 2026 13:28:59 +0000 Subject: [PATCH 01/12] feat: enable DSpark for DeepSeek-V4-Pro-0813 AgentX disagg on MI355X MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Add a DSpark arm to the AMD multi-node SGLang disaggregated path and switch dsv4-fp4-mi355x-sglang-disagg-agentic-hicache-mtp onto it. models.yaml: DeepSeek-V4-Pro-AgentX gains dspark_flags alongside (not replacing) mtp_flags, so recipes that stay on spec-decoding: mtp keep the EAGLE arm untouched. dp_flags gains --enable-dp-lm-head, which SGLang requires for DSpark under DP attention and which is harmless for EAGLE. server_sglang.sh: the config loader exports MODEL_DSPARK_FLAGS, and build_server_config branches on SPEC_DECODING == draft_model. DECODE_MTP_SIZE carries the draft length for both algorithms but they spend it differently: EAGLE runs that many sequential draft passes, while DSpark emits a whole block in one pass, so num-steps is pinned to 1 and the draft length becomes --speculative-dspark-block-size. The verify window (num-draft-tokens = draft length + 1) is identical, which is why the MORI decode dispatch scaling by (DECODE_MTP_SIZE + 1) needs no change. A model configured with draft_model but no dspark_flags now fails hard instead of silently falling back to EAGLE. amd-master.yaml: all four arms move to spec-decoding: draft_model at DECODE_MTP_SIZE=3 (golden AL 3.01), and the image moves to v0.5.19-rocm720-mi35x-20260913, the first tag carrying the DSpark optimizations. CLIENT_IMAGE stays at 20260907; it only runs the load generator. 在 AMD 多节点 SGLang 分离式路径上启用 DSpark,并将 dsv4-fp4-mi355x-sglang-disagg-agentic-hicache-mtp 切换过去。 models.yaml:DeepSeek-V4-Pro-AgentX 新增 dspark_flags,与 mtp_flags 并列而非 替换,因此仍使用 spec-decoding: mtp 的配方的 EAGLE 分支不受影响。dp_flags 增加 --enable-dp-lm-head,这是 SGLang 在 DP attention 下运行 DSpark 的必需项,对 EAGLE 无副作用。 server_sglang.sh:配置加载器新导出 MODEL_DSPARK_FLAGS,build_server_config 按 SPEC_DECODING == draft_model 分流。DECODE_MTP_SIZE 对两种算法都表示草稿长度, 但花法不同:EAGLE 顺序执行同样次数的草稿前向,而 DSpark 一次前向产出一整块, 因此 num-steps 固定为 1,草稿长度改为 --speculative-dspark-block-size。验证窗口 (num-draft-tokens = 草稿长度 + 1)两者相同,这正是 MORI 解码 dispatch 的 ×(DECODE_MTP_SIZE + 1) 缩放无需改动的原因。若模型配置了 draft_model 却没有 dspark_flags,现在直接硬失败,而不是静默退回 EAGLE。 amd-master.yaml:四条臂全部切到 spec-decoding: draft_model,DECODE_MTP_SIZE=3 (黄金 AL 3.01),镜像升级到 v0.5.19-rocm720-mi35x-20260913,即首个带 DSpark 优化的标签。CLIENT_IMAGE 保持 20260907,该容器只运行负载生成器。 --- benchmarks/multi_node/amd_utils/models.yaml | 21 ++++++++++++--- .../multi_node/amd_utils/server_sglang.sh | 27 +++++++++++++++++-- configs/amd-master.yaml | 14 ++++++---- 3 files changed, 52 insertions(+), 10 deletions(-) diff --git a/benchmarks/multi_node/amd_utils/models.yaml b/benchmarks/multi_node/amd_utils/models.yaml index f8507cf0f0..53666ce487 100644 --- a/benchmarks/multi_node/amd_utils/models.yaml +++ b/benchmarks/multi_node/amd_utils/models.yaml @@ -354,9 +354,21 @@ DeepSeek-R1-0528-MXFP4-v2: DeepSeek-V4-Pro-AgentX: &DeepSeek-V4-Pro-AgentX base_flags: "--enable-deepseek-v4-fp4-indexer --watchdog-timeout 3600 --load-balance-method round_robin --kv-cache-dtype fp8_e4m3 --attention-backend dsv4 --page-size 256 --swa-full-tokens-ratio 0.1 --enforce-shared-experts-fusion --tool-call-parser deepseekv4 --reasoning-parser deepseek-v4 --disaggregation-transfer-backend mori --tokenizer-worker-num 8 --stream-interval 20 --log-level info --log-level-http error" - dp_flags: "--enable-dp-attention --swa-full-tokens-ratio 0.15 --enable-dp-attention-local-control-broadcast" + # --enable-dp-lm-head is required by SGLang for DSpark under DP attention; it + # is harmless for the EAGLE/MTP arms, so it stays unconditional here rather + # than needing a second DP flag string. + dp_flags: "--enable-dp-attention --enable-dp-lm-head --swa-full-tokens-ratio 0.15 --enable-dp-attention-local-control-broadcast" ep_flags: "--ep-dispatch-algorithm fake --moe-a2a-backend mori --deepep-mode normal" mtp_flags: "--speculative-algorithm EAGLE --speculative-eagle-topk 1" + # DSpark draft head, selected when the sweep sets spec-decoding: draft_model. + # Kept alongside mtp_flags rather than replacing it, so recipes that stay on + # spec-decoding: mtp keep the EAGLE arm untouched. The -0813 checkpoint + # bundles the draft head (dspark_block_size / dspark_markov_rank / + # dspark_target_layer_ids in config.json), so --speculative-draft-model-path + # defaults to --model-path and no separate draft checkpoint is needed. + # server_sglang.sh appends the block size and the verify window from + # DECODE_MTP_SIZE (= gamma); unlike EAGLE, num-steps is pinned to 1. + dspark_flags: "--speculative-algorithm DSPARK --speculative-eagle-topk 1" prefill: disable_radix_cache: false disable_cuda_graph: true @@ -384,8 +396,11 @@ DeepSeek-V4-Pro-AgentX: &DeepSeek-V4-Pro-AgentX max_running_requests: "BENCH_MAX_CONC_VALUE*2" cuda_graph_bs_range: "1-BENCH_MAX_CONC_VALUE*2" -# Pro-0813 retains EAGLE 3-1-4 for PD compatibility. Synthetic acceptance is -# checkpoint-specific in server_sglang.sh: thinking-on, length 3 uses AL 3.01. +# Pro-0813 serves the PD path with DSPARK (spec-decoding: draft_model), which +# measured clean across c4-c256; the EAGLE 3-1-4 arm remains available through +# spec-decoding: mtp. Synthetic acceptance is checkpoint-specific in +# server_sglang.sh and does not depend on which of the two runs: thinking-on, +# draft length 3 uses AL 3.01. DeepSeek-V4-Pro-0813-AgentX: *DeepSeek-V4-Pro-AgentX DeepSeek-V4-Pro-DI: diff --git a/benchmarks/multi_node/amd_utils/server_sglang.sh b/benchmarks/multi_node/amd_utils/server_sglang.sh index b44ea1a11f..672a13f9bf 100755 --- a/benchmarks/multi_node/amd_utils/server_sglang.sh +++ b/benchmarks/multi_node/amd_utils/server_sglang.sh @@ -100,6 +100,7 @@ def parse_range(cuda_range, default_start, default_end): # Output shell variables print(f'MODEL_BASE_FLAGS=\"{m.get(\"base_flags\", \"\")}\"') print(f'MODEL_MTP_FLAGS=\"{m.get(\"mtp_flags\", \"\")}\"') +print(f'MODEL_DSPARK_FLAGS=\"{m.get(\"dspark_flags\", \"\")}\"') print(f'MODEL_DP_FLAGS=\"{m.get(\"dp_flags\", \"\")}\"') print(f'MODEL_EP_FLAGS=\"{m.get(\"ep_flags\", \"\")}\"') @@ -381,8 +382,30 @@ build_server_config() { local ep_config="" local specific_config="" + # Speculative-decoding config (only if a draft length is set). + # + # DECODE_MTP_SIZE carries the draft length for BOTH algorithms, but the two + # spend it differently: + # EAGLE/MTP -- num-steps = draft length, i.e. that many sequential draft + # forward passes, each producing one token. + # DSPARK -- one draft pass emits a whole block, so num-steps is + # pinned to 1 and the draft length becomes the block size + # (gamma). Passing gamma as num-steps here would ask for + # gamma sequential DSpark passes instead of one gamma-token + # block. + # The verify window (num-draft-tokens = draft length + 1) is the same for + # both, which is also what makes the MORI decode dispatch scaling + # (x (DECODE_MTP_SIZE + 1)) correct for DSpark without further change. if [ "$decode_mtp_size" -gt 0 ]; then - mtp_config="${MODEL_MTP_FLAGS} --speculative-num-steps ${decode_mtp_size} --speculative-num-draft-tokens $((decode_mtp_size + 1))" + if [[ "${SPEC_DECODING:-}" == "draft_model" ]]; then + if [[ -z "${MODEL_DSPARK_FLAGS// }" ]]; then + echo "FATAL: SPEC_DECODING=draft_model but model '${model_name}' has no dspark_flags in models.yaml." >&2 + exit 1 + fi + mtp_config="${MODEL_DSPARK_FLAGS} --speculative-dspark-block-size ${decode_mtp_size} --speculative-num-steps 1 --speculative-num-draft-tokens $((decode_mtp_size + 1))" + else + mtp_config="${MODEL_MTP_FLAGS} --speculative-num-steps ${decode_mtp_size} --speculative-num-draft-tokens $((decode_mtp_size + 1))" + fi fi if [[ "$enable_dp" == "true" ]]; then @@ -1393,7 +1416,7 @@ else if [[ -n "$DSV4_GOLDEN_AL" ]]; then DECODE_SIM_ACC_ENV="SGLANG_SIMULATE_ACC_LEN=${DSV4_GOLDEN_AL} SGLANG_SIMULATE_ACC_METHOD=match-expected SGLANG_SIMULATE_ACC_TOKEN_MODE=real-draft-token" else - echo "WARNING: agentic MTP run (model=${MODEL_NAME}, DECODE_MTP_SIZE=${DECODE_MTP_SIZE}) has no golden AL wired in server_sglang.sh -- falling back to real (unsimulated, non-representative) acceptance. Add a case in server_sglang.sh and golden_al_distribution/ before shipping this arm. See golden_al_distribution/README.md." >&2 + echo "WARNING: agentic spec-decoding run (model=${MODEL_NAME}, algorithm=${SPEC_DECODING:-mtp}, DECODE_MTP_SIZE=${DECODE_MTP_SIZE}) has no golden AL wired in server_sglang.sh -- falling back to real (unsimulated, non-representative) acceptance. Add a case in server_sglang.sh and golden_al_distribution/ before shipping this arm. See golden_al_distribution/README.md." >&2 fi fi fi diff --git a/configs/amd-master.yaml b/configs/amd-master.yaml index e1d48355a8..a4ea77306f 100644 --- a/configs/amd-master.yaml +++ b/configs/amd-master.yaml @@ -1194,7 +1194,11 @@ minimaxm3-fp8-mi325x-vllm-agentic-mtp: - { tp: 8, spec-decoding: mtp, kv-offloading: none, conc-list: [1, 2, 4, 8, 10, 12, 14, 16, 18] } dsv4-fp4-mi355x-sglang-disagg-agentic-hicache-mtp: - image: lmsysorg/sglang-rocm:v0.5.19-rocm720-mi35x-20260911 + # 20260913 rather than 20260911: it is the first tag carrying the DSpark + # optimizations, which is the whole point of this arm. CLIENT_IMAGE stays at + # 20260907 -- that container only runs the load generator, so the + # server-side DSpark work does not reach it. + image: lmsysorg/sglang-rocm:v0.5.19-rocm720-mi35x-20260913 model: deepseek-ai/DeepSeek-V4-Pro-0813 model-prefix: dsv4 runner: cluster:mi355x-amds @@ -1207,7 +1211,7 @@ dsv4-fp4-mi355x-sglang-disagg-agentic-hicache-mtp: agentic-coding: - dram-utilization: 0.80 search-space: - - spec-decoding: "mtp" + - spec-decoding: "draft_model" conc-list: [ 4 ] kv-offloading: none prefill: @@ -1226,7 +1230,7 @@ dsv4-fp4-mi355x-sglang-disagg-agentic-hicache-mtp: additional-settings: - "DECODE_NODES=1" - "DECODE_MTP_SIZE=3" - - spec-decoding: "mtp" + - spec-decoding: "draft_model" conc-list: [ 16 ] kv-offloading: none prefill: @@ -1245,7 +1249,7 @@ dsv4-fp4-mi355x-sglang-disagg-agentic-hicache-mtp: additional-settings: - "DECODE_NODES=1" - "DECODE_MTP_SIZE=3" - - spec-decoding: "mtp" + - spec-decoding: "draft_model" conc-list: [ 32, 48 ] kv-offloading: dram kv-offload-backend: { name: hicache } @@ -1266,7 +1270,7 @@ dsv4-fp4-mi355x-sglang-disagg-agentic-hicache-mtp: additional-settings: - "DECODE_NODES=1" - "DECODE_MTP_SIZE=3" - - spec-decoding: "mtp" + - spec-decoding: "draft_model" conc-list: [ 128, 192, 256 ] kv-offloading: dram kv-offload-backend: { name: umbp-linker } From 68df08efb8aa5a55d5677fbffbda50d0ee1490e9 Mon Sep 17 00:00:00 2001 From: Duyi-Wang Date: Wed, 16 Sep 2026 13:31:01 +0000 Subject: [PATCH 02/12] chore: record the DSpark switch in perf-changelog MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Append the perf-changelog entry for dsv4-fp4-mi355x-sglang-disagg-agentic-hicache-mtp moving to spec-decoding: draft_model on image -20260913. 为 dsv4-fp4-mi355x-sglang-disagg-agentic-hicache-mtp 切换到 spec-decoding: draft_model 及镜像 -20260913 追加 perf-changelog 条目。 --- perf-changelog.yaml | 12 ++++++++++++ 1 file changed, 12 insertions(+) diff --git a/perf-changelog.yaml b/perf-changelog.yaml index 8d11fa35b6..00bda14b41 100644 --- a/perf-changelog.yaml +++ b/perf-changelog.yaml @@ -7923,3 +7923,15 @@ - "Bump image to lmsysorg/sglang-rocm:v0.5.19-rocm720-mi35x-20260915." - "Switch HiCache defaults to --hicache-io-backend kernel and --hicache-mem-layout page_first (from direct / page_first_direct). Ratio 1.5 and write_through are unchanged." pr-link: https://github.com/SemiAnalysisAI/InferenceX/pull/3118 + +- config-keys: + - dsv4-fp4-mi355x-sglang-disagg-agentic-hicache-mtp + scenario-type: + - agentic-coding + description: + - "Switch all four arms from EAGLE (spec-decoding: mtp) to DSpark (spec-decoding: draft_model) at DECODE_MTP_SIZE=3, and bump the image from lmsysorg/sglang-rocm:v0.5.19-rocm720-mi35x-20260911 to -20260913, the first tag carrying the DSpark optimizations. CLIENT_IMAGE stays at 20260907 (load generator only). Topology, concurrency list, HiCache/UMBP settings, and precision are unchanged." + - "DeepSeek-V4-Pro-AgentX in models.yaml gains dspark_flags (--speculative-algorithm DSPARK --speculative-eagle-topk 1) alongside mtp_flags rather than replacing it, so recipes remaining on spec-decoding: mtp keep their EAGLE arm unchanged; dp_flags gains --enable-dp-lm-head, required by SGLang for DSpark under DP attention and inert for EAGLE. The -0813 checkpoint bundles the draft head, so no separate draft checkpoint is staged." + - "server_sglang.sh builds DSpark flags as --speculative-dspark-block-size DECODE_MTP_SIZE --speculative-num-steps 1 --speculative-num-draft-tokens (DECODE_MTP_SIZE + 1). DECODE_MTP_SIZE is the draft length for both algorithms but EAGLE spends it as sequential draft passes while DSpark emits one block, so num-steps is pinned to 1. The verify window is identical for both, so the MORI decode dispatch scaling by (DECODE_MTP_SIZE + 1) is unchanged. A model set to draft_model without dspark_flags now fails hard instead of silently serving EAGLE." + - "Golden AL is unchanged: the AgentX curve is keyed on checkpoint, thinking mode, and draft length, not on the algorithm, so DeepSeek-V4-Pro-0813 at draft length 3 keeps AL 3.01 from golden_al_distribution/dsv4-pro-0813-dspark.yaml (thinking_on)." + - "将四条臂全部从 EAGLE(spec-decoding: mtp)切换为 DSpark(spec-decoding: draft_model),DECODE_MTP_SIZE=3;镜像由 lmsysorg/sglang-rocm:v0.5.19-rocm720-mi35x-20260911 升级到 -20260913(首个带 DSpark 优化的标签),CLIENT_IMAGE 保持 20260907(仅负载生成器)。拓扑、并发列表、HiCache/UMBP 设置与精度均不变。models.yaml 中 DeepSeek-V4-Pro-AgentX 新增 dspark_flags(与 mtp_flags 并列而非替换,因此仍用 mtp 的配方 EAGLE 分支不受影响),dp_flags 增加 --enable-dp-lm-head。server_sglang.sh 按 SPEC_DECODING == draft_model 分流:DSpark 一次前向出一整块,故 num-steps 钉死为 1,草稿长度转为 --speculative-dspark-block-size;验证窗口(长度+1)与 EAGLE 相同,所以 MORI 解码 dispatch 的 ×(DECODE_MTP_SIZE+1) 缩放无需改动。黄金 AL 不变:按检查点而非算法选表,DeepSeek-V4-Pro-0813 在草稿长度 3 仍为 AL 3.01。" + pr-link: https://github.com/SemiAnalysisAI/InferenceX/pull/3188 From c68184dc459ed7ac72e3b3e5939b3dea34f55461 Mon Sep 17 00:00:00 2001 From: Duyi-Wang Date: Wed, 16 Sep 2026 13:47:16 +0000 Subject: [PATCH 03/12] refactor: rename the MI355X DSV4 disagg AgentX key to -umbp-dspark MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Rename dsv4-fp4-mi355x-sglang-disagg-agentic-hicache-mtp to dsv4-fp4-mi355x-sglang-disagg-agentic-umbp-dspark per review, so the key names the algorithm the arm now serves (DSpark) instead of the retired MTP label, matching the dsv41flash-*-agentic-dspark precedent. The rename touches only the config key. The emitted matrix is byte-identical to the pre-rename run: 7 points, same exp-names (derived from model-prefix and topology, not the key), image, spec-decoding, topology, concurrency, and offload backends. Historical perf-changelog entries keep the old key, which the append-only contract requires; the new entry records the rename so the lineage stays traceable. 按评审意见,将 dsv4-fp4-mi355x-sglang-disagg-agentic-hicache-mtp 改名为 dsv4-fp4-mi355x-sglang-disagg-agentic-umbp-dspark,使配置键反映该分支实际 使用的算法(DSpark),而非已停用的 MTP 标签,与 dsv41flash-*-agentic-dspark 的命名惯例一致。 改名仅涉及配置键。生成的 matrix 与改名前逐字节一致:7 个点,exp-name (由 model-prefix 与拓扑推导,与键名无关)、镜像、spec-decoding、拓扑、 并发列表、offload 后端均不变。perf-changelog 的历史条目保留旧键名(append-only 契约的要求),新条目记录了此次改名以保持谱系可追溯。 --- configs/amd-master.yaml | 5 ++++- perf-changelog.yaml | 4 +++- 2 files changed, 7 insertions(+), 2 deletions(-) diff --git a/configs/amd-master.yaml b/configs/amd-master.yaml index a4ea77306f..f7895ef786 100644 --- a/configs/amd-master.yaml +++ b/configs/amd-master.yaml @@ -1193,7 +1193,10 @@ minimaxm3-fp8-mi325x-vllm-agentic-mtp: search-space: - { tp: 8, spec-decoding: mtp, kv-offloading: none, conc-list: [1, 2, 4, 8, 10, 12, 14, 16, 18] } -dsv4-fp4-mi355x-sglang-disagg-agentic-hicache-mtp: +dsv4-fp4-mi355x-sglang-disagg-agentic-umbp-dspark: + # Renamed from dsv4-fp4-mi355x-sglang-disagg-agentic-hicache-mtp when this arm + # moved from EAGLE/MTP to DSpark; earlier perf-changelog entries are recorded + # under the old key. # 20260913 rather than 20260911: it is the first tag carrying the DSpark # optimizations, which is the whole point of this arm. CLIENT_IMAGE stays at # 20260907 -- that container only runs the load generator, so the diff --git a/perf-changelog.yaml b/perf-changelog.yaml index 00bda14b41..ab49a74f53 100644 --- a/perf-changelog.yaml +++ b/perf-changelog.yaml @@ -7925,7 +7925,7 @@ pr-link: https://github.com/SemiAnalysisAI/InferenceX/pull/3118 - config-keys: - - dsv4-fp4-mi355x-sglang-disagg-agentic-hicache-mtp + - dsv4-fp4-mi355x-sglang-disagg-agentic-umbp-dspark scenario-type: - agentic-coding description: @@ -7934,4 +7934,6 @@ - "server_sglang.sh builds DSpark flags as --speculative-dspark-block-size DECODE_MTP_SIZE --speculative-num-steps 1 --speculative-num-draft-tokens (DECODE_MTP_SIZE + 1). DECODE_MTP_SIZE is the draft length for both algorithms but EAGLE spends it as sequential draft passes while DSpark emits one block, so num-steps is pinned to 1. The verify window is identical for both, so the MORI decode dispatch scaling by (DECODE_MTP_SIZE + 1) is unchanged. A model set to draft_model without dspark_flags now fails hard instead of silently serving EAGLE." - "Golden AL is unchanged: the AgentX curve is keyed on checkpoint, thinking mode, and draft length, not on the algorithm, so DeepSeek-V4-Pro-0813 at draft length 3 keeps AL 3.01 from golden_al_distribution/dsv4-pro-0813-dspark.yaml (thinking_on)." - "将四条臂全部从 EAGLE(spec-decoding: mtp)切换为 DSpark(spec-decoding: draft_model),DECODE_MTP_SIZE=3;镜像由 lmsysorg/sglang-rocm:v0.5.19-rocm720-mi35x-20260911 升级到 -20260913(首个带 DSpark 优化的标签),CLIENT_IMAGE 保持 20260907(仅负载生成器)。拓扑、并发列表、HiCache/UMBP 设置与精度均不变。models.yaml 中 DeepSeek-V4-Pro-AgentX 新增 dspark_flags(与 mtp_flags 并列而非替换,因此仍用 mtp 的配方 EAGLE 分支不受影响),dp_flags 增加 --enable-dp-lm-head。server_sglang.sh 按 SPEC_DECODING == draft_model 分流:DSpark 一次前向出一整块,故 num-steps 钉死为 1,草稿长度转为 --speculative-dspark-block-size;验证窗口(长度+1)与 EAGLE 相同,所以 MORI 解码 dispatch 的 ×(DECODE_MTP_SIZE+1) 缩放无需改动。黄金 AL 不变:按检查点而非算法选表,DeepSeek-V4-Pro-0813 在草稿长度 3 仍为 AL 3.01。" + - "Rename the config key from dsv4-fp4-mi355x-sglang-disagg-agentic-hicache-mtp to dsv4-fp4-mi355x-sglang-disagg-agentic-umbp-dspark to reflect the DSpark algorithm and the UMBP-linker arm; earlier entries for this recipe are recorded under the old key. No other field changes with the rename." + - "配置键由 dsv4-fp4-mi355x-sglang-disagg-agentic-hicache-mtp 改名为 dsv4-fp4-mi355x-sglang-disagg-agentic-umbp-dspark,以反映 DSpark 算法与 UMBP-linker 分支;该配方此前的条目仍记录在旧键名下。改名不伴随任何其他字段变化。" pr-link: https://github.com/SemiAnalysisAI/InferenceX/pull/3188 From 64bf729ccb0b47307e817438574177dd7c88e075 Mon Sep 17 00:00:00 2001 From: Duyi-Wang Date: Wed, 16 Sep 2026 13:51:13 +0000 Subject: [PATCH 04/12] fix: abort the launch when draft_model is set without dspark_flags MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The guard lived inside build_server_config, which is only ever invoked as PREFILL_SERVER_CONFIG=$(build_server_config ...). exit there terminates the command-substitution subshell, and with no set -e and no status check the script carried on with an empty config string -- dropping base, ep, dp, and parallel flags for that mode, a worse failure than the silent EAGLE fallback the guard was meant to prevent. Move the check to top level, immediately before both build_server_config calls, where exit 1 actually stops the script. The in-function branch keeps a comment explaining why the validation is not co-located with the flag it guards. Verified: with SPEC_DECODING=draft_model and empty MODEL_DSPARK_FLAGS the old shape exits 0 and continues with an empty config, the new shape exits 1 before either assignment. Populated dspark_flags, spec-decoding mtp, empty SPEC_DECODING, and DECODE_MTP_SIZE=0 all still pass through unchanged. 原先的守卫位于 build_server_config 内,而该函数只以 PREFILL_SERVER_CONFIG=$(build_server_config ...) 的形式调用。子 shell 里的 exit 只会终止命令替换本身;脚本没有 set -e,也没有检查退出码,因此会带着空配置字符串 继续执行——该模式下的 base、ep、dp、parallel 参数全部丢失,比守卫本要防止的 "静默退回 EAGLE" 更糟。 将检查移到顶层、两次 build_server_config 调用之前,此处 exit 1 才真正终止脚本。 函数内分支保留注释,说明校验为何没有与它所保护的 flag 放在一起。 已验证:当 SPEC_DECODING=draft_model 且 MODEL_DSPARK_FLAGS 为空时,旧写法退出码 为 0 并带空配置继续,新写法在两次赋值之前即以退出码 1 终止。dspark_flags 非空、 spec-decoding 为 mtp、SPEC_DECODING 为空、以及 DECODE_MTP_SIZE=0 四种情况均保持 原有行为。 --- .../multi_node/amd_utils/server_sglang.sh | 17 +++++++++++++---- 1 file changed, 13 insertions(+), 4 deletions(-) diff --git a/benchmarks/multi_node/amd_utils/server_sglang.sh b/benchmarks/multi_node/amd_utils/server_sglang.sh index 672a13f9bf..8a25249b61 100755 --- a/benchmarks/multi_node/amd_utils/server_sglang.sh +++ b/benchmarks/multi_node/amd_utils/server_sglang.sh @@ -398,10 +398,10 @@ build_server_config() { # (x (DECODE_MTP_SIZE + 1)) correct for DSpark without further change. if [ "$decode_mtp_size" -gt 0 ]; then if [[ "${SPEC_DECODING:-}" == "draft_model" ]]; then - if [[ -z "${MODEL_DSPARK_FLAGS// }" ]]; then - echo "FATAL: SPEC_DECODING=draft_model but model '${model_name}' has no dspark_flags in models.yaml." >&2 - exit 1 - fi + # MODEL_DSPARK_FLAGS is validated at the call site, not here: this + # function is only ever invoked inside $( ), where an exit would + # terminate the subshell and leave the caller with an empty config + # rather than aborting the launch. mtp_config="${MODEL_DSPARK_FLAGS} --speculative-dspark-block-size ${decode_mtp_size} --speculative-num-steps 1 --speculative-num-draft-tokens $((decode_mtp_size + 1))" else mtp_config="${MODEL_MTP_FLAGS} --speculative-num-steps ${decode_mtp_size} --speculative-num-draft-tokens $((decode_mtp_size + 1))" @@ -462,6 +462,15 @@ build_server_config() { echo "$full_config" } +# Validate the DSpark path before building either config. This has to happen at +# top level: build_server_config only ever runs inside $( ), so an exit there +# would kill the subshell and hand the caller an empty config string instead of +# stopping the launch. +if [[ "$DECODE_MTP_SIZE" -gt 0 ]] && [[ "${SPEC_DECODING:-}" == "draft_model" ]] && [[ -z "${MODEL_DSPARK_FLAGS// }" ]]; then + echo "FATAL: SPEC_DECODING=draft_model but model '${MODEL_NAME}' has no dspark_flags in models.yaml." >&2 + exit 1 +fi + PREFILL_SERVER_CONFIG=$(build_server_config "prefill" "$MODEL_NAME" "$PREFILL_TP_SIZE" "$PREFILL_ENABLE_EP" "$PREFILL_ENABLE_DP" "$DECODE_MTP_SIZE") DECODE_SERVER_CONFIG=$(build_server_config "decode" "$MODEL_NAME" "$DECODE_TP_SIZE" "$DECODE_ENABLE_EP" "$DECODE_ENABLE_DP" "$DECODE_MTP_SIZE") From 0bf6ff208b98f0a9b2cbedf338e7a3bd6991ce6a Mon Sep 17 00:00:00 2001 From: Cam Quilici Date: Wed, 16 Sep 2026 11:53:16 -0500 Subject: [PATCH 05/12] fix(amd): forward required client inputs and collect all node logs --- benchmarks/multi_node/amd_utils/job.slurm | 27 ++++++ .../multi_node/amd_utils/server_sglang.sh | 6 +- .../multi_node/amd_utils/stage_node_logs.sh | 23 +++++ benchmarks/runtime_settings.sh | 2 +- docs/architecture.md | 2 +- docs/architecture_zh.md | 2 +- perf-changelog.yaml | 21 ++--- utils/test_amd_node_log_staging.py | 84 +++++++++++++++++++ 8 files changed, 152 insertions(+), 15 deletions(-) create mode 100755 benchmarks/multi_node/amd_utils/stage_node_logs.sh create mode 100644 utils/test_amd_node_log_staging.py diff --git a/benchmarks/multi_node/amd_utils/job.slurm b/benchmarks/multi_node/amd_utils/job.slurm index c12535776c..bb151f2039 100755 --- a/benchmarks/multi_node/amd_utils/job.slurm +++ b/benchmarks/multi_node/amd_utils/job.slurm @@ -775,6 +775,7 @@ echo \"[rank 0] Main container exited (rc=\$DOCKER_EXIT_CODE). Stopping vllm-rou \$DOCKER_CMD rm -f \"$ROUTER_CONT_NAME\" 2>/dev/null || true exit \$DOCKER_EXIT_CODE " +SERVER_SRUN_RC=$? if [[ "${KEEP_CONTAINERS}" != "1" ]]; then srun --nodelist="$SELECTED_NODELIST_SRUN" bash -c 'eval "$DOCKER_CMD_DETECT"; $DOCKER_CMD rm -f '"$DOCKER_CONT_NAME"' '"$CLIENT_CONT_NAME"' 2>/dev/null || true' @@ -786,3 +787,29 @@ if [[ "${KEEP_CONTAINERS}" != "1" ]]; then ' fi fi + +# /run_logs is backed by each compute node's local /tmp, so the node-0 copy +# performed by the engine launcher cannot see prefill/decode logs written on +# other nodes. Collect after the server step and container cleanup so failed +# runs also include shutdown output. KEEP_CONTAINERS=1 retains a snapshot of +# any containers left running for debugging. +# Use sudo because the container-created source and the existing node-0 +# destination can be root-owned. Restore ownership after the fan-in so a +# subsequent runner job can clean the workspace normally. +SHARED_JOB_LOGS="${BENCHMARK_LOGS_DIR}/logs/slurm_job-${SLURM_JOB_ID}" +if ! srun --nodelist="$SELECTED_NODELIST_SRUN" \ + --nodes="$NUM_NODES" --ntasks="$NUM_NODES" --ntasks-per-node=1 \ + bash "$DI_REPO_DIR/benchmarks/multi_node/amd_utils/stage_node_logs.sh" \ + "/tmp/slurm_job-${SLURM_JOB_ID}" "$SHARED_JOB_LOGS"; then + echo "[logs][ERROR] failed to stage logs from one or more Slurm nodes" >&2 + if [[ "$SERVER_SRUN_RC" -eq 0 ]]; then + SERVER_SRUN_RC=1 + fi +fi + +if [[ -d "$SHARED_JOB_LOGS" ]]; then + sudo chown -R "$(id -u):$(id -g)" "$SHARED_JOB_LOGS" 2>/dev/null || true + chmod -R a+rwX "$SHARED_JOB_LOGS" 2>/dev/null || true +fi + +exit "$SERVER_SRUN_RC" diff --git a/benchmarks/multi_node/amd_utils/server_sglang.sh b/benchmarks/multi_node/amd_utils/server_sglang.sh index 8a25249b61..b38035e73d 100755 --- a/benchmarks/multi_node/amd_utils/server_sglang.sh +++ b/benchmarks/multi_node/amd_utils/server_sglang.sh @@ -1215,6 +1215,8 @@ print(json.dumps(json.loads(sys.stdin.read())))' <<<"$_val")" || { # Must run from repo root so infx/evals/gsm8k.yaml resolves pushd /workspace + # Match the disaggregation router launched above. + export PORT=30000 source /workspace/benchmarks/benchmark_lib.sh # CONC must be exported before run_eval so meta_env.json matches validate_scores.py. @@ -1236,9 +1238,9 @@ print(json.dumps(json.loads(sys.stdin.read())))' <<<"$_val")" || { # arrive via Docker -e flags from job.slurm. if [[ "$DRY_RUN" -eq 1 ]]; then - echo "DRY RUN: run_eval --port 30000 (framework=${EVAL_FRAMEWORK}, conc=${EVAL_CONCURRENT_REQUESTS}, ctx=${EVAL_MAX_MODEL_LEN:-auto})" + echo "DRY RUN: run_eval --port ${PORT} (framework=${EVAL_FRAMEWORK}, conc=${EVAL_CONCURRENT_REQUESTS}, ctx=${EVAL_MAX_MODEL_LEN:-auto})" else - run_eval --port 30000 + run_eval --port "$PORT" eval_rc=$? if [[ $eval_rc -ne 0 ]]; then diff --git a/benchmarks/multi_node/amd_utils/stage_node_logs.sh b/benchmarks/multi_node/amd_utils/stage_node_logs.sh new file mode 100755 index 0000000000..f1e67814e4 --- /dev/null +++ b/benchmarks/multi_node/amd_utils/stage_node_logs.sh @@ -0,0 +1,23 @@ +#!/usr/bin/env bash + +set -eo pipefail + +if [[ $# -ne 2 ]]; then + echo "Usage: $0 " >&2 + exit 2 +fi + +SOURCE_LOGS=$1 +SHARED_LOGS=$2 + +if [[ ! -d "$SOURCE_LOGS" ]]; then + echo "[logs][ERROR] no node-local logs found on $(hostname): $SOURCE_LOGS" >&2 + exit 1 +fi + +# Server containers create the source tree as root, and node 0 may have already +# created the shared destination as root. The Slurm nodes provide passwordless +# sudo for the same Docker lifecycle used by job.slurm. +sudo mkdir -p "$SHARED_LOGS" +sudo cp -r "$SOURCE_LOGS"/. "$SHARED_LOGS"/ +echo "[logs] staged $(hostname):$SOURCE_LOGS -> $SHARED_LOGS" diff --git a/benchmarks/runtime_settings.sh b/benchmarks/runtime_settings.sh index 0b04efe9ba..a369898f8b 100644 --- a/benchmarks/runtime_settings.sh +++ b/benchmarks/runtime_settings.sh @@ -32,4 +32,4 @@ export SGLANG_TORCH_PROFILER_DIR='/workspace' export VLLM_TORCH_PROFILER_DIR='/workspace' # Explicitly forward these settings across container boundaries. -export INFERENCEX_RUNTIME_ENV_VARS="OPENAI_API_KEY SWEBENCH_EXPECTED_INSTANCES SWEBENCH_AGENT_STEP_LIMIT SWEBENCH_AGENT_TIMEOUT SWEBENCH_AGENT_EXIT_GRACE SWEBENCH_WATCHDOG_POLL SWEBENCH_SANDBOX_SWEEP SWEBENCH_SKIP_SCORE SWEBENCH_EVAL_TIMEOUT SWEBENCH_SCORE_TIMEOUT SWEBENCH_MAX_WORKERS EVAL_ENDPOINT_READY_TIMEOUT_SECONDS EVAL_MODEL_STABILIZATION_SECONDS AIPERF_FAILED_REQUEST_THRESHOLD AIPERF_LIVE_FAILED_REQUEST_THRESHOLD AIPERF_TRACE_IDLE_GAP_CAP_SECONDS AIPERF_PYTHON_VERSION AIPERF_WARMUP_REQUESTS_PER_LANE AIPERF_DATASET_WEKA_LIVE_ASSISTANT_RESPONSES AGENTIC_WARMUP_GRACE_PERIOD AIPERF_USE_DYNAMO_CONV_AWARE_ROUTING AIPERF_HTTP_X_DYNAMO_SESSION_ID_FROM_CORRELATION_ID AIPERF_DYNAMO_SESSION_TIMEOUT_SECONDS AIPERF_UNSAFE_OVERRIDE ENABLE_AGENTX_POWER VLLM_ENGINE_READY_TIMEOUT_S SGLANG_TORCH_PROFILER_DIR VLLM_TORCH_PROFILER_DIR" +export INFERENCEX_RUNTIME_ENV_VARS="OPENAI_API_KEY SWEBENCH_EXPECTED_INSTANCES SWEBENCH_AGENT_STEP_LIMIT SWEBENCH_AGENT_TIMEOUT SWEBENCH_AGENT_EXIT_GRACE SWEBENCH_WATCHDOG_POLL SWEBENCH_SANDBOX_SWEEP SWEBENCH_SKIP_SCORE SWEBENCH_EVAL_TIMEOUT SWEBENCH_SCORE_TIMEOUT SWEBENCH_MAX_WORKERS EVAL_ENDPOINT_READY_TIMEOUT_SECONDS EVAL_MODEL_STABILIZATION_SECONDS AIPERF_FAILED_REQUEST_THRESHOLD AIPERF_LIVE_FAILED_REQUEST_THRESHOLD AIPERF_TRACE_IDLE_GAP_CAP_SECONDS AIPERF_PYTHON_VERSION AIPERF_WARMUP_REQUESTS_PER_LANE AIPERF_DATASET_WEKA_LIVE_ASSISTANT_RESPONSES AGENTIC_WARMUP_GRACE_PERIOD AIPERF_USE_DYNAMO_CONV_AWARE_ROUTING AIPERF_HTTP_X_DYNAMO_SESSION_ID_FROM_CORRELATION_ID AIPERF_DYNAMO_SESSION_TIMEOUT_SECONDS AIPERF_EXPERIMENTAL_FAST AIPERF_UNSAFE_OVERRIDE ENABLE_AGENTX_POWER REQUIRE_POWER VLLM_ENGINE_READY_TIMEOUT_S SGLANG_TORCH_PROFILER_DIR VLLM_TORCH_PROFILER_DIR" diff --git a/docs/architecture.md b/docs/architecture.md index a462f7215b..9e4b5c743a 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -277,7 +277,7 @@ The collector and reusable-artifact validator share format recognition, concurre Agentic throughput jobs have a different contract. They validate AIPerf output with [`infx/results/agentic/validate_agentic_result.py`](../infx/results/agentic/validate_agentic_result.py), upload an aggregate `bmk_agentic_` artifact, and upload the raw `agentic_` sibling containing trace-replay material. InferenceX-app pairs those siblings by their shared suffix. Agentic eval-only jobs follow the eval output contract instead and do not require a throughput result. -Server logs and GPU metrics are diagnostic side artifacts. They are uploaded with `always()` so a failed run can still be investigated. Their presence does not turn a failed benchmark into a valid result. +Server logs and GPU metrics are diagnostic side artifacts. They are uploaded with `always()` so a failed run can still be investigated. Their presence does not turn a failed benchmark into a valid result. On the AMD Slurm fleet, `/run_logs` is node-local; after the server step finishes, `job.slurm` merges the closed log tree from every allocated node into shared storage so the diagnostic artifact includes prefill and decode logs from the full deployment. ## Stage 6: artifact collection and handoff diff --git a/docs/architecture_zh.md b/docs/architecture_zh.md index 3e8fde840f..377a62cba4 100644 --- a/docs/architecture_zh.md +++ b/docs/architecture_zh.md @@ -277,7 +277,7 @@ rows = build_rows(raw_eval, metadata, source="eval_job/results.json") 智能体吞吐量作业采用不同的契约。它们使用 [`infx/results/agentic/validate_agentic_result.py`](../infx/results/agentic/validate_agentic_result.py) 验证 AIPerf 输出,上传聚合的 `bmk_agentic_` 工件,并上传包含追踪重放材料的原始 `agentic_` 同级工件。InferenceX-app 通过它们共享的后缀对这些同级工件进行配对。智能体仅评测作业改为遵循评测输出契约,不要求吞吐量结果。 -服务器日志和 GPU 指标是诊断辅助工件。它们通过 `always()` 上传,因此失败的运行仍可供调查。它们的存在不会将失败的基准测试转变为有效结果。 +服务器日志和 GPU 指标是诊断辅助工件。它们通过 `always()` 上传,因此失败的运行仍可供调查。它们的存在不会将失败的基准测试转变为有效结果。在 AMD Slurm 机群上,`/run_logs` 是节点本地目录;服务器步骤结束后,`job.slurm` 会把每个已分配节点上已经关闭的日志树合并到共享存储中,使诊断工件包含整个部署的 Prefill 和 Decode 日志。 ## 阶段 6:工件收集与交接 diff --git a/perf-changelog.yaml b/perf-changelog.yaml index 6b0fdc1667..e7841bc1ba 100644 --- a/perf-changelog.yaml +++ b/perf-changelog.yaml @@ -7924,6 +7924,16 @@ - "Switch HiCache defaults to --hicache-io-backend kernel and --hicache-mem-layout page_first (from direct / page_first_direct). Ratio 1.5 and write_through are unchanged." pr-link: https://github.com/SemiAnalysisAI/InferenceX/pull/3118 +- config-keys: + - dsv4-fp4-mi355x-sglang-agentic-mtp + scenario-type: + - agentic-coding + description: + - "Bump the image to lmsysorg/sglang-rocm:v0.5.19-rocm720-mi35x-20260914." + - "Raise --prefill-decode-interval from 10 to 20 on both the TP-only and DP-attention serving paths, so the scheduler runs more decode steps between prefill admissions." + - "Switch the sglang-router policy for the DP-attention arms from consistent_hashing to cache_aware, and set the cache-aware absolute balance threshold to 32, down from the router default of 64, at concurrency above 160." + pr-link: https://github.com/SemiAnalysisAI/InferenceX/pull/3120 + - config-keys: - dsv4-fp4-mi355x-sglang-disagg-agentic-umbp-dspark scenario-type: @@ -7936,14 +7946,5 @@ - "将四条臂全部从 EAGLE(spec-decoding: mtp)切换为 DSpark(spec-decoding: draft_model),DECODE_MTP_SIZE=3;镜像由 lmsysorg/sglang-rocm:v0.5.19-rocm720-mi35x-20260911 升级到 -20260913(首个带 DSpark 优化的标签),CLIENT_IMAGE 保持 20260907(仅负载生成器)。拓扑、并发列表、HiCache/UMBP 设置与精度均不变。models.yaml 中 DeepSeek-V4-Pro-AgentX 新增 dspark_flags(与 mtp_flags 并列而非替换,因此仍用 mtp 的配方 EAGLE 分支不受影响),dp_flags 增加 --enable-dp-lm-head。server_sglang.sh 按 SPEC_DECODING == draft_model 分流:DSpark 一次前向出一整块,故 num-steps 钉死为 1,草稿长度转为 --speculative-dspark-block-size;验证窗口(长度+1)与 EAGLE 相同,所以 MORI 解码 dispatch 的 ×(DECODE_MTP_SIZE+1) 缩放无需改动。黄金 AL 不变:按检查点而非算法选表,DeepSeek-V4-Pro-0813 在草稿长度 3 仍为 AL 3.01。" - "Rename the config key from dsv4-fp4-mi355x-sglang-disagg-agentic-hicache-mtp to dsv4-fp4-mi355x-sglang-disagg-agentic-umbp-dspark to reflect the DSpark algorithm and the UMBP-linker arm; earlier entries for this recipe are recorded under the old key. No other field changes with the rename." - "配置键由 dsv4-fp4-mi355x-sglang-disagg-agentic-hicache-mtp 改名为 dsv4-fp4-mi355x-sglang-disagg-agentic-umbp-dspark,以反映 DSpark 算法与 UMBP-linker 分支;该配方此前的条目仍记录在旧键名下。改名不伴随任何其他字段变化。" + - "Collect completed logs from every Slurm node so both prefill and decode artifacts are uploaded; forward workflow-owned AgentX fast-mode and power inputs, and supply the existing router port to evaluation." pr-link: https://github.com/SemiAnalysisAI/InferenceX/pull/3188 - -- config-keys: - - dsv4-fp4-mi355x-sglang-agentic-mtp - scenario-type: - - agentic-coding - description: - - "Bump the image to lmsysorg/sglang-rocm:v0.5.19-rocm720-mi35x-20260914." - - "Raise --prefill-decode-interval from 10 to 20 on both the TP-only and DP-attention serving paths, so the scheduler runs more decode steps between prefill admissions." - - "Switch the sglang-router policy for the DP-attention arms from consistent_hashing to cache_aware, and set the cache-aware absolute balance threshold to 32, down from the router default of 64, at concurrency above 160." - pr-link: https://github.com/SemiAnalysisAI/InferenceX/pull/3120 diff --git a/utils/test_amd_node_log_staging.py b/utils/test_amd_node_log_staging.py new file mode 100644 index 0000000000..d194665588 --- /dev/null +++ b/utils/test_amd_node_log_staging.py @@ -0,0 +1,84 @@ +from __future__ import annotations + +import os +import subprocess +from pathlib import Path + +REPO_ROOT = Path(__file__).resolve().parents[1] +STAGE_SCRIPT = REPO_ROOT / "benchmarks/multi_node/amd_utils/stage_node_logs.sh" + + +def _stub_sudo(tmp_path: Path) -> dict[str, str]: + bin_dir = tmp_path / "bin" + bin_dir.mkdir() + sudo = bin_dir / "sudo" + sudo.write_text('#!/bin/sh\nexec "$@"\n') + sudo.chmod(0o755) + return {**os.environ, "PATH": f"{bin_dir}:{os.environ['PATH']}"} + + +def test_stage_node_logs_merges_prefill_and_decode_nodes(tmp_path: Path) -> None: + prefill = tmp_path / "prefill-node" + decode = tmp_path / "decode-node" + shared = tmp_path / "shared" + prefill.mkdir() + decode.mkdir() + (prefill / "prefill_host-a.log").write_text("prefill output\n") + (prefill / "server_host-a.log").write_text("frontend output\n") + (decode / "decode_host-b.log").write_text("decode output\n") + (decode / "server_host-b.log").write_text("decode wrapper output\n") + env = _stub_sudo(tmp_path) + + for node_logs in (prefill, decode): + subprocess.run( + ["bash", str(STAGE_SCRIPT), str(node_logs), str(shared)], + check=True, + env=env, + ) + + assert sorted(path.name for path in shared.iterdir()) == [ + "decode_host-b.log", + "prefill_host-a.log", + "server_host-a.log", + "server_host-b.log", + ] + assert (shared / "decode_host-b.log").read_text() == "decode output\n" + + +def test_stage_node_logs_rejects_a_node_without_logs(tmp_path: Path) -> None: + missing = tmp_path / "missing" + shared = tmp_path / "shared" + + completed = subprocess.run( + ["bash", str(STAGE_SCRIPT), str(missing), str(shared)], + check=False, + capture_output=True, + text=True, + ) + + assert completed.returncode == 1 + assert "no node-local logs found" in completed.stderr + assert not shared.exists() + + +def test_stage_node_logs_propagates_copy_failure(tmp_path: Path) -> None: + source = tmp_path / "node" + shared = tmp_path / "shared" + source.mkdir() + (source / "decode_host-b.log").write_text("decode output\n") + env = _stub_sudo(tmp_path) + cp = tmp_path / "bin" / "cp" + cp.write_text("#!/bin/sh\nexit 19\n") + cp.chmod(0o755) + + completed = subprocess.run( + ["bash", str(STAGE_SCRIPT), str(source), str(shared)], + check=False, + capture_output=True, + text=True, + env=env, + ) + + assert completed.returncode == 19 + assert "[logs] staged" not in completed.stdout + assert not (shared / "decode_host-b.log").exists() From cbcaed56af86b1774488c472f7c7a50aaf6bea2c Mon Sep 17 00:00:00 2001 From: billishyahao Date: Thu, 17 Sep 2026 00:26:56 +0000 Subject: [PATCH 06/12] fix lint --- perf-changelog.yaml | 30 +++++++++++++++--------------- 1 file changed, 15 insertions(+), 15 deletions(-) diff --git a/perf-changelog.yaml b/perf-changelog.yaml index e6a57521aa..2c8423ba45 100644 --- a/perf-changelog.yaml +++ b/perf-changelog.yaml @@ -7942,21 +7942,6 @@ pr-link: https://github.com/SemiAnalysisAI/InferenceX/pull/3176 - config-keys: - - dsv4-fp4-mi355x-sglang-disagg-agentic-umbp-dspark - scenario-type: - - agentic-coding - description: - - "Switch all four arms from EAGLE (spec-decoding: mtp) to DSpark (spec-decoding: draft_model) at DECODE_MTP_SIZE=3, and bump the image from lmsysorg/sglang-rocm:v0.5.19-rocm720-mi35x-20260911 to -20260913, the first tag carrying the DSpark optimizations. CLIENT_IMAGE stays at 20260907 (load generator only). Topology, concurrency list, HiCache/UMBP settings, and precision are unchanged." - - "DeepSeek-V4-Pro-AgentX in models.yaml gains dspark_flags (--speculative-algorithm DSPARK --speculative-eagle-topk 1) alongside mtp_flags rather than replacing it, so recipes remaining on spec-decoding: mtp keep their EAGLE arm unchanged; dp_flags gains --enable-dp-lm-head, required by SGLang for DSpark under DP attention and inert for EAGLE. The -0813 checkpoint bundles the draft head, so no separate draft checkpoint is staged." - - "server_sglang.sh builds DSpark flags as --speculative-dspark-block-size DECODE_MTP_SIZE --speculative-num-steps 1 --speculative-num-draft-tokens (DECODE_MTP_SIZE + 1). DECODE_MTP_SIZE is the draft length for both algorithms but EAGLE spends it as sequential draft passes while DSpark emits one block, so num-steps is pinned to 1. The verify window is identical for both, so the MORI decode dispatch scaling by (DECODE_MTP_SIZE + 1) is unchanged. A model set to draft_model without dspark_flags now fails hard instead of silently serving EAGLE." - - "Golden AL is unchanged: the AgentX curve is keyed on checkpoint, thinking mode, and draft length, not on the algorithm, so DeepSeek-V4-Pro-0813 at draft length 3 keeps AL 3.01 from golden_al_distribution/dsv4-pro-0813-dspark.yaml (thinking_on)." - - "将四条臂全部从 EAGLE(spec-decoding: mtp)切换为 DSpark(spec-decoding: draft_model),DECODE_MTP_SIZE=3;镜像由 lmsysorg/sglang-rocm:v0.5.19-rocm720-mi35x-20260911 升级到 -20260913(首个带 DSpark 优化的标签),CLIENT_IMAGE 保持 20260907(仅负载生成器)。拓扑、并发列表、HiCache/UMBP 设置与精度均不变。models.yaml 中 DeepSeek-V4-Pro-AgentX 新增 dspark_flags(与 mtp_flags 并列而非替换,因此仍用 mtp 的配方 EAGLE 分支不受影响),dp_flags 增加 --enable-dp-lm-head。server_sglang.sh 按 SPEC_DECODING == draft_model 分流:DSpark 一次前向出一整块,故 num-steps 钉死为 1,草稿长度转为 --speculative-dspark-block-size;验证窗口(长度+1)与 EAGLE 相同,所以 MORI 解码 dispatch 的 ×(DECODE_MTP_SIZE+1) 缩放无需改动。黄金 AL 不变:按检查点而非算法选表,DeepSeek-V4-Pro-0813 在草稿长度 3 仍为 AL 3.01。" - - "Rename the config key from dsv4-fp4-mi355x-sglang-disagg-agentic-hicache-mtp to dsv4-fp4-mi355x-sglang-disagg-agentic-umbp-dspark to reflect the DSpark algorithm and the UMBP-linker arm; earlier entries for this recipe are recorded under the old key. No other field changes with the rename." - - "配置键由 dsv4-fp4-mi355x-sglang-disagg-agentic-hicache-mtp 改名为 dsv4-fp4-mi355x-sglang-disagg-agentic-umbp-dspark,以反映 DSpark 算法与 UMBP-linker 分支;该配方此前的条目仍记录在旧键名下。改名不伴随任何其他字段变化。" - - "Collect completed logs from every Slurm node so both prefill and decode artifacts are uploaded; forward workflow-owned AgentX fast-mode and power inputs, and supply the existing router port to evaluation." - pr-link: https://github.com/SemiAnalysisAI/InferenceX/pull/3188 - -- config-keys: - dsv4-fp4-b200-sglang-agentic-hicache-mtp scenario-type: - agentic-coding @@ -8003,3 +7988,18 @@ - "Use a B300-specific vLLM AgentX script with explicit FULL_AND_PIECEWISE CUDA graph capture sizes, concurrency-tiered batched-token limits, and an explicit maximum sequence count." - "Add a TP2 search-space variant at concurrency 2-128 while retaining the TP4 concurrency 1-128 variant." pr-link: https://github.com/SemiAnalysisAI/InferenceX/pull/3109 + +- config-keys: + - dsv4-fp4-mi355x-sglang-disagg-agentic-umbp-dspark + scenario-type: + - agentic-coding + description: + - "Switch all four arms from EAGLE (spec-decoding: mtp) to DSpark (spec-decoding: draft_model) at DECODE_MTP_SIZE=3, and bump the image from lmsysorg/sglang-rocm:v0.5.19-rocm720-mi35x-20260911 to -20260913, the first tag carrying the DSpark optimizations. CLIENT_IMAGE stays at 20260907 (load generator only). Topology, concurrency list, HiCache/UMBP settings, and precision are unchanged." + - "DeepSeek-V4-Pro-AgentX in models.yaml gains dspark_flags (--speculative-algorithm DSPARK --speculative-eagle-topk 1) alongside mtp_flags rather than replacing it, so recipes remaining on spec-decoding: mtp keep their EAGLE arm unchanged; dp_flags gains --enable-dp-lm-head, required by SGLang for DSpark under DP attention and inert for EAGLE. The -0813 checkpoint bundles the draft head, so no separate draft checkpoint is staged." + - "server_sglang.sh builds DSpark flags as --speculative-dspark-block-size DECODE_MTP_SIZE --speculative-num-steps 1 --speculative-num-draft-tokens (DECODE_MTP_SIZE + 1). DECODE_MTP_SIZE is the draft length for both algorithms but EAGLE spends it as sequential draft passes while DSpark emits one block, so num-steps is pinned to 1. The verify window is identical for both, so the MORI decode dispatch scaling by (DECODE_MTP_SIZE + 1) is unchanged. A model set to draft_model without dspark_flags now fails hard instead of silently serving EAGLE." + - "Golden AL is unchanged: the AgentX curve is keyed on checkpoint, thinking mode, and draft length, not on the algorithm, so DeepSeek-V4-Pro-0813 at draft length 3 keeps AL 3.01 from golden_al_distribution/dsv4-pro-0813-dspark.yaml (thinking_on)." + - "将四条臂全部从 EAGLE(spec-decoding: mtp)切换为 DSpark(spec-decoding: draft_model),DECODE_MTP_SIZE=3;镜像由 lmsysorg/sglang-rocm:v0.5.19-rocm720-mi35x-20260911 升级到 -20260913(首个带 DSpark 优化的标签),CLIENT_IMAGE 保持 20260907(仅负载生成器)。拓扑、并发列表、HiCache/UMBP 设置与精度均不变。models.yaml 中 DeepSeek-V4-Pro-AgentX 新增 dspark_flags(与 mtp_flags 并列而非替换,因此仍用 mtp 的配方 EAGLE 分支不受影响),dp_flags 增加 --enable-dp-lm-head。server_sglang.sh 按 SPEC_DECODING == draft_model 分流:DSpark 一次前向出一整块,故 num-steps 钉死为 1,草稿长度转为 --speculative-dspark-block-size;验证窗口(长度+1)与 EAGLE 相同,所以 MORI 解码 dispatch 的 ×(DECODE_MTP_SIZE+1) 缩放无需改动。黄金 AL 不变:按检查点而非算法选表,DeepSeek-V4-Pro-0813 在草稿长度 3 仍为 AL 3.01。" + - "Rename the config key from dsv4-fp4-mi355x-sglang-disagg-agentic-hicache-mtp to dsv4-fp4-mi355x-sglang-disagg-agentic-umbp-dspark to reflect the DSpark algorithm and the UMBP-linker arm; earlier entries for this recipe are recorded under the old key. No other field changes with the rename." + - "配置键由 dsv4-fp4-mi355x-sglang-disagg-agentic-hicache-mtp 改名为 dsv4-fp4-mi355x-sglang-disagg-agentic-umbp-dspark,以反映 DSpark 算法与 UMBP-linker 分支;该配方此前的条目仍记录在旧键名下。改名不伴随任何其他字段变化。" + - "Collect completed logs from every Slurm node so both prefill and decode artifacts are uploaded; forward workflow-owned AgentX fast-mode and power inputs, and supply the existing router port to evaluation." + pr-link: https://github.com/SemiAnalysisAI/InferenceX/pull/3188 From e112620aa7bd756c4b73ca406e47db5411f54a9e Mon Sep 17 00:00:00 2001 From: Cam Quilici Date: Wed, 16 Sep 2026 15:44:41 -0500 Subject: [PATCH 07/12] fix: bound teardown of owned AMD SGLang process groups --- benchmarks/benchmark_lib.sh | 68 +++++++ .../multi_node/amd_utils/server_sglang.sh | 46 ++--- docs/recovery-results-procedures.md | 12 ++ docs/recovery-results-procedures_zh.md | 10 + utils/test_process_group_cleanup.py | 191 ++++++++++++++++++ 5 files changed, 301 insertions(+), 26 deletions(-) create mode 100644 utils/test_process_group_cleanup.py diff --git a/benchmarks/benchmark_lib.sh b/benchmarks/benchmark_lib.sh index 74dce63495..5deb069840 100644 --- a/benchmarks/benchmark_lib.sh +++ b/benchmarks/benchmark_lib.sh @@ -19,6 +19,74 @@ check_env_vars() { fi } +# Report live members of explicitly owned process groups. Zombies cannot hold +# output pipes open. Do not use leader liveness: a router can orphan its workers. +_background_process_groups_alive() { + local groups=" $* " + local listing + listing=$(ps -eo pgid=,stat=) || return 1 + awk -v groups="$groups" ' + index(groups, " " $1 " ") && $2 !~ /^[ZX]/ { alive[$1] = 1 } + END { for (group in alive) print group } + ' <<< "$listing" +} + +# Called only after benchmark/eval work ends. Preserve its exit status while +# bounding teardown of the setsid groups recorded by the launcher. Grace periods +# are explicit arguments, independent of benchmark duration and server readiness. +stop_background_process_groups() { + local work_status="$1" term_grace="$2" kill_grace="$3" + shift 3 + local pgid own_pgid remaining deadline cleanup_status=0 + local groups=("$@") + if [[ ! "$work_status" =~ ^[0-9]+$ || ! "$term_grace" =~ ^[0-9]+$ || ! "$kill_grace" =~ ^[0-9]+$ ]]; then + echo "ERROR: invalid process-group cleanup status or grace period" >&2 + return 1 + fi + if ! own_pgid=$(ps -o pgid= -p "$$"); then + if [[ "$work_status" -ne 0 ]]; then return "$work_status"; fi + return 1 + fi + own_pgid="${own_pgid//[[:space:]]/}" + for pgid in "${groups[@]}"; do + if [[ ! "$pgid" =~ ^[1-9][0-9]*$ || "$pgid" -le 1 || "$pgid" == "$own_pgid" ]]; then + echo "ERROR: refusing unsafe process-group cleanup: '$pgid'" >&2 + if [[ "$work_status" -ne 0 ]]; then return "$work_status"; fi + return 1 + fi + done + if [[ ${#groups[@]} -eq 0 ]]; then return "$work_status"; fi + + echo "Stopping owned process groups: ${groups[*]}" + for pgid in "${groups[@]}"; do + kill -TERM -- "-$pgid" 2>/dev/null || true + done + deadline=$((SECONDS + term_grace)) + while true; do + remaining=$(_background_process_groups_alive "${groups[@]}") || { cleanup_status=1; break; } + [[ -n "$remaining" && $SECONDS -lt $deadline ]] || break + sleep 1 + done + if [[ -n "$remaining" ]]; then + echo "TERM grace expired; force-stopping owned process groups: $remaining" + for pgid in $remaining; do + kill -KILL -- "-$pgid" 2>/dev/null || true + done + deadline=$((SECONDS + kill_grace)) + while true; do + remaining=$(_background_process_groups_alive "${groups[@]}") || { cleanup_status=1; break; } + [[ -n "$remaining" && $SECONDS -lt $deadline ]] || break + sleep 1 + done + if [[ -n "$remaining" ]]; then + echo "ERROR: process groups still alive after KILL grace: $remaining" >&2 + cleanup_status=1 + fi + fi + if [[ "$work_status" -ne 0 ]]; then return "$work_status"; fi + return "$cleanup_status" +} + # Launchers may load only input validation, without benchmark initialization. if [[ "${1-}" == "--validation-only" ]]; then return 0 diff --git a/benchmarks/multi_node/amd_utils/server_sglang.sh b/benchmarks/multi_node/amd_utils/server_sglang.sh index b38035e73d..f247ac4158 100755 --- a/benchmarks/multi_node/amd_utils/server_sglang.sh +++ b/benchmarks/multi_node/amd_utils/server_sglang.sh @@ -959,8 +959,7 @@ if [ "$NODE_RANK" -eq 0 ]; then > >(tee /run_logs/slurm_job-${SLURM_JOB_ID}/prefill_${host_name}.log >/dev/null) 2>&1 & set +x prefill0_pid=$! - prefill0_pgid=$(ps -o pgid= -p "$prefill0_pid" 2>/dev/null | tr -d ' ') - : "${prefill0_pgid:=$prefill0_pid}" + prefill0_pgid=$prefill0_pid fi echo "Waiting for all prefill and decode servers to be up . . ." @@ -1017,8 +1016,7 @@ if [ "$NODE_RANK" -eq 0 ]; then fi set +x proxy_pid=$! - proxy_pgid=$(ps -o pgid= -p "$proxy_pid" 2>/dev/null | tr -d ' ') - : "${proxy_pgid:=$proxy_pid}" + proxy_pgid=$proxy_pid HEALTH_BARRIER_CMD="python3 $SGLANG_WS_PATH/sync.py barrier \ --node-ips ${NODE0_ADDR} \ @@ -1114,6 +1112,7 @@ if [ "$NODE_RANK" -eq 0 ]; then IS_AGENTIC_RUN=1 fi + BENCHMARK_EXIT_CODE=0 if [[ "${EVAL_ONLY}" == "true" ]]; then echo "EVAL_ONLY mode: skipping throughput benchmark" elif [[ "$DRY_RUN" -eq 1 ]]; then @@ -1188,10 +1187,12 @@ print(json.dumps(json.loads(sys.stdin.read())))' <<<"$_val")" || { --entrypoint "" \ "${CLIENT_IMAGE}" \ bash -lc "cd /workspace/benchmarks/multi_node/amd_utils && bash trace_replay.sh /models ${MODEL_NAME} \"${BENCH_MAX_CONCURRENCY}\" /run_logs/slurm_job-${SLURM_JOB_ID}" + BENCHMARK_EXIT_CODE=$? set +x else set -x eval "$BENCH_CMD" + BENCHMARK_EXIT_CODE=$? set +x fi @@ -1281,22 +1282,17 @@ print(json.dumps(json.loads(sys.stdin.read())))' <<<"$_val")" || { echo "Copied results to $LOGS_OUTPUT/slurm_job-${SLURM_JOB_ID}" fi - echo "Killing the proxy server and prefill server" - - if [[ "$DRY_RUN" -eq 0 ]]; then - # Group-kill the router (setsid at launch): the python launcher has usually - # exited after spawning the Rust worker, which reparents to init but stays in - # this group; kill $proxy_pid alone misses it and :30000 stays open. - kill -TERM -"${proxy_pgid:-$proxy_pid}" 2>/dev/null || true - # Group-kill the prefill tree so TP-scheduler children release the tee pipe - # and the container can exit. - kill -TERM -"${prefill0_pgid:-$prefill0_pid}" 2>/dev/null || true + node_exit_status=$BENCHMARK_EXIT_CODE + if [[ "${EVAL_FAILED:-0}" -eq 1 && "$node_exit_status" -eq 0 ]]; then + node_exit_status=1 fi - - if [[ "${EVAL_FAILED:-0}" -eq 1 ]]; then - echo "ERROR: eval failed; exiting node-0 with rc=1" - exit 1 + if [[ "$DRY_RUN" -eq 0 ]]; then + # The router and prefill may retain TERM-resistant tokenizer workers. + # Keep the benchmark/eval status even when teardown also fails. + stop_background_process_groups "$node_exit_status" 30 5 "$proxy_pgid" "$prefill0_pgid" + exit $? fi + exit "$node_exit_status" elif [ "$NODE_RANK" -gt 0 ] && [ "$NODE_RANK" -lt "$NODE_OFFSET" ]; then echo "${host_name}:${host_ip} is Prefill Node (Model: ${MODEL_NAME})" @@ -1340,8 +1336,7 @@ elif [ "$NODE_RANK" -gt 0 ] && [ "$NODE_RANK" -lt "$NODE_OFFSET" ]; then > >(tee /run_logs/slurm_job-${SLURM_JOB_ID}/prefill_${host_name}.log >/dev/null) 2>&1 & set +x prefill_pid=$! - prefill_pgid=$(ps -o pgid= -p "$prefill_pid" 2>/dev/null | tr -d ' ') - : "${prefill_pgid:=$prefill_pid}" + prefill_pgid=$prefill_pid fi echo "Waiting for proxy server to be up..." @@ -1371,8 +1366,8 @@ elif [ "$NODE_RANK" -gt 0 ] && [ "$NODE_RANK" -lt "$NODE_OFFSET" ]; then echo "Killing the rank $NODE_RANK prefill server" if [[ "$DRY_RUN" -eq 0 ]]; then - # Group-kill so TP-scheduler children release the tee pipe and the container exits. - kill -TERM -"${prefill_pgid:-$prefill_pid}" 2>/dev/null || true + stop_background_process_groups 0 30 5 "$prefill_pgid" + exit $? fi else @@ -1459,8 +1454,7 @@ else set +x decode_pid=$! - decode_pgid=$(ps -o pgid= -p "$decode_pid" 2>/dev/null | tr -d ' ') - : "${decode_pgid:=$decode_pid}" + decode_pgid=$decode_pid fi echo "Waiting for proxy server to be up..." @@ -1489,8 +1483,8 @@ else echo "Killing the rank $RANK decode server" if [[ "$DRY_RUN" -eq 0 ]]; then - # Group-kill so TP-scheduler children release the tee pipe and the container exits. - kill -TERM -"${decode_pgid:-$decode_pid}" 2>/dev/null || true + stop_background_process_groups 0 30 5 "$decode_pgid" + exit $? fi fi diff --git a/docs/recovery-results-procedures.md b/docs/recovery-results-procedures.md index 0d140d4c84..4ba1ccf578 100644 --- a/docs/recovery-results-procedures.md +++ b/docs/recovery-results-procedures.md @@ -425,3 +425,15 @@ Remaining durable fix: ``` This evidence is the completion gate. “Workflow green” without artifact identity, source/merge identity, and ingest counts is not a verified result recovery. + +### AMD multi-node SGLang teardown + +After benchmark/eval work and result staging, the AMD SGLang launcher sends TERM +only to its recorded `setsid` process groups. It allows 30 seconds for graceful +exit, then sends KILL to surviving groups and checks for exit for another five +seconds. This handles orphaned or TERM-resistant workers that otherwise hold log +pipes open. These cleanup deadlines do not change profiling, evaluation, or server +readiness deadlines. A failed client retains its exit status; unresolved cleanup +fails an otherwise successful node. Kernel-blocked processes may still require +separately authorized node repair. Do not change or discard completed metrics to +work around teardown failures. diff --git a/docs/recovery-results-procedures_zh.md b/docs/recovery-results-procedures_zh.md index df3ed5e045..af5d2b1e21 100644 --- a/docs/recovery-results-procedures_zh.md +++ b/docs/recovery-results-procedures_zh.md @@ -425,3 +425,13 @@ Remaining durable fix: ``` 这些证据就是完成关卡。如果没有制品身份、source/merge 身份和摄取数量,仅仅“工作流绿色”并不代表结果恢复已经验证。 + +### AMD 多节点 SGLang 清理 + +基准测试或评估完成并暂存结果后,AMD SGLang 启动器仅向其记录的 `setsid` +进程组发送 TERM,等待最多 30 秒。随后向仍存活的进程组发送 KILL,再等待最多 +5 秒并检查退出状态。这可以清理已成为孤儿进程或忽略 TERM 的工作进程,避免其 +持续占用日志管道。这些清理期限不会改变性能采集、评估或服务器就绪检查的期限。 +客户端失败时保留原退出码;若客户端成功但清理仍未完成,则节点任务失败。 +内核阻塞的进程仍可能需要另行授权的节点修复。不要为绕过清理失败而修改或丢弃 +已完成的指标。 diff --git a/utils/test_process_group_cleanup.py b/utils/test_process_group_cleanup.py new file mode 100644 index 0000000000..186fb38bf0 --- /dev/null +++ b/utils/test_process_group_cleanup.py @@ -0,0 +1,191 @@ +"""Real process/pipe regressions for post-benchmark group teardown.""" + +import os +import signal +import subprocess +import sys +import time +from pathlib import Path + +import pytest + +LIBRARY = Path(__file__).resolve().parents[1] / "benchmarks/benchmark_lib.sh" + + +def stop_groups(status: int, *groups: int) -> subprocess.CompletedProcess[str]: + # The client exits independently; the actual cleanup implementation must + # preserve its code even when signaling and polling report success. + return subprocess.run( + [ + "bash", + "-c", + ( + 'source "$1" --validation-only; shift; ' + 'client_status=$1; shift; bash -c "exit $client_status"; status=$?; ' + 'stop_background_process_groups "$status" 1 1 "$@"' + ), + "cleanup", + str(LIBRARY), + str(status), + *map(str, groups), + ], + check=False, + capture_output=True, + text=True, + timeout=8, + ) + + +def await_file(path: Path) -> None: + deadline = time.monotonic() + 5 + while not path.exists(): + assert time.monotonic() < deadline, "Child did not start" + time.sleep(0.01) + + +@pytest.mark.parametrize("leader_exits, client_status", [(True, 0), (False, 19)]) +def test_stubborn_descendant_releases_pipe_and_preserves_client_status( + tmp_path: Path, + leader_exits: bool, + client_status: int, +) -> None: + ready = tmp_path / "child-ready" + child_code = ( + "import signal,pathlib,sys,time; " + "signal.signal(signal.SIGTERM, signal.SIG_IGN); " + 'pathlib.Path(sys.argv[1]).touch(); print("worker output", flush=True); time.sleep(60)' + ) + leader_code = ( + "import subprocess,sys,time; " + 'subprocess.Popen([sys.executable,"-c",sys.argv[1],sys.argv[2]]); ' + 'time.sleep(0 if sys.argv[3]=="True" else 60)' + ) + with ( + subprocess.Popen( + [ + sys.executable, + "-c", + leader_code, + child_code, + str(ready), + str(leader_exits), + ], + start_new_session=True, + stdout=subprocess.PIPE, + stderr=subprocess.STDOUT, + text=True, + ) as leader, + subprocess.Popen( + [sys.executable, "-c", "import time; time.sleep(60)"] + ) as unrelated, + ): + try: + await_file(ready) + if leader_exits: + assert leader.wait(timeout=3) == 0 + result = stop_groups(client_status, leader.pid) + assert result.returncode == client_status, result.stderr + assert "force-stopping owned process groups" in result.stdout + # An orphan that retains stdout makes communicate hang even after + # its leader exited. This tests the original tee-pipe failure. + output, _ = leader.communicate(timeout=3) + assert "worker output" in output + assert unrelated.poll() is None + finally: + try: + os.killpg(leader.pid, signal.SIGKILL) + except ProcessLookupError: + pass + except PermissionError: + # macOS can retain a zombie-only process group owned by init. + pass + unrelated.terminate() + + +def test_graceful_group_gets_term_without_kill(tmp_path: Path) -> None: + ready = tmp_path / "ready" + stopped = tmp_path / "stopped" + code = ( + "import signal,pathlib,sys,time; " + "signal.signal(signal.SIGTERM, lambda *_: (pathlib.Path(sys.argv[2]).touch(), sys.exit(0))); " + "pathlib.Path(sys.argv[1]).touch(); time.sleep(60)" + ) + with subprocess.Popen( + [sys.executable, "-c", code, str(ready), str(stopped)], start_new_session=True + ) as leader: + try: + await_file(ready) + result = stop_groups(0, leader.pid) + assert result.returncode == 0, result.stderr + assert leader.wait(timeout=2) == 0 + assert stopped.exists() + assert "force-stopping" not in result.stdout + finally: + if leader.poll() is None: + leader.kill() + + +@pytest.mark.parametrize("client_status, expected", [(0, 1), (23, 23)]) +def test_refused_cleanup_cannot_hide_work_failure_or_report_success( + client_status: int, expected: int +) -> None: + result = stop_groups(client_status, 1) + assert result.returncode == expected + assert "refusing unsafe process-group cleanup" in result.stderr + + +def test_refuses_callers_own_group() -> None: + result = subprocess.run( + [ + "bash", + "-c", + ( + 'source "$1" --validation-only; ' + 'group=$(ps -o pgid= -p $$ | tr -d " "); stop_background_process_groups 0 1 1 "$group"' + ), + "cleanup", + str(LIBRARY), + ], + start_new_session=True, + check=False, + capture_output=True, + text=True, + timeout=5, + ) + assert result.returncode == 1 + assert "refusing unsafe process-group cleanup" in result.stderr + + +@pytest.mark.parametrize("client_status, expected", [(0, 1), (23, 23)]) +def test_signal_failure_is_bounded_and_preserves_work_failure( + client_status: int, expected: int +) -> None: + with subprocess.Popen( + [sys.executable, "-c", "import time; time.sleep(60)"], start_new_session=True + ) as leader: + try: + # Only mock the OS signal collaborator: the real liveness checks, + # grace deadlines, and status selection all run against a live group. + result = subprocess.run( + [ + "bash", + "-c", + ( + 'source "$1" --validation-only; ' + 'kill() { return 1; }; stop_background_process_groups "$2" 1 1 "$3"' + ), + "cleanup", + str(LIBRARY), + str(client_status), + str(leader.pid), + ], + check=False, + capture_output=True, + text=True, + timeout=6, + ) + assert result.returncode == expected + assert "still alive after KILL grace" in result.stderr + assert leader.poll() is None + finally: + leader.kill() From 3450090826b228d9f70f10f8d3df95f0ff00dbfe Mon Sep 17 00:00:00 2001 From: Cam Quilici Date: Wed, 16 Sep 2026 23:36:44 -0500 Subject: [PATCH 08/12] fix: coordinate AMD node preflight before server startup --- benchmarks/benchmark_lib.sh | 22 +++ benchmarks/multi_node/amd_utils/job.slurm | 26 +-- .../multi_node/amd_utils/preflight_node.sh | 31 ++++ docs/recovery-results-procedures.md | 11 ++ docs/recovery-results-procedures_zh.md | 9 + utils/test_amd_multinode_preflight.py | 169 ++++++++++++++++++ 6 files changed, 246 insertions(+), 22 deletions(-) create mode 100644 benchmarks/multi_node/amd_utils/preflight_node.sh create mode 100644 utils/test_amd_multinode_preflight.py diff --git a/benchmarks/benchmark_lib.sh b/benchmarks/benchmark_lib.sh index 5deb069840..eb039ef8f3 100644 --- a/benchmarks/benchmark_lib.sh +++ b/benchmarks/benchmark_lib.sh @@ -87,6 +87,28 @@ stop_background_process_groups() { return "$cleanup_status" } +# Finish preflight on every allocated node before any server container starts its +# peer-readiness deadline. A failed node prevents the entire serving step. +run_amd_multinode_after_preflight() { + local nodelist="$1" node_count="$2" preflight_script="$3" + local container_filter="$4" skip_gpu_sanity="$5" + shift 5 + local preflight_rc + if srun --nodelist="$nodelist" \ + --nodes="$node_count" --ntasks="$node_count" --ntasks-per-node=1 \ + --kill-on-bad-exit=1 --unbuffered \ + bash "$preflight_script" "$container_filter" "$skip_gpu_sanity"; then + echo "[preflight] all nodes ready; launching server containers" + else + preflight_rc=$? + echo "[preflight][ERROR] node preflight failed; no server containers launched" >&2 + return "$preflight_rc" + fi + srun --nodelist="$nodelist" \ + --nodes="$node_count" --ntasks="$node_count" --ntasks-per-node=1 \ + --kill-on-bad-exit=1 --signal=TERM@30 --unbuffered "$@" +} + # Launchers may load only input validation, without benchmark initialization. if [[ "${1-}" == "--validation-only" ]]; then return 0 diff --git a/benchmarks/multi_node/amd_utils/job.slurm b/benchmarks/multi_node/amd_utils/job.slurm index bb151f2039..d1f5ea11a7 100755 --- a/benchmarks/multi_node/amd_utils/job.slurm +++ b/benchmarks/multi_node/amd_utils/job.slurm @@ -583,11 +583,10 @@ if [[ -n "${CLIENT_IMAGE:-}" ]]; then srun --nodelist="$SELECTED_NODELIST_SRUN" bash -c 'eval "$DOCKER_CMD_DETECT"; $DOCKER_CMD pull '"$CLIENT_IMAGE"' >/dev/null 2>&1 || true' 2>/dev/null || true fi -srun \ - --nodelist="$SELECTED_NODELIST_SRUN" \ - --kill-on-bad-exit=1 \ - --signal=TERM@30 \ - --unbuffered \ +run_amd_multinode_after_preflight \ + "$SELECTED_NODELIST_SRUN" "$NUM_NODES" \ + "$DI_REPO_DIR/benchmarks/multi_node/amd_utils/preflight_node.sh" \ + "$CONT_FILTER" "$SKIP_GPU_SANITY" \ bash -lc " set -eo pipefail @@ -675,23 +674,6 @@ else fi fi # end: if ENGINE == atom-disagg -# Pre-clean (idempotent): stop then force-remove so GPU VRAM is released -# before the drain gate. stop-only left containers in Created/Exited state -# on some nodes. -\$DOCKER_CMD ps -aq --filter \"$CONT_FILTER\" | xargs -r \$DOCKER_CMD rm -f || true -\$DOCKER_CMD ps -aq | xargs -r \$DOCKER_CMD stop -t 15 || true -\$DOCKER_CMD ps -aq | xargs -r \$DOCKER_CMD rm -f || true -sleep 2 - -# GPU drain gate: fail fast on leftover VRAM use instead of OOMing in model -# load ~15 min later. Reuses wait_for_amd_gpu_clean from benchmark_lib.sh. -if [[ \"${SKIP_GPU_SANITY}\" == \"1\" ]]; then - echo \"[INFO] SKIP_GPU_SANITY=1 set; skipping GPU pre-flight drain check\" -else - # Unset so benchmark_lib.sh's unrelated agentic KV_OFFLOADING check doesn't exit 1 here. - bash -c \"unset IS_AGENTIC SCENARIO_TYPE; source $DI_REPO_DIR/benchmarks/benchmark_lib.sh && wait_for_amd_gpu_clean\" -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 diff --git a/benchmarks/multi_node/amd_utils/preflight_node.sh b/benchmarks/multi_node/amd_utils/preflight_node.sh new file mode 100644 index 0000000000..34759bab89 --- /dev/null +++ b/benchmarks/multi_node/amd_utils/preflight_node.sh @@ -0,0 +1,31 @@ +#!/usr/bin/env bash +set -eo pipefail + +SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +source "$SCRIPT_DIR/../../benchmark_lib.sh" --validation-only +check_env_vars DOCKER_CMD_DETECT DI_REPO_DIR SLURM_JOB_ID +CONT_FILTER="$1" +SKIP_GPU_SANITY="$2" +check_env_vars CONT_FILTER SKIP_GPU_SANITY + +preflight_node() { + eval "$DOCKER_CMD_DETECT" + + # Preserve the existing pre-clean scope and ordering. Moving it into this + # separate Slurm step prevents one node starting while another still drains. + $DOCKER_CMD ps -aq --filter "$CONT_FILTER" | xargs -r $DOCKER_CMD rm -f || true + $DOCKER_CMD ps -aq | xargs -r $DOCKER_CMD stop -t 15 || true + $DOCKER_CMD ps -aq | xargs -r $DOCKER_CMD rm -f || true + sleep 2 + + if [[ "$SKIP_GPU_SANITY" == "1" ]]; then + echo "[INFO] SKIP_GPU_SANITY=1 set; skipping GPU pre-flight drain check" + else + # Avoid benchmark-only agentic initialization on the host, as before. + bash -c 'unset IS_AGENTIC SCENARIO_TYPE; source "$DI_REPO_DIR/benchmarks/benchmark_lib.sh" && wait_for_amd_gpu_clean' + fi +} + +NODE_LOG_DIR="/tmp/slurm_job-${SLURM_JOB_ID}" +mkdir -p "$NODE_LOG_DIR" +preflight_node 2>&1 | tee "$NODE_LOG_DIR/preflight_$(hostname).log" diff --git a/docs/recovery-results-procedures.md b/docs/recovery-results-procedures.md index 4ba1ccf578..438fba05f6 100644 --- a/docs/recovery-results-procedures.md +++ b/docs/recovery-results-procedures.md @@ -437,3 +437,14 @@ readiness deadlines. A failed client retains its exit status; unresolved cleanup fails an otherwise successful node. Kernel-blocked processes may still require separately authorized node repair. Do not change or discard completed metrics to work around teardown failures. + +### AMD multi-node GPU preflight coordination + +The Slurm launcher completes Docker pre-clean and the existing GPU VRAM drain +check on every selected node in a separate Slurm step before launching any server +container. A failed preflight prevents the serving step; it does not consume a +healthy peer's container-readiness deadline. Node-local `preflight_.log` +files are included in the normal log fan-in, including failures. The VRAM threshold, +15-minute GPU guard, and container/server readiness deadlines remain unchanged. +This coordination prevents a peer-barrier race; it does not repair a GPU driver +that fails to reclaim memory. The existing Docker pre-clean scope is unchanged. diff --git a/docs/recovery-results-procedures_zh.md b/docs/recovery-results-procedures_zh.md index af5d2b1e21..31ec736403 100644 --- a/docs/recovery-results-procedures_zh.md +++ b/docs/recovery-results-procedures_zh.md @@ -435,3 +435,12 @@ Remaining durable fix: 客户端失败时保留原退出码;若客户端成功但清理仍未完成,则节点任务失败。 内核阻塞的进程仍可能需要另行授权的节点修复。不要为绕过清理失败而修改或丢弃 已完成的指标。 + +### AMD 多节点 GPU 预检协调 + +Slurm 启动器先在独立步骤中完成所有选定节点的 Docker 预清理和现有 GPU VRAM +回收检查,然后才启动服务器容器。任一节点预检失败都会阻止服务步骤启动,不会 +消耗健康节点等待容器就绪的期限。节点本地的 `preflight_.log` 文件 +通过常规日志汇总流程收集,包括失败日志。VRAM 阈值、15 分钟 GPU 检查期限及 +容器和服务器就绪期限均保持不变。这一协调消除了节点间等待的竞态,但无法修复 +不能回收显存的 GPU 驱动。现有 Docker 预清理范围保持不变。 diff --git a/utils/test_amd_multinode_preflight.py b/utils/test_amd_multinode_preflight.py new file mode 100644 index 0000000000..035d431bde --- /dev/null +++ b/utils/test_amd_multinode_preflight.py @@ -0,0 +1,169 @@ +from __future__ import annotations + +import json +import os +import shutil +import subprocess +import time +import uuid +from pathlib import Path + +import pytest + +REPO_ROOT = Path(__file__).resolve().parents[1] +LIBRARY = REPO_ROOT / "benchmarks/benchmark_lib.sh" +PREFLIGHT = REPO_ROOT / "benchmarks/multi_node/amd_utils/preflight_node.sh" + + +@pytest.fixture +def cluster(tmp_path: Path): + """Only Slurm, Docker, GPU telemetry, host naming, and the sleep clock are fake.""" + bin_dir = tmp_path / "bin" + bin_dir.mkdir() + scripts = { + "srun": """#!/usr/bin/env python3 +import os, subprocess, sys +args = sys.argv[1:] +while args and args[0].startswith("--"): + args.pop(0) +children = [subprocess.Popen(args, env={**os.environ, "SLURM_PROCID": str(rank)}) + for rank in range(2)] +statuses = [child.wait() for child in children] +sys.exit(next((status for status in statuses if status), 0)) +""", + "docker": """#!/usr/bin/env python3 +import os, pathlib, sys +with (pathlib.Path(os.environ["TEST_STATE"]) / ("docker-" + os.environ["SLURM_PROCID"])).open("a") as f: + f.write(" ".join(sys.argv[1:]) + "\\n") +""", + "hostname": '#!/bin/sh\necho "node-$SLURM_PROCID"\n', + "sleep": "#!/bin/sh\nexit 0\n", + "rocm-smi": """#!/usr/bin/env python3 +import os, pathlib, time +state = pathlib.Path(os.environ["TEST_STATE"]) +rank = os.environ["SLURM_PROCID"] +mode = os.environ["TEST_MODE"] +with (state / ("probes-" + rank)).open("a") as f: + f.write("probe\\n") +if rank == "1" and mode == "delayed" and not (state / "release").exists(): + (state / "waiting").touch() + time.sleep(0.02) + used = 92 +elif rank == "1" and mode == "failure": + used = 92 +else: + used = 0 + (state / ("clean-" + rank)).touch() +print(f"GPU[0] : GPU Memory Allocated (VRAM%): {used}") +""", + } + for name, source in scripts.items(): + path = bin_dir / name + path.write_text(source) + path.chmod(0o755) + job_id = f"infx-preflight-test-{uuid.uuid4().hex}" + env = { + **os.environ, + "PATH": f"{bin_dir}:{os.environ['PATH']}", + "DOCKER_CMD_DETECT": "DOCKER_CMD=docker", + "DI_REPO_DIR": str(REPO_ROOT), + "SLURM_JOB_ID": job_id, + "TEST_STATE": str(tmp_path), + "TEST_MODE": "delayed", + } + yield tmp_path, env, Path("/tmp") / f"slurm_job-{job_id}" + shutil.rmtree(Path("/tmp") / f"slurm_job-{job_id}", ignore_errors=True) + + +def _command(skip: str = "0", server_status: int = 0) -> list[str]: + # The actual orchestration and node preflight run unchanged. The server is + # an external stand-in that records both observed clean markers at launch. + server = """import json, os, pathlib +p = pathlib.Path(os.environ["TEST_STATE"]) +(p / ("server-" + os.environ["SLURM_PROCID"])).write_text(json.dumps([ + (p / "clean-0").exists(), (p / "clean-1").exists()])) +raise SystemExit(int(os.environ["TEST_SERVER_STATUS"])) +""" + return [ + "bash", + "-c", + 'source "$1" --validation-only; shift; run_amd_multinode_after_preflight "$@"', + "test", + str(LIBRARY), + "node-0,node-1", + "2", + str(PREFLIGHT), + "name=^container_test_", + skip, + "env", + f"TEST_SERVER_STATUS={server_status}", + "python3", + "-c", + server, + ] + + +def test_slow_peer_finishes_gpu_preflight_before_any_server_starts(cluster) -> None: + state, env, logs = cluster + with (state / "output").open("w+") as output: + proc = subprocess.Popen(_command(), env=env, stdout=output, stderr=output) + try: + deadline = time.monotonic() + 10 + while not ((state / "waiting").exists() and (state / "clean-0").exists()): + assert proc.poll() is None, (state / "output").read_text() + assert time.monotonic() < deadline + time.sleep(0.02) + assert not list(state.glob("server-*")) + (state / "release").touch() + assert proc.wait(timeout=10) == 0 + finally: + if proc.poll() is None: + proc.kill() + proc.wait() + for rank in range(2): + assert json.loads((state / f"server-{rank}").read_text()) == [True, True] + assert "GPUs clean" in (logs / f"preflight_node-{rank}.log").read_text() + + +def test_failed_gpu_guard_prevents_all_servers_and_retains_diagnostics(cluster) -> None: + state, env, logs = cluster + completed = subprocess.run( + _command(), + env={**env, "TEST_MODE": "failure"}, + capture_output=True, + text=True, + check=False, + timeout=30, + ) + assert completed.returncode == 1 + assert not list(state.glob("server-*")) + assert len((state / "probes-1").read_text().splitlines()) == 90 + assert "GPUs still draining" in (logs / "preflight_node-1.log").read_text() + assert "no server containers launched" in completed.stderr + + +def test_explicit_gpu_skip_still_precleans_and_preserves_server_failure( + cluster, +) -> None: + state, env, logs = cluster + completed = subprocess.run( + _command(skip="1", server_status=23), + env=env, + capture_output=True, + text=True, + check=False, + timeout=10, + ) + assert completed.returncode == 23 + assert not list(state.glob("probes-*")) + assert len(list(state.glob("server-*"))) == 2 + for rank in range(2): + assert (state / f"docker-{rank}").read_text().splitlines() == [ + "ps -aq --filter name=^container_test_", + "ps -aq", + "ps -aq", + ] + assert ( + "skipping GPU pre-flight" + in (logs / f"preflight_node-{rank}.log").read_text() + ) From 37731636833fac0c38144bae4082b13d44ee5f77 Mon Sep 17 00:00:00 2001 From: Cam Quilici Date: Wed, 16 Sep 2026 23:49:46 -0500 Subject: [PATCH 09/12] fix: clean owned SGLang processes on startup failure --- benchmarks/benchmark_lib.sh | 19 +++ .../multi_node/amd_utils/server_sglang.sh | 23 ++-- docs/recovery-results-procedures.md | 8 +- docs/recovery-results-procedures_zh.md | 7 +- utils/test_process_group_cleanup.py | 111 ++++++++++++++++++ 5 files changed, 151 insertions(+), 17 deletions(-) diff --git a/benchmarks/benchmark_lib.sh b/benchmarks/benchmark_lib.sh index eb039ef8f3..fb219570fc 100644 --- a/benchmarks/benchmark_lib.sh +++ b/benchmarks/benchmark_lib.sh @@ -87,6 +87,25 @@ stop_background_process_groups() { return "$cleanup_status" } +# EXIT handler for launchers that own explicit setsid groups and optionally one +# auxiliary daemon PID. Disable this handler before exiting to avoid recursion. +# Run auxiliary cleanup even when group cleanup fails, preserving the work code. +exit_after_background_process_cleanup() { + local work_status="$1" term_grace="$2" kill_grace="$3" auxiliary_pid="$4" + shift 4 + local final_status + trap - EXIT + if stop_background_process_groups "$work_status" "$term_grace" "$kill_grace" "$@"; then + final_status=0 + else + final_status=$? + fi + if [[ -n "$auxiliary_pid" ]]; then + kill "$auxiliary_pid" 2>/dev/null || true + fi + exit "$final_status" +} + # Finish preflight on every allocated node before any server container starts its # peer-readiness deadline. A failed node prevents the entire serving step. run_amd_multinode_after_preflight() { diff --git a/benchmarks/multi_node/amd_utils/server_sglang.sh b/benchmarks/multi_node/amd_utils/server_sglang.sh index f247ac4158..39ab32be20 100755 --- a/benchmarks/multi_node/amd_utils/server_sglang.sh +++ b/benchmarks/multi_node/amd_utils/server_sglang.sh @@ -23,6 +23,12 @@ export BENCH_MAX_CONC_VALUE source $SGLANG_WS_PATH/setup_deps.sh source $SGLANG_WS_PATH/env.sh +# Install before starting UMBP or serving processes. Early readiness failures must +# close the same owned groups as normal completion, including orphaned workers. +SGLANG_OWNED_PGIDS=() +UMBP_SA_PID="" +trap 'exit_after_background_process_cleanup "$?" 30 5 "$UMBP_SA_PID" "${SGLANG_OWNED_PGIDS[@]}"' EXIT + host_ip=$(ip route get 1.1.1.1 | awk '/src/ {print $7}') host_name=$(hostname) @@ -686,7 +692,6 @@ elif [[ "$KV_OFFLOADING" != "none" && "$KV_OFFLOAD_BACKEND" == umbp-linker* ]]; "$UMBP_SA_BIN" "$UMBP_STANDALONE_ADDRESS" > "$UMBP_SA_LOG" 2>&1 & UMBP_SA_PID=$! echo "[UMBP] standalone server PID: $UMBP_SA_PID" - trap '[[ -n "${UMBP_SA_PID:-}" ]] && kill "$UMBP_SA_PID" 2>/dev/null || true' EXIT # Three waits, all bounded by wall time rather than by a guess at how # fast this node is. Bind time for a 549 GB tier measured 120 s on @@ -960,6 +965,7 @@ if [ "$NODE_RANK" -eq 0 ]; then set +x prefill0_pid=$! prefill0_pgid=$prefill0_pid + SGLANG_OWNED_PGIDS+=("$prefill0_pgid") fi echo "Waiting for all prefill and decode servers to be up . . ." @@ -1017,6 +1023,7 @@ if [ "$NODE_RANK" -eq 0 ]; then set +x proxy_pid=$! proxy_pgid=$proxy_pid + SGLANG_OWNED_PGIDS+=("$proxy_pgid") HEALTH_BARRIER_CMD="python3 $SGLANG_WS_PATH/sync.py barrier \ --node-ips ${NODE0_ADDR} \ @@ -1286,12 +1293,6 @@ print(json.dumps(json.loads(sys.stdin.read())))' <<<"$_val")" || { if [[ "${EVAL_FAILED:-0}" -eq 1 && "$node_exit_status" -eq 0 ]]; then node_exit_status=1 fi - if [[ "$DRY_RUN" -eq 0 ]]; then - # The router and prefill may retain TERM-resistant tokenizer workers. - # Keep the benchmark/eval status even when teardown also fails. - stop_background_process_groups "$node_exit_status" 30 5 "$proxy_pgid" "$prefill0_pgid" - exit $? - fi exit "$node_exit_status" elif [ "$NODE_RANK" -gt 0 ] && [ "$NODE_RANK" -lt "$NODE_OFFSET" ]; then @@ -1337,6 +1338,7 @@ elif [ "$NODE_RANK" -gt 0 ] && [ "$NODE_RANK" -lt "$NODE_OFFSET" ]; then set +x prefill_pid=$! prefill_pgid=$prefill_pid + SGLANG_OWNED_PGIDS+=("$prefill_pgid") fi echo "Waiting for proxy server to be up..." @@ -1366,8 +1368,7 @@ elif [ "$NODE_RANK" -gt 0 ] && [ "$NODE_RANK" -lt "$NODE_OFFSET" ]; then echo "Killing the rank $NODE_RANK prefill server" if [[ "$DRY_RUN" -eq 0 ]]; then - stop_background_process_groups 0 30 5 "$prefill_pgid" - exit $? + exit 0 fi else @@ -1455,6 +1456,7 @@ else set +x decode_pid=$! decode_pgid=$decode_pid + SGLANG_OWNED_PGIDS+=("$decode_pgid") fi echo "Waiting for proxy server to be up..." @@ -1483,8 +1485,7 @@ else echo "Killing the rank $RANK decode server" if [[ "$DRY_RUN" -eq 0 ]]; then - stop_background_process_groups 0 30 5 "$decode_pgid" - exit $? + exit 0 fi fi diff --git a/docs/recovery-results-procedures.md b/docs/recovery-results-procedures.md index 438fba05f6..8b8ac153ac 100644 --- a/docs/recovery-results-procedures.md +++ b/docs/recovery-results-procedures.md @@ -428,15 +428,17 @@ This evidence is the completion gate. “Workflow green” without artifact iden ### AMD multi-node SGLang teardown -After benchmark/eval work and result staging, the AMD SGLang launcher sends TERM -only to its recorded `setsid` process groups. It allows 30 seconds for graceful +On exit, including a failed startup/readiness check, the AMD SGLang launcher sends +TERM only to its recorded `setsid` process groups. Normal completion stages results +before this cleanup. It allows 30 seconds for graceful exit, then sends KILL to surviving groups and checks for exit for another five seconds. This handles orphaned or TERM-resistant workers that otherwise hold log pipes open. These cleanup deadlines do not change profiling, evaluation, or server readiness deadlines. A failed client retains its exit status; unresolved cleanup fails an otherwise successful node. Kernel-blocked processes may still require separately authorized node repair. Do not change or discard completed metrics to -work around teardown failures. +work around teardown failures. A single EXIT handler owns group cleanup and the +existing UMBP standalone PID cleanup; the latter still runs if group cleanup fails. ### AMD multi-node GPU preflight coordination diff --git a/docs/recovery-results-procedures_zh.md b/docs/recovery-results-procedures_zh.md index 31ec736403..1427077bb6 100644 --- a/docs/recovery-results-procedures_zh.md +++ b/docs/recovery-results-procedures_zh.md @@ -428,13 +428,14 @@ Remaining durable fix: ### AMD 多节点 SGLang 清理 -基准测试或评估完成并暂存结果后,AMD SGLang 启动器仅向其记录的 `setsid` -进程组发送 TERM,等待最多 30 秒。随后向仍存活的进程组发送 KILL,再等待最多 +退出时(包括启动或就绪检查失败),AMD SGLang 启动器仅向其记录的 `setsid` +进程组发送 TERM,等待最多 30 秒。正常完成时,先暂存结果再进行清理。随后向仍存活的进程组发送 KILL,再等待最多 5 秒并检查退出状态。这可以清理已成为孤儿进程或忽略 TERM 的工作进程,避免其 持续占用日志管道。这些清理期限不会改变性能采集、评估或服务器就绪检查的期限。 客户端失败时保留原退出码;若客户端成功但清理仍未完成,则节点任务失败。 内核阻塞的进程仍可能需要另行授权的节点修复。不要为绕过清理失败而修改或丢弃 -已完成的指标。 +已完成的指标。单一 EXIT 处理器统一负责进程组清理和现有 UMBP 独立进程 PID +清理;即使进程组清理失败,后者仍会执行。 ### AMD 多节点 GPU 预检协调 diff --git a/utils/test_process_group_cleanup.py b/utils/test_process_group_cleanup.py index 186fb38bf0..d01e6ae796 100644 --- a/utils/test_process_group_cleanup.py +++ b/utils/test_process_group_cleanup.py @@ -189,3 +189,114 @@ def test_signal_failure_is_bounded_and_preserves_work_failure( assert leader.poll() is None finally: leader.kill() + + +@pytest.mark.parametrize("work_status", [0, 17]) +def test_exit_handler_cleans_failed_startup_or_normal_work_and_umbp( + tmp_path: Path, work_status: int +) -> None: + ready = tmp_path / "group-ready" + auxiliary_ready = tmp_path / "auxiliary-ready" + auxiliary_stopped = tmp_path / "auxiliary-stopped" + worker_code = ( + "import signal,pathlib,sys,time; " + "signal.signal(signal.SIGTERM, signal.SIG_IGN); " + "pathlib.Path(sys.argv[1]).touch(); " + 'print("retained server output", flush=True); time.sleep(60)' + ) + auxiliary_code = ( + "import signal,pathlib,sys,time; " + "signal.signal(signal.SIGTERM, lambda *_: (pathlib.Path(sys.argv[2]).touch(), sys.exit(0))); " + "pathlib.Path(sys.argv[1]).touch(); time.sleep(60)" + ) + with ( + subprocess.Popen( + [sys.executable, "-c", worker_code, str(ready)], + start_new_session=True, + stdout=subprocess.PIPE, + stderr=subprocess.STDOUT, + text=True, + ) as worker, + subprocess.Popen( + [ + sys.executable, + "-c", + auxiliary_code, + str(auxiliary_ready), + str(auxiliary_stopped), + ] + ) as auxiliary, + ): + try: + await_file(ready) + await_file(auxiliary_ready) + completed = subprocess.run( + [ + "bash", + "-c", + ( + 'source "$1" --validation-only; ' + 'owned_groups=("$3"); auxiliary_pid=$4; ' + 'trap \'exit_after_background_process_cleanup "$?" 1 1 ' + '"$auxiliary_pid" "${owned_groups[@]}"\' EXIT; ' + 'bash -c "exit $2"; exit $?' + ), + "startup", + str(LIBRARY), + str(work_status), + str(worker.pid), + str(auxiliary.pid), + ], + capture_output=True, + text=True, + check=False, + timeout=7, + ) + assert completed.returncode == work_status, completed.stderr + assert "force-stopping owned process groups" in completed.stdout + # EOF proves the failed startup cannot retain the outer tee pipe. + output, _ = worker.communicate(timeout=2) + assert output == "retained server output\n" + assert worker.returncode == -signal.SIGKILL + assert auxiliary.wait(timeout=2) == 0 + assert auxiliary_stopped.exists() + finally: + if worker.poll() is None: + worker.kill() + if auxiliary.poll() is None: + auxiliary.kill() + + +@pytest.mark.parametrize("work_status, expected", [(0, 1), (17, 17)]) +def test_exit_handler_runs_auxiliary_cleanup_even_if_group_cleanup_fails( + work_status: int, expected: int +) -> None: + with subprocess.Popen( + [sys.executable, "-c", "import time; time.sleep(60)"] + ) as auxiliary: + try: + completed = subprocess.run( + [ + "bash", + "-c", + ( + 'source "$1" --validation-only; auxiliary_pid=$3; ' + 'trap \'exit_after_background_process_cleanup "$?" 1 1 ' + '"$auxiliary_pid" 1\' EXIT; exit "$2"' + ), + "startup", + str(LIBRARY), + str(work_status), + str(auxiliary.pid), + ], + capture_output=True, + text=True, + check=False, + timeout=5, + ) + assert completed.returncode == expected + assert "refusing unsafe process-group cleanup" in completed.stderr + assert auxiliary.wait(timeout=2) == -signal.SIGTERM + finally: + if auxiliary.poll() is None: + auxiliary.kill() From 3b4751cd1c81df55e20a526187c8ecf03ffefe15 Mon Sep 17 00:00:00 2001 From: Cam Quilici Date: Thu, 17 Sep 2026 09:37:17 -0500 Subject: [PATCH 10/12] fix: exclude failed AgentX rows from reusable artifact identities --- docs/results-and-ingestion.md | 2 + docs/results-and-ingestion_zh.md | 2 + infx/results/artifacts.py | 47 +++++++- utils/test_failed_agentic_artifacts.py | 150 +++++++++++++++++++++++++ 4 files changed, 197 insertions(+), 4 deletions(-) create mode 100644 utils/test_failed_agentic_artifacts.py diff --git a/docs/results-and-ingestion.md b/docs/results-and-ingestion.md index b6791cf6b2..6ec5934814 100644 --- a/docs/results-and-ingestion.md +++ b/docs/results-and-ingestion.md @@ -178,6 +178,8 @@ raw tree: results/**, excluding inputs.json and profile_export_raw.jso The aggregate artifact matches the `bmk_*` collection pattern and therefore also appears as a row in `results_bmk/agg_bmk.json`. The raw sibling is not fed to `collect_results.py`. InferenceX-app pairs `bmk_agentic_` with `agentic_` after stripping `bmk_` and `agentic_`. For files named with `_concN.json`, the concurrency is part of trace-sibling lookup. +AgentX reuse validation excludes explicitly failed rows with numeric zero successful requests and a finite, nonnegative integer-valued numeric total, consistent with ingestion excluding failed runs. A failed-only point artifact may lack its raw sibling; neither artifact is deleted. Successful, mixed, empty, unknown or malformed results retain existing identity and raw-sibling checks. A wholly failed source with no reusable results is still rejected, and excluding diagnostic failure rows does not establish full-sweep coverage. + Server logs are separate `server_logs_` artifacts. The app uses the fully stripped suffix fallback so AgentX rows can find a server log even though the log artifact has no `agentic_` prefix. Ordinary single-node AgentX runs enable the shared GPU power monitor by default. diff --git a/docs/results-and-ingestion_zh.md b/docs/results-and-ingestion_zh.md index 38917f66d7..01b2aeca4d 100644 --- a/docs/results-and-ingestion_zh.md +++ b/docs/results-and-ingestion_zh.md @@ -177,6 +177,8 @@ raw tree: results/**, excluding inputs.json and profile_export_raw.jso 聚合工件匹配 `bmk_*` 收集模式,因此也会成为 `results_bmk/agg_bmk.json` 中的一条记录。原始同级工件不会交给 `collect_results.py`。InferenceX-app 移除 `bmk_` 和 `agentic_` 后缀前缀,将 `bmk_agentic_` 与 `agentic_` 配对。对于以 `_concN.json` 命名的文件,并发也参与 trace 同级工件查找。 +AgentX 复用校验会排除明确失败的结果行:成功请求数必须是数值零,总请求数必须是有限、非负、整数值的数值。这与入库时排除失败运行的规则一致。只包含此类失败结果的点工件可以没有原始同级工件;校验不会删除任何工件。成功、混合、空、未知或格式不合法的结果仍接受原有身份及原始同级工件检查。没有任何可复用结果的全失败来源仍会被拒绝,排除失败诊断行本身不证明完整 sweep 已覆盖。 + 服务器日志是单独的 `server_logs_` 工件。应用会使用完全移除前缀后的后缀作为回退,从而让 AgentX 记录找到不含 `agentic_` 前缀的日志工件。 普通单节点 AgentX 提交默认启用共享 GPU 功耗监控。 diff --git a/infx/results/artifacts.py b/infx/results/artifacts.py index 3859dd4c08..764b233484 100644 --- a/infx/results/artifacts.py +++ b/infx/results/artifacts.py @@ -3,6 +3,7 @@ from __future__ import annotations import json +import math from collections import Counter from collections.abc import Iterable from pathlib import Path @@ -187,10 +188,42 @@ def agentic_keys_from_paths(paths: Iterable[Path]) -> list[tuple[Any, ...]]: return [ agentic_key(row) for _, row in json_rows(paths) - if row.get("scenario_type") == "agentic-coding" + if row.get("scenario_type") == "agentic-coding" and not failed_agentic_row(row) ] +def failed_agentic_row(row: Any) -> bool: + """Identify explicit zero-success results without treating unknown counts as failures.""" + if not isinstance(row, dict) or row.get("scenario_type") != "agentic-coding": + return False + successful = row.get("num_requests_successful") + total = row.get("num_requests_total") + return ( + isinstance(successful, (int, float)) + and not isinstance(successful, bool) + and successful == 0 + and isinstance(total, (int, float)) + and not isinstance(total, bool) + and total >= 0 + and (isinstance(total, int) or (math.isfinite(total) and total.is_integer())) + ) + + +def failed_agentic_point_names(artifacts_dir: Path, paths: Iterable[Path]) -> set[str]: + """Find failed-only artifacts; mixed, empty and unknown payloads stay strict.""" + failed_names: set[str] = set() + other_names: set[str] = set() + for path in paths: + name = path.relative_to(artifacts_dir).parts[0].removeprefix("bmk_") + data = load_json(path) + rows = data if isinstance(data, list) else [data] + if rows and all(failed_agentic_row(row) for row in rows): + failed_names.add(name) + else: + other_names.add(name) + return failed_names - other_names + + def actual_agentic_keys(artifacts_dir: Path) -> set[tuple[Any, ...]]: """Build actual agentic identities from aggregate and point results.""" paths = list((artifacts_dir / "results_bmk").glob("*.json")) @@ -254,7 +287,8 @@ def validate_agentic_artifacts( artifacts_dir: Path, ) -> list[str]: """Validate agentic point, raw, and aggregate artifacts agree.""" - point_rows = agentic_keys_from_paths(agentic_point_files(artifacts_dir)) + point_paths = agentic_point_files(artifacts_dir) + point_rows = agentic_keys_from_paths(point_paths) errors = duplicate_identity_errors("agentic point", point_rows) results_bmk = artifacts_dir / "results_bmk" @@ -270,14 +304,19 @@ def validate_agentic_artifacts( ) point_names = { - path.relative_to(artifacts_dir).parts[0].removeprefix("bmk_") - for path in agentic_point_files(artifacts_dir) + path.relative_to(artifacts_dir).parts[0].removeprefix("bmk_") for path in point_paths } + # Failed attempts can upload a point stub before a raw profile exists. Keep + # their files, but neither require nor reject a raw partner for failed-only + # artifacts. A successful or unknown row still requires its original partner. + failed_names = failed_agentic_point_names(artifacts_dir, point_paths) + point_names -= failed_names raw_names = { path.name for path in artifacts_dir.iterdir() if path.is_dir() and path.name.startswith("agentic_") } + raw_names -= failed_names if point_names != raw_names: missing_raw = point_names - raw_names extra_raw = raw_names - point_names diff --git a/utils/test_failed_agentic_artifacts.py b/utils/test_failed_agentic_artifacts.py new file mode 100644 index 0000000000..f0140f8238 --- /dev/null +++ b/utils/test_failed_agentic_artifacts.py @@ -0,0 +1,150 @@ +from __future__ import annotations + +import json +import sys +from pathlib import Path + +import pytest + +from infx.workflows.validate_reusable_sweep_artifacts import validate_agentic_artifacts + + +def agentic_result() -> dict: + return { + "scenario_type": "agentic-coding", + "hw": "test-gpu", + "infmax_model_prefix": "test-model", + "framework": "sglang", + "precision": "fp8", + "tp": 1, + "conc": 16, + } + + +def write_agentic_artifacts(root: Path) -> None: + point = root / "bmk_agentic_success" + point.mkdir() + (point / "result.json").write_text(json.dumps(agentic_result())) + (root / "agentic_success").mkdir() + + +@pytest.mark.parametrize("total,raw_present", [(0, False), (349.0, True)]) +def test_agentic_reuse_keeps_success_and_excludes_failed_attempt_without_deleting_files( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, + capsys: pytest.CaptureFixture, + total: float, + raw_present: bool, +) -> None: + from infx.workflows.validate_reusable_sweep_artifacts import main + + write_agentic_artifacts(tmp_path) + failed = { + **agentic_result(), + "num_requests_successful": 0.0, + "num_requests_total": total, + } + point = tmp_path / "bmk_agentic_failed" + point.mkdir() + payload = json.dumps(failed) + (point / "result.json").write_text(payload) + if raw_present: + (tmp_path / "agentic_failed").mkdir() + aggregate = tmp_path / "results_bmk" + aggregate.mkdir() + aggregate_payload = json.dumps([agentic_result(), failed]) + (aggregate / "results.json").write_text(aggregate_payload) + monkeypatch.setattr(sys, "argv", ["validate", "--artifacts-dir", str(tmp_path)]) + + assert main() == 0 + assert capsys.readouterr().out == ( + "Reusable sweep artifacts validated: " + "0 fixed-sequence row(s), 1 agentic row(s), 0 eval row(s).\n" + ) + assert (point / "result.json").read_text() == payload + assert (aggregate / "results.json").read_text() == aggregate_payload + assert (tmp_path / "agentic_failed").exists() is raw_present + + +@pytest.mark.parametrize( + "counts", + [ + {"num_requests_successful": 1, "num_requests_total": 349}, + {"num_requests_successful": "0", "num_requests_total": 349}, + {"num_requests_successful": False, "num_requests_total": 349}, + {"num_requests_successful": None, "num_requests_total": 349}, + {"num_requests_successful": -1, "num_requests_total": 349}, + {"num_requests_successful": 0}, + {"num_requests_successful": 0, "num_requests_total": "349"}, + {"num_requests_successful": 0, "num_requests_total": False}, + {"num_requests_successful": 0, "num_requests_total": None}, + {"num_requests_successful": 0, "num_requests_total": -1}, + {"num_requests_successful": 0, "num_requests_total": 0.5}, + {"num_requests_successful": 0, "num_requests_total": float("inf")}, + {"num_requests_successful": 0, "num_requests_total": float("nan")}, + ], +) +def test_agentic_nonfailed_or_unknown_counts_still_require_raw_artifacts( + tmp_path: Path, + counts: dict, +) -> None: + point = tmp_path / "bmk_agentic_unverified" + point.mkdir() + (point / "result.json").write_text(json.dumps({**agentic_result(), **counts})) + + assert validate_agentic_artifacts(tmp_path) == [ + "missing raw agentic artifact dir: agentic_unverified" + ] + + +@pytest.mark.parametrize( + "other_payload", + [ + [], + None, + ["unrecognized"], + [agentic_result()], + ], +) +def test_failed_agentic_directory_with_other_payload_remains_strict( + tmp_path: Path, + other_payload: object, +) -> None: + point = tmp_path / "bmk_agentic_mixed" + point.mkdir() + failed = { + **agentic_result(), + "num_requests_successful": 0, + "num_requests_total": 12, + } + (point / "failed.json").write_text(json.dumps(failed)) + (point / "other.json").write_text(json.dumps(other_payload)) + + assert validate_agentic_artifacts(tmp_path) == [ + "missing raw agentic artifact dir: agentic_mixed" + ] + + +def test_wholly_failed_agentic_sweep_has_no_reusable_results( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, + capsys: pytest.CaptureFixture, +) -> None: + from infx.workflows.validate_reusable_sweep_artifacts import main + + point = tmp_path / "bmk_agentic_failed" + point.mkdir() + failed = { + **agentic_result(), + "num_requests_successful": 0, + "num_requests_total": 12, + } + (point / "result.json").write_text(json.dumps(failed)) + monkeypatch.setattr(sys, "argv", ["validate", "--artifacts-dir", str(tmp_path)]) + + assert main() == 1 + assert ( + "no reusable benchmark, agentic, or eval result rows found" + in capsys.readouterr().err + ) + assert (point / "result.json").is_file() From f7533264587120d79cbe3e7dfc2312c59c058978 Mon Sep 17 00:00:00 2001 From: Cam Quilici Date: Thu, 17 Sep 2026 09:41:46 -0500 Subject: [PATCH 11/12] Revert "fix: exclude failed AgentX rows from reusable artifact identities" This reverts commit 3b4751cd1c81df55e20a526187c8ecf03ffefe15. --- docs/results-and-ingestion.md | 2 - docs/results-and-ingestion_zh.md | 2 - infx/results/artifacts.py | 47 +------- utils/test_failed_agentic_artifacts.py | 150 ------------------------- 4 files changed, 4 insertions(+), 197 deletions(-) delete mode 100644 utils/test_failed_agentic_artifacts.py diff --git a/docs/results-and-ingestion.md b/docs/results-and-ingestion.md index 6ec5934814..b6791cf6b2 100644 --- a/docs/results-and-ingestion.md +++ b/docs/results-and-ingestion.md @@ -178,8 +178,6 @@ raw tree: results/**, excluding inputs.json and profile_export_raw.jso The aggregate artifact matches the `bmk_*` collection pattern and therefore also appears as a row in `results_bmk/agg_bmk.json`. The raw sibling is not fed to `collect_results.py`. InferenceX-app pairs `bmk_agentic_` with `agentic_` after stripping `bmk_` and `agentic_`. For files named with `_concN.json`, the concurrency is part of trace-sibling lookup. -AgentX reuse validation excludes explicitly failed rows with numeric zero successful requests and a finite, nonnegative integer-valued numeric total, consistent with ingestion excluding failed runs. A failed-only point artifact may lack its raw sibling; neither artifact is deleted. Successful, mixed, empty, unknown or malformed results retain existing identity and raw-sibling checks. A wholly failed source with no reusable results is still rejected, and excluding diagnostic failure rows does not establish full-sweep coverage. - Server logs are separate `server_logs_` artifacts. The app uses the fully stripped suffix fallback so AgentX rows can find a server log even though the log artifact has no `agentic_` prefix. Ordinary single-node AgentX runs enable the shared GPU power monitor by default. diff --git a/docs/results-and-ingestion_zh.md b/docs/results-and-ingestion_zh.md index 01b2aeca4d..38917f66d7 100644 --- a/docs/results-and-ingestion_zh.md +++ b/docs/results-and-ingestion_zh.md @@ -177,8 +177,6 @@ raw tree: results/**, excluding inputs.json and profile_export_raw.jso 聚合工件匹配 `bmk_*` 收集模式,因此也会成为 `results_bmk/agg_bmk.json` 中的一条记录。原始同级工件不会交给 `collect_results.py`。InferenceX-app 移除 `bmk_` 和 `agentic_` 后缀前缀,将 `bmk_agentic_` 与 `agentic_` 配对。对于以 `_concN.json` 命名的文件,并发也参与 trace 同级工件查找。 -AgentX 复用校验会排除明确失败的结果行:成功请求数必须是数值零,总请求数必须是有限、非负、整数值的数值。这与入库时排除失败运行的规则一致。只包含此类失败结果的点工件可以没有原始同级工件;校验不会删除任何工件。成功、混合、空、未知或格式不合法的结果仍接受原有身份及原始同级工件检查。没有任何可复用结果的全失败来源仍会被拒绝,排除失败诊断行本身不证明完整 sweep 已覆盖。 - 服务器日志是单独的 `server_logs_` 工件。应用会使用完全移除前缀后的后缀作为回退,从而让 AgentX 记录找到不含 `agentic_` 前缀的日志工件。 普通单节点 AgentX 提交默认启用共享 GPU 功耗监控。 diff --git a/infx/results/artifacts.py b/infx/results/artifacts.py index 764b233484..3859dd4c08 100644 --- a/infx/results/artifacts.py +++ b/infx/results/artifacts.py @@ -3,7 +3,6 @@ from __future__ import annotations import json -import math from collections import Counter from collections.abc import Iterable from pathlib import Path @@ -188,42 +187,10 @@ def agentic_keys_from_paths(paths: Iterable[Path]) -> list[tuple[Any, ...]]: return [ agentic_key(row) for _, row in json_rows(paths) - if row.get("scenario_type") == "agentic-coding" and not failed_agentic_row(row) + if row.get("scenario_type") == "agentic-coding" ] -def failed_agentic_row(row: Any) -> bool: - """Identify explicit zero-success results without treating unknown counts as failures.""" - if not isinstance(row, dict) or row.get("scenario_type") != "agentic-coding": - return False - successful = row.get("num_requests_successful") - total = row.get("num_requests_total") - return ( - isinstance(successful, (int, float)) - and not isinstance(successful, bool) - and successful == 0 - and isinstance(total, (int, float)) - and not isinstance(total, bool) - and total >= 0 - and (isinstance(total, int) or (math.isfinite(total) and total.is_integer())) - ) - - -def failed_agentic_point_names(artifacts_dir: Path, paths: Iterable[Path]) -> set[str]: - """Find failed-only artifacts; mixed, empty and unknown payloads stay strict.""" - failed_names: set[str] = set() - other_names: set[str] = set() - for path in paths: - name = path.relative_to(artifacts_dir).parts[0].removeprefix("bmk_") - data = load_json(path) - rows = data if isinstance(data, list) else [data] - if rows and all(failed_agentic_row(row) for row in rows): - failed_names.add(name) - else: - other_names.add(name) - return failed_names - other_names - - def actual_agentic_keys(artifacts_dir: Path) -> set[tuple[Any, ...]]: """Build actual agentic identities from aggregate and point results.""" paths = list((artifacts_dir / "results_bmk").glob("*.json")) @@ -287,8 +254,7 @@ def validate_agentic_artifacts( artifacts_dir: Path, ) -> list[str]: """Validate agentic point, raw, and aggregate artifacts agree.""" - point_paths = agentic_point_files(artifacts_dir) - point_rows = agentic_keys_from_paths(point_paths) + point_rows = agentic_keys_from_paths(agentic_point_files(artifacts_dir)) errors = duplicate_identity_errors("agentic point", point_rows) results_bmk = artifacts_dir / "results_bmk" @@ -304,19 +270,14 @@ def validate_agentic_artifacts( ) point_names = { - path.relative_to(artifacts_dir).parts[0].removeprefix("bmk_") for path in point_paths + path.relative_to(artifacts_dir).parts[0].removeprefix("bmk_") + for path in agentic_point_files(artifacts_dir) } - # Failed attempts can upload a point stub before a raw profile exists. Keep - # their files, but neither require nor reject a raw partner for failed-only - # artifacts. A successful or unknown row still requires its original partner. - failed_names = failed_agentic_point_names(artifacts_dir, point_paths) - point_names -= failed_names raw_names = { path.name for path in artifacts_dir.iterdir() if path.is_dir() and path.name.startswith("agentic_") } - raw_names -= failed_names if point_names != raw_names: missing_raw = point_names - raw_names extra_raw = raw_names - point_names diff --git a/utils/test_failed_agentic_artifacts.py b/utils/test_failed_agentic_artifacts.py deleted file mode 100644 index f0140f8238..0000000000 --- a/utils/test_failed_agentic_artifacts.py +++ /dev/null @@ -1,150 +0,0 @@ -from __future__ import annotations - -import json -import sys -from pathlib import Path - -import pytest - -from infx.workflows.validate_reusable_sweep_artifacts import validate_agentic_artifacts - - -def agentic_result() -> dict: - return { - "scenario_type": "agentic-coding", - "hw": "test-gpu", - "infmax_model_prefix": "test-model", - "framework": "sglang", - "precision": "fp8", - "tp": 1, - "conc": 16, - } - - -def write_agentic_artifacts(root: Path) -> None: - point = root / "bmk_agentic_success" - point.mkdir() - (point / "result.json").write_text(json.dumps(agentic_result())) - (root / "agentic_success").mkdir() - - -@pytest.mark.parametrize("total,raw_present", [(0, False), (349.0, True)]) -def test_agentic_reuse_keeps_success_and_excludes_failed_attempt_without_deleting_files( - tmp_path: Path, - monkeypatch: pytest.MonkeyPatch, - capsys: pytest.CaptureFixture, - total: float, - raw_present: bool, -) -> None: - from infx.workflows.validate_reusable_sweep_artifacts import main - - write_agentic_artifacts(tmp_path) - failed = { - **agentic_result(), - "num_requests_successful": 0.0, - "num_requests_total": total, - } - point = tmp_path / "bmk_agentic_failed" - point.mkdir() - payload = json.dumps(failed) - (point / "result.json").write_text(payload) - if raw_present: - (tmp_path / "agentic_failed").mkdir() - aggregate = tmp_path / "results_bmk" - aggregate.mkdir() - aggregate_payload = json.dumps([agentic_result(), failed]) - (aggregate / "results.json").write_text(aggregate_payload) - monkeypatch.setattr(sys, "argv", ["validate", "--artifacts-dir", str(tmp_path)]) - - assert main() == 0 - assert capsys.readouterr().out == ( - "Reusable sweep artifacts validated: " - "0 fixed-sequence row(s), 1 agentic row(s), 0 eval row(s).\n" - ) - assert (point / "result.json").read_text() == payload - assert (aggregate / "results.json").read_text() == aggregate_payload - assert (tmp_path / "agentic_failed").exists() is raw_present - - -@pytest.mark.parametrize( - "counts", - [ - {"num_requests_successful": 1, "num_requests_total": 349}, - {"num_requests_successful": "0", "num_requests_total": 349}, - {"num_requests_successful": False, "num_requests_total": 349}, - {"num_requests_successful": None, "num_requests_total": 349}, - {"num_requests_successful": -1, "num_requests_total": 349}, - {"num_requests_successful": 0}, - {"num_requests_successful": 0, "num_requests_total": "349"}, - {"num_requests_successful": 0, "num_requests_total": False}, - {"num_requests_successful": 0, "num_requests_total": None}, - {"num_requests_successful": 0, "num_requests_total": -1}, - {"num_requests_successful": 0, "num_requests_total": 0.5}, - {"num_requests_successful": 0, "num_requests_total": float("inf")}, - {"num_requests_successful": 0, "num_requests_total": float("nan")}, - ], -) -def test_agentic_nonfailed_or_unknown_counts_still_require_raw_artifacts( - tmp_path: Path, - counts: dict, -) -> None: - point = tmp_path / "bmk_agentic_unverified" - point.mkdir() - (point / "result.json").write_text(json.dumps({**agentic_result(), **counts})) - - assert validate_agentic_artifacts(tmp_path) == [ - "missing raw agentic artifact dir: agentic_unverified" - ] - - -@pytest.mark.parametrize( - "other_payload", - [ - [], - None, - ["unrecognized"], - [agentic_result()], - ], -) -def test_failed_agentic_directory_with_other_payload_remains_strict( - tmp_path: Path, - other_payload: object, -) -> None: - point = tmp_path / "bmk_agentic_mixed" - point.mkdir() - failed = { - **agentic_result(), - "num_requests_successful": 0, - "num_requests_total": 12, - } - (point / "failed.json").write_text(json.dumps(failed)) - (point / "other.json").write_text(json.dumps(other_payload)) - - assert validate_agentic_artifacts(tmp_path) == [ - "missing raw agentic artifact dir: agentic_mixed" - ] - - -def test_wholly_failed_agentic_sweep_has_no_reusable_results( - tmp_path: Path, - monkeypatch: pytest.MonkeyPatch, - capsys: pytest.CaptureFixture, -) -> None: - from infx.workflows.validate_reusable_sweep_artifacts import main - - point = tmp_path / "bmk_agentic_failed" - point.mkdir() - failed = { - **agentic_result(), - "num_requests_successful": 0, - "num_requests_total": 12, - } - (point / "result.json").write_text(json.dumps(failed)) - monkeypatch.setattr(sys, "argv", ["validate", "--artifacts-dir", str(tmp_path)]) - - assert main() == 1 - assert ( - "no reusable benchmark, agentic, or eval result rows found" - in capsys.readouterr().err - ) - assert (point / "result.json").is_file() From f4a3fbaaa4c58b3662b7a4ddba28e086b07db5d9 Mon Sep 17 00:00:00 2001 From: Cam Quilici Date: Thu, 17 Sep 2026 09:47:11 -0500 Subject: [PATCH 12/12] test: remove added AMD lifecycle and log staging tests --- utils/test_amd_multinode_preflight.py | 169 -------------- utils/test_amd_node_log_staging.py | 84 ------- utils/test_process_group_cleanup.py | 302 -------------------------- 3 files changed, 555 deletions(-) delete mode 100644 utils/test_amd_multinode_preflight.py delete mode 100644 utils/test_amd_node_log_staging.py delete mode 100644 utils/test_process_group_cleanup.py diff --git a/utils/test_amd_multinode_preflight.py b/utils/test_amd_multinode_preflight.py deleted file mode 100644 index 035d431bde..0000000000 --- a/utils/test_amd_multinode_preflight.py +++ /dev/null @@ -1,169 +0,0 @@ -from __future__ import annotations - -import json -import os -import shutil -import subprocess -import time -import uuid -from pathlib import Path - -import pytest - -REPO_ROOT = Path(__file__).resolve().parents[1] -LIBRARY = REPO_ROOT / "benchmarks/benchmark_lib.sh" -PREFLIGHT = REPO_ROOT / "benchmarks/multi_node/amd_utils/preflight_node.sh" - - -@pytest.fixture -def cluster(tmp_path: Path): - """Only Slurm, Docker, GPU telemetry, host naming, and the sleep clock are fake.""" - bin_dir = tmp_path / "bin" - bin_dir.mkdir() - scripts = { - "srun": """#!/usr/bin/env python3 -import os, subprocess, sys -args = sys.argv[1:] -while args and args[0].startswith("--"): - args.pop(0) -children = [subprocess.Popen(args, env={**os.environ, "SLURM_PROCID": str(rank)}) - for rank in range(2)] -statuses = [child.wait() for child in children] -sys.exit(next((status for status in statuses if status), 0)) -""", - "docker": """#!/usr/bin/env python3 -import os, pathlib, sys -with (pathlib.Path(os.environ["TEST_STATE"]) / ("docker-" + os.environ["SLURM_PROCID"])).open("a") as f: - f.write(" ".join(sys.argv[1:]) + "\\n") -""", - "hostname": '#!/bin/sh\necho "node-$SLURM_PROCID"\n', - "sleep": "#!/bin/sh\nexit 0\n", - "rocm-smi": """#!/usr/bin/env python3 -import os, pathlib, time -state = pathlib.Path(os.environ["TEST_STATE"]) -rank = os.environ["SLURM_PROCID"] -mode = os.environ["TEST_MODE"] -with (state / ("probes-" + rank)).open("a") as f: - f.write("probe\\n") -if rank == "1" and mode == "delayed" and not (state / "release").exists(): - (state / "waiting").touch() - time.sleep(0.02) - used = 92 -elif rank == "1" and mode == "failure": - used = 92 -else: - used = 0 - (state / ("clean-" + rank)).touch() -print(f"GPU[0] : GPU Memory Allocated (VRAM%): {used}") -""", - } - for name, source in scripts.items(): - path = bin_dir / name - path.write_text(source) - path.chmod(0o755) - job_id = f"infx-preflight-test-{uuid.uuid4().hex}" - env = { - **os.environ, - "PATH": f"{bin_dir}:{os.environ['PATH']}", - "DOCKER_CMD_DETECT": "DOCKER_CMD=docker", - "DI_REPO_DIR": str(REPO_ROOT), - "SLURM_JOB_ID": job_id, - "TEST_STATE": str(tmp_path), - "TEST_MODE": "delayed", - } - yield tmp_path, env, Path("/tmp") / f"slurm_job-{job_id}" - shutil.rmtree(Path("/tmp") / f"slurm_job-{job_id}", ignore_errors=True) - - -def _command(skip: str = "0", server_status: int = 0) -> list[str]: - # The actual orchestration and node preflight run unchanged. The server is - # an external stand-in that records both observed clean markers at launch. - server = """import json, os, pathlib -p = pathlib.Path(os.environ["TEST_STATE"]) -(p / ("server-" + os.environ["SLURM_PROCID"])).write_text(json.dumps([ - (p / "clean-0").exists(), (p / "clean-1").exists()])) -raise SystemExit(int(os.environ["TEST_SERVER_STATUS"])) -""" - return [ - "bash", - "-c", - 'source "$1" --validation-only; shift; run_amd_multinode_after_preflight "$@"', - "test", - str(LIBRARY), - "node-0,node-1", - "2", - str(PREFLIGHT), - "name=^container_test_", - skip, - "env", - f"TEST_SERVER_STATUS={server_status}", - "python3", - "-c", - server, - ] - - -def test_slow_peer_finishes_gpu_preflight_before_any_server_starts(cluster) -> None: - state, env, logs = cluster - with (state / "output").open("w+") as output: - proc = subprocess.Popen(_command(), env=env, stdout=output, stderr=output) - try: - deadline = time.monotonic() + 10 - while not ((state / "waiting").exists() and (state / "clean-0").exists()): - assert proc.poll() is None, (state / "output").read_text() - assert time.monotonic() < deadline - time.sleep(0.02) - assert not list(state.glob("server-*")) - (state / "release").touch() - assert proc.wait(timeout=10) == 0 - finally: - if proc.poll() is None: - proc.kill() - proc.wait() - for rank in range(2): - assert json.loads((state / f"server-{rank}").read_text()) == [True, True] - assert "GPUs clean" in (logs / f"preflight_node-{rank}.log").read_text() - - -def test_failed_gpu_guard_prevents_all_servers_and_retains_diagnostics(cluster) -> None: - state, env, logs = cluster - completed = subprocess.run( - _command(), - env={**env, "TEST_MODE": "failure"}, - capture_output=True, - text=True, - check=False, - timeout=30, - ) - assert completed.returncode == 1 - assert not list(state.glob("server-*")) - assert len((state / "probes-1").read_text().splitlines()) == 90 - assert "GPUs still draining" in (logs / "preflight_node-1.log").read_text() - assert "no server containers launched" in completed.stderr - - -def test_explicit_gpu_skip_still_precleans_and_preserves_server_failure( - cluster, -) -> None: - state, env, logs = cluster - completed = subprocess.run( - _command(skip="1", server_status=23), - env=env, - capture_output=True, - text=True, - check=False, - timeout=10, - ) - assert completed.returncode == 23 - assert not list(state.glob("probes-*")) - assert len(list(state.glob("server-*"))) == 2 - for rank in range(2): - assert (state / f"docker-{rank}").read_text().splitlines() == [ - "ps -aq --filter name=^container_test_", - "ps -aq", - "ps -aq", - ] - assert ( - "skipping GPU pre-flight" - in (logs / f"preflight_node-{rank}.log").read_text() - ) diff --git a/utils/test_amd_node_log_staging.py b/utils/test_amd_node_log_staging.py deleted file mode 100644 index d194665588..0000000000 --- a/utils/test_amd_node_log_staging.py +++ /dev/null @@ -1,84 +0,0 @@ -from __future__ import annotations - -import os -import subprocess -from pathlib import Path - -REPO_ROOT = Path(__file__).resolve().parents[1] -STAGE_SCRIPT = REPO_ROOT / "benchmarks/multi_node/amd_utils/stage_node_logs.sh" - - -def _stub_sudo(tmp_path: Path) -> dict[str, str]: - bin_dir = tmp_path / "bin" - bin_dir.mkdir() - sudo = bin_dir / "sudo" - sudo.write_text('#!/bin/sh\nexec "$@"\n') - sudo.chmod(0o755) - return {**os.environ, "PATH": f"{bin_dir}:{os.environ['PATH']}"} - - -def test_stage_node_logs_merges_prefill_and_decode_nodes(tmp_path: Path) -> None: - prefill = tmp_path / "prefill-node" - decode = tmp_path / "decode-node" - shared = tmp_path / "shared" - prefill.mkdir() - decode.mkdir() - (prefill / "prefill_host-a.log").write_text("prefill output\n") - (prefill / "server_host-a.log").write_text("frontend output\n") - (decode / "decode_host-b.log").write_text("decode output\n") - (decode / "server_host-b.log").write_text("decode wrapper output\n") - env = _stub_sudo(tmp_path) - - for node_logs in (prefill, decode): - subprocess.run( - ["bash", str(STAGE_SCRIPT), str(node_logs), str(shared)], - check=True, - env=env, - ) - - assert sorted(path.name for path in shared.iterdir()) == [ - "decode_host-b.log", - "prefill_host-a.log", - "server_host-a.log", - "server_host-b.log", - ] - assert (shared / "decode_host-b.log").read_text() == "decode output\n" - - -def test_stage_node_logs_rejects_a_node_without_logs(tmp_path: Path) -> None: - missing = tmp_path / "missing" - shared = tmp_path / "shared" - - completed = subprocess.run( - ["bash", str(STAGE_SCRIPT), str(missing), str(shared)], - check=False, - capture_output=True, - text=True, - ) - - assert completed.returncode == 1 - assert "no node-local logs found" in completed.stderr - assert not shared.exists() - - -def test_stage_node_logs_propagates_copy_failure(tmp_path: Path) -> None: - source = tmp_path / "node" - shared = tmp_path / "shared" - source.mkdir() - (source / "decode_host-b.log").write_text("decode output\n") - env = _stub_sudo(tmp_path) - cp = tmp_path / "bin" / "cp" - cp.write_text("#!/bin/sh\nexit 19\n") - cp.chmod(0o755) - - completed = subprocess.run( - ["bash", str(STAGE_SCRIPT), str(source), str(shared)], - check=False, - capture_output=True, - text=True, - env=env, - ) - - assert completed.returncode == 19 - assert "[logs] staged" not in completed.stdout - assert not (shared / "decode_host-b.log").exists() diff --git a/utils/test_process_group_cleanup.py b/utils/test_process_group_cleanup.py deleted file mode 100644 index d01e6ae796..0000000000 --- a/utils/test_process_group_cleanup.py +++ /dev/null @@ -1,302 +0,0 @@ -"""Real process/pipe regressions for post-benchmark group teardown.""" - -import os -import signal -import subprocess -import sys -import time -from pathlib import Path - -import pytest - -LIBRARY = Path(__file__).resolve().parents[1] / "benchmarks/benchmark_lib.sh" - - -def stop_groups(status: int, *groups: int) -> subprocess.CompletedProcess[str]: - # The client exits independently; the actual cleanup implementation must - # preserve its code even when signaling and polling report success. - return subprocess.run( - [ - "bash", - "-c", - ( - 'source "$1" --validation-only; shift; ' - 'client_status=$1; shift; bash -c "exit $client_status"; status=$?; ' - 'stop_background_process_groups "$status" 1 1 "$@"' - ), - "cleanup", - str(LIBRARY), - str(status), - *map(str, groups), - ], - check=False, - capture_output=True, - text=True, - timeout=8, - ) - - -def await_file(path: Path) -> None: - deadline = time.monotonic() + 5 - while not path.exists(): - assert time.monotonic() < deadline, "Child did not start" - time.sleep(0.01) - - -@pytest.mark.parametrize("leader_exits, client_status", [(True, 0), (False, 19)]) -def test_stubborn_descendant_releases_pipe_and_preserves_client_status( - tmp_path: Path, - leader_exits: bool, - client_status: int, -) -> None: - ready = tmp_path / "child-ready" - child_code = ( - "import signal,pathlib,sys,time; " - "signal.signal(signal.SIGTERM, signal.SIG_IGN); " - 'pathlib.Path(sys.argv[1]).touch(); print("worker output", flush=True); time.sleep(60)' - ) - leader_code = ( - "import subprocess,sys,time; " - 'subprocess.Popen([sys.executable,"-c",sys.argv[1],sys.argv[2]]); ' - 'time.sleep(0 if sys.argv[3]=="True" else 60)' - ) - with ( - subprocess.Popen( - [ - sys.executable, - "-c", - leader_code, - child_code, - str(ready), - str(leader_exits), - ], - start_new_session=True, - stdout=subprocess.PIPE, - stderr=subprocess.STDOUT, - text=True, - ) as leader, - subprocess.Popen( - [sys.executable, "-c", "import time; time.sleep(60)"] - ) as unrelated, - ): - try: - await_file(ready) - if leader_exits: - assert leader.wait(timeout=3) == 0 - result = stop_groups(client_status, leader.pid) - assert result.returncode == client_status, result.stderr - assert "force-stopping owned process groups" in result.stdout - # An orphan that retains stdout makes communicate hang even after - # its leader exited. This tests the original tee-pipe failure. - output, _ = leader.communicate(timeout=3) - assert "worker output" in output - assert unrelated.poll() is None - finally: - try: - os.killpg(leader.pid, signal.SIGKILL) - except ProcessLookupError: - pass - except PermissionError: - # macOS can retain a zombie-only process group owned by init. - pass - unrelated.terminate() - - -def test_graceful_group_gets_term_without_kill(tmp_path: Path) -> None: - ready = tmp_path / "ready" - stopped = tmp_path / "stopped" - code = ( - "import signal,pathlib,sys,time; " - "signal.signal(signal.SIGTERM, lambda *_: (pathlib.Path(sys.argv[2]).touch(), sys.exit(0))); " - "pathlib.Path(sys.argv[1]).touch(); time.sleep(60)" - ) - with subprocess.Popen( - [sys.executable, "-c", code, str(ready), str(stopped)], start_new_session=True - ) as leader: - try: - await_file(ready) - result = stop_groups(0, leader.pid) - assert result.returncode == 0, result.stderr - assert leader.wait(timeout=2) == 0 - assert stopped.exists() - assert "force-stopping" not in result.stdout - finally: - if leader.poll() is None: - leader.kill() - - -@pytest.mark.parametrize("client_status, expected", [(0, 1), (23, 23)]) -def test_refused_cleanup_cannot_hide_work_failure_or_report_success( - client_status: int, expected: int -) -> None: - result = stop_groups(client_status, 1) - assert result.returncode == expected - assert "refusing unsafe process-group cleanup" in result.stderr - - -def test_refuses_callers_own_group() -> None: - result = subprocess.run( - [ - "bash", - "-c", - ( - 'source "$1" --validation-only; ' - 'group=$(ps -o pgid= -p $$ | tr -d " "); stop_background_process_groups 0 1 1 "$group"' - ), - "cleanup", - str(LIBRARY), - ], - start_new_session=True, - check=False, - capture_output=True, - text=True, - timeout=5, - ) - assert result.returncode == 1 - assert "refusing unsafe process-group cleanup" in result.stderr - - -@pytest.mark.parametrize("client_status, expected", [(0, 1), (23, 23)]) -def test_signal_failure_is_bounded_and_preserves_work_failure( - client_status: int, expected: int -) -> None: - with subprocess.Popen( - [sys.executable, "-c", "import time; time.sleep(60)"], start_new_session=True - ) as leader: - try: - # Only mock the OS signal collaborator: the real liveness checks, - # grace deadlines, and status selection all run against a live group. - result = subprocess.run( - [ - "bash", - "-c", - ( - 'source "$1" --validation-only; ' - 'kill() { return 1; }; stop_background_process_groups "$2" 1 1 "$3"' - ), - "cleanup", - str(LIBRARY), - str(client_status), - str(leader.pid), - ], - check=False, - capture_output=True, - text=True, - timeout=6, - ) - assert result.returncode == expected - assert "still alive after KILL grace" in result.stderr - assert leader.poll() is None - finally: - leader.kill() - - -@pytest.mark.parametrize("work_status", [0, 17]) -def test_exit_handler_cleans_failed_startup_or_normal_work_and_umbp( - tmp_path: Path, work_status: int -) -> None: - ready = tmp_path / "group-ready" - auxiliary_ready = tmp_path / "auxiliary-ready" - auxiliary_stopped = tmp_path / "auxiliary-stopped" - worker_code = ( - "import signal,pathlib,sys,time; " - "signal.signal(signal.SIGTERM, signal.SIG_IGN); " - "pathlib.Path(sys.argv[1]).touch(); " - 'print("retained server output", flush=True); time.sleep(60)' - ) - auxiliary_code = ( - "import signal,pathlib,sys,time; " - "signal.signal(signal.SIGTERM, lambda *_: (pathlib.Path(sys.argv[2]).touch(), sys.exit(0))); " - "pathlib.Path(sys.argv[1]).touch(); time.sleep(60)" - ) - with ( - subprocess.Popen( - [sys.executable, "-c", worker_code, str(ready)], - start_new_session=True, - stdout=subprocess.PIPE, - stderr=subprocess.STDOUT, - text=True, - ) as worker, - subprocess.Popen( - [ - sys.executable, - "-c", - auxiliary_code, - str(auxiliary_ready), - str(auxiliary_stopped), - ] - ) as auxiliary, - ): - try: - await_file(ready) - await_file(auxiliary_ready) - completed = subprocess.run( - [ - "bash", - "-c", - ( - 'source "$1" --validation-only; ' - 'owned_groups=("$3"); auxiliary_pid=$4; ' - 'trap \'exit_after_background_process_cleanup "$?" 1 1 ' - '"$auxiliary_pid" "${owned_groups[@]}"\' EXIT; ' - 'bash -c "exit $2"; exit $?' - ), - "startup", - str(LIBRARY), - str(work_status), - str(worker.pid), - str(auxiliary.pid), - ], - capture_output=True, - text=True, - check=False, - timeout=7, - ) - assert completed.returncode == work_status, completed.stderr - assert "force-stopping owned process groups" in completed.stdout - # EOF proves the failed startup cannot retain the outer tee pipe. - output, _ = worker.communicate(timeout=2) - assert output == "retained server output\n" - assert worker.returncode == -signal.SIGKILL - assert auxiliary.wait(timeout=2) == 0 - assert auxiliary_stopped.exists() - finally: - if worker.poll() is None: - worker.kill() - if auxiliary.poll() is None: - auxiliary.kill() - - -@pytest.mark.parametrize("work_status, expected", [(0, 1), (17, 17)]) -def test_exit_handler_runs_auxiliary_cleanup_even_if_group_cleanup_fails( - work_status: int, expected: int -) -> None: - with subprocess.Popen( - [sys.executable, "-c", "import time; time.sleep(60)"] - ) as auxiliary: - try: - completed = subprocess.run( - [ - "bash", - "-c", - ( - 'source "$1" --validation-only; auxiliary_pid=$3; ' - 'trap \'exit_after_background_process_cleanup "$?" 1 1 ' - '"$auxiliary_pid" 1\' EXIT; exit "$2"' - ), - "startup", - str(LIBRARY), - str(work_status), - str(auxiliary.pid), - ], - capture_output=True, - text=True, - check=False, - timeout=5, - ) - assert completed.returncode == expected - assert "refusing unsafe process-group cleanup" in completed.stderr - assert auxiliary.wait(timeout=2) == -signal.SIGTERM - finally: - if auxiliary.poll() is None: - auxiliary.kill()