Skip to content

vllm.model_executor.model_loader.weight_cache.daemon

Weight cache daemon for fast engine restarts.

One daemon process per GPU holds the post-quantized, TP-sharded weights of its rank in GPU memory and serves CUDA IPC handles to vLLM engines over a Unix domain socket. Restarting engines map the weights via zero-copy IPC instead of reloading from disk.

Launch one daemon per TP rank with a single command:

python -m vllm.model_executor.model_loader.weight_cache.daemon \
    --model /path/to/model --tensor-parallel-size 4

Engines then load from the daemons with:

vllm serve /path/to/model --tensor-parallel-size 4 \
    --load-format ipc_cache

Only tensor and expert parallelism are supported; pipeline and data parallelism are rejected at launch.

For multi-node tensor parallelism, run one launcher per node with a shared rendezvous so the global TP group forms across nodes (CUDA IPC handles are node-local, so each node serves only its local GPUs' shards). Reuse the same --nnodes/--node-rank/--master-addr flags you pass the engine, plus a --weight-cache-master-port distinct from the engine's --master-port:

# node 0 (8 local GPUs)
python -m vllm.model_executor.model_loader.weight_cache.daemon \
    --model /path/to/model --tensor-parallel-size 16 \
    --nnodes 2 --node-rank 0 --master-addr 10.0.0.1 \
    --weight-cache-master-port 29600
# node 1 (8 local GPUs)
python -m vllm.model_executor.model_loader.weight_cache.daemon \
    --model /path/to/model --tensor-parallel-size 16 \
    --nnodes 2 --node-rank 1 --master-addr 10.0.0.1 \
    --weight-cache-master-port 29600

The global TP rank of local GPU i on node r is r * (tp_size // nnodes) + i, matching vLLM's contiguous per-node rank assignment, so each engine worker maps its shard from the daemon on its own node.

Classes:

  • WeightCacheDaemon –

    Per-GPU process that loads one TP shard and serves CUDA IPC handles.

Functions:

  • export_entries –

    Export a model's tensors, preserving tied-parameter aliases.

  • get_daemon_model –

    Load the daemon's model, composed from the configured loader.

WeightCacheDaemon

Per-GPU process that loads one TP shard and serves CUDA IPC handles.

Methods:

Source code in vllm/model_executor/model_loader/weight_cache/daemon.py
class WeightCacheDaemon:
    """Per-GPU process that loads one TP shard and serves CUDA IPC handles."""

    def __init__(
        self,
        vllm_config: VllmConfig,
        tp_rank: int,
        local_rank: int,
        distributed_init_method: str,
        socket_dir: str | None = None,
    ):
        self.vllm_config = vllm_config
        self.tp_rank = tp_rank
        self.local_rank = local_rank
        self.distributed_init_method = distributed_init_method
        self.socket_dir = socket_dir
        self.model: torch.nn.Module | None = None
        # Fingerprint before loading: process_weights_after_loading may
        # mutate hf_config.quantization_config.
        self.cache_config = WeightCacheKey.from_model_config(
            vllm_config.model_config,
            tp_size=vllm_config.parallel_config.tensor_parallel_size,
            tp_rank=tp_rank,
        )

    def load_model(self) -> None:
        tp_size = self.cache_config.tp_size
        torch.accelerator.set_device_index(self.local_rank)
        init_distributed_environment(
            world_size=tp_size,
            rank=self.tp_rank,
            distributed_init_method=self.distributed_init_method,
            local_rank=self.local_rank,
            backend=current_platform.dist_backend,
        )
        with set_current_vllm_config(self.vllm_config):
            ensure_model_parallel_initialized(tp_size, 1)
            self.model = get_daemon_model(self.vllm_config)
        logger.info(
            "Weight cache daemon rank %d loaded model",
            self.tp_rank,
        )

    def serve_forever(self, ready_callback: Callable[[], None] | None = None) -> None:
        """Serve requests until terminated.

        The socket is only bound once the model is fully cached, so clients
        get a connection error (and fall back to disk) until the daemon is
        ready.

        Args:
            ready_callback: Invoked once the socket is bound and listening,
                so the launcher can report overall readiness.

        """
        socket_path = self._socket_path
        ensure_private_socket_dir(
            os.path.dirname(socket_path), strict_perms=self.socket_dir is None
        )
        # Hold an exclusive per-GPU lock for the daemon's lifetime so a second
        # daemon cannot remove this daemon's live socket and hijack the path.
        lock_fd = self._acquire_gpu_lock(socket_path)
        if os.path.exists(socket_path):
            os.unlink(socket_path)
        server = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM)
        server.bind(socket_path)
        os.chmod(socket_path, 0o600)
        server.listen()
        logger.info(
            "Weight cache daemon rank %d serving on %s", self.tp_rank, socket_path
        )
        if ready_callback is not None:
            ready_callback()
        try:
            while True:
                conn, _ = server.accept()
                with conn:
                    try:
                        verify_peer_is_owner(conn)
                        self._handle_connection(conn)
                    except (ConnectionError, EOFError):
                        logger.warning("Client disconnected mid-request")
                    except Exception as e:
                        # Report the error back instead of just closing the
                        # socket, but don't let it take the daemon down.
                        logger.exception(
                            "Error handling weight cache client; continuing"
                        )
                        with contextlib.suppress(OSError):
                            send_msg(conn, {"status": "error", "message": str(e)})
        finally:
            server.close()
            if os.path.exists(socket_path):
                os.unlink(socket_path)
            os.close(lock_fd)

    def _acquire_gpu_lock(self, socket_path: str) -> int:
        """Take an exclusive lock guarding this GPU's socket path.

        The lock is advisory and released automatically when the daemon exits
        (or crashes), so a stale socket is only ever removed by whoever owns
        the lock. A running daemon holding it makes a second daemon fail fast
        instead of clobbering the live socket.
        """
        lock_fd = os.open(f"{socket_path}.lock", os.O_CREAT | os.O_RDWR, 0o600)
        try:
            fcntl.flock(lock_fd, fcntl.LOCK_EX | fcntl.LOCK_NB)
        except OSError as e:
            os.close(lock_fd)
            raise WeightCacheUnavailableError(
                f"Another weight cache daemon already owns {socket_path}"
            ) from e
        return lock_fd

    @property
    def _socket_path(self) -> str:
        return get_socket_path(get_current_device_uuid(), self.socket_dir)

    def _handle_connection(self, conn: socket.socket) -> None:
        request = recv_msg(conn)
        cmd = request.get("cmd")
        if cmd == "get_state":
            self._handle_get_state(conn, request)
        elif cmd == "release":
            self._handle_release(conn)
        else:
            send_msg(conn, {"status": "error", "message": f"Unknown command {cmd!r}"})

    def _handle_get_state(self, conn: socket.socket, request: dict) -> None:
        client_config = request.get("cache_config")
        if not isinstance(client_config, WeightCacheKey):
            send_msg(conn, {"status": "error", "message": "Missing cache_config"})
            return
        mismatched = self.cache_config.mismatched_fields(client_config)
        if mismatched:
            logger.warning("WeightCacheKey mismatch on fields: %s", mismatched)
            send_msg(conn, {"status": "mismatch", "fields": mismatched})
            return
        if self.model is None:
            send_msg(conn, {"status": "error", "message": "Weights were released"})
            return
        gpu_uuid = get_current_device_uuid()
        entries, aliases = export_entries(self.model)
        send_msg(
            conn,
            {
                "status": "ok",
                "entries": entries,
                "aliases": aliases,
                "gpu_uuid": gpu_uuid,
            },
        )
        logger.info_once(
            "Weight cache daemon rank %d sent %d tensors (+%d aliases) to engine",
            self.tp_rank,
            len(entries),
            len(aliases),
        )

    def _handle_release(self, conn: socket.socket) -> None:
        self.model = None
        torch.accelerator.empty_cache()
        logger.info("Weight cache daemon rank %d released cached weights", self.tp_rank)
        send_msg(conn, {"status": "ok"})

_acquire_gpu_lock(socket_path)

Take an exclusive lock guarding this GPU's socket path.

The lock is advisory and released automatically when the daemon exits (or crashes), so a stale socket is only ever removed by whoever owns the lock. A running daemon holding it makes a second daemon fail fast instead of clobbering the live socket.

Source code in vllm/model_executor/model_loader/weight_cache/daemon.py
def _acquire_gpu_lock(self, socket_path: str) -> int:
    """Take an exclusive lock guarding this GPU's socket path.

    The lock is advisory and released automatically when the daemon exits
    (or crashes), so a stale socket is only ever removed by whoever owns
    the lock. A running daemon holding it makes a second daemon fail fast
    instead of clobbering the live socket.
    """
    lock_fd = os.open(f"{socket_path}.lock", os.O_CREAT | os.O_RDWR, 0o600)
    try:
        fcntl.flock(lock_fd, fcntl.LOCK_EX | fcntl.LOCK_NB)
    except OSError as e:
        os.close(lock_fd)
        raise WeightCacheUnavailableError(
            f"Another weight cache daemon already owns {socket_path}"
        ) from e
    return lock_fd

serve_forever(ready_callback=None)

Serve requests until terminated.

The socket is only bound once the model is fully cached, so clients get a connection error (and fall back to disk) until the daemon is ready.

Parameters:

  • ready_callback

    (Callable[[], None] | None, default: None ) –

    Invoked once the socket is bound and listening, so the launcher can report overall readiness.

Source code in vllm/model_executor/model_loader/weight_cache/daemon.py
def serve_forever(self, ready_callback: Callable[[], None] | None = None) -> None:
    """Serve requests until terminated.

    The socket is only bound once the model is fully cached, so clients
    get a connection error (and fall back to disk) until the daemon is
    ready.

    Args:
        ready_callback: Invoked once the socket is bound and listening,
            so the launcher can report overall readiness.

    """
    socket_path = self._socket_path
    ensure_private_socket_dir(
        os.path.dirname(socket_path), strict_perms=self.socket_dir is None
    )
    # Hold an exclusive per-GPU lock for the daemon's lifetime so a second
    # daemon cannot remove this daemon's live socket and hijack the path.
    lock_fd = self._acquire_gpu_lock(socket_path)
    if os.path.exists(socket_path):
        os.unlink(socket_path)
    server = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM)
    server.bind(socket_path)
    os.chmod(socket_path, 0o600)
    server.listen()
    logger.info(
        "Weight cache daemon rank %d serving on %s", self.tp_rank, socket_path
    )
    if ready_callback is not None:
        ready_callback()
    try:
        while True:
            conn, _ = server.accept()
            with conn:
                try:
                    verify_peer_is_owner(conn)
                    self._handle_connection(conn)
                except (ConnectionError, EOFError):
                    logger.warning("Client disconnected mid-request")
                except Exception as e:
                    # Report the error back instead of just closing the
                    # socket, but don't let it take the daemon down.
                    logger.exception(
                        "Error handling weight cache client; continuing"
                    )
                    with contextlib.suppress(OSError):
                        send_msg(conn, {"status": "error", "message": str(e)})
    finally:
        server.close()
        if os.path.exists(socket_path):
            os.unlink(socket_path)
        os.close(lock_fd)

_reject_unsupported_parallelism(parallel_config)

Reject parallelism modes other than tensor/expert parallelism.

Source code in vllm/model_executor/model_loader/weight_cache/daemon.py
def _reject_unsupported_parallelism(parallel_config: ParallelConfig) -> None:
    """Reject parallelism modes other than tensor/expert parallelism."""
    unsupported = {
        "pipeline parallelism": parallel_config.pipeline_parallel_size > 1,
        "data parallelism": parallel_config.data_parallel_size > 1,
    }
    for name, enabled in unsupported.items():
        if enabled:
            raise ValueError(
                f"The weight cache daemon only supports tensor and expert "
                f"parallelism; {name} is not supported"
            )

export_entries(model)

Export a model's tensors, preserving tied-parameter aliases.

named_parameters/named_buffers are iterated with remove_duplicate=False so tied weights (e.g. lm_head.weight sharing storage with embed_tokens.weight) are not silently dropped. Each unique tensor is exported once per call; every additional name that refers to the same tensor object is recorded in the returned alias map so the client can re-establish the shared identity instead of allocating uninitialized memory for it.

CUDA reduction arguments must be exported separately for each consumer so that PyTorch registers a reference for each IPC mapping's lifetime.

Returns:

Source code in vllm/model_executor/model_loader/weight_cache/daemon.py
def export_entries(
    model: torch.nn.Module,
) -> tuple[dict[str, TensorEntry], dict[str, str]]:
    """Export a model's tensors, preserving tied-parameter aliases.

    ``named_parameters``/``named_buffers`` are iterated with
    ``remove_duplicate=False`` so tied weights (e.g. ``lm_head.weight`` sharing
    storage with ``embed_tokens.weight``) are not silently dropped. Each unique
    tensor is exported once per call; every additional name that refers to the same
    tensor object is recorded in the returned alias map so the client can
    re-establish the shared identity instead of allocating uninitialized
    memory for it.

    CUDA reduction arguments must be exported separately for each consumer so
    that PyTorch registers a reference for each IPC mapping's lifetime.

    Returns:
        A ``(entries, aliases)`` pair where ``entries`` maps a canonical name to
        its ``TensorEntry`` and ``aliases`` maps each duplicate name to its
        canonical name.

    """
    entries: dict[str, TensorEntry] = {}
    aliases: dict[str, str] = {}
    canonical_by_id: dict[int, str] = {}

    def _add(name: str, tensor: torch.Tensor, kind: str) -> None:
        canonical = canonical_by_id.get(id(tensor))
        if canonical is not None:
            aliases[name] = canonical
            return
        canonical_by_id[id(tensor)] = name
        entries[name] = TensorEntry.from_tensor(tensor, kind)

    for name, param in model.named_parameters(remove_duplicate=False):
        _add(name, param, "param")
    # named_buffers includes non-persistent buffers (e.g. rotary embedding
    # caches) that state_dict would miss.
    for name, buffer in model.named_buffers(remove_duplicate=False):
        if name in entries or name in aliases:
            continue
        _add(name, buffer, "buffer")
    return entries, aliases

get_daemon_model(vllm_config)

Load the daemon's model, composed from the configured loader.

Runs the quantization check after model creation but before the slow weight load, so an unsupported method fails fast. Online quantization always fails the check, so load_model's finalize step for it is unnecessary here.

Source code in vllm/model_executor/model_loader/weight_cache/daemon.py
def get_daemon_model(vllm_config: VllmConfig) -> torch.nn.Module:
    """Load the daemon's model, composed from the configured loader.

    Runs the quantization check after model creation but before the slow
    weight load, so an unsupported method fails fast. Online quantization
    always fails the check, so load_model's finalize step for it is
    unnecessary here.
    """
    model_config = vllm_config.model_config
    load_config = vllm_config.load_config
    loader = get_model_loader(load_config)
    device_config = vllm_config.device_config
    target_device = torch.device(
        device_config.device if load_config.device is None else load_config.device
    )
    with set_default_torch_dtype(model_config.dtype):
        with target_device:
            model = loader.create_model(vllm_config, model_config)
        check_ipc_quant_support(model)
        loader.load_weights(model, model_config)
        process_weights_after_loading(model, model_config, target_device)
    return model.eval()