From 90871f1c79e4e47648900e38e55b9145dda38bbb Mon Sep 17 00:00:00 2001 From: Ishan Dhanani Date: Mon, 6 Jul 2026 12:48:07 +0000 Subject: [PATCH 01/12] feat(kv-router): consume Mooncake KV events Signed-off-by: Ishan Dhanani --- components/src/dynamo/sglang/register.py | 1 + .../backends/sglang/hicache.md | 30 +- lib/llm/src/discovery/model_manager.rs | 6 +- lib/llm/src/kv_router/shared_cache.rs | 340 ++++++++++-------- 4 files changed, 212 insertions(+), 165 deletions(-) diff --git a/components/src/dynamo/sglang/register.py b/components/src/dynamo/sglang/register.py index 794604f3135f..f067a792929f 100644 --- a/components/src/dynamo/sglang/register.py +++ b/components/src/dynamo/sglang/register.py @@ -314,6 +314,7 @@ def _get_mooncake_runtime_data(server_args: ServerArgs) -> Optional[dict[str, An "master_metrics_port": int( getattr(mooncake_config, "master_metrics_port", 9003) ), + "kv_events_endpoint": os.getenv("DYN_MOONCAKE_KV_EVENTS_ENDPOINT") or None, } diff --git a/docs/fern/pages/developer-guide/knowledge-base/modular-components/backends/sglang/hicache.md b/docs/fern/pages/developer-guide/knowledge-base/modular-components/backends/sglang/hicache.md index 35da2be0a3fb..3dbd2484b713 100644 --- a/docs/fern/pages/developer-guide/knowledge-base/modular-components/backends/sglang/hicache.md +++ b/docs/fern/pages/developer-guide/knowledge-base/modular-components/backends/sglang/hicache.md @@ -17,7 +17,7 @@ SGLang HiCache extends RadixAttention with a multi-tier KV cache that transparen What Dynamo adds on top of HiCache: - **Tier-aware routing.** The KV router tracks which cache tier each block lives on (GPU / Host / External) and uses that when scoring candidate workers — not just device overlap. -- **Shared-pool awareness.** When an external backend such as Mooncake is configured, the router queries the shared pool in parallel with its own indexer so it can discount prefill cost for blocks any worker can fetch, not just blocks the candidate holds locally. +- **Shared-pool awareness.** When an external backend such as Mooncake is configured, the router tracks Mooncake object events so it can discount prefill cost for blocks any worker can fetch, not just blocks the candidate holds locally. If you are running a single worker with HiCache and no shared pool, no Dynamo-side configuration is required — the worker reports KV events to the router as usual. @@ -84,15 +84,15 @@ flowchart LR Worker -- "KV events (store/remove + medium)" --> Router Worker -- "writes pages" --> Mooncake - Router -- "batch_query on each request" --> Mooncake + Mooncake -- "object events (stored/removed)" --> Router ``` -On every request the router runs two lookups in parallel: +The router maintains two indexes: - Its own radix tree, built from worker KV events (per-tier). -- A batch query to the Mooncake master for blocks reachable from the shared pool. +- A set of Mooncake object keys, built from the Mooncake master's event stream. -If the shared-pool query fails, the router falls back to indexer-only scoring and logs a warning. The request still succeeds. +Request routing checks both indexes locally. If the event stream starts after Mooncake already contains objects, the shared-pool index starts empty and learns about subsequent events. A sequence gap clears the shared-pool index to avoid stale hits. ### Scoring @@ -143,16 +143,18 @@ Earlier SGLang versions do not emit `medium=CPU_PINNED` for Host-tier residency, You also need: - Dynamo router started with `--shared-cache-type hicache` (see [Configuration](#configuration)). -- A Mooncake master reachable from the Dynamo frontend host. Worker-side Mooncake config (master address, page size, TP/PP layout, split-head layout) is published automatically via each worker's registration metadata when the worker is started with `--hicache-storage-backend mooncake`. +- A Mooncake master with KV events enabled. This requires the event publisher introduced by [kvcache-ai/Mooncake#2214](https://github.com/kvcache-ai/Mooncake/pull/2214) until that change is available in a Mooncake release. +- A Mooncake KV event endpoint reachable from the Dynamo frontend host. Worker-side Mooncake config (event endpoint, page size, TP/PP layout, and split-head layout) is published automatically through each worker's registration metadata. ## Setup > [!WARNING] > SGLang 0.5.11 through 0.5.12.x bundle Mooncake 0.3.10.post2, which is affected by an upstream `MemcpyWorkerPool` crash when both `--enable-metrics` and `--disable-piecewise-cuda-graph` are set on the worker. Upgrade to SGLang 0.5.13 or later, which bundles Mooncake 0.3.11.post1 with [the upstream fix](https://github.com/kvcache-ai/Mooncake/pull/2001), or omit both flags. The setup below omits them. -**SGLang worker** — HiCache with Mooncake storage: +**SGLang worker** — HiCache with Mooncake storage and the event endpoint advertised to Dynamo: ```bash +DYN_MOONCAKE_KV_EVENTS_ENDPOINT=tcp://mooncake-master.internal:5557 \ python -m dynamo.sglang \ --model-path Qwen/Qwen3-0.6B \ --page-size 64 \ @@ -180,12 +182,12 @@ python -m dynamo.frontend \ | Flag | Env var | Default | Description | | --------------------------- | ----------------------------- | ------- | ------------------------------------------------------------------------------------------------------------------------------------------------------------------ | -| `--shared-cache-type` | `DYN_SHARED_CACHE_TYPE` | `none` | `none` disables shared-pool lookups; `hicache` enables Mooncake queries. | -| `--shared-cache-multiplier` | `DYN_SHARED_CACHE_MULTIPLIER` | `0.5` | Discount factor for shared-pool hits. `0.0` queries but ignores them; `0.5` treats a shared hit as half a device hit; `1.0` treats shared and device hits equally. | +| `--shared-cache-type` | `DYN_SHARED_CACHE_TYPE` | `none` | `none` disables shared-pool tracking; `hicache` enables the Mooncake event-backed index. | +| `--shared-cache-multiplier` | `DYN_SHARED_CACHE_MULTIPLIER` | `0.5` | Discount factor for shared-pool hits. `0.0` ignores them; `0.5` treats a shared hit as half a device hit; `1.0` treats shared and device hits equally. | Per-request overrides are available via `RouterConfigOverride.shared_cache_multiplier` for A/B experimentation without restarting the router. -No extra flags are required on the worker. When `--hicache-storage-backend mooncake` is set, Dynamo publishes the required metadata (page size, TP/PP layout, master address) via the worker's `ModelRuntimeConfig.engine_specific` blob under the key `sglang_hicache_mooncake`. +Set `DYN_MOONCAKE_KV_EVENTS_ENDPOINT` on the worker to the Mooncake PUB endpoint, such as `tcp://mooncake-master.internal:5557`. The endpoint must be reachable from the frontend. When `--hicache-storage-backend mooncake` is set, Dynamo publishes the endpoint and required layout metadata through the worker's `ModelRuntimeConfig.engine_specific` blob under the key `sglang_hicache_mooncake`. ## Verification @@ -214,10 +216,10 @@ curl -s localhost:8000/metrics | grep shared_cache | Symptom | Likely cause | Fix | | -------------------------------------------------------- | ---------------------------------------------------------------------------- | --------------------------------------------------------------------------------------------- | -| `shared_cache_hit_rate` is always 0 | Mooncake master unreachable from the router host | Check network path; the router logs `Shared cache query failed` when it can't reach Mooncake. | -| Events only ever carry `medium=GPU` | SGLang older than 0.5.11 or a custom build missing [PR #22894](https://github.com/sgl-project/sglang/pull/22894) | Upgrade to SGLang 0.5.11 or later. | -| Workers registered but router never queries shared cache | `--shared-cache-type` left at default `none` | Set `--shared-cache-type hicache` on the frontend. | -| Queries issued but winning worker rarely changes | `--shared-cache-multiplier 0.0` | Raise the multiplier — typical starting range is `0.3`–`0.7`. | +| `shared_cache_hit_rate` is always 0 | Event endpoint missing, unreachable, or connected after objects were stored | Set `DYN_MOONCAKE_KV_EVENTS_ENDPOINT` on the worker and check the frontend subscriber log. | +| Events only ever carry `medium=GPU` | SGLang older than 0.5.11 or a custom build missing [PR #22894](https://github.com/sgl-project/sglang/pull/22894) | Upgrade to SGLang 0.5.11 or later. | +| Workers registered but router never tracks shared cache | `--shared-cache-type` left at default `none` | Set `--shared-cache-type hicache` on the frontend. | +| Shared hits do not affect worker selection | `--shared-cache-multiplier 0.0` | Raise the multiplier — typical starting range is `0.3`–`0.7`. | | Page-size mismatch warnings | Router `--page-size` doesn't match worker `--page-size` | They must agree; the router hashes pages using the worker's page size. | | Router logs "no workers have HiCache enabled" | No worker published `sglang_hicache_mooncake` metadata | Confirm workers started with `--hicache-storage-backend mooncake`. | diff --git a/lib/llm/src/discovery/model_manager.rs b/lib/llm/src/discovery/model_manager.rs index 5ebb92d8e14b..5236e1801efe 100644 --- a/lib/llm/src/discovery/model_manager.rs +++ b/lib/llm/src/discovery/model_manager.rs @@ -1127,9 +1127,9 @@ impl ModelManager { worker_component = worker_component_name, "Using HiCache shared KV cache" ); - Some(Box::new(HicacheSharedKvCache::new( - workers_with_configs.clone(), - ))) + let cache = HicacheSharedKvCache::new(workers_with_configs.clone()); + cache.start_subscriber(endpoint.component()); + Some(Box::new(cache)) } }; diff --git a/lib/llm/src/kv_router/shared_cache.rs b/lib/llm/src/kv_router/shared_cache.rs index 482605c5cfa8..428fdb13370a 100644 --- a/lib/llm/src/kv_router/shared_cache.rs +++ b/lib/llm/src/kv_router/shared_cache.rs @@ -3,34 +3,40 @@ //! HiCache shared KV cache client for SGLang + Mooncake. //! -//! Instead of querying a worker endpoint over the request plane, this client: +//! This client: //! 1. Reads Mooncake HiCache metadata published by SGLang workers in runtime config. //! 2. Recomputes the logical HiCache page hashes from request tokens using the //! same token -> page-hash logic as SGLang. //! 3. Expands those logical page hashes into the concrete Mooncake object keys //! SGLang uses for the configured TP/PP/MLA layout. -//! 4. Queries the Mooncake master HTTP service directly via `/batch_query_keys`. +//! 4. Tracks those object keys from the Mooncake master's KV event stream. -use std::collections::HashMap; -use std::time::Duration; +use std::sync::{ + Arc, + atomic::{AtomicU64, Ordering}, +}; use async_trait::async_trait; -use reqwest::Url; +use dashmap::DashSet; +use futures::StreamExt; use serde::{Deserialize, Serialize}; use sha2::{Digest, Sha256}; - -const MOONCAKE_HTTP_TIMEOUT: Duration = Duration::from_secs(2); +use tokio_util::sync::CancellationToken; use dynamo_kv_router::{ SharedKvCache, indexer::KvRouterError, protocols::{SharedCacheHits, WorkerId}, }; +use dynamo_runtime::{component::Component, traits::DistributedRuntimeProvider}; -use crate::{discovery::RuntimeConfigWatch, local_model::runtime_config::ModelRuntimeConfig}; +use crate::{ + discovery::RuntimeConfigWatch, + local_model::runtime_config::ModelRuntimeConfig, + utils::zmq::{connect_sub_socket, multipart_message}, +}; const SGLANG_HICACHE_MOONCAKE_RUNTIME_KEY: &str = "sglang_hicache_mooncake"; -const MOONCAKE_BATCH_QUERY_KEYS_CHUNK_SIZE: usize = 128; #[derive(Debug, Clone, Deserialize, Serialize, PartialEq, Eq)] struct SglangHicacheMooncakeConfig { @@ -45,45 +51,49 @@ struct SglangHicacheMooncakeConfig { extra_backend_tag: Option, master_server_address: Option, master_metrics_port: u16, + kv_events_endpoint: Option, } #[derive(Debug, Deserialize)] -struct MooncakeBatchQueryKeysResponse { - success: bool, +struct MooncakeObjectEvent { + event_type: String, #[serde(default)] - data: HashMap, -} - -#[derive(Debug, Deserialize, Default)] -struct MooncakeBatchQueryKeyResult { + object_key: Option, #[serde(default)] - ok: bool, + tenant_id: String, } +type MooncakeEventBatch = (i64, Vec, u32); + #[derive(Debug, Clone, Copy)] enum QueryToken { Single(u32), Bigram(u32, u32), } -/// Shared KV cache client that queries the Mooncake master HTTP service for -/// SGLang HiCache (L3) state. +/// Event-driven shared KV cache index for SGLang HiCache (L3) state. +#[derive(Clone)] pub struct HicacheSharedKvCache { runtime_configs: RuntimeConfigWatch, - http_client: reqwest::Client, + present_keys: Arc>, + last_sequence: Arc, } impl HicacheSharedKvCache { pub fn new(runtime_configs: RuntimeConfigWatch) -> Self { Self { runtime_configs, - http_client: reqwest::Client::builder() - .timeout(MOONCAKE_HTTP_TIMEOUT) - .build() - .expect("failed to build reqwest client"), + present_keys: Arc::new(DashSet::new()), + last_sequence: Arc::new(AtomicU64::new(0)), } } + pub fn start_subscriber(&self, component: &Component) { + let cache = self.clone(); + let cancellation_token = component.drt().child_token(); + tokio::spawn(async move { cache.run_subscriber(cancellation_token).await }); + } + fn resolve_mooncake_config(&self) -> Option { let workers = self.runtime_configs.borrow(); let mut configs = Vec::new(); @@ -107,57 +117,104 @@ impl HicacheSharedKvCache { Some(first.clone()) } - async fn fetch_key_presence( - &self, - endpoint: &Url, - actual_keys: &[String], - ) -> Result, KvRouterError> { - let mut key_presence = HashMap::with_capacity(actual_keys.len()); - - for chunk in actual_keys.chunks(MOONCAKE_BATCH_QUERY_KEYS_CHUNK_SIZE) { - let joined_keys = chunk.join(","); - - let mut url = endpoint.clone(); - // Mooncake expects a raw comma-separated `keys=` list. If commas are - // percent-encoded (`%2C`), Mooncake treats the entire value as one key. - url.set_query(Some(&format!("keys={joined_keys}"))); - - let response = self.http_client.get(url.clone()).send().await.map_err(|e| { - tracing::warn!(error = %e, url = %url, "Mooncake batch_query_keys request failed"); - KvRouterError::IndexerOffline - })?; - - let status = response.status(); - if !status.is_success() { - tracing::warn!( - status = %status, - url = %url, - "Mooncake batch_query_keys returned non-success status" - ); - return Err(KvRouterError::IndexerOffline); + fn apply_batch(&self, sequence: u64, events: Vec) { + let previous = self.last_sequence.swap(sequence, Ordering::AcqRel); + if (previous == 0 && sequence != 1) || (previous != 0 && sequence != previous + 1) { + self.present_keys.clear(); + tracing::warn!( + previous, + sequence, + "Mooncake KV event sequence gap; cleared shared-cache state" + ); + } + + for event in events { + // The shared-cache query contract currently has no tenant input and + // historically queried Mooncake's default tenant only. + if !event.tenant_id.is_empty() && event.tenant_id != "default" { + continue; + } + let Some(object_key) = event.object_key else { + continue; + }; + match event.event_type.as_str() { + "stored" => { + self.present_keys.insert(object_key); + } + "removed" => { + self.present_keys.remove(&object_key); + } + _ => {} + } + } + } + + fn clear(&self) { + self.present_keys.clear(); + self.last_sequence.store(0, Ordering::Release); + } + + async fn run_subscriber(mut self, cancellation_token: CancellationToken) { + let endpoint = loop { + if let Some(endpoint) = self + .resolve_mooncake_config() + .and_then(|config| config.kv_events_endpoint) + { + break endpoint; } - let body: MooncakeBatchQueryKeysResponse = response.json().await.map_err(|e| { - tracing::warn!( - error = %e, - url = %url, - "Failed to decode Mooncake batch_query_keys response" - ); - KvRouterError::IndexerOffline - })?; - - if !body.success { - tracing::warn!(url = %url, "Mooncake batch_query_keys reported failure"); - return Err(KvRouterError::IndexerOffline); + tokio::select! { + _ = cancellation_token.cancelled() => return, + result = self.runtime_configs.changed() => { + if result.is_err() { + return; + } + } } + }; - for key in chunk { - let exists = body.data.get(key).map(|entry| entry.ok).unwrap_or(false); - key_presence.insert(key.clone(), exists); + let mut socket = match connect_sub_socket(&endpoint, None).await { + Ok(socket) => socket, + Err(error) => { + tracing::warn!(%endpoint, %error, "Failed to connect to Mooncake KV events"); + return; + } + }; + tracing::info!(%endpoint, "Connected to Mooncake KV events"); + + loop { + tokio::select! { + _ = cancellation_token.cancelled() => return, + message = socket.next() => { + let frames = match message { + Some(Ok(frames)) => multipart_message(frames), + Some(Err(error)) => { + tracing::warn!(%endpoint, %error, "Mooncake KV event stream failed"); + self.clear(); + return; + } + None => { + tracing::warn!(%endpoint, "Mooncake KV event stream ended"); + self.clear(); + return; + } + }; + if frames.len() != 3 || frames[1].len() != 8 { + tracing::warn!(frame_count = frames.len(), "Invalid Mooncake KV event frame"); + self.clear(); + continue; + } + let sequence = u64::from_be_bytes(frames[1].as_slice().try_into().unwrap()); + match rmp_serde::from_slice::(&frames[2]) { + Ok((_, events, _)) => self.apply_batch(sequence, events), + Err(error) => { + tracing::warn!(%error, "Failed to decode Mooncake KV event batch"); + self.clear(); + } + } + } } } - - Ok(key_presence) } } @@ -205,29 +262,19 @@ impl SharedKvCache for HicacheSharedKvCache { return Ok(SharedCacheHits::default()); } - let Some(endpoint) = mooncake_batch_query_endpoint(&config) else { - tracing::debug!("Mooncake master HTTP endpoint is unavailable"); + if config.kv_events_endpoint.is_none() { + tracing::debug!("Mooncake KV event endpoint is unavailable"); return Ok(SharedCacheHits::default()); - }; + } let page_hashes = logical_page_hashes(tokens, config.page_size, config.is_eagle); if page_hashes.is_empty() { return Ok(SharedCacheHits::default()); } - let page_query_keys = build_page_query_keys(&page_hashes, &config); - let all_actual_keys = page_query_keys - .iter() - .flat_map(|keys| keys.iter().cloned()) - .collect::>(); - - let key_presence = self.fetch_key_presence(&endpoint, &all_actual_keys).await?; - let page_hits = page_query_keys + let page_hits = build_page_query_keys(&page_hashes, &config) .iter() - .map(|keys| { - keys.iter() - .all(|key| key_presence.get(key).copied().unwrap_or(false)) - }) + .map(|keys| keys.iter().all(|key| self.present_keys.contains(key))) .collect::>(); Ok(SharedCacheHits::from_hits(&page_hits)) @@ -255,33 +302,6 @@ fn mooncake_config_from_runtime( } } -fn mooncake_batch_query_endpoint(config: &SglangHicacheMooncakeConfig) -> Option { - let master_server_address = config.master_server_address.as_deref()?; - - let mut url = Url::parse(&format!("http://{master_server_address}")) - .inspect_err(|error| { - tracing::warn!( - master_server_address, - %error, - "Failed to parse Mooncake master address" - ); - }) - .ok()?; - - if url.set_port(Some(config.master_metrics_port)).is_err() { - tracing::warn!( - master_server_address, - master_metrics_port = config.master_metrics_port, - "Failed to set Mooncake master HTTP port" - ); - return None; - } - - url.set_path("/batch_query_keys"); - url.set_query(None); - Some(url) -} - fn logical_page_hashes(tokens: &[u32], page_size: u32, is_eagle: bool) -> Vec { let page_size = page_size as usize; if page_size == 0 { @@ -410,11 +430,9 @@ fn maybe_prefix_key(logical_key: &str, extra_backend_tag: Option<&str>) -> Strin #[cfg(test)] mod tests { - use std::ops::Range; + use std::{collections::HashMap, ops::Range}; use super::*; - use mockito::{Matcher, Server}; - use serde_json::json; use tokio::sync::watch; fn mooncake_config() -> SglangHicacheMooncakeConfig { @@ -430,6 +448,7 @@ mod tests { extra_backend_tag: None, master_server_address: Some("127.0.0.1:50051".to_string()), master_metrics_port: 9003, + kv_events_endpoint: Some("tcp://127.0.0.1:5557".to_string()), } } @@ -532,41 +551,30 @@ mod tests { } #[tokio::test] - async fn test_check_blocks_queries_mooncake_master() { - let mut server = Server::new_async().await; - let server_url = Url::parse(&server.url()).unwrap(); - + async fn test_check_blocks_uses_mooncake_events() { let hash0 = "cf97adeedb59e05bfd73a2b4c2a8885708c4f4f70c84c64b27120e72ab733b72".to_string(); let hash1 = "4ebfa8a1f3c341517621838c6e1b9aa350307e3f00b3cbd1a07ef740f54396d6".to_string(); - - let response = json!({ - "success": true, - "data": { - format!("{hash0}_0_k"): {"ok": true, "values": []}, - format!("{hash0}_0_v"): {"ok": true, "values": []}, - format!("{hash1}_0_k"): {"ok": true, "values": []}, - format!("{hash1}_0_v"): {"ok": false, "error": "not found"}, - } - }); - - let mock = server - .mock("GET", "/batch_query_keys") - .match_query(Matcher::Exact(format!( - "keys={hash0}_0_k,{hash0}_0_v,{hash1}_0_k,{hash1}_0_v" - ))) - .with_status(200) - .with_header("content-type", "application/json") - .with_body(response.to_string()) - .create_async() - .await; - - let config = SglangHicacheMooncakeConfig { - master_server_address: Some(format!("{}:50051", server_url.host_str().unwrap())), - master_metrics_port: server_url.port().unwrap(), - ..mooncake_config() - }; - - let cache = HicacheSharedKvCache::new(runtime_watch_with_config(config)); + let cache = HicacheSharedKvCache::new(runtime_watch_with_config(mooncake_config())); + cache.apply_batch( + 1, + vec![ + MooncakeObjectEvent { + event_type: "stored".to_string(), + object_key: Some(format!("{hash0}_0_k")), + tenant_id: "default".to_string(), + }, + MooncakeObjectEvent { + event_type: "stored".to_string(), + object_key: Some(format!("{hash0}_0_v")), + tenant_id: "default".to_string(), + }, + MooncakeObjectEvent { + event_type: "stored".to_string(), + object_key: Some(format!("{hash1}_0_k")), + tenant_id: "default".to_string(), + }, + ], + ); let hits = cache .check_blocks(&[1, 2, 3, 4, 5, 6, 7, 8], 4, None) .await @@ -575,7 +583,43 @@ mod tests { assert_eq!(hits.ranges, vec![Range { start: 0, end: 1 }]); assert_eq!(hits.total_hits, 1); - mock.assert_async().await; + cache.apply_batch( + 2, + vec![MooncakeObjectEvent { + event_type: "removed".to_string(), + object_key: Some(format!("{hash0}_0_v")), + tenant_id: "default".to_string(), + }], + ); + let hits = cache + .check_blocks(&[1, 2, 3, 4, 5, 6, 7, 8], 4) + .await + .unwrap(); + assert_eq!(hits.total_hits, 0); + } + + #[test] + fn test_sequence_gap_clears_stale_keys() { + let cache = HicacheSharedKvCache::new(runtime_watch_with_config(mooncake_config())); + cache.apply_batch( + 1, + vec![MooncakeObjectEvent { + event_type: "stored".to_string(), + object_key: Some("old_0_k".to_string()), + tenant_id: "default".to_string(), + }], + ); + cache.apply_batch( + 3, + vec![MooncakeObjectEvent { + event_type: "stored".to_string(), + object_key: Some("new_0_k".to_string()), + tenant_id: "default".to_string(), + }], + ); + + assert!(!cache.present_keys.contains("old_0_k")); + assert!(cache.present_keys.contains("new_0_k")); } #[tokio::test] From adb395ee3133a2b4fd8163bc8ca05addc661fc60 Mon Sep 17 00:00:00 2001 From: Ishan Dhanani Date: Mon, 6 Jul 2026 15:02:51 +0000 Subject: [PATCH 02/12] perf(kv-router): cache verified Mooncake groups Signed-off-by: Ishan Dhanani --- .../backends/sglang/hicache.md | 10 +- lib/llm/src/kv_router/shared_cache.rs | 101 ++++++++++++++++-- 2 files changed, 97 insertions(+), 14 deletions(-) diff --git a/docs/fern/pages/developer-guide/knowledge-base/modular-components/backends/sglang/hicache.md b/docs/fern/pages/developer-guide/knowledge-base/modular-components/backends/sglang/hicache.md index 3dbd2484b713..742e87554830 100644 --- a/docs/fern/pages/developer-guide/knowledge-base/modular-components/backends/sglang/hicache.md +++ b/docs/fern/pages/developer-guide/knowledge-base/modular-components/backends/sglang/hicache.md @@ -90,9 +90,11 @@ flowchart LR The router maintains two indexes: - Its own radix tree, built from worker KV events (per-tier). -- A set of Mooncake object keys, built from the Mooncake master's event stream. +- A set of Mooncake object keys and verified logical groups, built from the Mooncake master's event stream. -Request routing checks both indexes locally. If the event stream starts after Mooncake already contains objects, the shared-pool index starts empty and learns about subsequent events. A sequence gap clears the shared-pool index to avoid stale hits. +Request routing checks both indexes locally. When an event includes SGLang's `group_id`, the first lookup verifies every physical object in the group and caches the result. Later requests use one group lookup per page; any event for that group invalidates the cached result. Events without `group_id` use exact object-key checks. + +If the event stream starts after Mooncake already contains objects, the shared-pool index starts empty and learns about subsequent events. A sequence gap clears the shared-pool index to avoid stale hits. ### Scoring @@ -162,7 +164,7 @@ python -m dynamo.sglang \ --hicache-ratio 2 \ --hicache-write-policy write_through \ --hicache-storage-backend mooncake \ - --hicache-storage-backend-extra-config '{"master_server_address": "mooncake-master.internal:50051"}' \ + --hicache-storage-backend-extra-config '{"master_server_address": "mooncake-master.internal:50051", "enable_group_semantics": true}' \ --skip-tokenizer-init ``` @@ -189,6 +191,8 @@ Per-request overrides are available via `RouterConfigOverride.shared_cache_multi Set `DYN_MOONCAKE_KV_EVENTS_ENDPOINT` on the worker to the Mooncake PUB endpoint, such as `tcp://mooncake-master.internal:5557`. The endpoint must be reachable from the frontend. When `--hicache-storage-backend mooncake` is set, Dynamo publishes the endpoint and required layout metadata through the worker's `ModelRuntimeConfig.engine_specific` blob under the key `sglang_hicache_mooncake`. +Set `enable_group_semantics` to `true` in `--hicache-storage-backend-extra-config` to include SGLang logical group IDs in Mooncake metadata. Dynamo falls back to exact physical-key checks when group metadata is unavailable. + ## Verification **Events carry a medium.** Run the worker with `--log-level debug` and grep the log: diff --git a/lib/llm/src/kv_router/shared_cache.rs b/lib/llm/src/kv_router/shared_cache.rs index 428fdb13370a..99e99e63b217 100644 --- a/lib/llm/src/kv_router/shared_cache.rs +++ b/lib/llm/src/kv_router/shared_cache.rs @@ -17,7 +17,7 @@ use std::sync::{ }; use async_trait::async_trait; -use dashmap::DashSet; +use dashmap::{DashMap, DashSet}; use futures::StreamExt; use serde::{Deserialize, Serialize}; use sha2::{Digest, Sha256}; @@ -61,6 +61,8 @@ struct MooncakeObjectEvent { object_key: Option, #[serde(default)] tenant_id: String, + #[serde(default)] + group_id: Option, } type MooncakeEventBatch = (i64, Vec, u32); @@ -76,6 +78,7 @@ enum QueryToken { pub struct HicacheSharedKvCache { runtime_configs: RuntimeConfigWatch, present_keys: Arc>, + group_states: Arc>, last_sequence: Arc, } @@ -84,6 +87,7 @@ impl HicacheSharedKvCache { Self { runtime_configs, present_keys: Arc::new(DashSet::new()), + group_states: Arc::new(DashMap::new()), last_sequence: Arc::new(AtomicU64::new(0)), } } @@ -121,6 +125,7 @@ impl HicacheSharedKvCache { let previous = self.last_sequence.swap(sequence, Ordering::AcqRel); if (previous == 0 && sequence != 1) || (previous != 0 && sequence != previous + 1) { self.present_keys.clear(); + self.group_states.clear(); tracing::warn!( previous, sequence, @@ -137,6 +142,7 @@ impl HicacheSharedKvCache { let Some(object_key) = event.object_key else { continue; }; + let group_id = event.group_id.filter(|id| !id.is_empty()); match event.event_type.as_str() { "stored" => { self.present_keys.insert(object_key); @@ -146,11 +152,15 @@ impl HicacheSharedKvCache { } _ => {} } + if let Some(group_id) = group_id { + self.group_states.insert(group_id, (sequence, false)); + } } } fn clear(&self) { self.present_keys.clear(); + self.group_states.clear(); self.last_sequence.store(0, Ordering::Release); } @@ -272,9 +282,29 @@ impl SharedKvCache for HicacheSharedKvCache { return Ok(SharedCacheHits::default()); } - let page_hits = build_page_query_keys(&page_hashes, &config) + let page_hits = page_hashes .iter() - .map(|keys| keys.iter().all(|key| self.present_keys.contains(key))) + .map(|page_hash| { + let group_id = sglang_group_id(page_hash, &config); + let generation = self.group_states.get(&group_id).map(|state| *state); + if generation.is_some_and(|(_, verified)| verified) { + return true; + } + + let hit = expand_actual_query_keys(page_hash, &config) + .iter() + .all(|key| self.present_keys.contains(key)); + if hit + && let Some((generation, _)) = generation + && let Some(mut state) = self.group_states.get_mut(&group_id) + { + if state.0 != generation { + return false; + } + state.1 = true; + } + hit + }) .collect::>(); Ok(SharedCacheHits::from_hits(&page_hits)) @@ -367,14 +397,15 @@ fn hex_encode(bytes: &[u8]) -> String { output } -fn build_page_query_keys( - page_hashes: &[String], - config: &SglangHicacheMooncakeConfig, -) -> Vec> { - page_hashes - .iter() - .map(|page_hash| expand_actual_query_keys(page_hash, config)) - .collect() +fn sglang_group_id(logical_page_hash: &str, config: &SglangHicacheMooncakeConfig) -> String { + match config + .extra_backend_tag + .as_deref() + .filter(|tag| !tag.is_empty()) + { + Some(tag) => format!("sglang-hicache:{tag}_{logical_page_hash}"), + None => format!("sglang-hicache:{logical_page_hash}"), + } } fn expand_actual_query_keys( @@ -562,16 +593,19 @@ mod tests { event_type: "stored".to_string(), object_key: Some(format!("{hash0}_0_k")), tenant_id: "default".to_string(), + group_id: None, }, MooncakeObjectEvent { event_type: "stored".to_string(), object_key: Some(format!("{hash0}_0_v")), tenant_id: "default".to_string(), + group_id: None, }, MooncakeObjectEvent { event_type: "stored".to_string(), object_key: Some(format!("{hash1}_0_k")), tenant_id: "default".to_string(), + group_id: None, }, ], ); @@ -589,6 +623,7 @@ mod tests { event_type: "removed".to_string(), object_key: Some(format!("{hash0}_0_v")), tenant_id: "default".to_string(), + group_id: None, }], ); let hits = cache @@ -598,6 +633,46 @@ mod tests { assert_eq!(hits.total_hits, 0); } + #[tokio::test] + async fn test_check_blocks_uses_group_id() { + let hash = "cf97adeedb59e05bfd73a2b4c2a8885708c4f4f70c84c64b27120e72ab733b72"; + let group_id = format!("sglang-hicache:{hash}"); + let cache = HicacheSharedKvCache::new(runtime_watch_with_config(mooncake_config())); + cache.apply_batch( + 1, + vec![ + MooncakeObjectEvent { + event_type: "stored".to_string(), + object_key: Some(format!("{hash}_0_k")), + tenant_id: "default".to_string(), + group_id: Some(group_id.clone()), + }, + MooncakeObjectEvent { + event_type: "stored".to_string(), + object_key: Some(format!("{hash}_0_v")), + tenant_id: "default".to_string(), + group_id: Some(group_id.clone()), + }, + ], + ); + + let hits = cache.check_blocks(&[1, 2, 3, 4], 4).await.unwrap(); + assert_eq!(hits.total_hits, 1); + assert!(cache.group_states.get(&group_id).is_some_and(|v| v.1)); + + cache.apply_batch( + 2, + vec![MooncakeObjectEvent { + event_type: "removed".to_string(), + object_key: Some(format!("{hash}_0_v")), + tenant_id: "default".to_string(), + group_id: Some(group_id), + }], + ); + let hits = cache.check_blocks(&[1, 2, 3, 4], 4).await.unwrap(); + assert_eq!(hits.total_hits, 0); + } + #[test] fn test_sequence_gap_clears_stale_keys() { let cache = HicacheSharedKvCache::new(runtime_watch_with_config(mooncake_config())); @@ -607,18 +682,22 @@ mod tests { event_type: "stored".to_string(), object_key: Some("old_0_k".to_string()), tenant_id: "default".to_string(), + group_id: Some("old-group".to_string()), }], ); + assert!(cache.group_states.contains_key("old-group")); cache.apply_batch( 3, vec![MooncakeObjectEvent { event_type: "stored".to_string(), object_key: Some("new_0_k".to_string()), tenant_id: "default".to_string(), + group_id: None, }], ); assert!(!cache.present_keys.contains("old_0_k")); + assert!(cache.group_states.is_empty()); assert!(cache.present_keys.contains("new_0_k")); } From 0b92e88412feb0313a03b6036a13ab809a106e17 Mon Sep 17 00:00:00 2001 From: Ishan Dhanani Date: Tue, 28 Jul 2026 00:11:38 +0000 Subject: [PATCH 03/12] test(kv-router): update shared cache tests Signed-off-by: Ishan Dhanani --- lib/llm/src/kv_router/shared_cache.rs | 8 +++++--- 1 file changed, 5 insertions(+), 3 deletions(-) diff --git a/lib/llm/src/kv_router/shared_cache.rs b/lib/llm/src/kv_router/shared_cache.rs index 99e99e63b217..823113370c07 100644 --- a/lib/llm/src/kv_router/shared_cache.rs +++ b/lib/llm/src/kv_router/shared_cache.rs @@ -464,6 +464,8 @@ mod tests { use std::{collections::HashMap, ops::Range}; use super::*; + use mockito::Server; + use reqwest::Url; use tokio::sync::watch; fn mooncake_config() -> SglangHicacheMooncakeConfig { @@ -627,7 +629,7 @@ mod tests { }], ); let hits = cache - .check_blocks(&[1, 2, 3, 4, 5, 6, 7, 8], 4) + .check_blocks(&[1, 2, 3, 4, 5, 6, 7, 8], 4, None) .await .unwrap(); assert_eq!(hits.total_hits, 0); @@ -656,7 +658,7 @@ mod tests { ], ); - let hits = cache.check_blocks(&[1, 2, 3, 4], 4).await.unwrap(); + let hits = cache.check_blocks(&[1, 2, 3, 4], 4, None).await.unwrap(); assert_eq!(hits.total_hits, 1); assert!(cache.group_states.get(&group_id).is_some_and(|v| v.1)); @@ -669,7 +671,7 @@ mod tests { group_id: Some(group_id), }], ); - let hits = cache.check_blocks(&[1, 2, 3, 4], 4).await.unwrap(); + let hits = cache.check_blocks(&[1, 2, 3, 4], 4, None).await.unwrap(); assert_eq!(hits.total_hits, 0); } From 2637ade33337e7e44c16c6e8b2f88d6c55a85322 Mon Sep 17 00:00:00 2001 From: Ishan Dhanani Date: Tue, 28 Jul 2026 00:26:03 +0000 Subject: [PATCH 04/12] fix(kv-router): recover Mooncake cache subscription Signed-off-by: Ishan Dhanani --- lib/llm/src/discovery/model_manager.rs | 25 +++++- lib/llm/src/kv_router/shared_cache.rs | 111 ++++++++++++++++--------- 2 files changed, 93 insertions(+), 43 deletions(-) diff --git a/lib/llm/src/discovery/model_manager.rs b/lib/llm/src/discovery/model_manager.rs index 5236e1801efe..1e3afe356945 100644 --- a/lib/llm/src/discovery/model_manager.rs +++ b/lib/llm/src/discovery/model_manager.rs @@ -148,6 +148,9 @@ pub struct ModelManager { /// Per-endpoint runtime config watchers. Keyed by EndpointId (includes namespace). runtime_configs: DashMap, + /// Per-component HiCache state and its one Mooncake event subscriber. + hicache_caches: DashMap, + /// Shared KV-source membership coordinators, scoped by exact serving endpoint. /// Weak ownership lets the discovery loop stop when its last consumer goes away. kv_source_memberships: DashMap>, @@ -188,6 +191,7 @@ impl ModelManager { prefill_router_activators: DashMap::new(), encoder_router_activators: DashMap::new(), runtime_configs: DashMap::new(), + hicache_caches: DashMap::new(), kv_source_memberships: DashMap::new(), lora_domains: DashMap::new(), lora_enabled: crate::lora::lora_serving_enabled(), @@ -1050,6 +1054,21 @@ impl ModelManager { lora_enabled && worker_type == crate::protocols::common::timing::WORKER_TYPE_DECODE } + fn hicache_cache_for( + &self, + endpoint: &Endpoint, + runtime_configs: RuntimeConfigWatch, + ) -> HicacheSharedKvCache { + self.hicache_caches + .entry(endpoint.component().to_string()) + .or_insert_with(|| { + let cache = HicacheSharedKvCache::new(runtime_configs); + cache.start_subscriber(endpoint.component()); + cache + }) + .clone() + } + #[allow(clippy::too_many_arguments)] pub async fn kv_chooser_for( &self, @@ -1127,9 +1146,9 @@ impl ModelManager { worker_component = worker_component_name, "Using HiCache shared KV cache" ); - let cache = HicacheSharedKvCache::new(workers_with_configs.clone()); - cache.start_subscriber(endpoint.component()); - Some(Box::new(cache)) + Some(Box::new( + self.hicache_cache_for(endpoint, workers_with_configs.clone()), + )) } }; diff --git a/lib/llm/src/kv_router/shared_cache.rs b/lib/llm/src/kv_router/shared_cache.rs index 823113370c07..dffd315c93b8 100644 --- a/lib/llm/src/kv_router/shared_cache.rs +++ b/lib/llm/src/kv_router/shared_cache.rs @@ -15,6 +15,7 @@ use std::sync::{ Arc, atomic::{AtomicU64, Ordering}, }; +use std::time::Duration; use async_trait::async_trait; use dashmap::{DashMap, DashSet}; @@ -37,6 +38,7 @@ use crate::{ }; const SGLANG_HICACHE_MOONCAKE_RUNTIME_KEY: &str = "sglang_hicache_mooncake"; +const MOONCAKE_EVENT_RECONNECT_DELAY: Duration = Duration::from_secs(1); #[derive(Debug, Clone, Deserialize, Serialize, PartialEq, Eq)] struct SglangHicacheMooncakeConfig { @@ -146,15 +148,18 @@ impl HicacheSharedKvCache { match event.event_type.as_str() { "stored" => { self.present_keys.insert(object_key); + if let Some(group_id) = group_id { + self.group_states.insert(group_id, (sequence, false)); + } } "removed" => { self.present_keys.remove(&object_key); + // An older Mooncake publisher may omit `group_id` on removal. Clearing all + // verified groups is conservative and prevents a stale group fast-path hit. + self.group_states.clear(); } _ => {} } - if let Some(group_id) = group_id { - self.group_states.insert(group_id, (sequence, false)); - } } } @@ -183,47 +188,55 @@ impl HicacheSharedKvCache { } }; - let mut socket = match connect_sub_socket(&endpoint, None).await { - Ok(socket) => socket, - Err(error) => { - tracing::warn!(%endpoint, %error, "Failed to connect to Mooncake KV events"); - return; - } - }; - tracing::info!(%endpoint, "Connected to Mooncake KV events"); - loop { - tokio::select! { - _ = cancellation_token.cancelled() => return, - message = socket.next() => { - let frames = match message { - Some(Ok(frames)) => multipart_message(frames), - Some(Err(error)) => { - tracing::warn!(%endpoint, %error, "Mooncake KV event stream failed"); - self.clear(); - return; - } - None => { - tracing::warn!(%endpoint, "Mooncake KV event stream ended"); - self.clear(); - return; - } - }; - if frames.len() != 3 || frames[1].len() != 8 { - tracing::warn!(frame_count = frames.len(), "Invalid Mooncake KV event frame"); - self.clear(); - continue; + self.clear(); + let mut socket = match connect_sub_socket(&endpoint, None).await { + Ok(socket) => socket, + Err(error) => { + tracing::warn!(%endpoint, %error, "Failed to connect to Mooncake KV events; retrying"); + tokio::select! { + _ = cancellation_token.cancelled() => return, + _ = tokio::time::sleep(MOONCAKE_EVENT_RECONNECT_DELAY) => continue, } - let sequence = u64::from_be_bytes(frames[1].as_slice().try_into().unwrap()); - match rmp_serde::from_slice::(&frames[2]) { - Ok((_, events, _)) => self.apply_batch(sequence, events), - Err(error) => { - tracing::warn!(%error, "Failed to decode Mooncake KV event batch"); - self.clear(); + } + }; + tracing::info!(%endpoint, "Connected to Mooncake KV events"); + + loop { + tokio::select! { + _ = cancellation_token.cancelled() => return, + message = socket.next() => { + let frames = match message { + Some(Ok(frames)) => multipart_message(frames), + Some(Err(error)) => { + tracing::warn!(%endpoint, %error, "Mooncake KV event stream failed; reconnecting"); + break; + } + None => { + tracing::warn!(%endpoint, "Mooncake KV event stream ended; reconnecting"); + break; + } + }; + if frames.len() != 3 || frames[1].len() != 8 { + tracing::warn!(frame_count = frames.len(), "Invalid Mooncake KV event frame; reconnecting"); + break; + } + let sequence = u64::from_be_bytes(frames[1].as_slice().try_into().unwrap()); + match rmp_serde::from_slice::(&frames[2]) { + Ok((_, events, _)) => self.apply_batch(sequence, events), + Err(error) => { + tracing::warn!(%error, "Failed to decode Mooncake KV event batch; reconnecting"); + break; + } } } } } + + tokio::select! { + _ = cancellation_token.cancelled() => return, + _ = tokio::time::sleep(MOONCAKE_EVENT_RECONNECT_DELAY) => {} + } } } } @@ -636,7 +649,7 @@ mod tests { } #[tokio::test] - async fn test_check_blocks_uses_group_id() { + async fn test_check_blocks_invalidates_group_on_unlabeled_removal() { let hash = "cf97adeedb59e05bfd73a2b4c2a8885708c4f4f70c84c64b27120e72ab733b72"; let group_id = format!("sglang-hicache:{hash}"); let cache = HicacheSharedKvCache::new(runtime_watch_with_config(mooncake_config())); @@ -668,13 +681,31 @@ mod tests { event_type: "removed".to_string(), object_key: Some(format!("{hash}_0_v")), tenant_id: "default".to_string(), - group_id: Some(group_id), + group_id: None, }], ); + assert!(cache.group_states.is_empty()); let hits = cache.check_blocks(&[1, 2, 3, 4], 4, None).await.unwrap(); assert_eq!(hits.total_hits, 0); } + #[tokio::test] + async fn test_subscriber_retries_failed_connection_until_cancelled() { + let mut config = mooncake_config(); + config.kv_events_endpoint = Some("invalid://mooncake-events".to_string()); + let cache = HicacheSharedKvCache::new(runtime_watch_with_config(config)); + let cancellation_token = CancellationToken::new(); + let task = tokio::spawn(cache.run_subscriber(cancellation_token.clone())); + + tokio::time::sleep(Duration::from_millis(20)).await; + assert!(!task.is_finished()); + cancellation_token.cancel(); + tokio::time::timeout(Duration::from_secs(1), task) + .await + .unwrap() + .unwrap(); + } + #[test] fn test_sequence_gap_clears_stale_keys() { let cache = HicacheSharedKvCache::new(runtime_watch_with_config(mooncake_config())); From b775b6ddaf04f2a54defadb99069031099a224e8 Mon Sep 17 00:00:00 2001 From: Ishan Dhanani Date: Tue, 28 Jul 2026 16:36:00 +0000 Subject: [PATCH 05/12] fix(kv-router): scope HiCache group invalidation Signed-off-by: Ishan Dhanani --- lib/llm/src/kv_router/shared_cache.rs | 127 +++++++++++++++++++++----- 1 file changed, 103 insertions(+), 24 deletions(-) diff --git a/lib/llm/src/kv_router/shared_cache.rs b/lib/llm/src/kv_router/shared_cache.rs index dffd315c93b8..fca197acd2c6 100644 --- a/lib/llm/src/kv_router/shared_cache.rs +++ b/lib/llm/src/kv_router/shared_cache.rs @@ -13,7 +13,7 @@ use std::sync::{ Arc, - atomic::{AtomicU64, Ordering}, + atomic::{AtomicBool, AtomicU64, Ordering}, }; use std::time::Duration; @@ -82,6 +82,7 @@ pub struct HicacheSharedKvCache { present_keys: Arc>, group_states: Arc>, last_sequence: Arc, + has_sequence: Arc, } impl HicacheSharedKvCache { @@ -91,6 +92,7 @@ impl HicacheSharedKvCache { present_keys: Arc::new(DashSet::new()), group_states: Arc::new(DashMap::new()), last_sequence: Arc::new(AtomicU64::new(0)), + has_sequence: Arc::new(AtomicBool::new(false)), } } @@ -124,8 +126,12 @@ impl HicacheSharedKvCache { } fn apply_batch(&self, sequence: u64, events: Vec) { + let has_previous = self.has_sequence.swap(true, Ordering::AcqRel); let previous = self.last_sequence.swap(sequence, Ordering::AcqRel); - if (previous == 0 && sequence != 1) || (previous != 0 && sequence != previous + 1) { + if has_previous && sequence == previous { + return; + } + if has_previous && sequence != previous.wrapping_add(1) { self.present_keys.clear(); self.group_states.clear(); tracing::warn!( @@ -154,9 +160,13 @@ impl HicacheSharedKvCache { } "removed" => { self.present_keys.remove(&object_key); - // An older Mooncake publisher may omit `group_id` on removal. Clearing all - // verified groups is conservative and prevents a stale group fast-path hit. - self.group_states.clear(); + if let Some(group_id) = group_id { + self.group_states.remove(&group_id); + } else { + // An older Mooncake publisher may omit `group_id` on removal. Clearing all + // verified groups is conservative and prevents a stale group fast-path hit. + self.group_states.clear(); + } } _ => {} } @@ -167,6 +177,7 @@ impl HicacheSharedKvCache { self.present_keys.clear(); self.group_states.clear(); self.last_sequence.store(0, Ordering::Release); + self.has_sequence.store(false, Ordering::Release); } async fn run_subscriber(mut self, cancellation_token: CancellationToken) { @@ -307,14 +318,14 @@ impl SharedKvCache for HicacheSharedKvCache { let hit = expand_actual_query_keys(page_hash, &config) .iter() .all(|key| self.present_keys.contains(key)); - if hit - && let Some((generation, _)) = generation - && let Some(mut state) = self.group_states.get_mut(&group_id) - { - if state.0 != generation { - return false; + if hit && let Some((generation, _)) = generation { + if let Some(mut state) = self.group_states.get_mut(&group_id) { + // A concurrent stored event invalidates verification, not the physical + // key check that already proved this request is a hit. + if state.0 == generation { + state.1 = true; + } } - state.1 = true; } hit }) @@ -477,8 +488,6 @@ mod tests { use std::{collections::HashMap, ops::Range}; use super::*; - use mockito::Server; - use reqwest::Url; use tokio::sync::watch; fn mooncake_config() -> SglangHicacheMooncakeConfig { @@ -596,6 +605,16 @@ mod tests { ); } + #[test] + fn test_sglang_group_id_uses_extra_backend_tag() { + let config = SglangHicacheMooncakeConfig { + extra_backend_tag: Some("tag".to_string()), + ..mooncake_config() + }; + + assert_eq!(sglang_group_id("hash", &config), "sglang-hicache:tag_hash"); + } + #[tokio::test] async fn test_check_blocks_uses_mooncake_events() { let hash0 = "cf97adeedb59e05bfd73a2b4c2a8885708c4f4f70c84c64b27120e72ab733b72".to_string(); @@ -689,6 +708,75 @@ mod tests { assert_eq!(hits.total_hits, 0); } + #[tokio::test] + async fn test_labeled_removal_preserves_other_verified_groups() { + let hash0 = "cf97adeedb59e05bfd73a2b4c2a8885708c4f4f70c84c64b27120e72ab733b72"; + let hash1 = "4ebfa8a1f3c341517621838c6e1b9aa350307e3f00b3cbd1a07ef740f54396d6"; + let group0 = format!("sglang-hicache:{hash0}"); + let group1 = format!("sglang-hicache:{hash1}"); + let cache = HicacheSharedKvCache::new(runtime_watch_with_config(mooncake_config())); + cache.apply_batch( + 1, + [(&hash0, &group0), (&hash1, &group1)] + .into_iter() + .flat_map(|(hash, group_id)| { + ["k", "v"].into_iter().map(move |kind| MooncakeObjectEvent { + event_type: "stored".to_string(), + object_key: Some(format!("{hash}_0_{kind}")), + tenant_id: "default".to_string(), + group_id: Some(group_id.clone()), + }) + }) + .collect(), + ); + assert_eq!( + cache + .check_blocks(&[1, 2, 3, 4, 5, 6, 7, 8], 4, None) + .await + .unwrap() + .total_hits, + 2 + ); + + cache.apply_batch( + 2, + vec![MooncakeObjectEvent { + event_type: "removed".to_string(), + object_key: Some(format!("{hash0}_0_k")), + tenant_id: "default".to_string(), + group_id: Some(group0.clone()), + }], + ); + + assert!(!cache.group_states.contains_key(&group0)); + assert!(cache.group_states.get(&group1).is_some_and(|state| state.1)); + } + + #[test] + fn test_duplicate_sequence_preserves_shared_cache_state() { + let cache = HicacheSharedKvCache::new(runtime_watch_with_config(mooncake_config())); + cache.apply_batch( + 0, + vec![MooncakeObjectEvent { + event_type: "stored".to_string(), + object_key: Some("key-0".to_string()), + tenant_id: "default".to_string(), + group_id: None, + }], + ); + cache.apply_batch( + 0, + vec![MooncakeObjectEvent { + event_type: "removed".to_string(), + object_key: Some("key-0".to_string()), + tenant_id: "default".to_string(), + group_id: None, + }], + ); + + assert!(cache.present_keys.contains("key-0")); + } + #[tokio::test] async fn test_subscriber_retries_failed_connection_until_cancelled() { let mut config = mooncake_config(); @@ -736,16 +824,7 @@ mod tests { #[tokio::test] async fn test_check_blocks_skips_mooncake_for_cache_namespace() { - let server = Server::new_async().await; - let server_url = Url::parse(&server.url()).unwrap(); - - let config = SglangHicacheMooncakeConfig { - master_server_address: Some(format!("{}:50051", server_url.host_str().unwrap())), - master_metrics_port: server_url.port().unwrap(), - ..mooncake_config() - }; - - let cache = HicacheSharedKvCache::new(runtime_watch_with_config(config)); + let cache = HicacheSharedKvCache::new(runtime_watch_with_config(mooncake_config())); let hits = cache .check_blocks(&[1, 2, 3, 4, 5, 6, 7, 8], 4, Some("tenant-a")) .await From 9b491e8e905e46bf2698d38d2a153ed56aff7cbd Mon Sep 17 00:00:00 2001 From: Ishan Dhanani Date: Tue, 28 Jul 2026 16:41:44 +0000 Subject: [PATCH 06/12] fix(kv-router): recover Mooncake subscriber updates Signed-off-by: Ishan Dhanani --- lib/llm/src/kv_router/shared_cache.rs | 149 +++++++++++++++++++++----- 1 file changed, 120 insertions(+), 29 deletions(-) diff --git a/lib/llm/src/kv_router/shared_cache.rs b/lib/llm/src/kv_router/shared_cache.rs index fca197acd2c6..405413f9f1f0 100644 --- a/lib/llm/src/kv_router/shared_cache.rs +++ b/lib/llm/src/kv_router/shared_cache.rs @@ -33,6 +33,7 @@ use dynamo_runtime::{component::Component, traits::DistributedRuntimeProvider}; use crate::{ discovery::RuntimeConfigWatch, + kv_router::metrics::RoutingOverheadMetrics, local_model::runtime_config::ModelRuntimeConfig, utils::zmq::{connect_sub_socket, multipart_message}, }; @@ -56,7 +57,7 @@ struct SglangHicacheMooncakeConfig { kv_events_endpoint: Option, } -#[derive(Debug, Deserialize)] +#[derive(Debug, Deserialize, Serialize)] struct MooncakeObjectEvent { event_type: String, #[serde(default)] @@ -67,6 +68,8 @@ struct MooncakeObjectEvent { group_id: Option, } +/// Mooncake KV publisher payload: `(timestamp_ms, events, dp_rank)` after the empty-topic and +/// big-endian sequence ZMQ frames. type MooncakeEventBatch = (i64, Vec, u32); #[derive(Debug, Clone, Copy)] @@ -125,6 +128,11 @@ impl HicacheSharedKvCache { Some(first.clone()) } + fn kv_events_endpoint(&self) -> Option { + self.resolve_mooncake_config() + .and_then(|config| config.kv_events_endpoint) + } + fn apply_batch(&self, sequence: u64, events: Vec) { let has_previous = self.has_sequence.swap(true, Ordering::AcqRel); let previous = self.last_sequence.swap(sequence, Ordering::AcqRel); @@ -180,30 +188,34 @@ impl HicacheSharedKvCache { self.has_sequence.store(false, Ordering::Release); } + fn record_subscriber_error(&self) { + if let Some(metrics) = RoutingOverheadMetrics::get() { + metrics.inc_shared_cache_errors(); + } + } + async fn run_subscriber(mut self, cancellation_token: CancellationToken) { - let endpoint = loop { - if let Some(endpoint) = self - .resolve_mooncake_config() - .and_then(|config| config.kv_events_endpoint) - { - break endpoint; - } + loop { + let endpoint = loop { + if let Some(endpoint) = self.kv_events_endpoint() { + break endpoint; + } - tokio::select! { - _ = cancellation_token.cancelled() => return, - result = self.runtime_configs.changed() => { - if result.is_err() { - return; + tokio::select! { + _ = cancellation_token.cancelled() => return, + result = self.runtime_configs.changed() => { + if result.is_err() { + return; + } } } - } - }; + }; - loop { self.clear(); let mut socket = match connect_sub_socket(&endpoint, None).await { Ok(socket) => socket, Err(error) => { + self.record_subscriber_error(); tracing::warn!(%endpoint, %error, "Failed to connect to Mooncake KV events; retrying"); tokio::select! { _ = cancellation_token.cancelled() => return, @@ -216,28 +228,35 @@ impl HicacheSharedKvCache { loop { tokio::select! { _ = cancellation_token.cancelled() => return, + result = self.runtime_configs.changed() => { + if result.is_err() { + return; + } + let next_endpoint = self.kv_events_endpoint(); + if next_endpoint.as_deref() != Some(endpoint.as_str()) { + tracing::info!(%endpoint, next_endpoint = ?next_endpoint, "Mooncake KV event endpoint changed; reconnecting"); + break; + } + } message = socket.next() => { let frames = match message { Some(Ok(frames)) => multipart_message(frames), Some(Err(error)) => { + self.record_subscriber_error(); tracing::warn!(%endpoint, %error, "Mooncake KV event stream failed; reconnecting"); break; } None => { + self.record_subscriber_error(); tracing::warn!(%endpoint, "Mooncake KV event stream ended; reconnecting"); break; } }; - if frames.len() != 3 || frames[1].len() != 8 { - tracing::warn!(frame_count = frames.len(), "Invalid Mooncake KV event frame; reconnecting"); - break; - } - let sequence = u64::from_be_bytes(frames[1].as_slice().try_into().unwrap()); - match rmp_serde::from_slice::(&frames[2]) { - Ok((_, events, _)) => self.apply_batch(sequence, events), + match parse_mooncake_event_frames(&frames) { + Ok((sequence, events)) => self.apply_batch(sequence, events), Err(error) => { - tracing::warn!(%error, "Failed to decode Mooncake KV event batch; reconnecting"); - break; + self.record_subscriber_error(); + tracing::warn!(%error, "Dropping invalid Mooncake KV event frame"); } } } @@ -252,6 +271,22 @@ impl HicacheSharedKvCache { } } +fn parse_mooncake_event_frames( + frames: &[Vec], +) -> anyhow::Result<(u64, Vec)> { + let [_, sequence, payload] = frames else { + anyhow::bail!("expected three frames, got {}", frames.len()); + }; + let sequence = u64::from_be_bytes( + sequence + .as_slice() + .try_into() + .map_err(|_| anyhow::anyhow!("expected an 8-byte sequence frame"))?, + ); + let (_, events, _) = rmp_serde::from_slice::(payload)?; + Ok((sequence, events)) +} + #[async_trait] impl SharedKvCache for HicacheSharedKvCache { async fn check_blocks( @@ -507,7 +542,12 @@ mod tests { } } - fn runtime_watch_with_config(config: SglangHicacheMooncakeConfig) -> RuntimeConfigWatch { + fn runtime_watch_with_config_and_sender( + config: SglangHicacheMooncakeConfig, + ) -> ( + RuntimeConfigWatch, + watch::Sender>, + ) { let mut runtime_config = ModelRuntimeConfig::new(); runtime_config .set_engine_specific(SGLANG_HICACHE_MOONCAKE_RUNTIME_KEY, config) @@ -516,8 +556,12 @@ mod tests { let mut workers = HashMap::new(); workers.insert(1, runtime_config); - let (_tx, rx) = watch::channel(workers); - rx + let (tx, rx) = watch::channel(workers); + (rx, tx) + } + + fn runtime_watch_with_config(config: SglangHicacheMooncakeConfig) -> RuntimeConfigWatch { + runtime_watch_with_config_and_sender(config).0 } #[test] @@ -615,6 +659,52 @@ mod tests { assert_eq!(sglang_group_id("hash", &config), "sglang-hicache:tag_hash"); } + #[test] + fn test_parse_mooncake_event_frames() { + let payload = rmp_serde::to_vec(&( + 0_i64, + vec![MooncakeObjectEvent { + event_type: "stored".to_string(), + object_key: Some("key-0".to_string()), + tenant_id: "default".to_string(), + group_id: Some("group-0".to_string()), + }], + 0_u32, + )) + .unwrap(); + + let (sequence, events) = + parse_mooncake_event_frames(&[Vec::new(), 7_u64.to_be_bytes().to_vec(), payload]) + .unwrap(); + + assert_eq!(sequence, 7); + assert_eq!(events[0].object_key.as_deref(), Some("key-0")); + assert_eq!(events[0].group_id.as_deref(), Some("group-0")); + } + + #[test] + fn test_kv_events_endpoint_tracks_runtime_config_updates() { + let (runtime_configs, tx) = runtime_watch_with_config_and_sender(mooncake_config()); + let cache = HicacheSharedKvCache::new(runtime_configs); + assert_eq!( + cache.kv_events_endpoint().as_deref(), + Some("tcp://127.0.0.1:5557") + ); + + let mut updated = mooncake_config(); + updated.kv_events_endpoint = Some("tcp://127.0.0.1:5558".to_string()); + let mut runtime_config = ModelRuntimeConfig::new(); + runtime_config + .set_engine_specific(SGLANG_HICACHE_MOONCAKE_RUNTIME_KEY, updated) + .unwrap(); + tx.send(HashMap::from([(1, runtime_config)])).unwrap(); + + assert_eq!( + cache.kv_events_endpoint().as_deref(), + Some("tcp://127.0.0.1:5558") + ); + } + #[tokio::test] async fn test_check_blocks_uses_mooncake_events() { let hash0 = "cf97adeedb59e05bfd73a2b4c2a8885708c4f4f70c84c64b27120e72ab733b72".to_string(); @@ -781,7 +871,8 @@ mod tests { async fn test_subscriber_retries_failed_connection_until_cancelled() { let mut config = mooncake_config(); config.kv_events_endpoint = Some("invalid://mooncake-events".to_string()); - let cache = HicacheSharedKvCache::new(runtime_watch_with_config(config)); + let (runtime_configs, _tx) = runtime_watch_with_config_and_sender(config); + let cache = HicacheSharedKvCache::new(runtime_configs); let cancellation_token = CancellationToken::new(); let task = tokio::spawn(cache.run_subscriber(cancellation_token.clone())); From d6a34cac2e936c5a7228339e78f38e0140013ea1 Mon Sep 17 00:00:00 2001 From: Ishan Dhanani Date: Tue, 28 Jul 2026 16:44:40 +0000 Subject: [PATCH 07/12] fix(discovery): scope HiCache subscribers to endpoints Signed-off-by: Ishan Dhanani --- lib/llm/src/discovery/model_manager.rs | 62 +++++++++++++++++++++++--- lib/llm/src/discovery/watcher.rs | 8 ++++ lib/llm/src/kv_router/shared_cache.rs | 30 +++++++++---- 3 files changed, 86 insertions(+), 14 deletions(-) diff --git a/lib/llm/src/discovery/model_manager.rs b/lib/llm/src/discovery/model_manager.rs index 1e3afe356945..644cfd0f5764 100644 --- a/lib/llm/src/discovery/model_manager.rs +++ b/lib/llm/src/discovery/model_manager.rs @@ -148,8 +148,8 @@ pub struct ModelManager { /// Per-endpoint runtime config watchers. Keyed by EndpointId (includes namespace). runtime_configs: DashMap, - /// Per-component HiCache state and its one Mooncake event subscriber. - hicache_caches: DashMap, + /// Per-endpoint HiCache state and its one Mooncake event subscriber. + hicache_caches: DashMap, /// Shared KV-source membership coordinators, scoped by exact serving endpoint. /// Weak ownership lets the discovery loop stop when its last consumer goes away. @@ -1060,15 +1060,35 @@ impl ModelManager { runtime_configs: RuntimeConfigWatch, ) -> HicacheSharedKvCache { self.hicache_caches - .entry(endpoint.component().to_string()) + .entry(endpoint.id()) .or_insert_with(|| { - let cache = HicacheSharedKvCache::new(runtime_configs); - cache.start_subscriber(endpoint.component()); + let cache = HicacheSharedKvCache::new_with_cancellation( + runtime_configs, + endpoint.component().drt().child_token(), + ); + cache.start_subscriber(); cache }) .clone() } + pub fn remove_hicache_caches(&self, namespace: &str, component: &str) { + let endpoint_ids = self + .hicache_caches + .iter() + .filter(|entry| { + entry.key().namespace == namespace && entry.key().component == component + }) + .map(|entry| entry.key().clone()) + .collect::>(); + + for endpoint_id in endpoint_ids { + if let Some((_, cache)) = self.hicache_caches.remove(&endpoint_id) { + cache.shutdown(); + } + } + } + #[allow(clippy::too_many_arguments)] pub async fn kv_chooser_for( &self, @@ -2397,6 +2417,38 @@ mod tests { assert!(mm.get_model("llama").is_some()); } + #[test] + fn remove_hicache_caches_cancels_only_the_removed_component() { + let manager = ModelManager::new(); + let (_tx, runtime_configs) = + tokio::sync::watch::channel(HashMap::::new()); + let cancelled = tokio_util::sync::CancellationToken::new(); + let retained = tokio_util::sync::CancellationToken::new(); + manager.hicache_caches.insert( + EndpointId::from("ns.worker.generate"), + HicacheSharedKvCache::new_with_cancellation(runtime_configs.clone(), cancelled.clone()), + ); + manager.hicache_caches.insert( + EndpointId::from("ns.other.generate"), + HicacheSharedKvCache::new_with_cancellation(runtime_configs, retained.clone()), + ); + + manager.remove_hicache_caches("ns", "worker"); + + assert!(cancelled.is_cancelled()); + assert!(!retained.is_cancelled()); + assert!( + !manager + .hicache_caches + .contains_key(&EndpointId::from("ns.worker.generate")) + ); + assert!( + manager + .hicache_caches + .contains_key(&EndpointId::from("ns.other.generate")) + ); + } + #[test] fn test_alias_resolution_maps_to_primary() { let mm = ModelManager::new(); diff --git a/lib/llm/src/discovery/watcher.rs b/lib/llm/src/discovery/watcher.rs index 1bd5eb378840..20ac889e5317 100644 --- a/lib/llm/src/discovery/watcher.rs +++ b/lib/llm/src/discovery/watcher.rs @@ -1090,6 +1090,9 @@ impl ModelWatcher { && eid.component == *worker_component && worker_set_key(eid, other_card.model_type, other_card.worker_type) == ws_key }); + let endpoint_has_instances = active_instances.iter().any(|(eid, _)| { + eid.namespace == *worker_namespace && eid.component == *worker_component + }); if !component_has_instances { // No more workers of this component in this namespace — remove its WorkerSet @@ -1111,6 +1114,11 @@ impl ModelWatcher { } } + if !endpoint_has_instances { + self.manager + .remove_hicache_caches(worker_namespace, worker_component); + } + // Activator-state cleanup depends on which component just went away. // // PREFILL teardown (cached endpoint is stale): drop everything for diff --git a/lib/llm/src/kv_router/shared_cache.rs b/lib/llm/src/kv_router/shared_cache.rs index 405413f9f1f0..50f97827163d 100644 --- a/lib/llm/src/kv_router/shared_cache.rs +++ b/lib/llm/src/kv_router/shared_cache.rs @@ -24,19 +24,17 @@ use serde::{Deserialize, Serialize}; use sha2::{Digest, Sha256}; use tokio_util::sync::CancellationToken; -use dynamo_kv_router::{ - SharedKvCache, - indexer::KvRouterError, - protocols::{SharedCacheHits, WorkerId}, -}; -use dynamo_runtime::{component::Component, traits::DistributedRuntimeProvider}; - use crate::{ discovery::RuntimeConfigWatch, kv_router::metrics::RoutingOverheadMetrics, local_model::runtime_config::ModelRuntimeConfig, utils::zmq::{connect_sub_socket, multipart_message}, }; +use dynamo_kv_router::{ + SharedKvCache, + indexer::KvRouterError, + protocols::{SharedCacheHits, WorkerId}, +}; const SGLANG_HICACHE_MOONCAKE_RUNTIME_KEY: &str = "sglang_hicache_mooncake"; const MOONCAKE_EVENT_RECONNECT_DELAY: Duration = Duration::from_secs(1); @@ -86,25 +84,39 @@ pub struct HicacheSharedKvCache { group_states: Arc>, last_sequence: Arc, has_sequence: Arc, + cancellation_token: CancellationToken, } impl HicacheSharedKvCache { pub fn new(runtime_configs: RuntimeConfigWatch) -> Self { + Self::new_with_cancellation(runtime_configs, CancellationToken::new()) + } + + pub fn new_with_cancellation( + runtime_configs: RuntimeConfigWatch, + cancellation_token: CancellationToken, + ) -> Self { Self { runtime_configs, present_keys: Arc::new(DashSet::new()), group_states: Arc::new(DashMap::new()), last_sequence: Arc::new(AtomicU64::new(0)), has_sequence: Arc::new(AtomicBool::new(false)), + cancellation_token, } } - pub fn start_subscriber(&self, component: &Component) { + pub fn start_subscriber(&self) { let cache = self.clone(); - let cancellation_token = component.drt().child_token(); + let cancellation_token = self.cancellation_token.clone(); tokio::spawn(async move { cache.run_subscriber(cancellation_token).await }); } + pub fn shutdown(&self) { + self.cancellation_token.cancel(); + self.clear(); + } + fn resolve_mooncake_config(&self) -> Option { let workers = self.runtime_configs.borrow(); let mut configs = Vec::new(); From 0b075f479d7c832d9b04bdb48b8ed5a244d90664 Mon Sep 17 00:00:00 2001 From: Ishan Dhanani Date: Tue, 28 Jul 2026 16:55:11 +0000 Subject: [PATCH 08/12] fix(sglang): configure Mooncake events on frontend Signed-off-by: Ishan Dhanani --- components/src/dynamo/sglang/register.py | 36 ------ .../backends/sglang/hicache.md | 22 ++-- lib/llm/src/discovery/model_manager.rs | 6 +- lib/llm/src/kv_router/metrics.rs | 2 +- lib/llm/src/kv_router/shared_cache.rs | 105 ++++++++++++++++-- lib/runtime/src/metrics/prometheus_names.rs | 2 +- 6 files changed, 114 insertions(+), 59 deletions(-) diff --git a/components/src/dynamo/sglang/register.py b/components/src/dynamo/sglang/register.py index f067a792929f..275854c097e1 100644 --- a/components/src/dynamo/sglang/register.py +++ b/components/src/dynamo/sglang/register.py @@ -8,7 +8,6 @@ from typing import Any, List, Optional import sglang as sgl -from sglang.srt.environ import envs from sglang.srt.server_args import ServerArgs from sglang.srt.speculative.spec_info import SpeculativeAlgorithm @@ -232,33 +231,6 @@ def _get_mooncake_runtime_data(server_args: ServerArgs) -> Optional[dict[str, An getattr(server_args, "hicache_storage_backend_extra_config", None) ) - try: - from sglang.srt.mem_cache.storage.mooncake_store.mooncake_store import ( - MooncakeStoreConfig, - ) - except ImportError as e: - logging.warning(f"MooncakeStoreConfig import unavailable: {e}") - return None - - # Graceful degradation: Mooncake runtime metadata is optional. If config - # resolution fails for any reason (file not found, malformed env vars, - # upstream API change), skip publishing the metadata rather than crashing - # the worker -- the worker still serves requests, just without HiCache - # router hints. Broad catch is intentional per python-guidelines.md. - try: - if extra_config and ( - extra_config.get("master_server_address") is not None - or extra_config.get("client_server_address") is not None - ): - mooncake_config = MooncakeStoreConfig.load_from_extra_config(extra_config) - elif envs.SGLANG_HICACHE_MOONCAKE_CONFIG_PATH.is_set(): - mooncake_config = MooncakeStoreConfig.from_file() - else: - mooncake_config = MooncakeStoreConfig.load_from_env() - except Exception as e: - logging.warning(f"Failed to resolve Mooncake config for runtime metadata: {e}") - return None - tp_size = int(getattr(server_args, "tp_size", 1) or 1) pp_size = int(getattr(server_args, "pp_size", 1) or 1) @@ -296,10 +268,6 @@ def _get_mooncake_runtime_data(server_args: ServerArgs) -> Optional[dict[str, An if not isinstance(extra_backend_tag, str) or not extra_backend_tag: extra_backend_tag = None - master_server_address = getattr(mooncake_config, "master_server_address", None) - if not isinstance(master_server_address, str) or not master_server_address: - master_server_address = None - return { "backend": "mooncake", "page_size": int(getattr(server_args, "page_size", 1) or 1), @@ -310,10 +278,6 @@ def _get_mooncake_runtime_data(server_args: ServerArgs) -> Optional[dict[str, An "tp_lcm_size": tp_lcm_size, "should_split_heads": should_split_heads, "extra_backend_tag": extra_backend_tag, - "master_server_address": master_server_address, - "master_metrics_port": int( - getattr(mooncake_config, "master_metrics_port", 9003) - ), "kv_events_endpoint": os.getenv("DYN_MOONCAKE_KV_EVENTS_ENDPOINT") or None, } diff --git a/docs/fern/pages/developer-guide/knowledge-base/modular-components/backends/sglang/hicache.md b/docs/fern/pages/developer-guide/knowledge-base/modular-components/backends/sglang/hicache.md index 742e87554830..9956e8c398a1 100644 --- a/docs/fern/pages/developer-guide/knowledge-base/modular-components/backends/sglang/hicache.md +++ b/docs/fern/pages/developer-guide/knowledge-base/modular-components/backends/sglang/hicache.md @@ -153,10 +153,9 @@ You also need: > [!WARNING] > SGLang 0.5.11 through 0.5.12.x bundle Mooncake 0.3.10.post2, which is affected by an upstream `MemcpyWorkerPool` crash when both `--enable-metrics` and `--disable-piecewise-cuda-graph` are set on the worker. Upgrade to SGLang 0.5.13 or later, which bundles Mooncake 0.3.11.post1 with [the upstream fix](https://github.com/kvcache-ai/Mooncake/pull/2001), or omit both flags. The setup below omits them. -**SGLang worker** — HiCache with Mooncake storage and the event endpoint advertised to Dynamo: +**SGLang worker** — HiCache with Mooncake storage: ```bash -DYN_MOONCAKE_KV_EVENTS_ENDPOINT=tcp://mooncake-master.internal:5557 \ python -m dynamo.sglang \ --model-path Qwen/Qwen3-0.6B \ --page-size 64 \ @@ -170,9 +169,10 @@ python -m dynamo.sglang \ Launch additional workers on other GPUs / hosts with the same Mooncake config so they back to the same cluster. -**Dynamo frontend** — enable tier-aware routing: +**Dynamo frontend** — configure the Mooncake event endpoint and enable tier-aware routing: ```bash +DYN_MOONCAKE_KV_EVENTS_ENDPOINT=tcp://mooncake-master.internal:5557 \ python -m dynamo.frontend \ --http-port 8000 \ --router-mode kv \ @@ -186,10 +186,11 @@ python -m dynamo.frontend \ | --------------------------- | ----------------------------- | ------- | ------------------------------------------------------------------------------------------------------------------------------------------------------------------ | | `--shared-cache-type` | `DYN_SHARED_CACHE_TYPE` | `none` | `none` disables shared-pool tracking; `hicache` enables the Mooncake event-backed index. | | `--shared-cache-multiplier` | `DYN_SHARED_CACHE_MULTIPLIER` | `0.5` | Discount factor for shared-pool hits. `0.0` ignores them; `0.5` treats a shared hit as half a device hit; `1.0` treats shared and device hits equally. | +| — | `DYN_MOONCAKE_KV_EVENTS_ENDPOINT` | unset | Mooncake PUB endpoint consumed by the frontend. When unset, Dynamo uses one consistently advertised worker endpoint. | Per-request overrides are available via `RouterConfigOverride.shared_cache_multiplier` for A/B experimentation without restarting the router. -Set `DYN_MOONCAKE_KV_EVENTS_ENDPOINT` on the worker to the Mooncake PUB endpoint, such as `tcp://mooncake-master.internal:5557`. The endpoint must be reachable from the frontend. When `--hicache-storage-backend mooncake` is set, Dynamo publishes the endpoint and required layout metadata through the worker's `ModelRuntimeConfig.engine_specific` blob under the key `sglang_hicache_mooncake`. +Set `DYN_MOONCAKE_KV_EVENTS_ENDPOINT` on the frontend to the Mooncake PUB endpoint, such as `tcp://mooncake-master.internal:5557`. The endpoint must be reachable from the frontend. Workers can advertise the same variable as a fallback, but a missing worker value does not disable shared-cache routing. Set `enable_group_semantics` to `true` in `--hicache-storage-backend-extra-config` to include SGLang logical group IDs in Mooncake metadata. Dynamo falls back to exact physical-key checks when group metadata is unavailable. @@ -205,12 +206,13 @@ python -m dynamo.sglang ... --log-level debug 2>&1 | grep -E 'BlockStored|BlockR If `medium` is missing or Host-tier transitions never report `CPU_PINNED`, confirm that the worker runs SGLang 0.5.11 or later (or a custom build that includes PR #22894). -**Router sees the shared pool.** Two new histograms are exposed on the frontend's Prometheus endpoint: +**Router sees the shared pool.** Shared-cache metrics are exposed on the frontend's Prometheus endpoint: -| Metric | Meaning | -| ----------------------------------- | ------------------------------------------------------------------------ | -| `router_shared_cache_hit_rate` | Fraction of request blocks found in the shared pool (0.0–1.0). | -| `router_shared_cache_beyond_blocks` | Blocks in the shared pool _beyond_ the selected worker's device overlap. | +| Metric | Meaning | +| ----------------------------------------- | ------------------------------------------------------------------------ | +| `router_shared_cache_hit_rate` | Fraction of request blocks found in the shared pool (0.0–1.0). | +| `router_shared_cache_beyond_blocks` | Blocks in the shared pool _beyond_ the selected worker's device overlap. | +| `dynamo_router_shared_cache_errors_total` | Shared-cache query and Mooncake subscriber failures. | ```bash curl -s localhost:8000/metrics | grep shared_cache @@ -220,7 +222,7 @@ curl -s localhost:8000/metrics | grep shared_cache | Symptom | Likely cause | Fix | | -------------------------------------------------------- | ---------------------------------------------------------------------------- | --------------------------------------------------------------------------------------------- | -| `shared_cache_hit_rate` is always 0 | Event endpoint missing, unreachable, or connected after objects were stored | Set `DYN_MOONCAKE_KV_EVENTS_ENDPOINT` on the worker and check the frontend subscriber log. | +| `shared_cache_hit_rate` is always 0 | Event endpoint missing, unreachable, or connected after objects were stored | Set `DYN_MOONCAKE_KV_EVENTS_ENDPOINT` on the frontend and inspect `dynamo_router_shared_cache_errors_total`. | | Events only ever carry `medium=GPU` | SGLang older than 0.5.11 or a custom build missing [PR #22894](https://github.com/sgl-project/sglang/pull/22894) | Upgrade to SGLang 0.5.11 or later. | | Workers registered but router never tracks shared cache | `--shared-cache-type` left at default `none` | Set `--shared-cache-type hicache` on the frontend. | | Shared hits do not affect worker selection | `--shared-cache-multiplier 0.0` | Raise the multiplier — typical starting range is `0.3`–`0.7`. | diff --git a/lib/llm/src/discovery/model_manager.rs b/lib/llm/src/discovery/model_manager.rs index 644cfd0f5764..ee97ecb0e2fc 100644 --- a/lib/llm/src/discovery/model_manager.rs +++ b/lib/llm/src/discovery/model_manager.rs @@ -1062,9 +1062,13 @@ impl ModelManager { self.hicache_caches .entry(endpoint.id()) .or_insert_with(|| { - let cache = HicacheSharedKvCache::new_with_cancellation( + let frontend_kv_events_endpoint = std::env::var("DYN_MOONCAKE_KV_EVENTS_ENDPOINT") + .ok() + .filter(|endpoint| !endpoint.is_empty()); + let cache = HicacheSharedKvCache::new_with_cancellation_and_endpoint( runtime_configs, endpoint.component().drt().child_token(), + frontend_kv_events_endpoint, ); cache.start_subscriber(); cache diff --git a/lib/llm/src/kv_router/metrics.rs b/lib/llm/src/kv_router/metrics.rs index fdcaf168b114..a7c225636cd1 100644 --- a/lib/llm/src/kv_router/metrics.rs +++ b/lib/llm/src/kv_router/metrics.rs @@ -716,7 +716,7 @@ impl RoutingOverheadMetrics { routing_overhead::SHARED_CACHE_ERRORS_TOTAL ); prometheus::IntCounter::with_opts( - Opts::new(name, "Total shared cache query errors") + Opts::new(name, "Total shared cache failures") .const_label(labels::ROUTER_ID, &router_id), ) .expect("shared_cache_errors_total") diff --git a/lib/llm/src/kv_router/shared_cache.rs b/lib/llm/src/kv_router/shared_cache.rs index 50f97827163d..6c33402add02 100644 --- a/lib/llm/src/kv_router/shared_cache.rs +++ b/lib/llm/src/kv_router/shared_cache.rs @@ -39,7 +39,7 @@ use dynamo_kv_router::{ const SGLANG_HICACHE_MOONCAKE_RUNTIME_KEY: &str = "sglang_hicache_mooncake"; const MOONCAKE_EVENT_RECONNECT_DELAY: Duration = Duration::from_secs(1); -#[derive(Debug, Clone, Deserialize, Serialize, PartialEq, Eq)] +#[derive(Debug, Clone, Deserialize, Serialize)] struct SglangHicacheMooncakeConfig { backend: String, page_size: u32, @@ -47,14 +47,29 @@ struct SglangHicacheMooncakeConfig { pp_size: u32, is_mla_model: bool, is_eagle: bool, + #[serde(default)] tp_lcm_size: Option, should_split_heads: bool, + #[serde(default)] extra_backend_tag: Option, - master_server_address: Option, - master_metrics_port: u16, + #[serde(default)] kv_events_endpoint: Option, } +impl SglangHicacheMooncakeConfig { + fn has_same_layout(&self, other: &Self) -> bool { + self.backend == other.backend + && self.page_size == other.page_size + && self.tp_size == other.tp_size + && self.pp_size == other.pp_size + && self.is_mla_model == other.is_mla_model + && self.is_eagle == other.is_eagle + && self.tp_lcm_size == other.tp_lcm_size + && self.should_split_heads == other.should_split_heads + && self.extra_backend_tag == other.extra_backend_tag + } +} + #[derive(Debug, Deserialize, Serialize)] struct MooncakeObjectEvent { event_type: String, @@ -85,16 +100,25 @@ pub struct HicacheSharedKvCache { last_sequence: Arc, has_sequence: Arc, cancellation_token: CancellationToken, + frontend_kv_events_endpoint: Option, } impl HicacheSharedKvCache { pub fn new(runtime_configs: RuntimeConfigWatch) -> Self { - Self::new_with_cancellation(runtime_configs, CancellationToken::new()) + Self::new_with_cancellation_and_endpoint(runtime_configs, CancellationToken::new(), None) } pub fn new_with_cancellation( runtime_configs: RuntimeConfigWatch, cancellation_token: CancellationToken, + ) -> Self { + Self::new_with_cancellation_and_endpoint(runtime_configs, cancellation_token, None) + } + + pub fn new_with_cancellation_and_endpoint( + runtime_configs: RuntimeConfigWatch, + cancellation_token: CancellationToken, + frontend_kv_events_endpoint: Option, ) -> Self { Self { runtime_configs, @@ -103,6 +127,7 @@ impl HicacheSharedKvCache { last_sequence: Arc::new(AtomicU64::new(0)), has_sequence: Arc::new(AtomicBool::new(false)), cancellation_token, + frontend_kv_events_endpoint, } } @@ -129,7 +154,10 @@ impl HicacheSharedKvCache { let (_, first) = configs.first()?; - if configs.iter().any(|(_, config)| config != first) { + if configs + .iter() + .any(|(_, config)| !config.has_same_layout(first)) + { tracing::warn!( workers = ?configs.iter().map(|(worker_id, _)| *worker_id).collect::>(), "SGLang Mooncake HiCache runtime configs differ across workers; skipping shared-cache lookup" @@ -141,8 +169,25 @@ impl HicacheSharedKvCache { } fn kv_events_endpoint(&self) -> Option { - self.resolve_mooncake_config() - .and_then(|config| config.kv_events_endpoint) + self.resolve_mooncake_config()?; + if let Some(endpoint) = &self.frontend_kv_events_endpoint { + return Some(endpoint.clone()); + } + + let workers = self.runtime_configs.borrow(); + let mut endpoints = workers.iter().filter_map(|(worker_id, runtime_config)| { + mooncake_config_from_runtime(*worker_id, runtime_config) + .and_then(|config| config.kv_events_endpoint) + .filter(|endpoint| !endpoint.is_empty()) + }); + let endpoint = endpoints.next()?; + if endpoints.any(|candidate| candidate != endpoint) { + tracing::warn!( + "SGLang Mooncake KV event endpoints differ across workers; skipping shared-cache lookup" + ); + return None; + } + Some(endpoint) } fn apply_batch(&self, sequence: u64, events: Vec) { @@ -343,7 +388,7 @@ impl SharedKvCache for HicacheSharedKvCache { return Ok(SharedCacheHits::default()); } - if config.kv_events_endpoint.is_none() { + if self.kv_events_endpoint().is_none() { tracing::debug!("Mooncake KV event endpoint is unavailable"); return Ok(SharedCacheHits::default()); } @@ -548,8 +593,6 @@ mod tests { tp_lcm_size: None, should_split_heads: false, extra_backend_tag: None, - master_server_address: Some("127.0.0.1:50051".to_string()), - master_metrics_port: 9003, kv_events_endpoint: Some("tcp://127.0.0.1:5557".to_string()), } } @@ -717,6 +760,48 @@ mod tests { ); } + #[test] + fn test_kv_events_endpoint_tolerates_worker_metadata_omission() { + let mut advertised = mooncake_config(); + advertised.kv_events_endpoint = Some("tcp://127.0.0.1:5557".to_string()); + let mut missing = advertised.clone(); + missing.kv_events_endpoint = None; + let mut advertised_runtime = ModelRuntimeConfig::new(); + advertised_runtime + .set_engine_specific(SGLANG_HICACHE_MOONCAKE_RUNTIME_KEY, advertised) + .unwrap(); + let mut missing_runtime = ModelRuntimeConfig::new(); + missing_runtime + .set_engine_specific(SGLANG_HICACHE_MOONCAKE_RUNTIME_KEY, missing) + .unwrap(); + let (_tx, runtime_configs) = watch::channel(HashMap::from([ + (1, advertised_runtime), + (2, missing_runtime), + ])); + let cache = HicacheSharedKvCache::new(runtime_configs); + + assert_eq!( + cache.kv_events_endpoint().as_deref(), + Some("tcp://127.0.0.1:5557") + ); + } + + #[test] + fn test_frontend_kv_events_endpoint_overrides_worker_metadata() { + let mut worker_config = mooncake_config(); + worker_config.kv_events_endpoint = None; + let cache = HicacheSharedKvCache::new_with_cancellation_and_endpoint( + runtime_watch_with_config(worker_config), + CancellationToken::new(), + Some("tcp://frontend-config:5557".to_string()), + ); + + assert_eq!( + cache.kv_events_endpoint().as_deref(), + Some("tcp://frontend-config:5557") + ); + } + #[tokio::test] async fn test_check_blocks_uses_mooncake_events() { let hash0 = "cf97adeedb59e05bfd73a2b4c2a8885708c4f4f70c84c64b27120e72ab733b72".to_string(); diff --git a/lib/runtime/src/metrics/prometheus_names.rs b/lib/runtime/src/metrics/prometheus_names.rs index 42148314deb0..000cda5c8400 100644 --- a/lib/runtime/src/metrics/prometheus_names.rs +++ b/lib/runtime/src/metrics/prometheus_names.rs @@ -592,7 +592,7 @@ pub mod routing_overhead { /// Time spent querying the shared KV cache (Mooncake) pub const SHARED_CACHE_QUERY_MS: &str = "overhead_shared_cache_query_ms"; - /// Total shared cache query errors (timeouts, HTTP failures) + /// Total shared cache failures (query and subscriber failures) pub const SHARED_CACHE_ERRORS_TOTAL: &str = "shared_cache_errors_total"; } From c597a2f125d722ff6a1b3a1f86eb24ab0dceff35 Mon Sep 17 00:00:00 2001 From: Ishan Dhanani Date: Tue, 28 Jul 2026 16:58:13 +0000 Subject: [PATCH 09/12] style(kv-router): simplify group verification guard Signed-off-by: Ishan Dhanani --- lib/llm/src/kv_router/shared_cache.rs | 15 ++++++++------- 1 file changed, 8 insertions(+), 7 deletions(-) diff --git a/lib/llm/src/kv_router/shared_cache.rs b/lib/llm/src/kv_router/shared_cache.rs index 6c33402add02..998c717076b6 100644 --- a/lib/llm/src/kv_router/shared_cache.rs +++ b/lib/llm/src/kv_router/shared_cache.rs @@ -410,13 +410,14 @@ impl SharedKvCache for HicacheSharedKvCache { let hit = expand_actual_query_keys(page_hash, &config) .iter() .all(|key| self.present_keys.contains(key)); - if hit && let Some((generation, _)) = generation { - if let Some(mut state) = self.group_states.get_mut(&group_id) { - // A concurrent stored event invalidates verification, not the physical - // key check that already proved this request is a hit. - if state.0 == generation { - state.1 = true; - } + if hit + && let Some((generation, _)) = generation + && let Some(mut state) = self.group_states.get_mut(&group_id) + { + // A concurrent stored event invalidates verification, not the physical key + // check that already proved this request is a hit. + if state.0 == generation { + state.1 = true; } } hit From 3db6754052cb0c8fa09a60323326f3e83be4036f Mon Sep 17 00:00:00 2001 From: Ishan Dhanani Date: Tue, 28 Jul 2026 17:40:54 +0000 Subject: [PATCH 10/12] perf(kv-router): reuse Mooncake config resolution Signed-off-by: Ishan Dhanani --- lib/llm/src/kv_router/shared_cache.rs | 35 ++++++++++++--------------- 1 file changed, 15 insertions(+), 20 deletions(-) diff --git a/lib/llm/src/kv_router/shared_cache.rs b/lib/llm/src/kv_router/shared_cache.rs index 998c717076b6..c0b9bcff89b3 100644 --- a/lib/llm/src/kv_router/shared_cache.rs +++ b/lib/llm/src/kv_router/shared_cache.rs @@ -142,7 +142,9 @@ impl HicacheSharedKvCache { self.clear(); } - fn resolve_mooncake_config(&self) -> Option { + fn resolve_mooncake_config_and_endpoint( + &self, + ) -> Option<(SglangHicacheMooncakeConfig, String)> { let workers = self.runtime_configs.borrow(); let mut configs = Vec::new(); @@ -165,21 +167,14 @@ impl HicacheSharedKvCache { return None; } - Some(first.clone()) - } - - fn kv_events_endpoint(&self) -> Option { - self.resolve_mooncake_config()?; if let Some(endpoint) = &self.frontend_kv_events_endpoint { - return Some(endpoint.clone()); + return Some((first.clone(), endpoint.clone())); } - let workers = self.runtime_configs.borrow(); - let mut endpoints = workers.iter().filter_map(|(worker_id, runtime_config)| { - mooncake_config_from_runtime(*worker_id, runtime_config) - .and_then(|config| config.kv_events_endpoint) - .filter(|endpoint| !endpoint.is_empty()) - }); + let mut endpoints = configs + .iter() + .filter_map(|(_, config)| config.kv_events_endpoint.as_deref()) + .filter(|endpoint| !endpoint.is_empty()); let endpoint = endpoints.next()?; if endpoints.any(|candidate| candidate != endpoint) { tracing::warn!( @@ -187,7 +182,12 @@ impl HicacheSharedKvCache { ); return None; } - Some(endpoint) + Some((first.clone(), endpoint.to_string())) + } + + fn kv_events_endpoint(&self) -> Option { + self.resolve_mooncake_config_and_endpoint() + .map(|(_, endpoint)| endpoint) } fn apply_batch(&self, sequence: u64, events: Vec) { @@ -360,7 +360,7 @@ impl SharedKvCache for HicacheSharedKvCache { return Ok(SharedCacheHits::default()); } - let Some(config) = self.resolve_mooncake_config() else { + let Some((config, _endpoint)) = self.resolve_mooncake_config_and_endpoint() else { tracing::debug!("No SGLang Mooncake HiCache runtime config available"); return Ok(SharedCacheHits::default()); }; @@ -388,11 +388,6 @@ impl SharedKvCache for HicacheSharedKvCache { return Ok(SharedCacheHits::default()); } - if self.kv_events_endpoint().is_none() { - tracing::debug!("Mooncake KV event endpoint is unavailable"); - return Ok(SharedCacheHits::default()); - } - let page_hashes = logical_page_hashes(tokens, config.page_size, config.is_eagle); if page_hashes.is_empty() { return Ok(SharedCacheHits::default()); From 7d782e2f69d00d44137128f35e997da4dd3f7103 Mon Sep 17 00:00:00 2001 From: Ishan Dhanani Date: Fri, 31 Jul 2026 20:25:08 +0000 Subject: [PATCH 11/12] fix(kv-router): harden Mooncake event index Signed-off-by: Ishan Dhanani --- lib/llm/src/discovery/watcher.rs | 38 ++++++-- lib/llm/src/kv_router/shared_cache.rs | 134 ++++++++++++++++++++++++++ 2 files changed, 165 insertions(+), 7 deletions(-) diff --git a/lib/llm/src/discovery/watcher.rs b/lib/llm/src/discovery/watcher.rs index 20ac889e5317..075073af74cb 100644 --- a/lib/llm/src/discovery/watcher.rs +++ b/lib/llm/src/discovery/watcher.rs @@ -105,6 +105,16 @@ fn model_card_endpoint_id(mcid: &ModelCardInstanceId) -> EndpointId { } } +fn has_live_endpoint_card( + cards: &[(EndpointId, ModelDeploymentCard)], + namespace: &str, + component: &str, +) -> bool { + cards + .iter() + .any(|(endpoint, _)| endpoint.namespace == namespace && endpoint.component == component) +} + fn model_card_instance_id(instance: &DiscoveryInstance) -> anyhow::Result { match instance { DiscoveryInstance::Model { @@ -1025,10 +1035,13 @@ impl ModelWatcher { // state. If discovery is temporarily unavailable, retaining the card // lets the next reconciliation pass retry this stale entry instead of // losing the key while leaving its WorkerSet behind. - let active_instances = self - .cards_for_model_with_endpoints(&model_name, namespace_filter) - .await - .with_context(|| model_name.clone())?; + let all_cards = self.all_cards().await.with_context(|| model_name.clone())?; + let active_instances = all_cards + .iter() + .filter(|(endpoint_id, card)| { + card.name() == model_name && namespace_filter.matches(&endpoint_id.namespace) + }) + .collect::>(); let card = match self.manager.remove_model_card(&key) { Some(card) => card, @@ -1090,9 +1103,8 @@ impl ModelWatcher { && eid.component == *worker_component && worker_set_key(eid, other_card.model_type, other_card.worker_type) == ws_key }); - let endpoint_has_instances = active_instances.iter().any(|(eid, _)| { - eid.namespace == *worker_namespace && eid.component == *worker_component - }); + let endpoint_has_instances = + has_live_endpoint_card(&all_cards, worker_namespace, worker_component); if !component_has_instances { // No more workers of this component in this namespace — remove its WorkerSet @@ -2193,6 +2205,18 @@ mod tests { } } + #[test] + fn endpoint_liveness_considers_all_models() { + // This is the discovery snapshot after the last adapter card was removed: its base + // model remains on the same worker endpoint and must keep its HiCache subscriber. + let cards = vec![( + test_endpoint_id("generate"), + ModelDeploymentCard::with_name_only("base-model"), + )]; + + assert!(has_live_endpoint_card(&cards, "ns1", "workers")); + } + #[test] fn vllm_generate_requires_explicit_worker_capability() { let mut card = ModelDeploymentCard::with_name_only("model"); diff --git a/lib/llm/src/kv_router/shared_cache.rs b/lib/llm/src/kv_router/shared_cache.rs index c0b9bcff89b3..b9c4ceb00e2d 100644 --- a/lib/llm/src/kv_router/shared_cache.rs +++ b/lib/llm/src/kv_router/shared_cache.rs @@ -17,6 +17,7 @@ use std::sync::{ }; use std::time::Duration; +use arc_swap::ArcSwapOption; use async_trait::async_trait; use dashmap::{DashMap, DashSet}; use futures::StreamExt; @@ -38,6 +39,7 @@ use dynamo_kv_router::{ const SGLANG_HICACHE_MOONCAKE_RUNTIME_KEY: &str = "sglang_hicache_mooncake"; const MOONCAKE_EVENT_RECONNECT_DELAY: Duration = Duration::from_secs(1); +const MAX_MOONCAKE_INDEX_ENTRIES: usize = 1_000_000; #[derive(Debug, Clone, Deserialize, Serialize)] struct SglangHicacheMooncakeConfig { @@ -99,6 +101,7 @@ pub struct HicacheSharedKvCache { group_states: Arc>, last_sequence: Arc, has_sequence: Arc, + last_layout: Arc>, cancellation_token: CancellationToken, frontend_kv_events_endpoint: Option, } @@ -126,6 +129,7 @@ impl HicacheSharedKvCache { group_states: Arc::new(DashMap::new()), last_sequence: Arc::new(AtomicU64::new(0)), has_sequence: Arc::new(AtomicBool::new(false)), + last_layout: Arc::new(ArcSwapOption::empty()), cancellation_token, frontend_kv_events_endpoint, } @@ -142,6 +146,21 @@ impl HicacheSharedKvCache { self.clear(); } + fn clear_on_layout_change(&self, layout: &SglangHicacheMooncakeConfig) { + let last_layout = self.last_layout.load(); + if last_layout + .as_ref() + .is_some_and(|previous| previous.has_same_layout(layout)) + { + return; + } + if last_layout.is_some() { + self.clear(); + tracing::warn!("SGLang Mooncake HiCache layout changed; cleared shared-cache state"); + } + self.last_layout.store(Some(Arc::new(layout.clone()))); + } + fn resolve_mooncake_config_and_endpoint( &self, ) -> Option<(SglangHicacheMooncakeConfig, String)> { @@ -167,6 +186,8 @@ impl HicacheSharedKvCache { return None; } + self.clear_on_layout_change(first); + if let Some(endpoint) = &self.frontend_kv_events_endpoint { return Some((first.clone(), endpoint.clone())); } @@ -191,6 +212,8 @@ impl HicacheSharedKvCache { } fn apply_batch(&self, sequence: u64, events: Vec) { + // SGLang's ZmqEventPublisher increments this sequence once per published batch, so a + // non-consecutive value means one or more whole batches were missed. let has_previous = self.has_sequence.swap(true, Ordering::AcqRel); let previous = self.last_sequence.swap(sequence, Ordering::AcqRel); if has_previous && sequence == previous { @@ -236,6 +259,8 @@ impl HicacheSharedKvCache { _ => {} } } + + self.clear_if_index_too_large(MAX_MOONCAKE_INDEX_ENTRIES); } fn clear(&self) { @@ -245,6 +270,20 @@ impl HicacheSharedKvCache { self.has_sequence.store(false, Ordering::Release); } + fn clear_if_index_too_large(&self, max_entries: usize) { + let present_keys = self.present_keys.len(); + let group_states = self.group_states.len(); + if present_keys.saturating_add(group_states) > max_entries { + self.clear(); + tracing::warn!( + present_keys, + group_states, + max_entries, + "Mooncake KV event index exceeded its size limit; cleared shared-cache state" + ); + } + } + fn record_subscriber_error(&self) { if let Some(metrics) = RoutingOverheadMetrics::get() { metrics.inc_shared_cache_errors(); @@ -262,6 +301,7 @@ impl HicacheSharedKvCache { _ = cancellation_token.cancelled() => return, result = self.runtime_configs.changed() => { if result.is_err() { + self.clear(); return; } } @@ -287,6 +327,7 @@ impl HicacheSharedKvCache { _ = cancellation_token.cancelled() => return, result = self.runtime_configs.changed() => { if result.is_err() { + self.clear(); return; } let next_endpoint = self.kv_events_endpoint(); @@ -798,6 +839,56 @@ mod tests { ); } + #[tokio::test] + async fn test_layout_change_clears_cached_hits() { + let config = mooncake_config(); + let (runtime_configs, tx) = runtime_watch_with_config_and_sender(config.clone()); + let cache = HicacheSharedKvCache::new(runtime_configs); + let hash = logical_page_hashes(&[1, 2, 3, 4], config.page_size, config.is_eagle) + .pop() + .unwrap(); + let group_id = sglang_group_id(&hash, &config); + cache.apply_batch( + 1, + expand_actual_query_keys(&hash, &config) + .into_iter() + .map(|object_key| MooncakeObjectEvent { + event_type: "stored".to_string(), + object_key: Some(object_key), + tenant_id: "default".to_string(), + group_id: Some(group_id.clone()), + }) + .collect(), + ); + assert_eq!( + cache + .check_blocks(&[1, 2, 3, 4], config.page_size, None) + .await + .unwrap() + .total_hits, + 1 + ); + + let mut new_layout = config; + new_layout.tp_size = 2; + let mut runtime_config = ModelRuntimeConfig::new(); + runtime_config + .set_engine_specific(SGLANG_HICACHE_MOONCAKE_RUNTIME_KEY, new_layout) + .unwrap(); + tx.send(HashMap::from([(1, runtime_config)])).unwrap(); + + assert_eq!( + cache + .check_blocks(&[1, 2, 3, 4], 4, None) + .await + .unwrap() + .total_hits, + 0 + ); + assert!(cache.present_keys.is_empty()); + assert!(cache.group_states.is_empty()); + } + #[tokio::test] async fn test_check_blocks_uses_mooncake_events() { let hash0 = "cf97adeedb59e05bfd73a2b4c2a8885708c4f4f70c84c64b27120e72ab733b72".to_string(); @@ -960,6 +1051,25 @@ mod tests { assert!(cache.present_keys.contains("key-0")); } + #[test] + fn test_index_size_limit_clears_shared_cache_state() { + let cache = HicacheSharedKvCache::new(runtime_watch_with_config(mooncake_config())); + cache.apply_batch( + 1, + vec![MooncakeObjectEvent { + event_type: "stored".to_string(), + object_key: Some("key-0".to_string()), + tenant_id: "default".to_string(), + group_id: Some("group-0".to_string()), + }], + ); + + cache.clear_if_index_too_large(1); + + assert!(cache.present_keys.is_empty()); + assert!(cache.group_states.is_empty()); + } + #[tokio::test] async fn test_subscriber_retries_failed_connection_until_cancelled() { let mut config = mooncake_config(); @@ -978,6 +1088,30 @@ mod tests { .unwrap(); } + #[tokio::test] + async fn test_subscriber_clears_state_when_runtime_config_watch_closes() { + let (tx, runtime_configs) = watch::channel(HashMap::::new()); + let cache = HicacheSharedKvCache::new(runtime_configs); + cache.apply_batch( + 1, + vec![MooncakeObjectEvent { + event_type: "stored".to_string(), + object_key: Some("key-0".to_string()), + tenant_id: "default".to_string(), + group_id: None, + }], + ); + let task = tokio::spawn(cache.clone().run_subscriber(CancellationToken::new())); + + drop(tx); + tokio::time::timeout(Duration::from_secs(1), task) + .await + .unwrap() + .unwrap(); + + assert!(cache.present_keys.is_empty()); + } + #[test] fn test_sequence_gap_clears_stale_keys() { let cache = HicacheSharedKvCache::new(runtime_watch_with_config(mooncake_config())); From 73a58c4c23bff08f016f901445b3befd1afa40fe Mon Sep 17 00:00:00 2001 From: Ishan Dhanani Date: Fri, 31 Jul 2026 20:31:01 +0000 Subject: [PATCH 12/12] style(discovery): keep watcher liveness generic Signed-off-by: Ishan Dhanani --- lib/llm/src/discovery/watcher.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/lib/llm/src/discovery/watcher.rs b/lib/llm/src/discovery/watcher.rs index 075073af74cb..dcfd364bab93 100644 --- a/lib/llm/src/discovery/watcher.rs +++ b/lib/llm/src/discovery/watcher.rs @@ -2208,7 +2208,7 @@ mod tests { #[test] fn endpoint_liveness_considers_all_models() { // This is the discovery snapshot after the last adapter card was removed: its base - // model remains on the same worker endpoint and must keep its HiCache subscriber. + // model remains on the same worker endpoint, so the endpoint is still live. let cards = vec![( test_endpoint_id("generate"), ModelDeploymentCard::with_name_only("base-model"),