Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions experimental/sgl-router/monitoring/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,8 @@ Families the router emits. The dashboard graphs all of them except the
| `sgl_router_kv_event_batches_lost_total` | Counter | KV-event batches dropped in transit, from gaps in each publisher's sequence |
| `sgl_router_kv_tree_accounting_errors_total` | Counter | Occupancy-bookkeeping contradictions, by `reason`. Always 0 on a correct tree |
| `sgl_router_kv_tree_maintained` | Gauge | 1 when this router maintains its own KV tree, 0 under an external Indexer |
| `sgl_router_kv_bootstrap_peers` | Gauge | Ready sibling router replicas peer bootstrap could pull a tree snapshot from. Emitted only with `--kv-peer-selector` |
| `sgl_router_kv_bootstrap_peers_synced` | Gauge | 1 once peer discovery has reported at least once; a 0 that never becomes 1 means the peer watch is not delivering (check EndpointSlice RBAC). Emitted only with `--kv-peer-selector` |

The legacy `sgl_router_overlap_blocks` metric was removed with the
`cache_aware_zmq` policy and has no direct replacement. Remove queries, alerts,
Expand Down
201 changes: 199 additions & 2 deletions experimental/sgl-router/src/config/cli.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,11 +3,12 @@

//! Grouped CLI options and conversion into a validated [`Config`].

use anyhow::{anyhow, ensure, Context, Result};
use anyhow::{anyhow, bail, ensure, Context, Result};
use clap::Parser;
use std::num::NonZeroU32;

use crate::config::sampling::{parse_sampling_overrides, ConflictPolicy};
use crate::config::types::is_selector_empty;
use crate::config::{
default_cb_cool_down, default_host, default_port, default_proxy_request_timeout_secs,
default_shutdown_drain_secs, default_stale_request_timeout_secs, resolve_mode, AffinityConfig,
Expand All @@ -16,6 +17,7 @@ use crate::config::{
InflightLoadConfig, K8sDiscoveryConfig, KvIndexerEndpointConfig, LogFormat, ModelConfig,
ObservabilityConfig, PolicyKind, ProxyConfig, ServerConfig, SessionAffinityMode,
StaticUrlsDiscoveryConfig, StickyConfig, StickyFallbackKind, DEFAULT_FUSE,
DEFAULT_KV_BOOTSTRAP_TIMEOUT_MS,
};

const DEFAULT_KV_INDEXER_QUERY_TIMEOUT_MS: u64 = 100;
Expand Down Expand Up @@ -164,6 +166,14 @@ pub struct DiscoveryArgs {
/// Decode equality selector terms (key=value or key==value). Requires --prefill-selector.
#[arg(long, num_args = 1..)]
pub decode_selector: Vec<String>,

/// Label selector matching the EndpointSlices of this router's OWN Service
/// (the Service's labels, not the pods'), so a booting replica can find
/// siblings to pull a cache-aware tree snapshot from, e.g.
/// `kubernetes.io/service-name=sgl-router`. Unset disables peer bootstrap
/// and every replica starts cold.
#[arg(long)]
pub kv_peer_selector: Option<String>,
}

#[derive(clap::Args, Debug)]
Expand Down Expand Up @@ -238,6 +248,11 @@ pub struct CacheArgs {
#[arg(long)]
pub kv_indexer_query_max_inflight: Option<usize>,

/// Peer-bootstrap deadline in milliseconds. Requires --kv-peer-selector.
/// Defaults to 600000 (10 minutes).
#[arg(long)]
pub kv_bootstrap_timeout_ms: Option<u64>,

/// Minimum cache-hit tokens for a candidate. Defaults to 1024.
#[arg(long)]
pub cache_affinity_min_matched_tokens: Option<u64>,
Expand Down Expand Up @@ -371,7 +386,27 @@ impl Cli {
.map(load_bucket_config)
.transpose()?;
let circuit_breaker = self.routing.build_circuit_breaker()?;
let kv_bootstrap_timeout_ms = self.cache.kv_bootstrap_timeout_ms;
let cache_aware = self.cache.into_config(self.routing.policy)?;
// Peer bootstrap grafts into this router's own radix tree, so it needs
// cache-aware over a local tree; checked here because it spans groups.
let has_peer_selector = discovery.peer_selector().is_some();
ensure!(
has_peer_selector || kv_bootstrap_timeout_ms.is_none(),
"--kv-bootstrap-timeout-ms requires --kv-peer-selector, which is what \
enables peer bootstrap"
);
if has_peer_selector {
match cache_aware.as_ref().map(|c| c.prefix_provider) {
None => bail!("--kv-peer-selector requires --policy cache_aware"),
Some(p) if p != CachePrefixProvider::RadixTree => bail!(
"--kv-peer-selector requires --cache-prefix-provider radix_tree \
(there is no local tree to bootstrap when an external Indexer is \
the prefix source)"
),
_ => {}
}
}
let fused = self.routing.build_fused()?;
let eligibility = self.routing.build_eligibility()?;
let sticky = self.affinity.into_sticky_config(self.routing.policy)?;
Expand Down Expand Up @@ -467,6 +502,14 @@ impl DiscoveryArgs {
"--service-discovery-namespace / --selector / --prefill-selector / \
--decode-selector require --service-discovery"
);
// Peer replicas are found through EndpointSlices too, so with
// any other backend the flag would be accepted and then
// silently ignored.
ensure!(
self.kv_peer_selector.is_none(),
"--kv-peer-selector requires --service-discovery (peer replicas are \
found via Kubernetes EndpointSlices)"
);
DiscoveryBackend::StaticUrls(StaticUrlsDiscoveryConfig {
urls: self.worker_urls,
})
Expand All @@ -477,9 +520,19 @@ impl DiscoveryArgs {
join_selector(&self.prefill_selector).as_deref(),
join_selector(&self.decode_selector).as_deref(),
)?;
// An empty selector matches every EndpointSlice in the watched
// namespace(s), turning unrelated Services into "siblings".
ensure!(
self.kv_peer_selector
.as_deref()
.is_none_or(|s| !is_selector_empty(s)),
"--kv-peer-selector must not be empty: an empty selector matches \
every EndpointSlice in the watched namespace"
);
DiscoveryBackend::K8s(K8sDiscoveryConfig {
namespace: self.service_discovery_namespace.unwrap_or_default(),
mode,
peer_selector: self.kv_peer_selector,
})
}
};
Expand Down Expand Up @@ -630,6 +683,9 @@ impl CacheArgs {
Ok(Some(CacheAwareConfig {
prefix_provider: cache_prefix_provider,
kv_indexer_endpoint,
bootstrap_timeout_ms: self
.kv_bootstrap_timeout_ms
.unwrap_or(DEFAULT_KV_BOOTSTRAP_TIMEOUT_MS),
}))
}
}
Expand Down Expand Up @@ -860,7 +916,9 @@ fn join_selector(terms: &[String]) -> Option<String> {
#[cfg(test)]
mod tests {
use super::*;
use crate::config::{DiscoveryBackend, K8sDiscoveryMode, ScoreTermKind};
use crate::config::{
DiscoveryBackend, K8sDiscoveryMode, ScoreTermKind, MAX_KV_BOOTSTRAP_TIMEOUT_MS,
};

#[test]
fn chat_routing_selects_policy_implementation_without_another_config() {
Expand Down Expand Up @@ -1166,6 +1224,145 @@ mod tests {
assert!(err.contains("require --service-discovery"), "got: {err}");
}

// ---- peer bootstrap ----

/// k8s discovery + cache-aware, the base every peer-bootstrap config needs.
fn peer_cfg(extra: &[&str]) -> Result<Config> {
let mut args = vec![
"--service-discovery",
"--selector",
"app=sglang",
"--policy",
"cache_aware",
];
args.extend_from_slice(extra);
into_config_owned(with_model(&args))
}

#[test]
fn peer_selector_rides_on_the_k8s_backend() {
let c = peer_cfg(&[
"--kv-peer-selector",
"app=sgl-router",
"--kv-bootstrap-timeout-ms",
"12000",
])
.unwrap();
match &c.discovery {
DiscoveryBackend::K8s(k) => {
assert_eq!(k.peer_selector.as_deref(), Some("app=sgl-router"))
}
_ => panic!("expected k8s backend"),
}
assert_eq!(
c.model.cache_aware.as_ref().unwrap().bootstrap_timeout_ms,
12_000,
);
}

#[test]
fn peer_bootstrap_is_off_by_default() {
let c = peer_cfg(&[]).unwrap();
match &c.discovery {
DiscoveryBackend::K8s(k) => assert!(k.peer_selector.is_none()),
_ => panic!("expected k8s backend"),
}
assert_eq!(
c.model.cache_aware.as_ref().unwrap().bootstrap_timeout_ms,
DEFAULT_KV_BOOTSTRAP_TIMEOUT_MS,
);
}

/// Peers are found through EndpointSlices, so without k8s discovery the
/// flag would be accepted and then silently ignored.
#[test]
fn rejects_peer_selector_without_service_discovery() {
let err = into_config_owned(with_model(&[
"--worker-urls",
"http://w:30000",
"--policy",
"cache_aware",
"--kv-peer-selector",
"app=sgl-router",
]))
.unwrap_err()
.to_string();
assert!(err.contains("requires --service-discovery"), "got: {err}");
}

#[test]
fn rejects_peer_selector_without_cache_aware_policy() {
let err = into_config_owned(with_model(&[
"--service-discovery",
"--selector",
"app=sglang",
"--kv-peer-selector",
"app=sgl-router",
]))
.unwrap_err()
.to_string();
assert!(err.contains("--policy cache_aware"), "got: {err}");
}

/// With an external Indexer there is no local tree to graft into, so the
/// flag would be accepted and then do nothing.
#[test]
fn rejects_peer_selector_with_an_external_indexer() {
let err = peer_cfg(&[
"--cache-prefix-provider",
"indexer",
"--kv-indexer-endpoint",
"http://indexer:50051",
"--kv-peer-selector",
"app=sgl-router",
])
.unwrap_err()
.to_string();
assert!(
err.contains("--cache-prefix-provider radix_tree"),
"got: {err}",
);
}

#[test]
fn kv_bootstrap_timeout_bounds() {
let parse = |ms: u64| {
peer_cfg(&[
"--kv-peer-selector",
"app=sgl-router",
"--kv-bootstrap-timeout-ms",
&ms.to_string(),
])
};
assert!(parse(MAX_KV_BOOTSTRAP_TIMEOUT_MS).is_ok());
let err = parse(MAX_KV_BOOTSTRAP_TIMEOUT_MS + 1)
.unwrap_err()
.to_string();
assert!(err.contains("ceiling"), "got: {err}");
let err = parse(0).unwrap_err().to_string();
assert!(err.contains("greater than zero"), "got: {err}");
}

/// Without a selector the timeout tunes a bootstrap that never runs.
#[test]
fn rejects_bootstrap_timeout_without_a_peer_selector() {
let err = peer_cfg(&["--kv-bootstrap-timeout-ms", "20000"])
.unwrap_err()
.to_string();
assert!(err.contains("requires --kv-peer-selector"), "got: {err}");
}

/// An empty selector matches every EndpointSlice in the namespace.
#[test]
fn rejects_an_empty_peer_selector() {
for selector in ["", " , "] {
let err = peer_cfg(&["--kv-peer-selector", selector])
.unwrap_err()
.to_string();
assert!(err.contains("must not be empty"), "{selector:?} got: {err}");
}
}

#[test]
fn rejects_static_urls_duplicate() {
let err = into_config_owned(with_model(&[
Expand Down
12 changes: 12 additions & 0 deletions experimental/sgl-router/src/config/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -60,6 +60,18 @@ impl Config {
can check it against the pod's real budget",
self.server.shutdown_drain_secs,
);
if let Some(cache) = self.model.cache_aware.as_ref() {
let ms = cache.bootstrap_timeout_ms;
ensure!(
ms > 0,
"--kv-bootstrap-timeout-ms must be greater than zero; to boot cold, \
leave --kv-peer-selector unset"
);
ensure!(
ms <= MAX_KV_BOOTSTRAP_TIMEOUT_MS,
"--kv-bootstrap-timeout-ms {ms} exceeds the {MAX_KV_BOOTSTRAP_TIMEOUT_MS}ms ceiling"
);
}
match &self.discovery {
DiscoveryBackend::StaticUrls(s) => {
ensure!(
Expand Down
57 changes: 55 additions & 2 deletions experimental/sgl-router/src/config/types.rs
Original file line number Diff line number Diff line change
Expand Up @@ -401,14 +401,38 @@ pub enum CachePrefixProvider {
}

/// Per-model Cache-Aware configuration.
#[derive(Debug, Clone, Default)]
#[derive(Debug, Clone)]
pub struct CacheAwareConfig {
/// Prefix-match source for native Cache-Aware.
pub prefix_provider: CachePrefixProvider,
/// External Indexer configuration when `prefix_provider = indexer`.
pub kv_indexer_endpoint: Option<KvIndexerEndpointConfig>,
/// Deadline for peer bootstrap, validated by `Config::validate`. Only
/// meaningful when a peer selector is set; see
/// [`K8sDiscoveryConfig::peer_selector`].
pub bootstrap_timeout_ms: u64,
}

impl Default for CacheAwareConfig {
fn default() -> Self {
Self {
prefix_provider: CachePrefixProvider::default(),
kv_indexer_endpoint: None,
bootstrap_timeout_ms: DEFAULT_KV_BOOTSTRAP_TIMEOUT_MS,
}
}
}

/// Budget for the whole peer-bootstrap sweep: 10 minutes, room for several
/// attempts at a fleet-sized snapshot (the producer's export build, the
/// transfer and the graft). Readiness waits on it, so a pod's startup or
/// readiness probe must tolerate a replica that stays unready this long.
pub const DEFAULT_KV_BOOTSTRAP_TIMEOUT_MS: u64 = 600_000;

/// Ceiling on `--kv-bootstrap-timeout-ms` (1 hour); past it, `Instant +
/// Duration` can overflow and panic.
pub const MAX_KV_BOOTSTRAP_TIMEOUT_MS: u64 = 3_600_000;

/// Default request header for sticky routing.
pub const DEFAULT_STICKY_HEADER: &str = "x-sgl-routing-key";

Expand Down Expand Up @@ -568,6 +592,17 @@ pub enum DiscoveryBackend {
K8s(K8sDiscoveryConfig),
}

impl DiscoveryBackend {
/// The peer-bootstrap selector, which only the Kubernetes backend carries.
/// `Some` is what enables peer bootstrap.
pub fn peer_selector(&self) -> Option<&str> {
match self {
Self::K8s(k) => k.peer_selector.as_deref(),
Self::StaticUrls(_) => None,
}
}
}

/// Workers registered at startup. Roles, models, and bootstrap ports come from
/// `/server_info`; topology changes require a restart.
#[derive(Debug, Clone)]
Expand All @@ -582,6 +617,24 @@ pub struct K8sDiscoveryConfig {
pub namespace: String,
/// Resolved + validated selector mode (plain vs PD).
pub mode: K8sDiscoveryMode,
/// Label selector matching the EndpointSlices of this router's OWN
/// Service, so a booting replica can find sibling replicas to pull a
/// cache-aware tree snapshot from.
///
/// Matched against EndpointSlice labels, exactly like the worker
/// selectors: those are the Service's labels plus
/// `kubernetes.io/service-name`, never the pods' labels. A pod-template
/// label matches nothing, which reads as "no siblings" and boots every
/// replica cold; `kubernetes.io/service-name=<router-service>` is the
/// unambiguous choice.
///
/// `None` disables peer bootstrap. Scoped to the same namespace as
/// `namespace` (so an empty `namespace` means every namespace), which is
/// why this lives here rather than on [`crate::config::CacheAwareConfig`].
///
/// Requires the router's ServiceAccount to have `list`/`watch` on
/// EndpointSlices in that namespace.
pub peer_selector: Option<String>,
}

/// Validated selector mode. Plain selectors run server-side; PD selectors
Expand Down Expand Up @@ -639,7 +692,7 @@ pub enum ConfigError {
}

/// An empty selector matches every slice, which is invalid for a PD role.
fn is_selector_empty(selector: &str) -> bool {
pub(crate) fn is_selector_empty(selector: &str) -> bool {
selector.split(',').all(|t| t.trim().is_empty())
}

Expand Down
Loading
Loading