From 0316f6f3970e0e140427e493bb90dbc5b25e4f1c Mon Sep 17 00:00:00 2001 From: TianZhuo <2770730562@qq.com> Date: Thu, 20 Aug 2026 16:04:58 +0800 Subject: [PATCH 1/2] [Refactor] Drop FT_COMMUNICATION_OPS_ABORT_TIMEOUT_MS env var Signed-off-by: TianZhuo <2770730562@qq.com> --- .../test_fault_tolerance_e2e.py | 4 ++-- vllm_ascend/ascend_config.py | 22 +++++++++---------- vllm_ascend/envs.py | 6 ----- vllm_ascend/worker/worker.py | 20 ++++++++++++++++- 4 files changed, 32 insertions(+), 20 deletions(-) diff --git a/tests/e2e/pull_request/four_card/fault_tolerance/test_fault_tolerance_e2e.py b/tests/e2e/pull_request/four_card/fault_tolerance/test_fault_tolerance_e2e.py index bcb53477323d..c105eaecbb3b 100644 --- a/tests/e2e/pull_request/four_card/fault_tolerance/test_fault_tolerance_e2e.py +++ b/tests/e2e/pull_request/four_card/fault_tolerance/test_fault_tolerance_e2e.py @@ -27,7 +27,7 @@ # - NPU: HCSP operator timeout detects the dead peer. # - Deadline (45s): slowest fallback (30s) + margin. CPU_DISTRIBUTED_TIMEOUT_S = 30 -FT_COMMUNICATION_OPS_ABORT_TIMEOUT_MS = 15 +FT_COMMUNICATION_ABORT_TIMEOUT_S = 15 FAULT_DETECTION_DEADLINE_S = 45 @@ -129,7 +129,7 @@ def _ft_server_args() -> list[str]: "--fault-tolerance-config", '{"engine_recovery_timeout_sec": 120}', "--additional-config", - f'{{"ft_communication_ops_abort_timeout_ms": {FT_COMMUNICATION_OPS_ABORT_TIMEOUT_MS}}}', + f'{{"ft_communication_abort_timeout": {FT_COMMUNICATION_ABORT_TIMEOUT_S}}}', ] diff --git a/vllm_ascend/ascend_config.py b/vllm_ascend/ascend_config.py index c5c6fc058b37..fe7e7e256569 100644 --- a/vllm_ascend/ascend_config.py +++ b/vllm_ascend/ascend_config.py @@ -200,20 +200,20 @@ def __init__(self, vllm_config: "VllmConfig"): "VLLM_ASCEND_FUSION_OP_TRANSPOSE_KV_CACHE_BY_BLOCK", ascend_envs.VLLM_ASCEND_FUSION_OP_TRANSPOSE_KV_CACHE_BY_BLOCK, ) - self.ft_communication_ops_abort_timeout_ms = self._get_config_value( - additional_config, - "ft_communication_ops_abort_timeout_ms", - "FT_COMMUNICATION_OPS_ABORT_TIMEOUT_MS", - ascend_envs.FT_COMMUNICATION_OPS_ABORT_TIMEOUT_MS, - ) - if ( - not isinstance(self.ft_communication_ops_abort_timeout_ms, int) - or self.ft_communication_ops_abort_timeout_ms < 0 + # Fault-tolerance communication abort timeout (seconds); when fault + # tolerance is enabled and > 0, a hung NPU comm op is aborted after + # this many seconds. Drives HCCL_EVENT_TIMEOUT / HCCL_EXEC_TIMEOUT + # (= timeout - 1) and set_op_timeout_ms(timeout * 1000). 0 disables. + abort_timeout = additional_config.get("ft_communication_abort_timeout", 0) + if not isinstance(abort_timeout, int) or not ( + abort_timeout == 0 or abort_timeout >= 2 ): raise ValueError( - "ft_communication_ops_abort_timeout_ms must be a non-negative " - f"integer, got {self.ft_communication_ops_abort_timeout_ms}" + "ft_communication_abort_timeout must be 0 (disabled) or an " + f"integer of at least 2 seconds, " + f"got {abort_timeout}" ) + self.ft_communication_abort_timeout = abort_timeout self.pd_tp_ratio = 1 self.pd_head_ratio = 1 diff --git a/vllm_ascend/envs.py b/vllm_ascend/envs.py index 918519159b81..f6a111d5c972 100644 --- a/vllm_ascend/envs.py +++ b/vllm_ascend/envs.py @@ -100,12 +100,6 @@ # Control the aclrtMemcpyBatchAsync compile path for KV cache offloading. # "1": force enable, "0": force disable, None: auto-detect from CANN headers. "VLLM_ASCEND_ENABLE_BATCH_MEMCPY": lambda: os.getenv("VLLM_ASCEND_ENABLE_BATCH_MEMCPY", None), - # Fault-tolerance timeout (in milliseconds) after which a hung NPU operator - # (e.g. a communication op) is aborted with a timeout exception, so fault - # tolerance can detect it and trigger recovery. 0 (default) disables the - # timeout. Applied via torch_npu.npu.set_op_timeout_ms when fault tolerance - # is enabled. - "FT_COMMUNICATION_OPS_ABORT_TIMEOUT_MS": lambda: int(os.getenv("FT_COMMUNICATION_OPS_ABORT_TIMEOUT_MS", "0")), } # end-env-vars-definition diff --git a/vllm_ascend/worker/worker.py b/vllm_ascend/worker/worker.py index 821966ead291..f6af89c155b2 100644 --- a/vllm_ascend/worker/worker.py +++ b/vllm_ascend/worker/worker.py @@ -20,6 +20,7 @@ import copy import gc import logging +import os from types import NoneType from typing import Any @@ -370,7 +371,24 @@ def _init_device(self): if self.parallel_config.enable_fault_tolerance: import torch_npu - torch_npu.npu.set_op_timeout_ms(get_ascend_config().ft_communication_ops_abort_timeout_ms) + abort_timeout = get_ascend_config().ft_communication_abort_timeout + if abort_timeout > 0: + # User-provided HCCL timeouts win; otherwise derive them + # from the config value. HCCL_EVENT_TIMEOUT must be + # greater than HCCL_EXEC_TIMEOUT, hence EXEC defaults to + # abort_timeout - 1. + os.environ.setdefault("HCCL_EVENT_TIMEOUT", str(abort_timeout)) + os.environ.setdefault("HCCL_EXEC_TIMEOUT", str(abort_timeout - 1)) + if ( + int(os.environ["HCCL_EVENT_TIMEOUT"]) + <= int(os.environ["HCCL_EXEC_TIMEOUT"]) + ): + raise ValueError( + f"HCCL_EVENT_TIMEOUT ({os.environ['HCCL_EVENT_TIMEOUT']}) " + "must be greater than HCCL_EXEC_TIMEOUT " + f"({os.environ['HCCL_EXEC_TIMEOUT']})" + ) + torch_npu.npu.set_op_timeout_ms(abort_timeout * 1000) if ( parallel_config.distributed_executor_backend not in ("ray", "external_launcher") and parallel_config.data_parallel_backend != "ray" From 24be8231226ceae1150b84c7fd2d83b5fff31671 Mon Sep 17 00:00:00 2001 From: "pre-commit-ci[bot]" <66853113+pre-commit-ci[bot]@users.noreply.github.com> Date: Fri, 21 Aug 2026 01:29:11 +0000 Subject: [PATCH 2/2] [pre-commit.ci] auto fixes from pre-commit.com hooks for more information, see https://pre-commit.ci --- vllm_ascend/ascend_config.py | 4 +--- vllm_ascend/worker/worker.py | 5 +---- 2 files changed, 2 insertions(+), 7 deletions(-) diff --git a/vllm_ascend/ascend_config.py b/vllm_ascend/ascend_config.py index fe7e7e256569..432e5ce80458 100644 --- a/vllm_ascend/ascend_config.py +++ b/vllm_ascend/ascend_config.py @@ -205,9 +205,7 @@ def __init__(self, vllm_config: "VllmConfig"): # this many seconds. Drives HCCL_EVENT_TIMEOUT / HCCL_EXEC_TIMEOUT # (= timeout - 1) and set_op_timeout_ms(timeout * 1000). 0 disables. abort_timeout = additional_config.get("ft_communication_abort_timeout", 0) - if not isinstance(abort_timeout, int) or not ( - abort_timeout == 0 or abort_timeout >= 2 - ): + if not isinstance(abort_timeout, int) or not (abort_timeout == 0 or abort_timeout >= 2): raise ValueError( "ft_communication_abort_timeout must be 0 (disabled) or an " f"integer of at least 2 seconds, " diff --git a/vllm_ascend/worker/worker.py b/vllm_ascend/worker/worker.py index f6af89c155b2..7e2d371aee82 100644 --- a/vllm_ascend/worker/worker.py +++ b/vllm_ascend/worker/worker.py @@ -379,10 +379,7 @@ def _init_device(self): # abort_timeout - 1. os.environ.setdefault("HCCL_EVENT_TIMEOUT", str(abort_timeout)) os.environ.setdefault("HCCL_EXEC_TIMEOUT", str(abort_timeout - 1)) - if ( - int(os.environ["HCCL_EVENT_TIMEOUT"]) - <= int(os.environ["HCCL_EXEC_TIMEOUT"]) - ): + if int(os.environ["HCCL_EVENT_TIMEOUT"]) <= int(os.environ["HCCL_EXEC_TIMEOUT"]): raise ValueError( f"HCCL_EVENT_TIMEOUT ({os.environ['HCCL_EVENT_TIMEOUT']}) " "must be greater than HCCL_EXEC_TIMEOUT "