Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
42 commits
Select commit Hold shift + click to select a range
1205202
test: Added Ray Serve SGLang PD disaggregation support
limarkdcunha Jul 2, 2026
2e85e8e
updates as per feedback
limarkdcunha Jul 3, 2026
0f0c2ec
cloud mirror model path fix
limarkdcunha Jul 3, 2026
0cba02a
bootstrap port offset fix
limarkdcunha Jul 3, 2026
0c1a7dd
prefill token clamp fix
limarkdcunha Jul 3, 2026
b22e455
DP gang scheduling explicit failure
limarkdcunha Jul 3, 2026
db26012
Keep vLLM request-shaping policy out of the neutral connector base
limarkdcunha Jul 8, 2026
ba9e702
validation step against mixed engine setup
limarkdcunha Jul 8, 2026
e7e9798
Stale engine config fix
limarkdcunha Jul 8, 2026
d4d5552
benchmark scripts for testing
limarkdcunha Aug 1, 2026
822cb17
feedback fix
limarkdcunha Aug 1, 2026
448cdde
fix in r_ud
limarkdcunha Aug 1, 2026
7f73224
request id fix
limarkdcunha Aug 1, 2026
9f6194e
fix as per feedback
limarkdcunha Aug 1, 2026
27d6e5e
Changes as per feedback
limarkdcunha Aug 1, 2026
a6646b2
Multi GPU collision fix
limarkdcunha Aug 1, 2026
a95a1f8
setdefault role correction skip fix
limarkdcunha Aug 1, 2026
cf895d2
Added role validator
limarkdcunha Aug 1, 2026
9cc692b
optimization attempt: 1 _concurrent_decode
limarkdcunha Aug 1, 2026
999501f
abort on prefill failures as before
limarkdcunha Aug 1, 2026
37060fa
Stagger native PD server launches to avoid PID limit spike
limarkdcunha Aug 2, 2026
8cc96f9
Added build_asgi_app for sglang
limarkdcunha Aug 2, 2026
1b7bd5d
added tracer logs
limarkdcunha Aug 8, 2026
f182e54
updated in benchmark code
limarkdcunha Aug 9, 2026
ee6c8fc
tracer code fix
limarkdcunha Aug 9, 2026
feee4e7
setting stream_batching_interval_ms to 0
limarkdcunha Aug 9, 2026
86a6633
testing skip slot reservation for thesis testing
limarkdcunha Aug 9, 2026
ea41f52
fix in router as per feedback
limarkdcunha Aug 9, 2026
b794ee4
forward thread caps to pd replicas always
limarkdcunha Aug 9, 2026
765a056
Undoing experimental changes
limarkdcunha Aug 9, 2026
ce0eaf5
reverted legit change
limarkdcunha Aug 9, 2026
f3b030e
added trace to measure one specific flow
limarkdcunha Aug 10, 2026
d8012a5
resolved as per feedback
limarkdcunha Aug 10, 2026
e467352
log queue wait and probe time as aggregates
limarkdcunha Aug 15, 2026
72f085e
Added more logging steps
limarkdcunha Sep 11, 2026
c1f01a3
Some more logging
limarkdcunha Sep 11, 2026
ff404ac
fix to last logging
limarkdcunha Sep 11, 2026
b201b25
Buffer PD trace logging off hot path; fix sweep harness
limarkdcunha Sep 26, 2026
0048e22
Write PD trace rows direct to file; drain on start and exit
limarkdcunha Sep 26, 2026
f4207a4
Report server token counts and support ignore_eos in PD benchmarks
limarkdcunha Sep 26, 2026
94f1194
Count output chunks on both arms, not server token counts
limarkdcunha Sep 26, 2026
b9d350a
Clean up benchmark code and comments
limarkdcunha Sep 28, 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
67 changes: 67 additions & 0 deletions python/ray/llm/_internal/serve/core/configs/engine_adapter.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,67 @@
"""Per-engine adapters.

Each adapter declares how an engine names its KV connector and which
engine-config class builds its ``EngineConfig``. This replaces ``if vLLM /
elif SGLang`` branching in ``llm_config.py`` so adding an engine is additive:
add an adapter and register it here.
"""

import abc
from typing import TYPE_CHECKING, Any, Dict, Optional, Type

if TYPE_CHECKING:
pass


class EngineAdapter(abc.ABC):
@abc.abstractmethod
def connector_name(self, engine_kwargs: Dict[str, Any]) -> Optional[str]:
"""The KV-connector registry name implied by engine_kwargs, or None."""
...

@abc.abstractmethod
def engine_config_cls(self) -> Type:
"""The engine-config class whose ``from_llm_config`` builds the config."""
...


class VLLMAdapter(EngineAdapter):
def connector_name(self, engine_kwargs: Dict[str, Any]) -> Optional[str]:
cfg = engine_kwargs.get("kv_transfer_config")
if not cfg:
return None
kv_connector = cfg.get("kv_connector")
if not kv_connector:
# Fail fast: a kv_transfer_config with no kv_connector is a
# misconfiguration, not a "no connector" case.
raise ValueError("Connector type is not specified.")
return kv_connector

def engine_config_cls(self) -> Type:
from ray.llm._internal.serve.engines.vllm.vllm_models import VLLMEngineConfig

return VLLMEngineConfig


class SGLangAdapter(EngineAdapter):
def connector_name(self, engine_kwargs: Dict[str, Any]) -> Optional[str]:
return (
"SGLang" if engine_kwargs.get("disaggregation_transfer_backend") else None
)

def engine_config_cls(self) -> Type:
from ray.llm._internal.serve.engines.sglang.sglang_engine import (
SGLangEngineConfig,
)

return SGLangEngineConfig


_ADAPTERS = {"vLLM": VLLMAdapter, "SGLang": SGLangAdapter}


def get_engine_adapter(llm_engine: str) -> EngineAdapter:
try:
return _ADAPTERS[llm_engine]()
except KeyError:
raise ValueError(f"Unsupported engine: {llm_engine}")
63 changes: 38 additions & 25 deletions python/ray/llm/_internal/serve/core/configs/llm_config.py
Original file line number Diff line number Diff line change
Expand Up @@ -45,14 +45,14 @@
TPUConfig,
infer_hardware_kind_from_bundles,
)
from ray.llm._internal.serve.engines.vllm.kv_transfer.factory import (
from ray.llm._internal.serve.engines.common.kv_transfer.factory import (
KVConnectorBackendFactory,
)
from ray.llm._internal.serve.observability.logging import get_logger
from ray.serve._private.config import DeploymentConfig, handle_num_replicas_auto

if TYPE_CHECKING:
from ray.llm._internal.serve.engines.vllm.kv_transfer.base import (
from ray.llm._internal.serve.engines.common.kv_transfer.base import (
BaseConnectorBackend,
)

Expand Down Expand Up @@ -86,6 +86,7 @@ class LLMEngine(str, Enum):
"""Enum that represents an LLMEngine."""

vLLM = "vLLM"
SGLang = "SGLang"


class LoraConfig(BaseModelExtended):
Expand Down Expand Up @@ -146,7 +147,7 @@ class ModelLoadingConfig(BaseModelExtended):
)


EngineConfigType = Union[None, "VLLMEngineConfig"] # noqa: F821
EngineConfigType = Union[None, "VLLMEngineConfig", "SGLangEngineConfig"] # noqa: F821


class LLMConfig(BaseModelExtended):
Expand Down Expand Up @@ -589,19 +590,23 @@ def get_engine_config(self) -> EngineConfigType:
if self._engine_config:
return self._engine_config

if self.llm_engine == LLMEngine.vLLM:
from ray.llm._internal.serve.engines.vllm.vllm_models import (
VLLMEngineConfig,
)

self._engine_config = VLLMEngineConfig.from_llm_config(self)
else:
# Note (genesu): This should never happen because we validate the engine
# in the config.
raise ValueError(f"Unsupported engine: {self.llm_engine}")
from ray.llm._internal.serve.core.configs.engine_adapter import (
get_engine_adapter,
)

adapter = get_engine_adapter(self.llm_engine)
self._engine_config = adapter.engine_config_cls().from_llm_config(self)
return self._engine_config

@property
def num_devices(self) -> int:
"""Total devices per replica (tensor-parallel × pipeline-parallel).

Neutral accessor used by connector backends for port spacing, so the
connector base never reaches into an engine-specific config type.
"""
return self.get_engine_config().num_devices

def update_engine_kwargs(self, **kwargs: Any) -> None:
"""Update the engine_kwargs and the engine_config engine_kwargs.

Expand All @@ -611,27 +616,35 @@ def update_engine_kwargs(self, **kwargs: Any) -> None:
self.engine_kwargs.update(kwargs)
# engine_config may be created before engine starts, this makes sure
# the engine_config is updated with the latest engine_kwargs.
if self._engine_config:
# Not every engine config mirrors engine_kwargs
# (the minimal SGLangEngineConfig does not — SGLangServer reads llm_config.engine_kwargs directly), so
# only sync configs that carry it.
if self._engine_config is not None and hasattr(
self._engine_config, "engine_kwargs"
):
self._engine_config.engine_kwargs.update(kwargs)

def setup_engine_backend(self):
self._setup_kv_connector_backend()

def _setup_kv_connector_backend(self):
"""Private method to setup kv connector depending on the local deployment state"""
Comment thread
cursor[bot] marked this conversation as resolved.
# 1. validate that the backend is one of the backends supported (Nixl or LMCache)
kv_transfer_config = self.engine_kwargs.get("kv_transfer_config")
if not kv_transfer_config:
return
"""Create and set up the KV connector backend for this engine, if any.

kv_connector = kv_transfer_config.get("kv_connector")
if not kv_connector:
raise ValueError("Connector type is not specified.")
The connector name is resolved through the per-engine adapter so this
path carries no engine-specific knowledge (vLLM reads ``kv_transfer_config.kv_connector``;
SGLang reads ``disaggregation_transfer_backend``).
"""

# 2. Setup the backend using factory
kv_connector_backend = KVConnectorBackendFactory.create_backend(
kv_connector, self
from ray.llm._internal.serve.core.configs.engine_adapter import (
get_engine_adapter,
)

adapter = get_engine_adapter(self.llm_engine)
name = adapter.connector_name(self.engine_kwargs)
if not name:
return
Comment thread
cursor[bot] marked this conversation as resolved.

kv_connector_backend = KVConnectorBackendFactory.create_backend(name, self)
kv_connector_backend.setup()
# 3. Stash the instance so the P/D orchestrator can reach the connector's
# coordination protocol (request shaping, peer binding, handoff
Expand Down
Empty file.
Empty file.
176 changes: 176 additions & 0 deletions python/ray/llm/_internal/serve/engines/common/kv_transfer/base.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,176 @@
import abc
import random
import string
from typing import TYPE_CHECKING, Any, Dict, Optional, Union

from ray import serve

if TYPE_CHECKING:
from ray.llm._internal.serve.core.configs.llm_config import LLMConfig
from ray.llm._internal.serve.core.configs.openai_api_models import (
ChatCompletionRequest,
CompletionRequest,
)

# The two OpenAI request models the P/D orchestrator shapes. Defined under
# TYPE_CHECKING (and used as a string annotation) to avoid an import cycle
# between this module and the config/openai-models modules.
RequestType = Union[ChatCompletionRequest, CompletionRequest]


def clamp_request_to_single_token(request: "RequestType") -> None:
"""Clamp a prefill request to a single, non-streaming token (in place)."""
request.max_tokens = 1
if hasattr(request, "max_completion_tokens"):
request.max_completion_tokens = 1
request.stream = False
if hasattr(request, "stream_options"):
request.stream_options = None


class BaseConnectorBackend(abc.ABC):
# ---- P/D coordination protocol ----
#
# These class attributes and methods let the P/D orchestrator
# (``PDOrchestratorMixin``) delegate request shaping, peer addressing, and
# handoff discipline to the connector. They are connector-agnostic: a
# connector picks a quadrant of (``requires_peer_binding``,
# ``concurrent_handoff``) and implements ``prepare_prefill_request`` /
# ``prepare_decode_request`` accordingly.
#
# ``requires_peer_binding``:
# * False -> the orchestrator dispatches prefill via the standard handle
# path; the peer (if any) is resolved post-hoc from the prefill response.
# * True -> the orchestrator selects the prefill replica first
# (``choose_replica``) and passes its ``replica_metadata`` to the backend
# as ``peer`` (pre-dispatch addressing).
#
# ``concurrent_handoff``:
# * False -> prefill runs to its first chunk before local decode starts
# (sequential handoff).
# * True -> prefill dispatch and local decode run concurrently.
#
# The two flags are independent; the known combos:
# * (False, False) — e.g. NixlConnector / LMCacheConnectorV1: decode learns
# everything it needs (remote engine id / block ids) from the prefill
# response, so it must wait for that response (sequential).
# * (True, True) — e.g. MoRIIO WRITE mode / SGLang: the peer address is
# bound up front and prefill *pushes* KV to decode, so decode needs
# nothing from prefill's response and can start immediately (concurrent).
# Concurrent handoff is only possible because peer binding happens up
# front.
# * (True, False) — e.g. MoRIIO READ mode: the peer is bound up front,
# but decode *pulls* KV using block ids returned in prefill's response,
# so it still waits for prefill to finish (sequential).
requires_peer_binding: bool = False
concurrent_handoff: bool = False

def __init__(self, llm_config: "LLMConfig"):
"""Base class for connector backends.

Args:
llm_config: The llm configuration for this engine
"""
self.llm_config = llm_config

def _get_unique_suffix(self, len: int = 6) -> str:
"""Generates unique alphanumeric suffix.

Args:
len: Length of the suffix to generate.
Returns:
A unique alphanumeric suffix string of specified length.
"""
return "".join(random.choices(string.ascii_letters + string.digits, k=len))

def _compute_port_offset(self) -> int:
"""Compute a deterministic port offset for this replica.

Uses data_parallel_rank if DP case, otherwise falls back to
the replica rank assigned by Ray Serve (TP/PP case).

For TP/PP cases, multiply by num_devices (tp × pp) to reserve
sufficient port space, since each worker needs a unique port.
Each TP worker adds its tp_rank (0, 1, ..., tp_size-1) to the
base port at bind time, and PP stages also need separate ports.

``num_devices`` is read from the neutral ``LLMConfig`` accessor, so this
base never reaches into an engine-specific config type.

Returns:
Non-negative integer offset to add to a base port.
"""
# Prefer explicit DP rank when available
dp_rank = self.llm_config.engine_kwargs.get("data_parallel_rank")
if isinstance(dp_rank, int) and dp_rank >= 0:
# vLLM already accounts for TP spacing in DP offset calculation
# (data_parallel_rank × tp_size), don't multiply here
return dp_rank

# NOTE (jeffreywang): A missing replica context must fail loudly, not
# silently return a 0 offset that collides colocated replicas on the
# same side-channel port. get_replica_context() raises RayServeException
# outside a replica.
rc = serve.get_replica_context()
num_devices = self.llm_config.num_devices
return rc.rank.rank * num_devices

@abc.abstractmethod
def prepare_prefill_request(
self, *, request: "RequestType", peer: Optional[Dict[str, Any]]
) -> "RequestType":
"""Shape the request sent to the remote prefill engine.

Args:
request: The incoming chat/completion request.
peer: The selected prefill replica's ``replica_metadata`` dict when
the connector opted into pre-dispatch peer binding
(``requires_peer_binding=True``), else None.

Returns:
A new request object to dispatch to the prefill engine.
"""
...

@abc.abstractmethod
def prepare_decode_request(
self,
*,
request: "RequestType",
peer: Optional[Dict[str, Any]],
prefill_response: Optional[Any],
) -> "RequestType":
"""Shape the request run on the local decode engine.

Args:
request: The incoming chat/completion request.
peer: The selected prefill replica's ``replica_metadata`` dict when
the connector opted into pre-dispatch peer binding, else None.
prefill_response: The captured prefill response chunk whose
``kv_transfer_params`` may be forwarded, or None when no chunk is
captured before decode starts (concurrent-handoff mode).

Returns:
A new request object to run on the local decode engine.
"""
...

def setup(self) -> None:
"""Setup the connector backend.

This method is called to setup the connector backend.
"""
pass

def replica_metadata(self) -> Dict[str, Any]:
"""Static per-replica coordination data published to the orchestrator.

Surfaced via the replica-metadata hook on ``ReplicaSelection`` so that a
connector opting into ``requires_peer_binding`` can address the selected
prefill peer. The default backend publishes nothing; connectors that need
to advertise an address (e.g. MoRIIO's zmq endpoint) override this.

Returns:
A JSON-serializable dict of per-replica metadata (empty by default).
"""
return {}
Loading
Loading