SWE-Race › Tasks › vllm-project-vllm-42585 ← prevnext →

vllm-project-vllm-42585

vllm-project/vllmhardsinglemerged 2026-05-26Apache-2.0fix: 6 files, +187 −492 fail-to-pass · 5 pass-to-pass
Results
Modelsolved / attemptsmedian stepsmedian costattempts
GPT-5.6 Luna6/824$0.0231✗ 2✗ 3✓ 4✓ 5✓ 6✓ 7✓ 8✓
DeepSeek V4 Flash2/424$0.0181✗ 2✗ 3✓ 4✓
GLM-5.3 Flash2/350$0.0291✗ 2✓ 3✓
The prompt the agent sees

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_gather_actual_addresses_child_crash_before_reporttest_gather_actual_addresses_end_to_end
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]

What each model did

Submitted diff first, then every command. Outputs show 12 lines; expand for the rest.
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]
step 1
$ 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:
step 2
$ 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
step 3
$ cd /repo && sed -n 220,320p vllm/v1/utils.py && grep -rn "APIServerProcessManager(" vllm -n
output · 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(
step 4
$ 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))
step 5
$ 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_config
23 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,
step 6
$ 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
step 7
$ 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,
step 8
$ 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(
step 9
$ 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
step 10
$ 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
step 11
$ 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
step 12
$ 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)
step 13
$ 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,
step 14
$ 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,
step 15
$ 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
step 16
$ 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_config
25 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
step 17
$ 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
step 18
$ 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(
step 19
$ 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)
step 20
$ 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.
step 21
$ 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,
step 22
$ 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"
                )
step 23
$ 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.
step 24
$ cd /repo && grep -rn "client_handshake_address" vllm -r --include=*.py | grep -v "engine/utils.py\|engine/core.py"
step 25
$ 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,
        )
step 26
$ 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
step 27
$ 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(
step 28
$ 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,
        )
step 29
$ 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
step 30
$ 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,
step 31
$ 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
step 32
$ 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
step 33
$ 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 scheduler
15 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)
step 34
$ cd /repo && grep -rn "EngineZmqAddresses(" vllm --include=*.py
output · 1 lines
vllm/v1/engine/utils.py:981:    return EngineZmqAddresses(
step 35
$ 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,
step 36
$ 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:
step 37
$ 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)
step 38
$ 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,
step 39
$ 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]
step 40
$ 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."""
step 41

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
step 42
$ cd /repo && python - <<'PY'
output · 1 lines
ok
step 43
$ 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:
step 44
$ cd /repo && python - <<'PY'
output · 1 lines
ok
step 45
$ cd /repo && python - <<'PY'
output · 1 lines
ok
step 46
$ cd /repo && python - <<'PY'
output · 1 lines
ok
step 47

Now the child-side reporting in `MPClient`:

$ cd /repo && python - <<'PY'
output · 2 lines
ok
parsed
step 48
$ 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
step 49
$ 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.
step 50
$ 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
step 51
$ 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']
step 52
$ 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
step 53
$ 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
step 54
$ 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
step 55
$ 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)
step 56
$ 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
step 57
$ 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
step 58
$ 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'
step 59
$ 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
step 60
$ 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)
step 61
$ 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:
step 62
$ 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
step 63
$ 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
step 64
$ 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(timeout
output · 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.