vllm-project-vllm-42585
Multi-API-server startup can intermittently fail with `zmq.error.ZMQError: Address already in use` when several API servers are launched concurrently, especially on multi-node data-parallel deployments or hosts with many other startup processes. The failure occurs because an API server is given a TCP port that was free when selected but has been claimed by another process before the child server binds to it, causing the server process to exit and the overall startup to fail.
API-server startup must reliably provide each child with usable, distinct TCP endpoints, even under heavy port contention. The parent process must not proceed until it has received the actual endpoints successfully bound by every child, so those addresses can be propagated to the engines. If a child exits before providing its endpoints, startup must fail promptly with a clear `RuntimeError` rather than hanging indefinitely or returning missing endpoint values.
Interface the hidden tests use. `get_engine_client_zmq_addr(local_only=False, host=host)` returns the placeholder `tcp://{host}:0`, so a child binds with a kernel-assigned port instead of a pre-selected one. `APIServerProcessManager` gains `gather_actual_addresses(timeout: float = 60.0) -> tuple[list[str], list[str]]`, returning the `(input_addresses, output_addresses)` each child actually bound, one entry per child in child order, all distinct `tcp://host:port` endpoints with a positive port; it raises `RuntimeError` whose message contains "reporting" when a child exits before reporting its endpoints, and also raises `RuntimeError` on timeout. Each child is started with the positional arguments `(listen_address, sock, args, client_config)`, and `client_config` carries an `actual_address_pipe` (a multiprocessing connection) on which the child sends a dict with keys `input_address` and `output_address` holding its bound endpoints and then closes the pipe; the tests supply their own child function that does exactly this, and one that exits without touching the pipe. `manager.processes` has one entry per server and `manager.shutdown()` still terminates them.
Hidden tests · 2 fail-to-pass, 5 pass-to-passrun after the agent submits, in a clean verifier
Test patch · 153 lines
diff --git a/tests/entrypoints/test_api_server_process_manager.py b/tests/entrypoints/test_api_server_process_manager.py
index 738ed4a22f11..ada7a8797fa2 100644
--- a/tests/entrypoints/test_api_server_process_manager.py
+++ b/tests/entrypoints/test_api_server_process_manager.py
@@ -8,8 +8,14 @@
from unittest.mock import patch
import pytest
+import zmq
-from vllm.v1.utils import APIServerProcessManager, wait_for_completion_or_failure
+from vllm.utils.network_utils import make_zmq_socket, split_zmq_path
+from vllm.v1.utils import (
+ APIServerProcessManager,
+ get_engine_client_zmq_addr,
+ wait_for_completion_or_failure,
+)
# Global variables to control worker behavior
WORKER_RUNTIME_SECONDS = 0.5
@@ -23,6 +29,39 @@ def mock_run_api_server_worker(listen_address, sock, args, client_config=None):
print("Mock worker completed successfully")
+# Module-level stub for the gather_actual_addresses test. Must be
+# importable by `multiprocessing.spawn` (no closures, no nesting).
+def defer_addresses_stub_worker(listen_address, sock, args, client_config):
+ """Bind ROUTER/PULL with a kernel-assigned port, report the actual
+ endpoints back via the pipe, then exit."""
+ ctx = zmq.Context()
+ try:
+ in_sock = make_zmq_socket(
+ ctx, client_config["input_address"], zmq.ROUTER, bind=True
+ )
+ out_sock = make_zmq_socket(
+ ctx, client_config["output_address"], zmq.PULL, bind=True
+ )
+ try:
+ pipe = client_config["actual_address_pipe"]
+ try:
+ pipe.send(
+ {
+ "input_address": in_sock.getsockopt(zmq.LAST_ENDPOINT).decode(),
+ "output_address": out_sock.getsockopt(
+ zmq.LAST_ENDPOINT
+ ).decode(),
+ }
+ )
+ finally:
+ pipe.close()
+ finally:
+ in_sock.close(linger=0)
+ out_sock.close(linger=0)
+ finally:
+ ctx.term()
+
+
@pytest.fixture
def api_server_args():
"""Fixture to provide arguments for APIServerProcessManager."""
@@ -268,3 +307,92 @@ def run_with_exception_capture():
manager.shutdown()
mock_coordinator.shutdown()
time.sleep(0.2)
+
+
+@pytest.mark.timeout(60)
+def test_gather_actual_addresses_end_to_end():
+ """Each child binds ROUTER/PULL with a kernel-picked port and reports
+ the bound endpoints back via its per-child pipe; the manager surfaces
+ them via :py:meth:`gather_actual_addresses`."""
+ host = "127.0.0.1"
+ num_servers = 4
+
+ placeholder_inputs = [
+ get_engine_client_zmq_addr(local_only=False, host=host)
+ for _ in range(num_servers)
+ ]
+ placeholder_outputs = [
+ get_engine_client_zmq_addr(local_only=False, host=host)
+ for _ in range(num_servers)
+ ]
+ for addr in placeholder_inputs + placeholder_outputs:
+ assert addr == f"tcp://{host}:0", addr
+
+ sock = socket.socket()
+ manager = APIServerProcessManager(
+ listen_address=f"tcp://{host}:0",
+ sock=sock,
+ args="test_args",
+ num_servers=num_servers,
+ input_addresses=placeholder_inputs,
+ output_addresses=placeholder_outputs,
+ target_server_fn=defer_addresses_stub_worker,
+ )
+
+ try:
+ assert len(manager.processes) == num_servers
+ actual_inputs, actual_outputs = manager.gather_actual_addresses(timeout=15.0)
+ finally:
+ manager.shutdown()
+ time.sleep(0.2)
+ sock.close()
+
+ assert len(actual_inputs) == num_servers
+ assert len(actual_outputs) == num_servers
+
+ for addr in actual_inputs + actual_outputs:
+ scheme, parsed_host, port = split_zmq_path(addr)
+ assert scheme == "tcp", addr
+ assert parsed_host == host, addr
+ assert port and int(port) > 0, addr
+
+ all_addrs = actual_inputs + actual_outputs
+ assert len(set(all_addrs)) == len(all_addrs), all_addrs
+
+
+@pytest.mark.timeout(30)
+def test_gather_actual_addresses_child_crash_before_report():
+ """A child that exits before sending its endpoints must surface a
+ clear ``RuntimeError`` rather than hang or return ``None`` slots."""
+ host = "127.0.0.1"
+ num_servers = 2
+ placeholder_inputs = [
+ get_engine_client_zmq_addr(local_only=False, host=host)
+ for _ in range(num_servers)
+ ]
+ placeholder_outputs = [
+ get_engine_client_zmq_addr(local_only=False, host=host)
+ for _ in range(num_servers)
+ ]
+
+ sock = socket.socket()
+ manager = APIServerProcessManager(
+ listen_address=f"tcp://{host}:0",
+ sock=sock,
+ args="test_args",
+ num_servers=num_servers,
+ input_addresses=placeholder_inputs,
+ output_addresses=placeholder_outputs,
+ # mock_run_api_server_worker exits without touching
+ # ``actual_address_pipe`` — simulates a child that dies before
+ # reporting its bound addresses.
+ target_server_fn=mock_run_api_server_worker,
+ )
+ try:
+ # Sentinel-first vs pipe-EOF-first both produce "reporting".
+ with pytest.raises(RuntimeError, match="reporting"):
+ manager.gather_actual_addresses(timeout=10.0)
+ finally:
+ manager.shutdown()
+ time.sleep(0.2)
+ sock.close()
Reference fix · 6 files, +187 −49the upstream merge, used only for grading calibration
The agent could not see this: the repository holds one commit and the sandbox has no network. Leak audit.
vllm/entrypoints/cli/serve.py, vllm/v1/engine/async_llm.py, vllm/v1/engine/coordinator.py, vllm/v1/engine/core_client.py, vllm/v1/engine/utils.py, vllm/v1/utils.py
diff --git a/vllm/entrypoints/cli/serve.py b/vllm/entrypoints/cli/serve.py
index 5972bd48aa46..cbd3a44724a9 100644
--- a/vllm/entrypoints/cli/serve.py
+++ b/vllm/entrypoints/cli/serve.py
@@ -308,7 +308,14 @@ def signal_handler(signum, frame):
from vllm.v1.engine.utils import get_engine_zmq_addresses
- addresses = get_engine_zmq_addresses(vllm_config, num_api_servers)
+ # Per-API-server ports are picked by the kernel at each child's bind()
+ # to avoid parent-probe vs child-bind TOCTOU; Rust front-end opts out
+ # because it has no port-report-back channel.
+ addresses = get_engine_zmq_addresses(
+ vllm_config,
+ num_api_servers,
+ defer_api_server_ports=not rust_frontend_path,
+ )
with launch_core_engines(
vllm_config, executor_class, log_stats, addresses, num_api_servers
@@ -341,6 +348,12 @@ def signal_handler(signum, frame):
tensor_queue=tensor_queue,
)
+ # Forward each child's bound endpoints to the engine handshake
+ # (runs on ``with`` exit).
+ actual_inputs, actual_outputs = api_server_manager.gather_actual_addresses()
+ addresses.inputs = actual_inputs
+ addresses.outputs = actual_outputs
+
# Wait for API servers.
try:
wait_for_completion_or_failure(
diff --git a/vllm/v1/engine/async_llm.py b/vllm/v1/engine/async_llm.py
index 160f148f5c59..419e15163a9f 100644
--- a/vllm/v1/engine/async_llm.py
+++ b/vllm/v1/engine/async_llm.py
@@ -81,7 +81,7 @@ def __init__(
start_engine_loop: bool = True,
stat_loggers: list[StatLoggerFactory] | None = None,
aggregate_engine_logging: bool = False,
- client_addresses: dict[str, str] | None = None,
+ client_addresses: dict[str, Any] | None = None,
client_count: int = 1,
client_index: int = 0,
) -> None:
@@ -209,7 +209,7 @@ def from_vllm_config(
enable_log_requests: bool = False,
aggregate_engine_logging: bool = False,
disable_log_stats: bool = False,
- client_addresses: dict[str, str] | None = None,
+ client_addresses: dict[str, Any] | None = None,
client_count: int = 1,
client_index: int = 0,
) -> "AsyncLLM":
diff --git a/vllm/v1/engine/coordinator.py b/vllm/v1/engine/coordinator.py
index 8ebf976c5fa1..87937cc6a870 100644
--- a/vllm/v1/engine/coordinator.py
+++ b/vllm/v1/engine/coordinator.py
@@ -11,7 +11,7 @@
from vllm.config import ParallelConfig
from vllm.logger import init_logger
-from vllm.utils.network_utils import get_tcp_uri, make_zmq_socket
+from vllm.utils.network_utils import make_zmq_socket
from vllm.utils.system_utils import get_mp_context, set_process_title
from vllm.v1.engine import EngineCoreOutputs, EngineCoreRequestType
from vllm.v1.serial_utils import MsgpackDecoder
@@ -91,16 +91,9 @@ def __init__(
if parallel_config.enable_elastic_ep:
local_only_eng = False
- def bind_address(local_only: bool) -> str:
- return (
- get_engine_client_zmq_addr(local_only=True, host=host)
- if local_only
- else get_tcp_uri(host, 0)
- )
-
- front_publish_address = bind_address(local_only)
- back_publish_address = bind_address(local_only_eng)
- back_output_address = bind_address(local_only_eng)
+ front_publish_address = get_engine_client_zmq_addr(local_only, host=host)
+ back_publish_address = get_engine_client_zmq_addr(local_only_eng, host=host)
+ back_output_address = get_engine_client_zmq_addr(local_only_eng, host=host)
context = get_mp_context()
parent_zmq_addr_pipe, child_zmq_addr_pipe = context.Pipe(duplex=False)
diff --git a/vllm/v1/engine/core_client.py b/vllm/v1/engine/core_client.py
index 2f2c15d246f8..c26380e6e151 100644
--- a/vllm/v1/engine/core_client.py
+++ b/vllm/v1/engine/core_client.py
@@ -11,6 +11,7 @@
from collections.abc import Awaitable, Callable, Sequence
from concurrent.futures import Future
from dataclasses import dataclass
+from multiprocessing.connection import Connection
from multiprocessing.queues import Queue
from threading import Thread
from typing import Any, TypeAlias, TypeVar
@@ -108,7 +109,7 @@ def make_async_mp_client(
vllm_config: VllmConfig,
executor_class: type[Executor],
log_stats: bool,
- client_addresses: dict[str, str] | None = None,
+ client_addresses: dict[str, Any] | None = None,
client_count: int = 1,
client_index: int = 0,
) -> "AsyncMPClient":
@@ -476,7 +477,7 @@ def __init__(
vllm_config: VllmConfig,
executor_class: type[Executor],
log_stats: bool,
- client_addresses: dict[str, str] | None = None,
+ client_addresses: dict[str, Any] | None = None,
):
self.vllm_config = vllm_config
@@ -507,7 +508,7 @@ def __init__(
output_address = client_addresses["output_address"]
self.stats_update_address = client_addresses.get("stats_update_address")
# Tensor queues passed via client_addresses for multi-API-server case
- tensor_queue = client_addresses.get("tensor_queue") # type: ignore[assignment]
+ tensor_queue = client_addresses.get("tensor_queue")
self.input_socket = self.resources.input_socket = make_zmq_socket(
self.ctx,
input_address,
@@ -518,6 +519,28 @@ def __init__(
self.resources.output_socket = make_zmq_socket(
self.ctx, output_address, zmq.PULL
)
+
+ # Report bound endpoints back so the parent can forward
+ # them to engines (mirrors the DPCoordinator pattern).
+ actual_address_pipe: Connection | None = client_addresses.get(
+ "actual_address_pipe"
+ )
+ if actual_address_pipe is not None:
+ try:
+ actual_input = self.input_socket.getsockopt(
+ zmq.LAST_ENDPOINT
+ ).decode()
+ actual_output = self.resources.output_socket.getsockopt(
+ zmq.LAST_ENDPOINT
+ ).decode()
+ actual_address_pipe.send(
+ {
+ "input_address": actual_input,
+ "output_address": actual_output,
+ }
+ )
+ finally:
+ actual_address_pipe.close()
else:
# Engines are managed by this client.
addresses = get_engine_zmq_addresses(vllm_config)
@@ -532,6 +555,15 @@ def __init__(
self.ctx, addresses.outputs[0], zmq.PULL
)
+ # Resolve ``tcp://host:0`` placeholders to bound endpoints
+ # before engines DEALER-connect. No-op for IPC.
+ addresses.inputs[0] = self.input_socket.getsockopt(
+ zmq.LAST_ENDPOINT
+ ).decode()
+ addresses.outputs[0] = self.resources.output_socket.getsockopt(
+ zmq.LAST_ENDPOINT
+ ).decode()
+
with launch_core_engines(
vllm_config, executor_class, log_stats, addresses
) as (engine_manager, coordinator, addresses, tensor_queue):
@@ -893,7 +925,7 @@ def __init__(
vllm_config: VllmConfig,
executor_class: type[Executor],
log_stats: bool,
- client_addresses: dict[str, str] | None = None,
+ client_addresses: dict[str, Any] | None = None,
client_count: int = 1,
client_index: int = 0,
):
@@ -1143,7 +1175,7 @@ def __init__(
vllm_conf
… [11253 more characters]GLM-5.3 Flash · attempt 3 · solved · 64 steps · $0.032
Submitted diff · 220 lines
diff --git a/vllm/entrypoints/cli/serve.py b/vllm/entrypoints/cli/serve.py
index 5972bd4..34534d1 100644
--- a/vllm/entrypoints/cli/serve.py
+++ b/vllm/entrypoints/cli/serve.py
@@ -340,6 +340,12 @@ def run_multi_api_server(args: argparse.Namespace):
stats_update_address=stats_update_address,
tensor_queue=tensor_queue,
)
+ # Wait until every API server has bound its (kernel-assigned)
+ # ZMQ endpoints and propagate the actual addresses so that
+ # they can be used by the engines.
+ actual_inputs, actual_outputs = api_server_manager.gather_actual_addresses()
+ addresses.inputs = actual_inputs
+ addresses.outputs = actual_outputs
# Wait for API servers.
try:
diff --git a/vllm/v1/engine/core_client.py b/vllm/v1/engine/core_client.py
index 2f2c15d..be4a3c7 100644
--- a/vllm/v1/engine/core_client.py
+++ b/vllm/v1/engine/core_client.py
@@ -457,6 +457,18 @@ class ElasticScalingCache:
pending_notifications: dict[EEPNotificationType, set[int]]
+def _resolve_bound_endpoint(socket: Any, address: str) -> str:
+ """Return the endpoint a ZMQ socket actually bound.
+
+ If the socket was bound to a TCP address with a kernel-assigned port
+ (port 0), query the socket for the real endpoint. Otherwise the
+ address is returned unchanged.
+ """
+ if address.startswith("tcp://") and address.rsplit(":", 1)[-1] == "0":
+ return socket.getsockopt(zmq.LAST_ENDPOINT).decode()
+ return address
+
+
class MPClient(EngineCoreClient):
"""
MPClient: base client for multi-proc EngineCore.
@@ -501,6 +513,10 @@ class MPClient(EngineCoreClient):
self.stats_update_address: str | None = None
tensor_queue: Queue | None = None
+ # Pipe on which this front-end process reports the ZMQ
+ # endpoints it actually bound (the configured addresses may be
+ # placeholders with a kernel-assigned port, i.e. port 0).
+ actual_address_pipe = client_addresses.pop("actual_address_pipe", None)
if client_addresses:
# Engines are managed externally to this client.
input_address = client_addresses["input_address"]
@@ -518,6 +534,26 @@ class MPClient(EngineCoreClient):
self.resources.output_socket = make_zmq_socket(
self.ctx, output_address, zmq.PULL
)
+ if actual_address_pipe is not None:
+ try:
+ actual_input_address = _resolve_bound_endpoint(
+ self.input_socket, input_address
+ )
+ actual_output_address = _resolve_bound_endpoint(
+ self.resources.output_socket, output_address
+ )
+ # Propagate the actual endpoints so that the rest of
+ # this process (e.g. engine handshakes) uses them.
+ client_addresses["input_address"] = actual_input_address
+ client_addresses["output_address"] = actual_output_address
+ actual_address_pipe.send(
+ {
+ "input_address": actual_input_address,
+ "output_address": actual_output_address,
+ }
+ )
+ finally:
+ actual_address_pipe.close()
else:
# Engines are managed by this client.
addresses = get_engine_zmq_addresses(vllm_config)
diff --git a/vllm/v1/utils.py b/vllm/v1/utils.py
index afa621a..dc15d3c 100644
--- a/vllm/v1/utils.py
+++ b/vllm/v1/utils.py
@@ -29,7 +29,7 @@ from torch.autograd.profiler import record_function
import vllm.envs as envs
from vllm.logger import init_logger
from vllm.usage.usage_lib import UsageContext, is_usage_stats_enabled, usage_message
-from vllm.utils.network_utils import get_open_port, get_open_zmq_ipc_path, get_tcp_uri
+from vllm.utils.network_utils import get_open_zmq_ipc_path, get_tcp_uri
from vllm.utils.system_utils import decorate_logs, kill_process_tree, set_process_title
from vllm.v1.core.sched.output import SchedulerOutput
@@ -156,7 +156,7 @@ def get_engine_client_zmq_addr(local_only: bool, host: str, port: int = 0) -> st
return (
get_open_zmq_ipc_path()
if local_only
- else (get_tcp_uri(host, port or get_open_port()))
+ else (get_tcp_uri(host, port))
)
@@ -199,6 +199,8 @@ class APIServerProcessManager:
# Start API servers
spawn_context = multiprocessing.get_context("spawn")
self.processes: list[BaseProcess] = []
+ self._address_pipes: list[connection.Connection] = []
+ self._gathered_addresses: tuple[list[str], list[str]] | None = None
for i, in_addr, out_addr in zip(
range(num_servers), input_addresses, output_addresses
@@ -214,6 +216,17 @@ class APIServerProcessManager:
if tensor_queue is not None:
client_config["tensor_queue"] = tensor_queue
+ # Each child binds its ZMQ sockets with kernel-assigned ports
+ # (the addresses may be placeholders with port 0) and reports
+ # the endpoints it actually bound back via this pipe. This
+ # avoids races where a pre-selected port is claimed by another
+ # process before the child binds to it.
+ actual_address_recv, actual_address_send = spawn_context.Pipe(
+ duplex=False
+ )
+ client_config["actual_address_pipe"] = actual_address_send
+ self._address_pipes.append(actual_address_recv)
+
proc = spawn_context.Process(
target=target_server_fn or run_api_server_worker_proc,
name=f"
… [4014 more characters]$ cd /repo && ls && grep -rn "get_engine_client_zmq_addr" --include=*.py | head -30
output · 38 lines
AGENTS.md CLAUDE.md CMakeLists.txt CODE_OF_CONDUCT.md CONTRIBUTING.md DCO LICENSE MANIFEST.in README.md RELEASE.md SECURITY.md benchmarks
26 more lines
build_rust.sh cmake codecov.yml csrc docker docs examples mkdocs.yaml pyproject.toml requirements rust rust-toolchain.toml scripts setup.py tests tools use_existing_torch.py vllm vllm.egg-info vllm/v1/engine/coordinator.py:18:from vllm.v1.utils import get_engine_client_zmq_addr, shutdown vllm/v1/engine/coordinator.py:96: get_engine_client_zmq_addr(local_only=True, host=host) vllm/v1/engine/utils.py:31:from vllm.v1.utils import get_engine_client_zmq_addr, shutdown vllm/v1/engine/utils.py:983: get_engine_client_zmq_addr(client_local_only, host) vllm/v1/engine/utils.py:987: get_engine_client_zmq_addr(client_local_only, host) vllm/v1/engine/utils.py:1098: handshake_address = get_engine_client_zmq_addr( vllm/v1/utils.py:147:def get_engine_client_zmq_addr(local_only: bool, host: str, port: int = 0) -> str:
$ cd /repo && sed -n 100,220p vllm/v1/utils.py && grep -rn "APIServerProcessManager" -l vllm | head
output · 39 lines
return len(self._x)
def __repr__(self):
return f"ConstantList({self._x})"
def copy(self) -> list[T]:
return self._x.copy()
class CpuGpuBuffer:
"""Buffer to easily copy tensors between CPU and GPU."""
27 more lines
def __init__(
self,
*size: int | torch.SymInt,
dtype: torch.dtype,
device: torch.device,
pin_memory: bool,
with_numpy: bool = True,
) -> None:
self.cpu = torch.zeros(*size, dtype=dtype, device="cpu", pin_memory=pin_memory)
self.gpu = torch.zeros_like(self.cpu, device=device)
self.np: np.ndarray
# To keep type hints simple (avoiding generics and subclasses), we
# only conditionally create the numpy array attribute. This can cause
# AttributeError if `self.np` is accessed when `with_numpy=False`.
if with_numpy:
if dtype == torch
…[3239 chars omitted]…
ate_address"] = stats_update_address
if tensor_queue is not None:
client_config["tensor_queue"] = tensor_queue
proc = spawn_context.Process(
target=target_server_fn or run_api_server_worker_proc,
name=f"ApiServer_{i}",
args=(listen_address, sock, args, client_config),
vllm/entrypoints/cli/serve.py
vllm/v1/utils.py$ cd /repo && sed -n 220,320p vllm/v1/utils.py && grep -rn "APIServerProcessManager(" vllm -noutput · 35 lines
args=(listen_address, sock, args, client_config),
)
self.processes.append(proc)
proc.start()
logger.info("Started %d API server processes", len(self.processes))
# Shutdown only the API server processes on garbage collection
# The extra processes are managed by their owners
self._finalizer = weakref.finalize(self, shutdown, self.processes)
def shutdown(self, timeout: float | None = None) -> None:23 more lines
"""Shutdown API server processes with configurable timeout"""
if self._finalizer.detach() is not None:
shutdown(self.processes, timeout=timeout)
class RustFrontendProcessManager:
"""Manages a single Rust frontend subprocess.
Launches the Rust vllm-rs binary in 'frontend' mode, passing the
listening socket fd and ZMQ transport addresses. Provides the same
int
…[2243 chars omitted]…
recv, send = connection.Pipe(duplex=False)
self._sentinel_conn = recv
self._sentinel_send = send
def monitor_subprocess() -> None:
try:
proc.wait()
finally:
with contextlib.suppress(Exception):
send.close()
vllm/entrypoints/cli/serve.py:333: api_server_manager = APIServerProcessManager($ cd /repo && sed -n 280,420p vllm/entrypoints/cli/serve.py
output · 36 lines
signal.signal(signal.SIGTERM, signal_handler)
signal.signal(signal.SIGINT, signal_handler)
listen_address, sock = setup_server(args)
engine_args = vllm.AsyncEngineArgs.from_cli_args(args)
engine_args._api_process_count = num_api_servers
engine_args._api_process_rank = -1
usage_context = UsageContext.OPENAI_API_SERVER
vllm_config = engine_args.create_engine_config(usage_context=usage_context)
24 more lines
if num_api_servers > 1 and envs.VLLM_ALLOW_RUNTIME_LORA_UPDATING:
raise ValueError(
"VLLM_ALLOW_RUNTIME_LORA_UPDATING cannot be used with api_server_count > 1"
)
executor_class = Executor.get_class(vllm_config)
log_stats = not engine_args.disable_log_stats
parallel_config = vllm_config.parallel_config
dp_rank = parallel_config.data_parallel_rank
assert parallel_config.local_engines_only or dp_rank == 0
…[2015 chars omitted]…
loat | None) -> float | None:
return (
deadline if deadline is None else max(deadline - time.monotonic(), 0.0)
)
api_server_manager.shutdown(timeout=timeout)
if local_engine_manager:
local_engine_manager.shutdown(timeout=to_timeout(shutdown_by))
if coordinator:
coordinator.shutdown(timeout=to_timeout(shutdown_by))$ cd /repo && grep -rn "get_engine_zmq_addresses\|class EngineAddresses\|addresses.inputs" vllm/v1/engine/utils.py | head -20 && sed -n 950,1120p vllm/v1/engine/utils.py
output · 35 lines
955:def get_engine_zmq_addresses(
ray.kill(actor)
for pg in self.created_placement_groups:
ray.util.remove_placement_group(pg)
def get_engine_zmq_addresses(
vllm_config: VllmConfig,
num_api_servers: int = 1,
) -> EngineZmqAddresses:
"""Allocate ZMQ addresses for engine-client communication."""
parallel_config = vllm_config.parallel_config23 more lines
local_engine_count = parallel_config.data_parallel_size_local
local_start_index = parallel_config.data_parallel_rank_local
dp_size = parallel_config.data_parallel_size
host = parallel_config.data_parallel_master_ip
local_engines_only = parallel_config.local_engines_only
# In offline mode there is an LLM instance per DP rank and
# one core engine per LLM, see
# examples/features/data_parallel/data_parallel_offline.py.
offline_mode = local_start_index is not None
…[5224 chars omitted]…
R, bind=True
) as handshake_socket:
# Start local engines.
if local_engine_count:
local_engine_manager = CoreEngineProcManager(
vllm_config=vllm_config,
executor_class=executor_class,
log_stats=log_stats,
handshake_address=handshake_address,
client_handshake_address=client_handshake_address,$ cd /repo && grep -rn "addresses.inputs\|addresses.outputs\|input_address" vllm/v1/engine/core.py vllm/v1/engine/utils.py | head -30 && grep -rn "EngineZmqAddresses" vllm/v1 --include=*.py | head
output · 20 lines
vllm/v1/engine/core.py:922: addresses.inputs, vllm/v1/engine/core.py:934: addresses.outputs, vllm/v1/engine/core.py:1008: addresses.inputs = client_addresses.inputs vllm/v1/engine/core.py:1009: addresses.outputs = client_addresses.outputs vllm/v1/engine/core.py:1402: input_addresses: list[str], vllm/v1/engine/core.py:1403: coord_input_address: str | None, vllm/v1/engine/core.py:1419: ctx, input_address, zmq.DEALER, identity=identity, bind=False vllm/v1/engine/core.py:1422: for input_address in input_addresses vllm/v1/engine/core.py:1424: if coord_input_address is None: vllm/v1/engine/core.py:1430: coord_input_address, vllm/v1/engine/core.py:1979: ctx, self.addresses.outputs[0], zmq.PUSH, linger=4000 vllm/v1/engine/
8 more lines
…[222 chars omitted]… EngineZmqAddresses, None, None]: vllm/v1/engine/core.py:1060: ) -> EngineZmqAddresses: vllm/v1/engine/core.py:2020: addresses: EngineZmqAddresses, vllm/v1/engine/core.py:2144: addresses: EngineZmqAddresses, vllm/v1/engine/core.py:2167: addresses: EngineZmqAddresses, vllm/v1/engine/utils.py:58:class EngineZmqAddresses: vllm/v1/engine/utils.py:80: addresses: EngineZmqAddresses
$ cd /repo && sed -n 900,1030p vllm/v1/engine/core.py
output · 37 lines
EEPNotificationType.NEW_CORE_ENGINES_INIT_READY,
vllm_config=vllm_config,
)
self._init_data_parallel(vllm_config)
super().__init__(
vllm_config,
executor_class,
log_stats,
executor_fail_callback,
internal_dp_balancing,
)25 more lines
# Background Threads and Queues for IO. These enable us to
# overlap ZMQ socket IO with GPU since they release the GIL,
# and to overlap some serialization/deserialization with the
# model forward pass.
# Threads handle Socket <-> Queues and core_busy_loop uses Queue.
ready_event = threading.Event()
input_thread = threading.Thread(
target=self.process_input_sockets,
args=(
…[3826 chars omitted]…
zmq.Context,
handshake_address: str,
identity: bytes,
local_client: bool,
headless: bool,
vllm_config: VllmConfig,
parallel_config_to_update: ParallelConfig | None = None,
) -> Generator[EngineZmqAddresses, None, None]:
with make_zmq_socket(
ctx,
handshake_address,
zmq.DEALER,
identity=identity,$ cd /repo && sed -n 1390,1445p vllm/v1/engine/core.py && grep -rn "process_input_sockets\|process_output_sockets" vllm/v1/engine/core.py | head
output · 32 lines
self.output_queue.put_nowait(EngineCoreProc.ENGINE_CORE_DEAD)
# Wait until msg sent by the daemon before shutdown.
self.output_thread.join(timeout=5.0)
if self.output_thread.is_alive():
logger.fatal(
"vLLM shutdown signal from EngineCore failed "
"to send. Please report this issue."
)
def process_input_sockets(
self,20 more lines
input_addresses: list[str],
coord_input_address: str | None,
identity: bytes,
ready_event: threading.Event,
):
"""Input socket IO thread."""
# Msgpack serialization decoding with optional tensor IPC receiver.
add_request_decoder = MsgpackDecoder(
EngineCoreRequest, oob_tensor_provider=self.tensor_ipc_receiver
)
generic_decoder = MsgpackDecoder(oob_tensor_provider=self.tensor_ipc_recei
…[1103 chars omitted]…
m_gpu_blocks=self.vllm_config.cache_config.num_gpu_blocks or 0,
dp_stats_address=self.frontend_stats_publish_address,
dtype=str(self.vllm_config.model_config.dtype).removeprefix("torch."),
920: target=self.process_input_sockets,
932: target=self.process_output_sockets,
1400: def process_input_sockets(
1495: def process_output_sockets($ cd /repo && grep -rn "input_address\|output_address\|actual_address" vllm/v1/engine/api_server.py | head -30
output · 1 lines
grep: vllm/v1/engine/api_server.py: No such file or directory
$ cd /repo && grep -rln "run_api_server_worker_proc" vllm && grep -rn "input_address" vllm/entrypoints/api_server.py | head -20
output · 1 lines
vllm/v1/utils.py
$ cd /repo && grep -rn "run_api_server_worker_proc" vllm -r | head && grep -rn "client_config" vllm --include=*.py -l
output · 4 lines
vllm/v1/utils.py:218: target=target_server_fn or run_api_server_worker_proc, vllm/v1/utils.py:377:def run_api_server_worker_proc( vllm/entrypoints/openai/api_server.py vllm/v1/utils.py
$ cd /repo && sed -n 370,470p vllm/v1/utils.py && grep -rn "client_config" vllm/entrypoints/openai/api_server.py | head
output · 37 lines
proc.join(remaining)
for proc in procs:
if proc.is_alive() and (pid := proc.pid) is not None:
kill_process_tree(pid)
def run_api_server_worker_proc(
listen_address, sock, args, client_config=None, **uvicorn_kwargs
) -> None:
"""Entrypoint for individual API server worker processes."""
25 more lines
from vllm.entrypoints.openai.api_server import run_server_worker
client_config = client_config or {}
server_index = client_config.get("client_index", 0)
# Set process title and add process-specific prefix to stdout and stderr.
set_process_title("APIServer", str(server_index))
decorate_logs()
uvloop.run(
run_server_worker(listen_address, sock, args, client_config, **uvicorn_kwargs)
)
def wait_for_completion_or_failure(
api_server_manager: "APIServerProcessManager | RustFrontendProcessManager",
engine_manag
…[3155 chars omitted]…
engine_args._api_process_rank = client_config.get("client_index", 0)
103: client_config=client_config,
113: client_config: dict[str, Any] | None = None,
130: # Don't mutate the input client_config
131: client_config = dict(client_config) if client_config else {}
132: client_count = client_config.pop("client_count", 1)
133: client_index = client_config.pop("client_index", 0)$ cd /repo && sed -n 70,200p vllm/entrypoints/openai/api_server.py
output · 38 lines
# Cannot use __name__ (https://github.com/vllm-project/vllm/pull/4765)
logger = init_logger("vllm.entrypoints.openai.api_server")
_FALLBACK_SUPPORTED_TASKS: tuple[SupportedTask, ...] = ("generate",)
@asynccontextmanager
async def build_async_engine_client(
args: Namespace,
*,
usage_context: UsageContext = UsageContext.OPENAI_API_SERVER,
client_config: dict[str, Any] | None = None,26 more lines
) -> AsyncIterator[EngineClient]:
if os.getenv("VLLM_WORKER_MULTIPROC_METHOD") == "forkserver":
# The executor is expected to be mp.
# Pre-import heavy modules in the forkserver process
logger.debug("Setup forkserver with pre-imports")
multiprocessing.set_start_method("forkserver")
multiprocessing.set_forkserver_preload(["vllm.v1.engine.async_llm"])
forkserver.ensure_running()
logger.debug("Forkserver setup complete!")
# Con
…[3236 chars omitted]…
ster_models_api_router,
)
register_models_api_router(app)
from vllm.entrypoints.sagemaker.api_router import (
attach_router as register_sagemaker_api_router,
)
register_sagemaker_api_router(app, supported_tasks, model_config)
if "generate" in supported_tasks:
from vllm.entrypoints.openai.generate.api_router import (
register_generate_api_routers,$ cd /repo && grep -rn "client_addresses\|input_address" vllm/v1/engine/async_llm.py | head -20 && grep -rn "def get_client_input_addresses\|input_address" vllm/v1/engine/utils.py | head -20
output · 4 lines
84: client_addresses: dict[str, str] | None = None, 150: client_addresses=client_addresses, 212: client_addresses: dict[str, str] | None = None, 226: client_addresses=client_addresses,
$ cd /repo && grep -rn "client_addresses" vllm/v1 --include=*.py | grep -v async_llm | head && grep -rn "client_addresses" vllm/v1/core/client.py 2>/dev/null | head; ls vllm/v1/core
output · 19 lines
vllm/v1/engine/core.py:1003: with handshake as addresses, local_handshake as client_addresses: vllm/v1/engine/core.py:1007: # (client_addresses). vllm/v1/engine/core.py:1008: addresses.inputs = client_addresses.inputs vllm/v1/engine/core.py:1009: addresses.outputs = client_addresses.outputs vllm/v1/engine/core_client.py:111: client_addresses: dict[str, str] | None = None, vllm/v1/engine/core_client.py:120: client_addresses, vllm/v1/engine/core_client.py:479: client_addresses: dict[str, str] | None = None, vllm/v1/engine/core_client.py:504: if client_addresses: vllm/v1/engine/core_client.py:506: input_address = client_addresses["input_address"] vllm/v1/engine/core_client.py:507: output_address = client_addresses["output_address"] __init__.py block_pool.py
7 more lines
encoder_cache_manager.py kv_cache_coordinator.py kv_cache_manager.py kv_cache_metrics.py kv_cache_utils.py sched single_type_kv_cache_manager.py
$ cd /repo && sed -n 470,560p vllm/v1/engine/core_client.py
output · 37 lines
* SyncMPClient subclass for LLM usage
"""
def __init__(
self,
asyncio_mode: bool,
vllm_config: VllmConfig,
executor_class: type[Executor],
log_stats: bool,
client_addresses: dict[str, str] | None = None,
):
self.vllm_config = vllm_config25 more lines
# ZMQ setup.
sync_ctx = zmq.Context(io_threads=2)
self.ctx = zmq.asyncio.Context(sync_ctx) if asyncio_mode else sync_ctx
# This will ensure resources created so far are closed
# when the client is garbage collected, even if an
# exception is raised mid-construction.
self.resources = BackgroundResources(ctx=sync_ctx)
self._finalizer = weakref.finalize(self, self.resources)
success = False
try:
# State used for data parallel.
self.engines_running = False
paralle
…[2989 chars omitted]…
None:
tensor_ipc_sender = TensorIpcSender(tensor_queue)
self.encoder = MsgpackEncoder(oob_tensor_consumer=tensor_ipc_sender)
self.decoder = MsgpackDecoder(EngineCoreOutputs)
dp_size = parallel_config.data_parallel_size
dp_rank = parallel_config.data_parallel_index
dp_local_size = parallel_config.data_parallel_size_local$ cd /repo && grep -rn "LAST_ENDPOINT" vllm -r | head && grep -rn "def make_zmq_socket" -A5 vllm/utils/zmq.py | head -20
output · 4 lines
vllm/v1/engine/coordinator.py:237: publish_front.getsockopt(zmq.LAST_ENDPOINT).decode(), vllm/v1/engine/coordinator.py:238: output_back.getsockopt(zmq.LAST_ENDPOINT).decode(), vllm/v1/engine/coordinator.py:239: publish_back.getsockopt(zmq.LAST_ENDPOINT).decode(), grep: vllm/utils/zmq.py: No such file or directory
$ cd /repo && grep -rn "def make_zmq_socket" vllm -r && grep -rn "get_open_port\b" vllm/utils/__init__.py vllm/utils/*.py | head
output · 7 lines
vllm/utils/network_utils.py:284:def make_zmq_socket( vllm/utils/network_utils.py:150:def get_open_port() -> int: vllm/utils/network_utils.py:163: candidate_port = _get_open_port() vllm/utils/network_utils.py:166: return _get_open_port() vllm/utils/network_utils.py:179: port = _get_open_port(start_port=next_port, max_attempts=1000) vllm/utils/network_utils.py:185: ports_set.add(get_open_port()) vllm/utils/network_utils.py:190:def _get_open_port(
$ cd /repo && sed -n 140,200p vllm/utils/network_utils.py && sed -n 280,400p vllm/utils/network_utils.py
output · 44 lines
def get_open_zmq_ipc_path() -> str:
base_rpc_path = envs.VLLM_RPC_BASE_PATH
return f"ipc://{base_rpc_path}/{uuid4()}"
def get_open_zmq_inproc_path() -> str:
return f"inproc://{uuid4()}"
def get_open_port() -> int:
"""
Get an open port for the vLLM process to listen on.32 more lines
An edge case to handle, is when we run data parallel,
we need to avoid ports that are potentially used by
the data parallel master process.
Right now we reserve 10 ports for the data parallel master
process. Currently it uses 2 ports.
"""
if "VLLM_DP_MASTER_PORT" in os.environ:
dp_master_port = envs.VLLM_DP_MASTER_PORT
reserved_port_range = range(dp_master_port, dp_master_port + 10)
while True:
candidate_port = _get_open_port()
if candidate_port not in reserved_port_range:
return candidate_port
return
…[3517 chars omitted]…
ZMQ socket"""
ctx = zmq.Context() # type: ignore[attr-defined]
try:
yield make_zmq_socket(
ctx,
path,
socket_type,
bind=bind,
identity=identity,
router_handover=router_handover,
)
except KeyboardInterrupt:
logger.debug("Got Keyboard Interrupt.")
finally:
ctx.destroy(linger=linger)$ cd /repo && grep -rn "handshake" vllm/v1/engine/core_client.py | head -20
output · 1 lines
686: # initial handshake with the rank 0 front-end.
$ cd /repo && grep -rn "handshake" vllm/v1/engine/*.py vllm/v1/metrics.py 2>/dev/null | grep -v core.py | head -30
output · 18 lines
vllm/v1/engine/core_client.py:686: # initial handshake with the rank 0 front-end. vllm/v1/engine/utils.py:68: # Not used by engine, just relayed to front-end in handshake response. vllm/v1/engine/utils.py:75: """Metadata sent to each engine process during startup handshake, vllm/v1/engine/utils.py:112: handshake_address: str, vllm/v1/engine/utils.py:115: client_handshake_address: str | None = None, vllm/v1/engine/utils.py:122: "handshake_address": handshake_address, vllm/v1/engine/utils.py:128: if client_handshake_address: vllm/v1/engine/utils.py:129: common_kwargs["client_handshake_address"] = client_handshake_address vllm/v1/engine/utils.py:1068: engines_to_handshake = [CoreEngine(index=dp_rank, local=True)] vllm/v1/engine/utils.py:1070: # Rank 0 holds Coordinator, so it handshakes with all Cores vllm/v1/engine/ut …[1194 chars omitted]…
6 more lines
local_handshake_address, zmq.ROUTER, bind=True vllm/v1/engine/utils.py:1112: ) as handshake_socket: vllm/v1/engine/utils.py:1119: handshake_address=handshake_address, vllm/v1/engine/utils.py:1120: client_handshake_address=client_handshake_address, vllm/v1/engine/utils.py:1134: handshake_socket, vllm/v1/engine/utils.py:1136: engines_to_handshake,
$ cd /repo && sed -n 1130,1230p vllm/v1/engine/utils.py
output · 39 lines
yield local_engine_manager, coordinator, addresses, tensor_queue
# Now wait for engines to start.
wait_for_engine_startup(
handshake_socket,
addresses,
engines_to_handshake,
parallel_config,
dp_size > 1 and vllm_config.model_config.is_moe,
vllm_config.cache_config,
local_engine_manager,
coordinator.proc if coordinator else None,27 more lines
)
def wait_for_engine_startup(
handshake_socket: zmq.Socket,
addresses: EngineZmqAddresses,
core_engines: list[CoreEngine],
parallel_config: ParallelConfig,
coordinated_dp: bool,
cache_config: CacheConfig,
proc_manager: CoreEngineProcManager | None,
coord_process: Process | None,
):
# Wait for engine core process(es) to send ready messages.
local_count = parallel_config.data_parallel_size_loc
…[2773 chars omitted]…
r(
f"Remote engine {eng_index} must not use "
f"--headless in external or hybrid dp lb "
f"mode"
)
else:
raise RuntimeError(
f"Remote engine {eng_index} must use "
f"--headless unless in external or hybrid "
f"dp lb mode"
)$ cd /repo && grep -rn "client_handshake\|_perform_handshake\|handshake" vllm/v1/engine/core_client.py vllm/entrypoints/openai/api_server.py | head
output · 1 lines
vllm/v1/engine/core_client.py:686: # initial handshake with the rank 0 front-end.
$ cd /repo && grep -rn "client_handshake_address" vllm -r --include=*.py | grep -v "engine/utils.py\|engine/core.py"
$ cd /repo && sed -n 1230,1330p vllm/v1/engine/utils.py
output · 36 lines
if status == "HELLO" and engine.state == CoreEngineState.NEW:
# Send init message with DP config info.
init_message = msgspec.msgpack.encode(
EngineHandshakeMetadata(
addresses=addresses,
parallel_config={
k: getattr(parallel_config, k)
for k in (
"data_parallel_master_ip",
"data_parallel_master_port",
"_data_parallel_master_port_list",
"data_parallel_size",24 more lines
)
}
if coordinated_dp
else {},
)
)
handshake_socket.send_multipart((eng_identity, init_message), copy=False)
conn_pending[0 if local else 1] -= 1
…[1234 chars omitted]…
.READY
else:
raise RuntimeError(
f"Unexpected {status} message for "
f"{'local' if local else 'remote'} engine "
f"{eng_index} in {engine.state} state."
)
logger.debug(
"%s from %s core engine process %s.",
status,
"local" if local else "remote",
eng_index,
)$ cd /repo && sed -n 640,760p vllm/v1/engine/core_client.py
output · 31 lines
def start_engine_core_monitor(self):
"""Start a monitor thread for engine core processes."""
engine_manager = self.resources.engine_manager
if engine_manager is None:
# No engine processes to monitor
return
self_ref = weakref.ref(self)
# Monitor engine core process liveness. If any die unexpectedly,
# marks the engine as dead, and shuts down the client.
def monitor_engine_cores():19 more lines
engine_manager.monitor_engine_liveness()
_self = self_ref()
if not _self or not _self._finalizer.alive or _self.resources.engine_dead:
return
_self.resources.engine_dead = True
_self.shutdown()
# Note: For MPClient, we don't have a failure callback mechanism
# like MultiprocExecutor, but we set engine_dead flag which will
…[3604 chars omitted]…
poller.register(shutdown_socket, zmq.POLLIN)
poller.register(out_socket, zmq.POLLIN)
while True:
socks = poller.poll()
if not socks:
continue
if len(socks) == 2 or socks[0][0] == shutdown_socket:
# shutdown signal, exit thread.
break$ cd /repo && grep -rn "EngineHandshakeMetadata\|client_handshake" vllm -r --include=*.py
output · 18 lines
vllm/v1/engine/core.py:72: EngineHandshakeMetadata, vllm/v1/engine/core.py:849: client_handshake_address: str | None = None, vllm/v1/engine/core.py:876: client_handshake_address, vllm/v1/engine/core.py:957: client_handshake_address: str | None, vllm/v1/engine/core.py:982: is_local = local_client and client_handshake_address is None vllm/v1/engine/core.py:993: if client_handshake_address is None: vllm/v1/engine/core.py:1001: input_ctx, client_handshake_address, identity, True, False, vllm_config vllm/v1/engine/core.py:1081: init_message: EngineHandshakeMetadata = msgspec.msgpack.decode( vllm/v1/engine/core.py:1082: init_bytes, type=EngineHandshakeMetadata vllm/v1/engine/core.py:1662: client_handshake_address: str | None = None, vllm/v1/engine/core.py:1695: client_handshake_address, vllm/v1/engi
6 more lines
…[293 chars omitted]… : common_kwargs["client_handshake_address"] = client_handshake_address vllm/v1/engine/utils.py:1105: client_handshake_address = local_handshake_address vllm/v1/engine/utils.py:1108: client_handshake_address = None vllm/v1/engine/utils.py:1120: client_handshake_address=client_handshake_address, vllm/v1/engine/utils.py:1234: EngineHandshakeMetadata(
$ cd /repo && sed -n 2090,2200p vllm/v1/engine/core.py
output · 43 lines
f"local range: [{local_dp_rank * world_size}, "
f"{(local_dp_rank + 1) * world_size}) "
f'base value: "{os.getenv(device_control_env_var)}"'
) from e
@contextmanager
def _perform_handshakes(
self,
handshake_address: str,
identity: bytes,
local_client: bool,
vllm_config: VllmConfig,31 more lines
client_handshake_address: str | None,
):
"""
For Ray, we don't need to actually perform handshake.
All addresses information is known before the actor creation.
Therefore, we simply yield these addresses.
"""
yield self.addresses
def wait_for_init(self):
"""
Wait until the engine core is initialized.
This is just an empty method. When ray.get() on this method
(or any other method of the actor) returns, i
…[1703 chars omitted]…
el_size_local = 1
vllm_config.parallel_config.data_parallel_rank = 0
EngineCoreActorMixin.__init__(
self, vllm_config, addresses, dp_rank, local_dp_rank
)
EngineCoreProc.__init__(
self,
vllm_config,
local_client,
"",
executor_class,
log_stats,
engine_index=dp_rank,
)$ cd /repo && grep -rn "handshake" vllm --include=*.py -l | grep -v "v1/engine/core.py\|v1/engine/utils.py"
output · 21 lines
vllm/benchmarks/serve.py vllm/distributed/kv_transfer/kv_connector/utils.py vllm/distributed/kv_transfer/kv_connector/v1/base.py vllm/distributed/kv_transfer/kv_connector/v1/mooncake/mooncake_connector.py vllm/distributed/kv_transfer/kv_connector/v1/moriio/moriio_common.py vllm/distributed/kv_transfer/kv_connector/v1/moriio/moriio_connector.py vllm/distributed/kv_transfer/kv_connector/v1/multi_connector.py vllm/distributed/kv_transfer/kv_connector/v1/nixl/connector.py vllm/distributed/kv_transfer/kv_connector/v1/nixl/metadata.py vllm/distributed/kv_transfer/kv_connector/v1/nixl/scheduler.py vllm/distributed/kv_transfer/kv_connector/v1/nixl/tp_mapping.py vllm/distributed/kv_transfer/kv_connector/v1/nixl/worker.py
9 more lines
vllm/distributed/kv_transfer/kv_connector/v1/ssm_conv_transfer_utils.py vllm/entrypoints/cli/serve.py vllm/envs.py vllm/v1/engine/core_client.py vllm/v1/executor/abstract.py vllm/v1/executor/multiproc_executor.py vllm/v1/outputs.py vllm/v1/worker/gpu_worker.py vllm/vllm_flash_attn/cute/block_sparse_utils.py
$ cd /repo && grep -rn "handshake" vllm/entrypoints/cli/serve.py
output · 3 lines
225: handshake_address = get_tcp_uri(host, port) 231: handshake_address, 241: handshake_address=handshake_address,
$ cd /repo && sed -n 180,280p vllm/entrypoints/cli/serve.py
output · 43 lines
vllm_config = engine_args.create_engine_config(
usage_context=usage_context, headless=True
)
if engine_args.data_parallel_hybrid_lb:
raise ValueError("data_parallel_hybrid_lb is not applicable in headless mode")
parallel_config = vllm_config.parallel_config
local_engine_count = parallel_config.data_parallel_size_local
if local_engine_count <= 0:
raise ValueError("data_parallel_size_local must be > 0 in headless mode")31 more lines
shutdown_requested = False
# Catch SIGTERM and SIGINT to allow graceful shutdown.
def signal_handler(signum, frame):
nonlocal shutdown_requested
logger.debug("Received %d signal.", signum)
if not shutdown_requested:
shutdown_requested = True
raise SystemExit
signal.signal(signal.SIGTERM, signal_handler)
signal.signal(signal.SIGINT, signal_handler)
…[2128 chars omitted]…
1"
)
if num_api_servers > 1:
setup_multiprocess_prometheus()
shutdown_requested = False
# Catch SIGTERM and SIGINT to allow graceful shutdown.
def signal_handler(signum, frame):
nonlocal shutdown_requested
logger.debug("Received %d signal.", signum)
if not shutdown_requested:
shutdown_requested = True
raise SystemExit$ cd /repo && grep -rn "handshake\|client_handshake" vllm/v1/engine/async_llm.py vllm/entrypoints/openai/api_server.py vllm/v1/engine/output_processor.py 2>/dev/null | head
$ cd /repo && grep -rn "client_handshake_address\|local_handshake" vllm/v1/engine/core.py | head; sed -n 1640,1720p vllm/v1/engine/core.py
output · 27 lines
849: client_handshake_address: str | None = None,
876: client_handshake_address,
957: client_handshake_address: str | None,
982: is_local = local_client and client_handshake_address is None
993: if client_handshake_address is None:
1000: local_handshake = self._perform_handshake(
1001: input_ctx, client_handshake_address, identity, True, False, vllm_config
1003: with handshake as addresses, local_handshake as client_addresses:
1662: client_handshake_address: str | None = None,
1695: client_handshake_address,
def _send_abort_outputs(self, aborted_reqs: list[tuple[str, int]]) -> None:
# TODO(nick) this will be moved inside the scheduler15 more lines
if aborted_reqs:
# Map client_index to list of request_ids that belong to that client.
by_client = defaultdict[int,
…[2455 chars omitted]…
rank < dp_size
self.dp_rank = dp_rank
self.dp_size = dp_size
dp_group, dp_store = parallel_config.stateless_init_dp_group(return_store=True)
self.dp_group, self.dp_store = dp_group, dp_store
def shutdown(self):
super().shutdown()
if dp_group := getattr(self, "dp_group", None):
stateless_destroy_torch_distributed_process_group(dp_group)$ cd /repo && grep -rn "EngineZmqAddresses(" vllm --include=*.pyoutput · 1 lines
vllm/v1/engine/utils.py:981: return EngineZmqAddresses(
$ cd /repo && grep -rn "handshake" vllm/v1/engine/core.py | sed -n 1,40p && sed -n 830,900p vllm/v1/engine/core.py
output · 27 lines
167: # handshake metadata from all workers so the connector in the scheduler 173: xfer_handshake_metadata = ( 174: self.model_executor.get_kv_connector_handshake_metadata() 177: if xfer_handshake_metadata: 178: # xfer_handshake_metadata is list of dicts from workers 182: for worker_dict in xfer_handshake_metadata: 185: kv_connector.set_xfer_handshake_metadata(content) 846: handshake_address: str, 849: client_handshake_address: str | None = None, 871: with self._perform_handshakes( 872: handshake_address, 876: client_handshake_address,
15 more lines
942: # Don't complete handshake until DP coordinator ready message is
951: def _perform_handshakes(
953: handshake_address: str,
957: client_handshake_address: str | None,
960: Perform star
…[3598 chars omitted]…
e stats to coordinator for "internal"
# and "hybrid" LB modes.
self.publish_dp_lb_stats = internal_dp_balancing
self.addresses = addresses
self.process_input_queue_block = True
if envs.VLLM_ELASTIC_EP_SCALE_UP_LAUNCH:
self._eep_send_engine_core_notification(
EEPNotificationType.NEW_CORE_ENGINES_INIT_READY,$ cd /repo && grep -rn '"HELLO"\|HELLO' vllm --include=*.py | grep -v test | head
output · 3 lines
vllm/v1/engine/core.py:1065: "status": "HELLO", vllm/v1/engine/utils.py:1198: # Receive HELLO and READY messages from the input socket. vllm/v1/engine/utils.py:1231: if status == "HELLO" and engine.state == CoreEngineState.NEW:
$ cd /repo && sed -n 1055,1140p vllm/v1/engine/core.py
output · 34 lines
def startup_handshake(
handshake_socket: zmq.Socket,
local_client: bool,
headless: bool,
parallel_config: ParallelConfig | None = None,
) -> EngineZmqAddresses:
# Send registration message.
handshake_socket.send(
msgspec.msgpack.encode(
{
"status": "HELLO",
"local": local_client,22 more lines
"headless": headless,
}
)
)
# Receive initialization message.
logger.debug("Waiting for init message from front-end.")
if not handshake_socket.poll(timeout=HANDSHAKE_TIMEOUT_MINS * 60_000):
raise RuntimeError(
"Did not receive response from front-end "
f"process within {HANDSHAKE_TIMEOUT_MINS} "
f"minutes"
)
init_bytes = handshak
…[2499 chars omitted]…
dent, so treat like DP=1.
# Note that parallel_config.data_parallel_index will still reflect
# the original DP rank.
parallel_config.data_parallel_size = 1
parallel_config.data_parallel_size_local = 1
parallel_config.data_parallel_rank = 0
engine_core = EngineCoreProc(*args, engine_index=dp_rank, **kwargs)$ cd /repo && grep -rn "ROUTER\|XSUB\|handover" vllm/v1/engine/core_client.py | head
output · 7 lines
498: # identity. The client input ROUTER needs handover to allow the new 500: enable_input_socket_handover = parallel_config.enable_elastic_ep 514: zmq.ROUTER, 516: router_handover=enable_input_socket_handover, 527: zmq.ROUTER, 529: router_handover=enable_input_socket_handover, 1190: make_zmq_socket(self.ctx, stats_addr, zmq.XSUB, linger=0) as socket,
$ cd /repo && sed -n 440,510p vllm/v1/engine/core_client.py
output · 30 lines
# We must ensure that the sync output socket is
# closed cleanly in its own thread.
with self.ctx.socket(zmq.PAIR) as shutdown_sender:
shutdown_sender.connect(self.shutdown_path)
# Send shutdown signal.
shutdown_sender.send(b"")
def validate_alive(self, frames: Sequence[zmq.Frame]):
if len(frames) == 1 and (frames[0].buffer == EngineCoreProc.ENGINE_CORE_DEAD):
self.engine_dead = True
raise EngineDeadError()
18 more lines
@dataclass
class ElasticScalingCache:
existing_core_engines: list[EngineIdentity]
num_new_core_engines: int
pending_notifications: dict[EEPNotificationType, set[int]]
class MPClient(EngineCoreClient):
"""
MPClient: base client for multi-proc EngineCore.
EngineCore runs in a background process busy loop, gett
…[1653 chars omitted]…
input_address = client_addresses["input_address"]
output_address = client_addresses["output_address"]
self.stats_update_address = client_addresses.get("stats_update_address")
# Tensor queues passed via client_addresses for multi-API-server case
tensor_queue = client_addresses.get("tensor_queue") # type: ignore[assignment]$ cd /repo && sed -n 75,160p vllm/v1/engine/async_llm.py && sed -n 200,260p vllm/v1/engine/async_llm.py
output · 36 lines
vllm_config: VllmConfig,
executor_class: type[Executor],
log_stats: bool,
usage_context: UsageContext = UsageContext.ENGINE_CONTEXT,
mm_registry: MultiModalRegistry = MULTIMODAL_REGISTRY,
log_requests: bool = True,
start_engine_loop: bool = True,
stat_loggers: list[StatLoggerFactory] | None = None,
aggregate_engine_logging: bool = False,
client_addresses: dict[str, str] | None = None,
client_count: int = 1,
client_index: int = 0,24 more lines
) -> None:
"""
Create an AsyncLLM.
Args:
vllm_config: global configuration.
executor_class: an Executor impl, e.g. MultiprocExecutor.
log_stats: Whether to log stats.
usage_context: Usage context of the LLM.
mm_registry: Multi-modal registry.
log_requests: Whether t
…[4359 chars omitted]…
engine_args.enable_log_requests,
log_stats=not engine_args.disable_log_stats,
start_engine_loop=start_engine_loop,
usage_context=usage_context,
stat_loggers=stat_loggers,
)
def __del__(self):
self.shutdown()
def shutdown(self, timeout: float | None = None) -> None:
"""Shutdown, cleaning up the background proc and IPC."""Now I'll implement the changes.
$ cd /repo && sed -n 140,175p vllm/v1/utils.py && sed -n 1,40p vllm/v1/utils.py
output · 37 lines
"""NOTE: Because this method is non-blocking, explicit synchronization
is needed to ensure the data is copied to CPU."""
if n is None:
return self.cpu.copy_(self.gpu, non_blocking=True)
return self.cpu[:n].copy_(self.gpu[:n], non_blocking=True)
def get_engine_client_zmq_addr(local_only: bool, host: str, port: int = 0) -> str:
"""Assign a new ZMQ socket address.
If local_only is True, participants are colocated and so a unique IPC
address will be returned.25 more lines
Otherwise, the provided host and port will be used to construct a TCP
address (port == 0 means assign an available port)."""
return (
get_open_zmq_ipc_path()
if local_only
else (get_tcp_uri(host, port or get_open_port()))
)
class APIServerProcessManager:
"""Manages a group of API server processes.
Handles creation, monitori
…[1076 chars omitted]…
lm.utils.network_utils import get_open_port, get_open_zmq_ipc_path, get_tcp_uri
from vllm.utils.system_utils import decorate_logs, kill_process_tree, set_process_title
from vllm.v1.core.sched.output import SchedulerOutput
if TYPE_CHECKING:
import numpy as np
from vllm.v1.engine.coordinator import DPCoordinator
from vllm.v1.engine.utils import CoreEngineActorManager, CoreEngineProcManager$ cd /repo && python - <<'PY'
output · 1 lines
ok
$ cd /repo && python - <<'PY'
output · 36 lines
def gather_actual_addresses(
self, timeout: float = 60.0
) -> tuple[list[str], list[str]]:
"""Wait for every API server process to report the ZMQ endpoints
it actually bound.
Args:
timeout: Maximum time to wait, in seconds.
Returns:
Tuple of (input_addresses, output_addresses), one entry per
API server process, in child order.24 more lines
Raises:
RuntimeError: If a child process exits before reporting its
endpoints, or if the timeout expires.
"""
if self._gathered_addresses is not None:
return self._gathered_addresses
num_servers = len(self.processes)
input_addresses: list[str | None] = [None] * num_servers
output_addresses: list[str | None] = [None] * num_servers
pipe_to_index = {pipe: i for i, pipe in enumerate(se
…[1997 chars omitted]…
f"API servers reported duplicate addresses: "
f"inputs={inputs}, outputs={outputs}"
)
self._gathered_addresses = (inputs, outputs)
for pipe in self._address_pipes:
with contextlib.suppress(Exception):
pipe.close()
return self._gathered_addresses
def shutdown(self, timeout: float | None = None) -> None:$ cd /repo && python - <<'PY'
output · 1 lines
ok
$ cd /repo && python - <<'PY'
output · 1 lines
ok
$ cd /repo && python - <<'PY'
output · 1 lines
ok
Now the child-side reporting in `MPClient`:
$ cd /repo && python - <<'PY'
output · 2 lines
ok parsed
$ cd /repo && grep -n "get_open_port" vllm/v1/utils.py; grep -n "^from\|^import" vllm/v1/engine/core_client.py | head -30
output · 31 lines
32:from vllm.utils.network_utils import get_open_port, get_open_zmq_ipc_path, get_tcp_uri 3:import asyncio 4:import contextlib 5:import queue 6:import sys 7:import uuid 8:import weakref 9:from abc import ABC, abstractmethod 10:from collections import defaultdict, deque 11:from collections.abc import Awaitable, Callable, Sequence 12:from concurrent.futures import Future 13:from dataclasses import dataclass
19 more lines
14:from multiprocessing.queues import Queue 15:from threading import Thread 16:from typing import Any, TypeAlias, TypeVar 18:import msgspec.msgpack 19:import zmq 20:import zmq.asyncio 22:from vllm.config import VllmConfig 23:from vllm.envs import VLLM_ENGINE_READY_TIMEOUT_S 24:from vllm.logger import init_logger 25:from vllm.lora.request import LoRARequest 26:from vllm.tasks import SupportedTask 27:from vllm.tracing import instrument 28:from vllm.utils.async_utils import in_loop 29:from vllm.utils.network_utils import ( 34:from vllm.v1.engine import ( 46:from vllm.v1.engine.coordinator import DPCoordinator 47:from vllm.v1.engine.core import EngineCore, EngineCoreProc 48:from vllm.v1.engine.exceptions import EngineDeadError 49:from vllm.v1.engine.tensor_ipc import TensorIpcSender
$ cd /repo && sed -i 's/from vllm.utils.network_utils import get_open_port, get_open_zmq_ipc_path, get_tcp_uri/from vllm.utils.network_utils import get_open_zmq_ipc_path, get_tcp_uri/' vllm/v1/utils.py && cat > /tmp/repro.py <<'EOF'
output · 27 lines
placeholder: tcp://127.0.0.1:0
INFO 10-01 19:51:18 [utils.py:241] Started 3 API server processes
placeholder: tcp://127.0.0.1:0
Traceback (most recent call last):
File "<string>", line 1, in <module>
File "/usr/local/lib/python3.12/multiprocessing/spawn.py", line 122, in spawn_main
exitcode = _main(fd, parent_sentinel)
^^^^^^^^^^^^^^^^^^^^^^^^^^
File "/usr/local/lib/python3.12/multiprocessing/spawn.py", line 131, in _main
prepare(preparation_data)
File "/usr/local/lib/python3.12/multiprocessing/spawn.py", line 246, in prepare
_fixup_main_from_path(data['init_main_from_path'])15 more lines
File "/usr/local/lib/python3.12/multiprocessing/spawn.py", line 297, in _fixup_main_from_path
main_content = runpy.run_path(main_path,
^^^^^^^^^^^^^^^^^^^^^^^^^
File "<frozen runpy>", line 287, in run_path
File "<frozen runpy>", line 98, in _run_modu
…[7715 chars omitted]…
or: API server process ApiServer_0 (PID: 953) exited with code 1 before reporting its addresses.
During handling of the above exception, another exception occurred:
Traceback (most recent call last):
File "/tmp/repro.py", line 35, in <module>
assert expect_fail, e
^^^^^^^^^^^
AssertionError: API server process ApiServer_0 (PID: 953) exited with code 1 before reporting its addresses.$ cd /repo && python - <<'PY'
output · 4 lines
File "/tmp/repro.py", line 11
pipe = client_config["actual_address_pipe"]
^^^^
IndentationError: expected an indented block after function definition on line 10$ cd /repo && cat > /tmp/repro.py <<'EOF'
output · 22 lines
placeholder: tcp://127.0.0.1:0
INFO 10-01 19:51:43 [utils.py:241] Started 3 API server processes
inputs: ['tcp://127.0.0.1:44195', 'tcp://127.0.0.1:37975', 'tcp://127.0.0.1:40303']
outputs: ['tcp://127.0.0.1:42913', 'tcp://127.0.0.1:39113', 'tcp://127.0.0.1:35183']
INFO 10-01 19:51:49 [utils.py:241] Started 3 API server processes
Traceback (most recent call last):
File "/tmp/repro.py", line 37, in main
ins, outs = m.gather_actual_addresses(timeout=15)
^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
File "/repo/vllm/v1/utils.py", line 277, in gather_actual_addresses
raise RuntimeError(
RuntimeError: Timed out after 15 seconds waiting for API server processes to report their addresses. Still waiting on: ['ApiServer_1', 'ApiServer_2', 'ApiServer_0']10 more lines
During handling of the above exception, another exception occurred:
Traceback (most recent call last):
File "/tmp/repro.py", line 54, in <module>
main()
File "/tmp/repro.py", line 40, in main
assert "reporting" in str(e), e
^^^^^^^^^^^^^^^^^^^^^
AssertionError: Timed out after 15 seconds waiting for API server processes to report their addresses. Still waiting on: ['ApiServer_1', 'ApiServer_2', 'ApiServer_0']$ cd /repo && sed -i 's/def bad_child(listen_address, sock, args, client_config):\n time.sleep(60)//' /tmp/repro.py && python - <<'PY'
output · 7 lines
placeholder: tcp://127.0.0.1:0 INFO 10-01 19:52:15 [utils.py:241] Started 3 API server processes inputs: ['tcp://127.0.0.1:41977', 'tcp://127.0.0.1:39479', 'tcp://127.0.0.1:35745'] outputs: ['tcp://127.0.0.1:39863', 'tcp://127.0.0.1:37789', 'tcp://127.0.0.1:33873'] INFO 10-01 19:52:22 [utils.py:241] Started 3 API server processes expected failure: API server process ApiServer_2 (PID: 1170) exited with code 0 before reporting its addresses. PASS
$ cd /repo && grep -rln "APIServerProcessManager\|get_engine_client_zmq_addr\|get_open_port" tests | head
output · 10 lines
tests/distributed/test_dcp_a2a.py tests/distributed/test_mnnvl_alltoall.py tests/distributed/test_node_count.py tests/distributed/test_rocm_quick_reduce.py tests/distributed/test_same_node.py tests/distributed/test_shm_broadcast.py tests/distributed/test_utils.py tests/distributed/test_weight_transfer.py tests/entrypoints/openai/completion/test_shutdown.py tests/entrypoints/serve/instrumentator/test_metrics.py
$ cd /repo && grep -rn "get_engine_client_zmq_addr\|APIServerProcessManager" tests | head; python -m pytest tests/distributed/test_utils.py -x -q 2>&1 | tail -5
output · 12 lines
tests/entrypoints/test_api_server_process_manager.py:12:from vllm.v1.utils import APIServerProcessManager, wait_for_completion_or_failure tests/entrypoints/test_api_server_process_manager.py:28: """Fixture to provide arguments for APIServerProcessManager.""" tests/entrypoints/test_api_server_process_manager.py:52: """Test initializing the APIServerProcessManager.""" tests/entrypoints/test_api_server_process_manager.py:62: manager = APIServerProcessManager(**args) tests/entrypoints/test_api_server_process_manager.py:99: manager = APIServerProcessManager(**api_server_args) tests/entrypoints/test_api_server_process_manager.py:157: manager = APIServerProcessManager(**api_server_args) tests/entrypoints/test_api_server_process_manager.py:216: manager = APIServerProcessManager(**api_server_args) =========================== short test summary info ============================ ERROR tests/distributed/test_utils.py !!!!!!!!!!!!!!!!!!!!!!!!!! stopping after 1 failures !!!!!!!!!!!!!!!!!!!!!!!!!!! !!!!!!!!!!!!!!!!!!!! Interrupted: 1 error during collection !!!!!!!!!!!!!!!!!!!! 1 error in 1.72s
$ cd /repo && sed -n 1,120p tests/entrypoints/test_api_server_process_manager.py
output · 44 lines
# SPDX-License-Identifier: Apache-2.0 # SPDX-FileCopyrightText: Copyright contributors to the vLLM project import multiprocessing import socket import threading import time from unittest.mock import patch import pytest from vllm.v1.utils import APIServerProcessManager, wait_for_completion_or_failure
32 more lines
# Global variables to control worker behavior
WORKER_RUNTIME_SECONDS = 0.5
# Mock implementation of run_api_server_worker
def mock_run_api_server_worker(listen_address, sock, args, client_config=None):
"""Mock run_api_server_worker that runs for a specific time."""
print(f"Mock worker started with client_config: {client_config}")
time.sleep(WORKER_RUNTIME_SECONDS)
print("Mock worker completed successfully")
@pytest.fixture
def api_server_args():
"""Fixture to provide arguments for APIServerProcessManager."""
sock = socket.socket()
return {
"targe
…[2508 chars omitted]…
ver_manager=manager)
except Exception as e:
result["exception"] = e
finally:
manager.shutdown()
# Start a thread to run wait_for_completion_or_failure
wait_thread = threading.Thread(target=run_with_exception_capture, daemon=True)
wait_thread.start()
# Let all processes run for a short time
time.sleep(0.2)$ cd /repo && timeout 300 python -m pytest tests/entrypoints/test_api_server_process_manager.py -q 2>&1 | tail -5
output · 5 lines
/repo/tests/entrypoints/test_api_server_process_manager.py:188: PytestUnknownMarkWarning: Unknown pytest.mark.timeout - is this a typo? You can register custom marks to avoid this warning - for details, see https://docs.pytest.org/en/stable/how-to/mark.html
@pytest.mark.timeout(30)
-- Docs: https://docs.pytest.org/en/stable/how-to/capture-warnings.html
5 passed, 2 warnings in 11.01s$ cd /repo && timeout 120 python -m pytest tests/distributed/test_utils.py -q 2>&1 | tail -8; git stash list
output · 8 lines
------------------------------- Captured stdout -------------------------------- INFO 10-01 19:54:00 [importing.py:46] Triton is installed but 0 active driver(s) found (expected 1). Disabling Triton to prevent runtime errors. WARNING 10-01 19:54:00 [importing.py:58] Triton is installed, but doesn't include CPU backend. Disabling Triton. INFO 10-01 19:54:00 [importing.py:81] Triton not installed or not compatible; certain GPU-related functions will not be available. =========================== short test summary info ============================ ERROR tests/distributed/test_utils.py !!!!!!!!!!!!!!!!!!!! Interrupted: 1 error during collection !!!!!!!!!!!!!!!!!!!! 1 error in 1.18s
$ cd /repo && timeout 120 python -m pytest tests/distributed/test_utils.py -q 2>&1 | grep -B5 Error | head -20
output · 10 lines
==================================== ERRORS ====================================
_______________ ERROR collecting tests/distributed/test_utils.py _______________
ImportError while importing test module '/repo/tests/distributed/test_utils.py'.
--
/usr/local/lib/python3.12/importlib/__init__.py:90: in import_module
return _bootstrap._gcd_import(name[level:], package, level)
^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
tests/distributed/test_utils.py:7: in <module>
import ray
E ModuleNotFoundError: No module named 'ray'$ cd /repo && (ruff check vllm/v1/utils.py vllm/v1/engine/core_client.py vllm/entrypoints/cli/serve.py 2>&1 | tail -5); python -m pytest tests/entrypoints/openai/completion/test_shutdown.py -q 2>&1 | tail -3
output · 4 lines
bash: line 537: ruff: command not found FAILED tests/entrypoints/openai/completion/test_shutdown.py::test_request_rejection_during_shutdown FAILED tests/entrypoints/openai/completion/test_shutdown.py::test_multi_api_server_shutdown 8 failed, 21 warnings in 16.05s
$ cd /repo && git stash && timeout 300 python -m pytest tests/entrypoints/openai/completion/test_shutdown.py -q 2>&1 | tail -3; git stash pop
output · 14 lines
Saved working directory and index state WIP on master: 1e54ea4 base FAILED tests/entrypoints/openai/completion/test_shutdown.py::test_request_rejection_during_shutdown FAILED tests/entrypoints/openai/completion/test_shutdown.py::test_multi_api_server_shutdown 8 failed, 21 warnings in 9.38s On branch master Changes not staged for commit: (use "git add <file>..." to update what will be committed) (use "git restore <file>..." to discard changes in working directory) modified: vllm/entrypoints/cli/serve.py modified: vllm/v1/engine/core_client.py modified: vllm/v1/utils.py
2 more lines
no changes added to commit (use "git add" and/or "git commit -a")
Dropped refs/stash@{0} (c0529829e6efa1562c66cbab4c67b5e6209dbe1f)$ cd /repo && git diff
output · 31 lines
diff --git a/vllm/entrypoints/cli/serve.py b/vllm/entrypoints/cli/serve.py
index 5972bd4..34534d1 100644
--- a/vllm/entrypoints/cli/serve.py
+++ b/vllm/entrypoints/cli/serve.py
@@ -340,6 +340,12 @@ def run_multi_api_server(args: argparse.Namespace):
stats_update_address=stats_update_address,
tensor_queue=tensor_queue,
)
+ # Wait until every API server has bound its (kernel-assigned)
+ # ZMQ endpoints and propagate the actual addresses so that
+ # they can be used by the engines.
+ actual_inputs, actual_outputs = api_server_manager.gather_actual_addresses()19 more lines
+ addresses.inputs = actual_inputs
+ addresses.outputs = actual_outputs
# Wait for API servers.
try:
diff --git a/vllm/v1/engine/core_client.py b/vllm/v1/engine/core_client.py
index 2f2c15d..be4a3c7 100644
---
…[8628 chars omitted]…
+
+ self._gathered_addresses = (inputs, outputs)
+ for pipe in self._address_pipes:
+ with contextlib.suppress(Exception):
+ pipe.close()
+ return self._gathered_addresses
+
def shutdown(self, timeout: float | None = None) -> None:
"""Shutdown API server processes with configurable timeout"""
if self._finalizer.detach() is not None:$ cd /repo && python - <<'PY'
output · 6 lines
outputs: ['tcp://127.0.0.1:34923', 'tcp://127.0.0.1:37045', 'tcp://127.0.0.1:42801'] INFO 10-01 19:56:14 [utils.py:241] Started 3 API server processes expected failure: API server process ApiServer_0 (PID: 1628) exited with code 0 before reporting its addresses. PASS -- Docs: https://docs.pytest.org/en/stable/how-to/capture-warnings.html 5 passed, 2 warnings in 10.04s
$ cd /repo && python -c "import vllm.v1.utils, vllm.v1.engine.core_client, vllm.entrypoints.cli.serve; print('imports ok')"output · 4 lines
INFO 10-01 19:56:54 [importing.py:46] Triton is installed but 0 active driver(s) found (expected 1). Disabling Triton to prevent runtime errors. WARNING 10-01 19:56:54 [importing.py:58] Triton is installed, but doesn't include CPU backend. Disabling Triton. INFO 10-01 19:56:54 [importing.py:81] Triton not installed or not compatible; certain GPU-related functions will not be available. imports ok
$ done: {"notes": "Fix multi-API-server port race: get_engine_client_zmq_addr now returns tcp://host:0 placeholder (kernel-assigned port) instead of pre-selecting a port; APIServerProcessManager creates a per-child pipe (actual_address_pipe in client_config), closes the parent send end after start, and gains gather_actual_addresses(timeoutoutput · 3 lines
…[280 chars omitted]… sses; MPClient resolves bound endpoints via zmq.LAST_ENDPOINT and reports them on the pipe. Reproducer and tests/entrypoints/test_api_server_process_manager.py pass."} submitted.