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
13 changes: 13 additions & 0 deletions bindings/python/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -516,6 +516,8 @@ struct Router {
capacity_credit_terminal_retention_secs: u64,
capacity_credit_required: bool,
priority_scheduler_adaptive_capacity: bool,
adaptive_admission_distribution_headroom_partitions: Vec<String>,
adaptive_admission_distribution_headroom_partition_seed_cap: u32,
}

impl Router {
Expand Down Expand Up @@ -831,6 +833,11 @@ impl Router {
feedback_max_token_usage: self.adaptive_admission_feedback_max_token_usage,
feedback_throughput_improvement_ratio: self
.adaptive_admission_feedback_throughput_improvement_ratio,
distribution_headroom_partitions: self
.adaptive_admission_distribution_headroom_partitions
.clone(),
distribution_headroom_partition_seed_cap: self
.adaptive_admission_distribution_headroom_partition_seed_cap,
})
.cors_allowed_origins(self.cors_allowed_origins.clone())
.retry_config(config::RetryConfig {
Expand Down Expand Up @@ -1070,6 +1077,8 @@ impl Router {
capacity_credit_terminal_retention_secs = 600,
capacity_credit_required = false,
priority_scheduler_adaptive_capacity = false,
adaptive_admission_distribution_headroom_partitions = vec![],
adaptive_admission_distribution_headroom_partition_seed_cap = 0,
))]
#[expect(clippy::too_many_arguments)]
#[expect(
Expand Down Expand Up @@ -1225,6 +1234,8 @@ impl Router {
capacity_credit_terminal_retention_secs: u64,
capacity_credit_required: bool,
priority_scheduler_adaptive_capacity: bool,
adaptive_admission_distribution_headroom_partitions: Vec<String>,
adaptive_admission_distribution_headroom_partition_seed_cap: u32,
) -> PyResult<Self> {
let mut all_urls = worker_urls.clone();

Expand Down Expand Up @@ -1394,6 +1405,8 @@ impl Router {
capacity_credit_terminal_retention_secs,
capacity_credit_required,
priority_scheduler_adaptive_capacity,
adaptive_admission_distribution_headroom_partitions,
adaptive_admission_distribution_headroom_partition_seed_cap,
})
}

Expand Down
17 changes: 17 additions & 0 deletions bindings/python/src/smg/router_args.py
Original file line number Diff line number Diff line change
Expand Up @@ -143,6 +143,10 @@ class RouterArgs:
adaptive_admission_feedback_max_waiting_requests_per_healthy_replica: int = 2
adaptive_admission_feedback_max_token_usage: float = 0.9
adaptive_admission_feedback_throughput_improvement_ratio: float = 0.02
adaptive_admission_distribution_headroom_partitions: list[str] = dataclasses.field(
default_factory=list
)
adaptive_admission_distribution_headroom_partition_seed_cap: int = 0
# Token bucket refill rate (tokens per second). If not set, defaults to max_concurrent_requests
rate_limit_tokens_per_second: int | None = None
# Cluster-wide requests-per-second ceiling. Requires mesh and the same value on every gateway.
Expand Down Expand Up @@ -1022,6 +1026,19 @@ def add_cli_args(
default=RouterArgs.adaptive_admission_feedback_throughput_improvement_ratio,
help="Relative throughput gain required to raise the learned concurrency knee",
)
adaptive_admission_group.add_argument(
f"--{prefix}adaptive-admission-distribution-headroom-partitions",
type=str,
nargs="+",
default=[],
help="Exact admission partitions allowed to expose clean-worker headroom",
)
adaptive_admission_group.add_argument(
f"--{prefix}adaptive-admission-distribution-headroom-partition-seed-cap",
type=int,
default=RouterArgs.adaptive_admission_distribution_headroom_partition_seed_cap,
help="Hard active distribution seed-lease cap per allowlisted partition",
)

# Retry configuration
retry_group.add_argument(
Expand Down
12 changes: 12 additions & 0 deletions bindings/python/tests/test_arg_parser.py
Original file line number Diff line number Diff line change
Expand Up @@ -72,6 +72,8 @@ def test_default_values(self):
assert args.adaptive_admission_feedback_max_waiting_requests_per_healthy_replica == 2
assert args.adaptive_admission_feedback_max_token_usage == 0.9
assert args.adaptive_admission_feedback_throughput_improvement_ratio == 0.02
assert args.adaptive_admission_distribution_headroom_partitions == []
assert args.adaptive_admission_distribution_headroom_partition_seed_cap == 0

def test_parse_priority_scheduler_options(self):
args = parse_router_args(
Expand Down Expand Up @@ -146,6 +148,11 @@ def test_parse_adaptive_admission_options(self):
"0.85",
"--adaptive-admission-feedback-throughput-improvement-ratio",
"0.03",
"--adaptive-admission-distribution-headroom-partitions",
"k3-prod",
"k3-canary",
"--adaptive-admission-distribution-headroom-partition-seed-cap",
"2",
]
)

Expand All @@ -162,6 +169,11 @@ def test_parse_adaptive_admission_options(self):
assert args.adaptive_admission_feedback_max_waiting_requests_per_healthy_replica == 4
assert args.adaptive_admission_feedback_max_token_usage == 0.85
assert args.adaptive_admission_feedback_throughput_improvement_ratio == 0.03
assert args.adaptive_admission_distribution_headroom_partitions == [
"k3-prod",
"k3-canary",
]
assert args.adaptive_admission_distribution_headroom_partition_seed_cap == 2

def test_parse_selector_valid(self):
"""Test parsing valid selector arguments."""
Expand Down
10 changes: 10 additions & 0 deletions model_gateway/src/config/types.rs
Original file line number Diff line number Diff line change
Expand Up @@ -150,6 +150,14 @@ pub struct AdaptiveAdmissionConfig {
/// concurrency knee upward. Near-equal throughput may move it downward.
#[serde(default = "default_feedback_throughput_improvement_ratio")]
pub feedback_throughput_improvement_ratio: f64,
/// Exact admission-partition names allowed to expose clean-worker
/// distribution headroom. An empty allowlist keeps the feature disabled.
#[serde(default)]
pub distribution_headroom_partitions: Vec<String>,
/// Hard process-local ceiling on active distribution-headroom seed leases
/// in any allowlisted partition. Zero keeps the feature disabled.
#[serde(default)]
pub distribution_headroom_partition_seed_cap: u32,
}

impl Default for AdaptiveAdmissionConfig {
Expand All @@ -169,6 +177,8 @@ impl Default for AdaptiveAdmissionConfig {
default_feedback_max_waiting_requests_per_healthy_replica(),
feedback_max_token_usage: default_feedback_max_token_usage(),
feedback_throughput_improvement_ratio: default_feedback_throughput_improvement_ratio(),
distribution_headroom_partitions: Vec::new(),
distribution_headroom_partition_seed_cap: 0,
}
}
}
Expand Down
106 changes: 106 additions & 0 deletions model_gateway/src/config/validation.rs
Original file line number Diff line number Diff line change
Expand Up @@ -772,6 +772,17 @@ impl ConfigValidator {
.to_string(),
});
}
if config.priority_scheduler_adaptive_capacity
&& config
.adaptive_admission
.distribution_headroom_partition_seed_cap
> 0
{
return Err(ConfigError::ValidationFailed {
reason: "distribution headroom requires a static scheduler ceiling; adaptive scheduler capacity may collapse before target selection"
.to_string(),
});
}

if config.worker_startup_timeout_secs == 0 {
return Err(ConfigError::InvalidValue {
Expand Down Expand Up @@ -865,6 +876,51 @@ impl ConfigValidator {
reason: "Must be finite and in [0, 1]".to_string(),
});
}
let mut distribution_partitions = std::collections::HashSet::new();
for partition in &adaptive.distribution_headroom_partitions {
if partition.is_empty() || partition.trim() != partition || partition.len() > 128 {
return Err(ConfigError::InvalidValue {
field: "adaptive_admission.distribution_headroom_partitions".to_string(),
value: partition.clone(),
reason: "Partition names must be non-empty, unpadded, and at most 128 bytes"
.to_string(),
});
}
if !distribution_partitions.insert(partition) {
return Err(ConfigError::InvalidValue {
field: "adaptive_admission.distribution_headroom_partitions".to_string(),
value: partition.clone(),
reason: "Partition names must be unique exact matches".to_string(),
});
}
}
if adaptive.distribution_headroom_partition_seed_cap > 0 {
if adaptive.distribution_headroom_partition_seed_cap != 1 {
return Err(ConfigError::InvalidValue {
field: "adaptive_admission.distribution_headroom_partition_seed_cap"
.to_string(),
value: adaptive
.distribution_headroom_partition_seed_cap
.to_string(),
reason: "Only a single in-flight seed per partition is currently supported"
.to_string(),
});
}
if adaptive.distribution_headroom_partitions.is_empty() {
return Err(ConfigError::ValidationFailed {
reason: "distribution headroom seed capacity requires a non-empty exact partition allowlist"
.to_string(),
});
}
if adaptive.mode != AdaptiveAdmissionMode::Enforce
|| adaptive.strategy != AdaptiveAdmissionStrategy::EngineFeedback
{
return Err(ConfigError::ValidationFailed {
reason: "distribution headroom seed capacity requires adaptive admission mode=enforce and strategy=engine_feedback"
.to_string(),
});
}
}

Ok(())
}
Expand Down Expand Up @@ -2045,9 +2101,59 @@ mod tests {
assert!(config.capacity_credit_generation.is_none());
assert!(!config.capacity_credit_required);
assert!(!config.priority_scheduler_adaptive_capacity);
assert!(config
.adaptive_admission
.distribution_headroom_partitions
.is_empty());
assert_eq!(
config
.adaptive_admission
.distribution_headroom_partition_seed_cap,
0
);
assert!(ConfigValidator::validate(&config).is_ok());
}

#[test]
fn distribution_headroom_requires_exact_allowlist_and_enforced_feedback() {
let mut config = RouterConfig::default();
config.adaptive_admission.distribution_headroom_partitions = vec!["k3".to_string()];
assert!(ConfigValidator::validate(&config).is_ok());

config
.adaptive_admission
.distribution_headroom_partition_seed_cap = 1;
assert!(ConfigValidator::validate(&config).is_err());
config.adaptive_admission.mode = AdaptiveAdmissionMode::Enforce;
config.adaptive_admission.strategy = AdaptiveAdmissionStrategy::EngineFeedback;
assert!(ConfigValidator::validate(&config).is_ok());

config
.adaptive_admission
.distribution_headroom_partition_seed_cap = 2;
assert!(ConfigValidator::validate(&config).is_err());
config
.adaptive_admission
.distribution_headroom_partition_seed_cap = 1;

config.priority_scheduler_enabled = true;
config.priority_scheduler_adaptive_capacity = true;
assert!(ConfigValidator::validate(&config).is_err());
config.priority_scheduler_adaptive_capacity = false;
assert!(ConfigValidator::validate(&config).is_ok());

config.adaptive_admission.distribution_headroom_partitions = vec![" k3".to_string()];
assert!(ConfigValidator::validate(&config).is_err());
config.adaptive_admission.distribution_headroom_partitions =
vec!["k3".to_string(), "k3".to_string()];
assert!(ConfigValidator::validate(&config).is_err());
config
.adaptive_admission
.distribution_headroom_partitions
.clear();
assert!(ConfigValidator::validate(&config).is_err());
}

#[test]
fn capacity_credit_generation_requires_scheduler_service_key_and_preferred_identity() {
let mut config = RouterConfig {
Expand Down
33 changes: 31 additions & 2 deletions model_gateway/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -631,6 +631,21 @@ struct CliArgs {
#[arg(long, default_value_t = 0.02, help_heading = "Adaptive Admission")]
adaptive_admission_feedback_throughput_improvement_ratio: f64,

/// Exact admission partitions allowed to expose clean-worker distribution
/// headroom. Empty by default.
#[arg(
long,
value_delimiter = ',',
num_args = 1..,
help_heading = "Adaptive Admission"
)]
adaptive_admission_distribution_headroom_partitions: Vec<String>,

/// Hard process-local active seed-lease cap per allowlisted partition.
/// Zero disables distribution headroom.
#[arg(long, default_value_t = 0, help_heading = "Adaptive Admission")]
adaptive_admission_distribution_headroom_partition_seed_cap: u32,

// ==================== Tenant Rate Limit ====================
/// Enable per-tenant LLM token/request rate limiting. When unset
/// (default), no rate limiter is constructed.
Expand Down Expand Up @@ -1654,6 +1669,11 @@ impl CliArgs {
feedback_max_token_usage: self.adaptive_admission_feedback_max_token_usage,
feedback_throughput_improvement_ratio: self
.adaptive_admission_feedback_throughput_improvement_ratio,
distribution_headroom_partitions: self
.adaptive_admission_distribution_headroom_partitions
.clone(),
distribution_headroom_partition_seed_cap: self
.adaptive_admission_distribution_headroom_partition_seed_cap,
})
.tenant_rate_limit_enabled(self.tenant_rate_limit_enabled)
.tenant_rate_limit_config(self.tenant_rate_limit_config.clone())
Expand Down Expand Up @@ -2027,7 +2047,7 @@ mod tests {
fn engine_feedback_admission_options_flow_into_router_config() {
let cli = cli_args_from(&[
"--adaptive-admission-mode",
"shadow",
"enforce",
"--adaptive-admission-strategy",
"engine-feedback",
"--adaptive-admission-feedback-probe-requests-per-healthy-replica",
Expand All @@ -2038,11 +2058,15 @@ mod tests {
"0.85",
"--adaptive-admission-feedback-throughput-improvement-ratio",
"0.03",
"--adaptive-admission-distribution-headroom-partitions",
"k3-prod,k3-canary",
"--adaptive-admission-distribution-headroom-partition-seed-cap",
"2",
]);

let router_config = cli.to_router_config(vec![], vec![]).unwrap();
let adaptive = router_config.adaptive_admission;
assert_eq!(adaptive.mode, AdaptiveAdmissionMode::Shadow);
assert_eq!(adaptive.mode, AdaptiveAdmissionMode::Enforce);
assert_eq!(adaptive.strategy, AdaptiveAdmissionStrategy::EngineFeedback);
assert_eq!(adaptive.feedback_probe_requests_per_healthy_replica, 3);
assert_eq!(
Expand All @@ -2051,6 +2075,11 @@ mod tests {
);
assert_eq!(adaptive.feedback_max_token_usage, 0.85);
assert_eq!(adaptive.feedback_throughput_improvement_ratio, 0.03);
assert_eq!(
adaptive.distribution_headroom_partitions,
["k3-prod", "k3-canary"]
);
assert_eq!(adaptive.distribution_headroom_partition_seed_cap, 2);
}

/// The multimodal transport flags must reach both `RouterConfig` and the
Expand Down
Loading
Loading