Skip to content
Merged
Show file tree
Hide file tree
Changes from all 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
12 changes: 12 additions & 0 deletions lib/bindings/python/rust/llm/kv.rs
Original file line number Diff line number Diff line change
Expand Up @@ -118,6 +118,14 @@ struct KvIndexerCli {
#[arg(long)]
watch_recover_port: Option<u16>,

/// Data-parallel ranks per discovered engine pod (vLLM
/// `--data-parallel-size`). Rank r publishes KV events on
/// `--watch-zmq-port + r` and serves recovery on `--watch-recover-port + r`;
/// every rank is subscribed under the pod's instance. Leaving this at 1 on
/// a DP engine silently indexes only rank 0's cache.
#[arg(long, default_value_t = 1)]
watch_dp_size: u32,

/// Model name whose engine pods to discover, e.g. "openai/gpt-oss-120b".
/// Unless --watch-label overrides it, the pod watch uses the selector
/// `di/model_name=<sanitized name>` (the stable label the backend stamps
Expand Down Expand Up @@ -161,11 +169,15 @@ where
let block_size = cli.block_size.ok_or_else(|| {
anyhow::anyhow!("--block-size is required when --watch-namespace is set")
})?;
if cli.watch_dp_size == 0 {
anyhow::bail!("--watch-dp-size must be at least 1");
}
Some(KubeDiscoveryConfig {
namespace,
label_selector,
zmq_port: cli.watch_zmq_port,
recover_port: cli.watch_recover_port,
dp_size: cli.watch_dp_size,
model_name: cli.watch_model_name.unwrap_or_else(|| cli.model_name.clone()),
tenant_id: cli.tenant_id.clone(),
block_size,
Expand Down
11 changes: 11 additions & 0 deletions lib/kv-router/src/standalone_indexer/docs.md
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,17 @@ If no `recover_endpoint` is configured, gaps are logged and the dropped batches
are lost. Implementation lives in `listener.rs` (`recover_gap`,
`apply_recover_response`, `apply_recovered_events`).

## Data-parallel engines (`--watch-dp-size`)

A vLLM engine with `--data-parallel-size N` runs N ranks per pod, each with its
own KV cache and its own event stream: rank `r` publishes on `zmq_port + r` and
serves `/kv_recover` on `kv_recover_port + r`. Pod discovery registers one
listener per rank under the pod's instance (`--watch-dp-size N`, default 1),
with `dp_rank = r` and the per-rank endpoints, so `/workers` shows N
`listeners` per pod. With the default on a DP engine only rank 0's cache is
indexed and every prefill scheduled on another rank is invisible to `/query`
while the engine still hits it (`pod_watcher.rs`, `rank_endpoints`).

## Audit logging (`--enable-logging`)

Pass `--enable-logging` to `python -m dynamo.indexer` to turn on verbose audit
Expand Down
11 changes: 8 additions & 3 deletions lib/kv-router/src/standalone_indexer/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -124,12 +124,17 @@ pub struct KubeDiscoveryConfig {
/// from the model name (`di/model_name=<sanitized>`, spans every
/// engine_hash) or supplied raw, e.g. `engine_hash=d4b7a85131172ca6`.
pub label_selector: String,
/// ZMQ KV-event port the engines publish on (e.g. 5557).
/// ZMQ KV-event port the engines publish on (e.g. 5557). Data-parallel
/// rank `r` publishes on `zmq_port + r`.
pub zmq_port: u16,
/// Optional HTTP port serving `GET /kv_recover` on the engines, used for
/// per-worker gap recovery. `http://<pod-ip>:<recover_port>` is the base
/// URL the indexer queries on a detected gap.
/// per-worker gap recovery. `http://<pod-ip>:<recover_port + r>` is the
/// base URL the indexer queries for rank `r` on a detected gap.
pub recover_port: Option<u16>,
/// Data-parallel ranks per engine pod. Each rank has its own KV cache and
/// its own event stream, so every rank is registered as a listener of the
/// same instance. Engines without data parallelism run one rank.
pub dp_size: u32,
/// Model name discovered pods are registered under.
pub model_name: String,
/// Tenant id discovered pods are registered under.
Expand Down
140 changes: 120 additions & 20 deletions lib/kv-router/src/standalone_indexer/pod_watcher.rs
Original file line number Diff line number Diff line change
Expand Up @@ -119,6 +119,7 @@ pub fn spawn_pod_watcher(
namespace = %config.namespace,
label_selector = %config.label_selector,
zmq_port = config.zmq_port,
dp_size = config.dp_size,
model_name = %config.model_name,
block_size = config.block_size,
"Starting Kubernetes pod watcher"
Expand Down Expand Up @@ -259,30 +260,35 @@ async fn reconcile(
let instance_id = instance_id_for(name);
let ip = &pod.ip;
let tenant_id = pod.tenant_id(&config.tenant_id);
let endpoint = format!("tcp://{ip}:{}", config.zmq_port);
let recover_endpoint = config
.recover_port
.map(|port| format!("http://{ip}:{port}"));

match registry
.register(
instance_id,
endpoint,
0, // dp_rank: single-rank engines
config.model_name.clone(),
tenant_id.clone(),
config.block_size,
recover_endpoint,
Some(name.clone()),
)
.await
{
Ok(()) => {
let mut failed: Option<(u32, anyhow::Error)> = None;
for rank in rank_endpoints(ip, config) {
if let Err(error) = registry
.register(
instance_id,
rank.endpoint,
rank.dp_rank,
config.model_name.clone(),
tenant_id.clone(),
config.block_size,
rank.recover_endpoint,
Some(name.clone()),
)
.await
{
failed = Some((rank.dp_rank, error));
break;
}
}

match failed {
None => {
tracing::info!(
pod = %name,
ip = %ip,
instance_id,
tenant_id = %tenant_id,
dp_size = config.dp_size,
"Subscribed to engine pod"
);
subscribed.insert(
Expand All @@ -294,19 +300,57 @@ async fn reconcile(
},
);
}
Err(error) => {
// Not recorded, so it is retried on the next reconcile.
Some((dp_rank, error)) => {
// Not recorded, so it is retried on the next reconcile. Drop
// any ranks that did register: register() rejects a rank that
// is already present, so a partial instance would never
// converge on retry.
tracing::warn!(
pod = %name,
ip = %ip,
dp_rank,
error = %error,
"Failed to register engine pod; will retry"
);
if let Err(error) = registry
.deregister(instance_id, &config.model_name, &tenant_id)
.await
{
tracing::debug!(
pod = %name,
error = %error,
"Cleanup deregister was a no-op"
);
}
}
}
}
}

/// One listener to register for a pod: its data-parallel rank and the
/// per-rank event and recovery endpoints.
#[derive(Debug, PartialEq)]
struct RankEndpoint {
dp_rank: u32,
endpoint: String,
recover_endpoint: Option<String>,
}

/// Endpoints for every data-parallel rank of a pod. vLLM offsets both the ZMQ
/// KV-event port and the `/kv_recover` port by the rank, so rank `r` of a pod
/// at `ip` publishes on `zmq_port + r` and recovers on `recover_port + r`.
fn rank_endpoints(ip: &str, config: &KubeDiscoveryConfig) -> Vec<RankEndpoint> {
(0..config.dp_size.max(1))
.map(|dp_rank| RankEndpoint {
dp_rank,
endpoint: format!("tcp://{ip}:{}", u32::from(config.zmq_port) + dp_rank),
recover_endpoint: config
.recover_port
.map(|port| format!("http://{ip}:{}", u32::from(port) + dp_rank)),
})
.collect()
}

/// Derive a stable `WorkerId` from the pod name so the same pod always maps to
/// the same registry entry across MODIFIED events.
fn instance_id_for(pod_name: &str) -> WorkerId {
Expand Down Expand Up @@ -456,4 +500,60 @@ mod tests {
assert_eq!(instance_id_for("engine-abc"), instance_id_for("engine-abc"));
assert_ne!(instance_id_for("engine-abc"), instance_id_for("engine-xyz"));
}

fn discovery(dp_size: u32, recover_port: Option<u16>) -> KubeDiscoveryConfig {
KubeDiscoveryConfig {
namespace: "deepinfra".to_string(),
label_selector: "di/model_name=m".to_string(),
zmq_port: 5557,
recover_port,
dp_size,
model_name: "m".to_string(),
tenant_id: "default".to_string(),
block_size: 128,
}
}

#[test]
fn single_rank_engine_registers_rank_zero_on_base_ports() {
assert_eq!(
rank_endpoints("10.0.0.1", &discovery(1, Some(5559))),
vec![RankEndpoint {
dp_rank: 0,
endpoint: "tcp://10.0.0.1:5557".to_string(),
recover_endpoint: Some("http://10.0.0.1:5559".to_string()),
}]
);
}

#[test]
fn data_parallel_ranks_offset_event_and_recovery_ports() {
assert_eq!(
rank_endpoints("10.0.0.1", &discovery(2, Some(5559))),
vec![
RankEndpoint {
dp_rank: 0,
endpoint: "tcp://10.0.0.1:5557".to_string(),
recover_endpoint: Some("http://10.0.0.1:5559".to_string()),
},
RankEndpoint {
dp_rank: 1,
endpoint: "tcp://10.0.0.1:5558".to_string(),
recover_endpoint: Some("http://10.0.0.1:5560".to_string()),
},
]
);
}

#[test]
fn no_recover_port_means_no_recovery_endpoint_on_any_rank() {
let ranks = rank_endpoints("10.0.0.1", &discovery(2, None));
assert_eq!(ranks.len(), 2);
assert!(ranks.iter().all(|r| r.recover_endpoint.is_none()));
}

#[test]
fn zero_dp_size_still_registers_rank_zero() {
assert_eq!(rank_endpoints("10.0.0.1", &discovery(0, None)).len(), 1);
}
}
Loading