Skip to content

vllm.utils.network_utils

Classes:

  • ZmqListener –

    A supervisor-bound listener that a ZMQ socket can adopt.

Functions:

ZmqListener dataclass

A supervisor-bound listener that a ZMQ socket can adopt.

Source code in vllm/utils/network_utils.py
@dataclass
class ZmqListener:
    """A supervisor-bound listener that a ZMQ socket can adopt."""

    address: str
    socket: socket.socket

    def close(self) -> None:
        self.socket.close()

    def cleanup(self) -> None:
        self.close()
        if self.address.startswith("ipc://"):
            with contextlib.suppress(FileNotFoundError):
                os.unlink(self.address.removeprefix("ipc://"))

_get_reserved_port_range()

Ports reserved for the data parallel master process (empty if unset).

Source code in vllm/utils/network_utils.py
def _get_reserved_port_range() -> range:
    """Ports reserved for the data parallel master process (empty if unset)."""
    if "VLLM_DP_MASTER_PORT" not in os.environ:
        return range(0)
    dp_master_port = envs.VLLM_DP_MASTER_PORT
    return range(dp_master_port, dp_master_port + 10)

_get_zmq_socket_buffer_size()

Choose the shared libzmq and inherited-listener buffer policy.

Source code in vllm/utils/network_utils.py
def _get_zmq_socket_buffer_size() -> int:
    """Choose the shared libzmq and inherited-listener buffer policy."""
    mem = psutil.virtual_memory()
    # For systems with substantial memory (>32GB total, >16GB available):
    # - Set a large 0.5GB buffer to improve throughput
    # For systems with less memory:
    # - Use system default (-1) to avoid excessive memory consumption
    if (
        mem.total > _ZMQ_LARGE_BUFFER_MIN_TOTAL_MEMORY
        and mem.available > _ZMQ_LARGE_BUFFER_MIN_AVAILABLE_MEMORY
    ):
        return _ZMQ_LARGE_BUFFER_SIZE
    return -1

aiter_requires_tcp_store()

AITER custom all-reduce requires a pure-TCP default store (its IPC metadata exchange asserts on TCPStore); the file:// rendezvous yields a FileStore and trips that assertion. Prefer the TCP rendezvous (pre-#50999) for ROCm + AITER custom AR until AITER accepts FileStore.

Source code in vllm/utils/network_utils.py
def aiter_requires_tcp_store() -> bool:
    """AITER custom all-reduce requires a pure-TCP default store (its IPC
    metadata exchange asserts on ``TCPStore``); the file:// rendezvous yields a
    ``FileStore`` and trips that assertion. Prefer the TCP rendezvous
    (pre-#50999) for ROCm + AITER custom AR until AITER accepts FileStore.
    """
    from vllm._aiter_ops import rocm_aiter_ops

    return rocm_aiter_ops.is_custom_all_reduce_enabled()

get_open_port()

Get an open port for the vLLM process to listen on. 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.

Source code in vllm/utils/network_utils.py
def get_open_port() -> int:
    """Get an open port for the vLLM process to listen on.
    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.
    """
    reserved_port_range = _get_reserved_port_range()
    port = _get_open_port()
    if port in reserved_port_range:
        port = _get_open_port(start_port=reserved_port_range.stop, max_attempts=1000)
    return port

get_open_ports_list(count=5)

Get a list of unique open ports.

When VLLM_PORT is set, scans upward from that port, advancing the start position after each find so every port is unique.

Source code in vllm/utils/network_utils.py
def get_open_ports_list(count: int = 5) -> list[int]:
    """Get a list of unique open ports.

    When VLLM_PORT is set, scans upward from that port, advancing
    the start position after each find so every port is unique.
    """
    ports_set = set[int]()
    if envs.VLLM_PORT is not None:
        reserved_port_range = _get_reserved_port_range()
        next_port = envs.VLLM_PORT
        for _ in range(count):
            port = _get_open_port(start_port=next_port, max_attempts=1000)
            if port in reserved_port_range:
                port = _get_open_port(
                    start_port=reserved_port_range.stop, max_attempts=1000
                )
            ports_set.add(port)
            next_port = port + 1
        return list(ports_set)
    else:
        while len(ports_set) < count:
            ports_set.add(get_open_port())

    return list(ports_set)

make_zmq_listener(path, socket_type)

Bind and listen on a raw socket for later ZMQ adoption.

Source code in vllm/utils/network_utils.py
def make_zmq_listener(path: str, socket_type: Any) -> ZmqListener:
    """Bind and listen on a raw socket for later ZMQ adoption."""
    scheme, host, port = split_zmq_path(path)
    if scheme == "tcp":
        family = socket.AF_INET6 if is_valid_ipv6_address(host) else socket.AF_INET
        bind_address: str | tuple[str, int] = (host, int(port))
    elif scheme == "ipc":
        family = socket.AF_UNIX
        bind_address = path.removeprefix("ipc://")
        os.makedirs(os.path.dirname(bind_address), exist_ok=True)
    else:
        raise ValueError(f"Cannot inherit a {scheme} ZMQ listener")

    listener = socket.socket(family, socket.SOCK_STREAM)
    bound = False
    try:
        buf_size = _get_zmq_socket_buffer_size()
        if buf_size >= 0:
            # Accepted stream sockets inherit these settings from the listener.
            # ZMQ_USE_FD adopts this pre-created socket, so the policy belongs
            # here as well as on the eventual libzmq socket.
            if socket_type in (zmq.PULL, zmq.DEALER, zmq.ROUTER):
                listener.setsockopt(socket.SOL_SOCKET, socket.SO_RCVBUF, buf_size)
            if socket_type in (zmq.PUSH, zmq.DEALER, zmq.ROUTER):
                listener.setsockopt(socket.SOL_SOCKET, socket.SO_SNDBUF, buf_size)

        listener.bind(bind_address)
        bound = True
        listener.listen()
        if scheme == "tcp":
            path = get_tcp_uri(host, listener.getsockname()[1])
        return ZmqListener(address=path, socket=listener)
    except BaseException:
        listener.close()
        # A failed bind leaves any existing pathname owned by its listener.
        if scheme == "ipc" and bound:
            assert isinstance(bind_address, str)
            with contextlib.suppress(FileNotFoundError):
                os.unlink(bind_address)
        raise

make_zmq_path(scheme, host, port=None)

Make a ZMQ path from its parts.

Parameters:

  • scheme

    (str) –

    The ZMQ transport scheme (e.g. tcp, ipc, inproc).

  • host

    (str) –

    The host - can be an IPv4 address, IPv6 address, or hostname.

  • port

    (int | None, default: None ) –

    Optional port number, only used for TCP sockets.

Returns:

  • str –

    A properly formatted ZMQ path string.

Source code in vllm/utils/network_utils.py
def make_zmq_path(scheme: str, host: str, port: int | None = None) -> str:
    """Make a ZMQ path from its parts.

    Args:
        scheme: The ZMQ transport scheme (e.g. tcp, ipc, inproc).
        host: The host - can be an IPv4 address, IPv6 address, or hostname.
        port: Optional port number, only used for TCP sockets.

    Returns:
        A properly formatted ZMQ path string.

    """
    if port is None:
        return f"{scheme}://{host}"
    if is_valid_ipv6_address(host):
        return f"{scheme}://[{host}]:{port}"
    return f"{scheme}://{host}:{port}"

make_zmq_socket(ctx, path, socket_type, bind=None, identity=None, linger=None, router_handover=False, listener=None)

Make a ZMQ socket with the proper bind/connect semantics.

When supplied, listener is detached and its fd moves to libzmq.

Source code in vllm/utils/network_utils.py
def make_zmq_socket(
    ctx: zmq.asyncio.Context | zmq.Context,  # type: ignore[name-defined]
    path: str,
    socket_type: Any,
    bind: bool | None = None,
    identity: bytes | None = None,
    linger: int | None = None,
    router_handover: bool = False,
    listener: socket.socket | None = None,
) -> zmq.Socket | zmq.asyncio.Socket:  # type: ignore[name-defined]
    """Make a ZMQ socket with the proper bind/connect semantics.

    When supplied, ``listener`` is detached and its fd moves to libzmq.
    """
    socket = ctx.socket(socket_type)
    buf_size = _get_zmq_socket_buffer_size()

    if bind is None:
        bind = socket_type not in (zmq.PUSH, zmq.SUB, zmq.XSUB)
    if listener is not None and not bind:
        raise ValueError("An inherited ZMQ listener requires bind=True")

    if socket_type in (zmq.PULL, zmq.DEALER, zmq.ROUTER):
        socket.setsockopt(zmq.RCVHWM, 0)
        socket.setsockopt(zmq.RCVBUF, buf_size)

    if socket_type in (zmq.PUSH, zmq.DEALER, zmq.ROUTER):
        socket.setsockopt(zmq.SNDHWM, 0)
        socket.setsockopt(zmq.SNDBUF, buf_size)

    if socket_type == zmq.ROUTER and router_handover:
        # Let a new connection take over an identity left behind by a dead one.
        socket.setsockopt(zmq.ROUTER_HANDOVER, 1)

    if identity is not None:
        socket.setsockopt(zmq.IDENTITY, identity)

    if linger is not None:
        socket.setsockopt(zmq.LINGER, linger)

    if socket_type == zmq.XPUB:
        socket.setsockopt(zmq.XPUB_VERBOSE, True)

    if listener is not None:
        socket.setsockopt(zmq.USE_FD, listener.detach())

    # Determine if the path is a TCP socket with an IPv6 address.
    # Enable IPv6 on the zmq socket if so.
    scheme, host, _ = split_zmq_path(path)
    if scheme == "tcp" and is_valid_ipv6_address(host):
        socket.setsockopt(zmq.IPV6, 1)

    if bind:
        socket.bind(path)
    else:
        socket.connect(path)

    return socket

split_zmq_path(path)

Split a zmq path into its parts.

Source code in vllm/utils/network_utils.py
def split_zmq_path(path: str) -> tuple[str, str, str]:
    """Split a zmq path into its parts."""
    parsed = parse_url(path)
    if not parsed.scheme:
        raise ValueError(f"Invalid zmq path: {path}")

    scheme = parsed.scheme
    host = parsed.hostname or ""
    port = "" if parsed.port is None else str(parsed.port)
    if host.startswith("[") and host.endswith("]"):
        host = host[1:-1]  # Remove brackets for IPv6 address

    if scheme == "tcp" and not all((host, port)):
        # The host and port fields are required for tcp
        raise ValueError(f"Invalid zmq path: {path}")

    if scheme != "tcp" and port:
        # port only makes sense with tcp
        raise ValueError(f"Invalid zmq path: {path}")

    return scheme, host, port

zmq_socket_ctx(path, socket_type, bind=None, linger=0, identity=None, router_handover=False)

Context manager for a ZMQ socket.

Source code in vllm/utils/network_utils.py
@contextlib.contextmanager
def zmq_socket_ctx(
    path: str,
    socket_type: Any,
    bind: bool | None = None,
    linger: int = 0,
    identity: bytes | None = None,
    router_handover: bool = False,
) -> Iterator[zmq.Socket]:
    """Context manager for a 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)