Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
24 commits
Select commit Hold shift + click to select a range
ed49ee0
[None][feat] Mooncake store part 1: pool, CLI, and V2 scheduler preem…
brb-nv Sep 15, 2026
dbbdfff
[None][chore] Address comments from CodeRabbit
brb-nv Sep 19, 2026
502cda7
[None][fix] Address comments from Thor
brb-nv Sep 19, 2026
a6b5124
[None][chore] Address follow-up comments from CodeRabbit
brb-nv Sep 20, 2026
6edc039
[None][fix] Address comments from zhaoyuanh-nvidia
brb-nv Sep 20, 2026
1cd820b
[None][fix] Offer --telemetry/--no-telemetry on the mooncake commands
brb-nv Sep 22, 2026
2189f02
[None][fix] Hand a signalled mooncake command to the telemetry boundary
brb-nv Sep 22, 2026
e23d334
[None][fix] Do not call a pending KV transfer a V2 scheduler deadlock
brb-nv Sep 22, 2026
e11bdc7
[None][fix] Reserve preempted KV pages for the request that freed them
brb-nv Sep 22, 2026
ad5a85e
[None][chore] Tighten comments on the review follow-ups
brb-nv Sep 22, 2026
4929e6e
[None][chore] Shorten the mooncake CMake package cleanup
brb-nv Sep 22, 2026
83793db
[None][chore] Defer the Mooncake LlmArgs surface to the connector MR
brb-nv Sep 23, 2026
93bd853
Address comments from nv-xtf
brb-nv Sep 23, 2026
fdfb2f3
Address comments from Coderabbit
brb-nv Sep 23, 2026
36c55d0
Logging for scheduler deadlock
brb-nv Sep 23, 2026
3f81b48
[None][chore] Share the mooncake client pin between the wheel and con…
brb-nv Sep 23, 2026
0284774
[None][fix] Drop the Mooncake import check from the container build
brb-nv Sep 24, 2026
ba8aca4
[None][fix] Install the Mooncake client in the release image
brb-nv Sep 29, 2026
929f0b8
[None][infra] Update CI image tags for the Mooncake client install
brb-nv Sep 29, 2026
27e6b0b
[None][chore] Keep the Mooncake client out of the container images
brb-nv Sep 29, 2026
a5142ca
[None][fix] Withdraw the donor's ready file along with the segment
brb-nv Sep 29, 2026
45bb57d
[None][fix] Match the install the donor's import error actually names
brb-nv Sep 29, 2026
4df8110
[None][chore] Install the Mooncake client with the default dependency…
brb-nv Sep 30, 2026
6b23512
[None][fix] Let each Mooncake client discover its own RDMA device
brb-nv Sep 30, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions requirements.txt
Original file line number Diff line number Diff line change
Expand Up @@ -116,3 +116,4 @@ cache-dit>=1.3.5
librosa
msgpack
uvloop>=0.19.0
mooncake-transfer-engine-cuda13==0.3.13
Original file line number Diff line number Diff line change
@@ -0,0 +1,69 @@
# SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved.
# SPDX-License-Identifier: Apache-2.0
#
# Licensed under the Apache License, Version 2.0 (the "License");
# you may not use this file except in compliance with the License.
# You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.
"""A Mooncake distributed store to back a KV cache connector.

The store is a shared CPU memory pool addressed by content, so a prefix computed
by one engine can be replayed by another, which regular block reuse cannot do
because it never leaves the instance that computed it.

This is a different component from the Mooncake transfer engine that the C++
cache transceiver uses for disaggregated prefill/decode handoff: that moves KV
point to point between two known peers, while this one publishes pages into a
pool addressed by content. The two compose, so a context server can write pages
here and still hand off over NIXL.

`master.py` brings a pool up: it resolves or launches the `mooncake_master`
and renders the client config the workers read. Capacity comes only from
processes that open a store handle, which in a disaggregated deployment is the
context servers alone, so `donor.py` lends a node's memory to the pool without
giving it a connector. `trtllm-serve mooncake_master` and `mooncake_donor`
expose both. Both need the Mooncake Python bindings, which `tensorrt-llm` pulls
in as `mooncake-transfer-engine-cuda13`.

`keys.py` and `staging.py` hold what the store side shares with the connector
that moves pages in and out of the pool: how a block of tokens becomes a store
key, and how pages reach the fabric on hosts without GPUDirect RDMA. The
connector itself, and the `LlmArgs` surface that selects it, land with the KV
cache manager V2 support it depends on.
"""

from .config import MooncakeStoreConnectorConfig, StoreRole, parse_size
from .donor import DEFAULT_DONOR_LOCAL_BUFFER_SIZE, donate_segment
from .master import (
PoolSpec,
local_address,
master_timeout,
provision_pool,
resolve_master_address,
running_master,
wait_for_master,
write_client_config,
)

__all__ = [
"DEFAULT_DONOR_LOCAL_BUFFER_SIZE",
"MooncakeStoreConnectorConfig",
"PoolSpec",
"StoreRole",
"donate_segment",
"local_address",
"master_timeout",
"parse_size",
"provision_pool",
"resolve_master_address",
"running_master",
"wait_for_master",
"write_client_config",
]
282 changes: 282 additions & 0 deletions tensorrt_llm/_torch/pyexecutor/connectors/mooncake_store/config.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,282 @@
# SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved.
# SPDX-License-Identifier: Apache-2.0
#
# Licensed under the Apache License, Version 2.0 (the "License");
# you may not use this file except in compliance with the License.
# You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.
"""Configuration for the Mooncake store KV cache connector.

Topology settings are read from the JSON file named by `MOONCAKE_CONFIG_PATH`,
the same file and environment variable the vLLM Mooncake store connector uses,
so one deployment can point both engines at the same pool.

`KvCacheConnectorConfig` carries no free-form dictionary, so the two settings
that are TensorRT-LLM's rather than Mooncake's, the read/write role and the key
prefix, are also taken from the environment.
"""

import json
import os
import re
from dataclasses import dataclass
from enum import Enum
from typing import Any, Optional

__all__ = [
"CLIENT_CONFIG_NAME",
"CONFIG_PATH_ENV",
"MooncakeStoreConnectorConfig",
"ROLE_ENV",
"RUN_DIR_ENV",
"STAGE_THROUGH_HOST_ENV",
"StoreRole",
"provisioned_config_path",
]

CONFIG_PATH_ENV = "MOONCAKE_CONFIG_PATH"
#: Where a server keeps the client config it renders and the master's log. Set
#: it to keep them after shutdown; otherwise they live in a temporary directory.
RUN_DIR_ENV = "TRTLLM_MOONCAKE_RUN_DIR"
#: Name the rendered client config takes in the run directory.
CLIENT_CONFIG_NAME = "mooncake.json"
ROLE_ENV = "TRTLLM_MOONCAKE_STORE_ROLE"
CACHE_PREFIX_ENV = "TRTLLM_MOONCAKE_STORE_PREFIX"
MODEL_KEY_ENV = "TRTLLM_MOONCAKE_STORE_MODEL_KEY"
STAGE_THROUGH_HOST_ENV = "TRTLLM_MOONCAKE_STORE_STAGE_THROUGH_HOST"

DEFAULT_GLOBAL_SEGMENT_SIZE = 3355443200
DEFAULT_LOCAL_BUFFER_SIZE = 1073741824
DEFAULT_CACHE_PREFIX = "trtllm"
DEFAULT_STAGING_BUFFER_SIZE = 536870912
#: Mooncake's own peer-to-peer handshake, which keeps a separate metadata
#: process out of the deployment. Nothing else is a sensible fallback: an empty
#: connstring is not one of the forms `store.setup` accepts, so a config that
#: leaves the field out means this rather than meaning no metadata service.
DEFAULT_METADATA_SERVER = "P2PHANDSHAKE"

_TRUE = {"1", "true", "yes", "on"}
_FALSE = {"0", "false", "no", "off"}

_SIZE_UNITS = {
"": 1,
"b": 1,
"k": 1000,
"kb": 1000,
"m": 1000**2,
"mb": 1000**2,
"g": 1000**3,
"gb": 1000**3,
"t": 1000**4,
"tb": 1000**4,
"kib": 1024,
"mib": 1024**2,
"gib": 1024**3,
"tib": 1024**4,
}
_SIZE_RE = re.compile(r"^\s*([0-9]+(?:\.[0-9]+)?)\s*([a-zA-Z]*)\s*$")


class StoreRole(Enum):
"""Which directions of traffic this engine is allowed to drive.

A disaggregated deployment typically runs context servers as `both` and
leaves generation servers unconfigured: generated tokens are rarely a reused
prefix, so writing them costs bandwidth for no hit rate.
"""

PRODUCER = "producer"
CONSUMER = "consumer"
BOTH = "both"

@property
def loads(self) -> bool:
"""Whether this role reads previously stored KV back onto the GPU."""
return self is not StoreRole.PRODUCER

@property
def saves(self) -> bool:
"""Whether this role writes newly computed KV into the store."""
return self is not StoreRole.CONSUMER


def parse_size(value: Any) -> int:
"""Accept either a byte count or a suffixed string such as `"4GiB"`."""
if isinstance(value, bool):
raise ValueError(f"expected a size, got {value!r}")
if isinstance(value, int):
return value
if isinstance(value, float):
return int(value)
match = _SIZE_RE.match(str(value))
if match is None:
raise ValueError(f"cannot parse size {value!r}")
magnitude, unit = match.groups()
scale = _SIZE_UNITS.get(unit.lower())
if scale is None:
raise ValueError(f"unknown size unit {unit!r} in {value!r}")
return int(float(magnitude) * scale)


def provisioned_config_path() -> Optional[str]:
"""The client config a server on this node rendered, if there is one.

`provision_pool` writes one and exports `MOONCAKE_CONFIG_PATH`, which the
ranks the LLM constructor spawns inherit. Ranks an external launcher
started, one task per rank, were already running by then and never see it,
so they read the config back from the run directory instead.

Only possible when the deployment named that directory, since it otherwise
defaults to a per-process temporary one that no other rank could read.
"""
run_dir = os.getenv(RUN_DIR_ENV)
if not run_dir:
return None
path = os.path.join(run_dir, CLIENT_CONFIG_NAME)
return path if os.path.exists(path) else None


@dataclass(frozen=True)
class MooncakeStoreConnectorConfig:
"""Everything needed to open a store handle and name keys in it."""

master_server_address: str
metadata_server: str = DEFAULT_METADATA_SERVER
protocol: str = "rdma"
device_name: str = ""
global_segment_size: int = DEFAULT_GLOBAL_SEGMENT_SIZE
local_buffer_size: int = DEFAULT_LOCAL_BUFFER_SIZE
local_hostname: Optional[str] = None
tenant_id: Optional[str] = None
role: StoreRole = StoreRole.BOTH
cache_prefix: str = DEFAULT_CACHE_PREFIX
#: Identity the keys are namespaced by. Two engines share cache only when
#: they agree on this, and two that disagree about what it names read each
#: other's pages, so it is required rather than defaulted. See
#: :meth:`resolve_model_key`.
model_key: Optional[str] = None
#: How many page keys go into one store call. Bounds the size of a single
#: RPC without bounding how much a request may transfer.
transfer_batch_size: int = 64
#: Pass pages through a pinned host buffer instead of registering the KV
#: pools with Mooncake. Costs a copy each way, but works without GPUDirect
#: RDMA, which registering device memory requires.
stage_through_host: bool = False
#: Ceiling on the pinned allocation per direction when staging. Slots are
#: sized from the layout's largest page, so this caps how many pages may be
#: in flight rather than how large one may be.
staging_buffer_bytes: int = DEFAULT_STAGING_BUFFER_SIZE

def __post_init__(self) -> None:
"""Reject settings that would fail later, inside a transfer."""
if not self.master_server_address:
raise ValueError("master_server_address is required")
if self.local_buffer_size <= 0:
raise ValueError("local_buffer_size must be > 0")
if self.global_segment_size < 0:
raise ValueError("global_segment_size must be >= 0")
if self.transfer_batch_size <= 0:
raise ValueError("transfer_batch_size must be > 0")
if self.stage_through_host and self.staging_buffer_bytes <= 0:
raise ValueError("staging_buffer_bytes must be > 0 when staging is on")

@staticmethod
def from_file(path: str) -> "MooncakeStoreConnectorConfig":
"""Read the topology from a vLLM-compatible Mooncake JSON config."""
with open(path) as handle:
raw = json.load(handle)
return MooncakeStoreConnectorConfig(
master_server_address=raw.get("master_server_address", ""),
metadata_server=raw.get("metadata_server") or DEFAULT_METADATA_SERVER,
protocol=raw.get("protocol", "rdma"),
device_name=raw.get("device_name", ""),
global_segment_size=parse_size(
raw.get("global_segment_size", DEFAULT_GLOBAL_SEGMENT_SIZE)
),
local_buffer_size=parse_size(raw.get("local_buffer_size", DEFAULT_LOCAL_BUFFER_SIZE)),
local_hostname=raw.get("local_hostname") or None,
tenant_id=raw.get("tenant_id") or None,
role=StoreRole(str(raw.get("role", StoreRole.BOTH.value)).strip().lower()),
cache_prefix=str(raw.get("cache_prefix", DEFAULT_CACHE_PREFIX)),
model_key=raw.get("model_key") or None,
transfer_batch_size=int(raw.get("transfer_batch_size", 64)),
stage_through_host=bool(raw.get("stage_through_host", False)),
staging_buffer_bytes=parse_size(
raw.get("staging_buffer_bytes", DEFAULT_STAGING_BUFFER_SIZE)
),
)

@staticmethod
def from_env() -> "MooncakeStoreConnectorConfig":
"""Load the JSON config, then apply the TensorRT-LLM env overrides."""
path = os.getenv(CONFIG_PATH_ENV) or provisioned_config_path()
if not path:
raise ValueError(
f"The mooncake-store connector needs {CONFIG_PATH_ENV} set to a "
"Mooncake JSON config (metadata_server, master_server_address, "
"protocol, device_name, global_segment_size, local_buffer_size), "
"or kv_connector_config.mooncake_store set so the server renders "
f"one, into ${RUN_DIR_ENV} if this rank was started by the "
"launcher rather than spawned by the server."
)
config = MooncakeStoreConnectorConfig.from_file(path)
return config.with_env_overrides()

def with_env_overrides(self) -> "MooncakeStoreConnectorConfig":
"""Apply `TRTLLM_MOONCAKE_STORE_*` on top of the file's settings."""
import dataclasses

updates: dict[str, Any] = {}
role = os.getenv(ROLE_ENV)
if role:
try:
updates["role"] = StoreRole(role.strip().lower())
except ValueError as exc:
known = ", ".join(member.value for member in StoreRole)
raise ValueError(f"{ROLE_ENV}={role!r} is not one of: {known}") from exc
prefix = os.getenv(CACHE_PREFIX_ENV)
if prefix:
updates["cache_prefix"] = prefix
model_key = os.getenv(MODEL_KEY_ENV)
if model_key:
updates["model_key"] = model_key
staging = os.getenv(STAGE_THROUGH_HOST_ENV)
if staging:
normalized = staging.strip().lower()
if normalized in _TRUE:
updates["stage_through_host"] = True
elif normalized in _FALSE:
updates["stage_through_host"] = False
else:
known = ", ".join(sorted(_TRUE | _FALSE))
raise ValueError(
f"{STAGE_THROUGH_HOST_ENV}={staging!r} is not a boolean; use one of: {known}"
)
return dataclasses.replace(self, **updates) if updates else self

def resolve_model_key(self, model: Any) -> str:
"""The model identity to namespace keys by.

Deliberately has no default. Deriving one from the model path would
make two checkpoints that happen to share a directory name, such as
`org-a/model` and `org-b/model` or two revisions mounted alike, agree
on a namespace while disagreeing on what the pages mean, and each would
read the other's KV as its own.
"""
if self.model_key:
return self.model_key
raise ValueError(
f"The mooncake-store connector needs a model key to namespace its "
f"pool keys by, and there is no safe default: set "
f"kv_connector_config.mooncake_store.model_key, the model_key field "
f"of the Mooncake JSON config, or ${MODEL_KEY_ENV}. Give it a value "
f"that separates this checkpoint from any other an engine sharing "
f"the pool might load, rather than one derived from {model!r}."
)
Loading
Loading