Skip to content
Merged
Show file tree
Hide file tree
Changes from 3 commits
Commits
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
8 changes: 4 additions & 4 deletions e2e_test/router/test_pd_mmlu.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,8 @@
workers for improved throughput and resource utilization.

Backends:
- "pd_http": HTTP mode (SGLang only - vLLM does not support HTTP)
- "pd_http": HTTP mode (SGLang parallel bootstrap dispatch; vLLM sequential
prefill-then-decode with kv_transfer_params relay)
- "pd_grpc": gRPC mode (both SGLang and vLLM)

Requirements:
Expand All @@ -16,7 +17,7 @@
# SGLang (runs both HTTP and gRPC)
pytest e2e_test/router/test_pd_mmlu.py -v

# vLLM (runs gRPC only, HTTP skipped)
# vLLM (runs both HTTP and gRPC)
E2E_RUNTIME=vllm pytest e2e_test/router/test_pd_mmlu.py -v
"""

Expand All @@ -31,11 +32,10 @@
logger = logging.getLogger(__name__)


@pytest.mark.engine("sglang")
@pytest.mark.engine("sglang", "vllm")
@pytest.mark.gpu(2)
@pytest.mark.model("meta-llama/Llama-3.1-8B-Instruct")
@pytest.mark.e2e
@pytest.mark.skip_for_runtime("vllm", reason="vLLM does not support HTTP mode")
@pytest.mark.parametrize("setup_backend", ["pd_http"], indirect=True)
class TestPDMMLUHttp:
"""MMLU evaluation tests using PD disaggregation (HTTP mode)."""
Expand Down
208 changes: 208 additions & 0 deletions model_gateway/src/routers/common/kv_transfer.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,208 @@
//! vLLM PD KV-transfer connector handling, shared by both transports.
//!
//! vLLM disaggregation is sequential: the prefill leg is tagged with
//! connector params, the engine returns (or the router mints) the handoff
//! params, and the decode leg carries them. The gRPC pipeline and the HTTP
//! PD router both implement that flow; the connector vocabulary lives here
//! so the two stay in lockstep.

use crate::{
observability::metrics::metrics_labels,
worker::{Worker, DEFAULT_BOOTSTRAP_PORT, MOONCAKE_CONNECTOR, NIXL_CONNECTOR},
};

/// KV-transfer params tagged onto the NIXL prefill leg so the engine pins its
/// KV blocks and returns the handoff params for the decode worker.
pub(crate) const NIXL_PREFILL_KV_PARAMS: &str =
r#"{"do_remote_decode":true,"do_remote_prefill":false}"#;

/// PD KV-transfer behavior derived from prefill worker metadata.
#[derive(Debug, Clone, PartialEq)]
pub(crate) enum KvConnectorMode {
/// MooncakeConnector: mint a transfer_id, tag both legs, synthesize decode
/// params from worker metadata; legacy host/port injection when the
/// servicer predates kv_engine_id reporting (or DP runs without a pinned rank).
Mooncake {
host: String,
port: u32,
engine_id: Option<String>,
},
/// NixlConnector: tag prefill with do_remote_decode, relay returned params to decode.
Nixl,
/// Unknown/absent connector: relay returned params opportunistically.
Passthrough,
}

impl KvConnectorMode {
pub(crate) fn metrics_label(&self) -> &'static str {
match self {
Self::Mooncake { .. } => metrics_labels::KV_CONNECTOR_MOONCAKE,
Self::Nixl => metrics_labels::KV_CONNECTOR_NIXL,
Self::Passthrough => metrics_labels::KV_CONNECTOR_PASSTHROUGH,
}
}
}

pub(crate) fn kv_connector_mode(
kv_connector: Option<&str>,
bootstrap_host: &str,
bootstrap_port: Option<u16>,
kv_engine_id: Option<&str>,
) -> KvConnectorMode {
match kv_connector {
Some(MOONCAKE_CONNECTOR) => KvConnectorMode::Mooncake {
host: bootstrap_host.to_string(),
port: u32::from(bootstrap_port.unwrap_or(DEFAULT_BOOTSTRAP_PORT)),
// Empty means unknown (forces the legacy fallback)
engine_id: kv_engine_id.filter(|s| !s.is_empty()).map(str::to_string),
},
Some(NIXL_CONNECTOR) => KvConnectorMode::Nixl,
_ => KvConnectorMode::Passthrough,
}
}

/// Connector id of the engine core serving the prefill leg. With DP the cores
/// suffix the configured id as `{base}_dp{rank}`, so minting needs a pinned
/// rank; unpinned DP>1 yields None (no mint — decode recomputes locally).
pub(crate) fn effective_kv_engine_id(
base: Option<&str>,
dp_size: Option<usize>,
dp_rank: Option<usize>,
) -> Option<String> {
let base = base.filter(|s| !s.is_empty())?;
if dp_size.unwrap_or(1) > 1 {
dp_rank.map(|rank| format!("{base}_dp{rank}"))
} else {
Some(base.to_string())
}
}

/// The connector mode for a PD pair, read off the prefill worker's metadata.
/// Discovered dp_size matters even without `--dp-aware` expansion: a DP>1
/// engine behind an unexpanded worker must not be minted for.
pub(crate) fn connector_mode_for_worker(worker: &dyn Worker) -> KvConnectorMode {
let meta = worker.metadata();
let dp_size = worker
.dp_size()
.or_else(|| meta.spec.labels.get("dp_size").and_then(|s| s.parse().ok()));
let engine_id =
effective_kv_engine_id(meta.spec.kv_engine_id.as_deref(), dp_size, worker.dp_rank());
Comment thread
coderabbitai[bot] marked this conversation as resolved.
Outdated
kv_connector_mode(
meta.spec.kv_connector.as_deref(),
&meta.spec.bootstrap_host,
meta.spec.bootstrap_port,
engine_id.as_deref(),
)
}

/// Prefill-leg params for Mooncake: the engine pins blocks under the minted id.
pub(crate) fn mooncake_prefill_params(transfer_id: &str) -> String {
serde_json::json!({
"do_remote_decode": true,
"do_remote_prefill": false,
"transfer_id": transfer_id,
})
.to_string()
}

/// Decode-leg params for Mooncake, synthesized from prefill worker metadata
/// (the engine returns nothing to relay; the connector is push-based).
pub(crate) fn mooncake_decode_params(
transfer_id: &str,
engine_id: &str,
host: &str,
port: u32,
) -> String {
serde_json::json!({
"do_remote_decode": false,
"do_remote_prefill": true,
"transfer_id": transfer_id,
"remote_engine_id": engine_id,
"remote_bootstrap_addr": format!("http://{host}:{port}"),
})
.to_string()
}

#[cfg(test)]
mod tests {
use super::*;

#[test]
fn kv_connector_mode_mooncake_uses_bootstrap_metadata() {
let mode = kv_connector_mode(
Some(MOONCAKE_CONNECTOR),
"prefill-host",
Some(9090),
Some("engine-1"),
);
assert_eq!(
mode,
KvConnectorMode::Mooncake {
host: "prefill-host".to_string(),
port: 9090,
engine_id: Some("engine-1".to_string()),
}
);
}

#[test]
fn kv_connector_mode_mooncake_defaults_port_and_tolerates_missing_engine_id() {
let mode = kv_connector_mode(Some(MOONCAKE_CONNECTOR), "prefill-host", None, None);
assert_eq!(
mode,
KvConnectorMode::Mooncake {
host: "prefill-host".to_string(),
port: u32::from(DEFAULT_BOOTSTRAP_PORT),
engine_id: None,
}
);
}

#[test]
fn kv_connector_mode_mooncake_empty_engine_id_means_legacy() {
let mode = kv_connector_mode(Some(MOONCAKE_CONNECTOR), "host", Some(9090), Some(""));
assert_eq!(
mode,
KvConnectorMode::Mooncake {
host: "host".to_string(),
port: 9090,
engine_id: None,
}
);
}

#[test]
fn kv_connector_mode_nixl() {
assert_eq!(
kv_connector_mode(Some(NIXL_CONNECTOR), "ignored", Some(9090), None),
KvConnectorMode::Nixl
);
}

#[test]
fn kv_connector_mode_unknown_or_missing_is_passthrough() {
assert_eq!(
kv_connector_mode(Some("LMCacheConnector"), "host", None, None),
KvConnectorMode::Passthrough
);
assert_eq!(
kv_connector_mode(None, "host", None, None),
KvConnectorMode::Passthrough
);
}

#[test]
fn effective_engine_id_requires_pinned_rank_under_dp() {
assert_eq!(
effective_kv_engine_id(Some("eng"), Some(2), Some(1)),
Some("eng_dp1".to_string())
);
assert_eq!(effective_kv_engine_id(Some("eng"), Some(2), None), None);
assert_eq!(
effective_kv_engine_id(Some("eng"), None, None),
Some("eng".to_string())
);
assert_eq!(effective_kv_engine_id(Some(""), None, None), None);
assert_eq!(effective_kv_engine_id(None, Some(2), Some(0)), None);
}
}
1 change: 1 addition & 0 deletions model_gateway/src/routers/common/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@
//! responses to clients and parsing upstream SSE byte streams

pub mod header_utils;
pub(crate) mod kv_transfer;
pub mod mcp_utils;
pub mod openai_bridge;
pub mod overload;
Expand Down
Loading
Loading