Rate Limit Planner — Rust source
Turn RPM/TPM limits into a concrete request schedule — batch size, spacing, binding limit, and total run time, with a safety factor for retries. 100% client-side.
This is the Rust implementation — the same logic the interactive tool runs, in a shareable, citable form.
//! Rate Limit Planner — turn provider rate limits plus a workload into a
//! concrete schedule: batch size, spacing, timeline and wall time.
//!
//! Language: Rust (edition 2021, standard library only)
//! Port of src/lib/rateLimitPlanner.ts (the canonical TypeScript
//! implementation).
//! Tool page: https://dev.cosmolabs.org/tools/rate-limit-planner
//!
//! Deterministic — no time reads. The TS lib throws RangeError; this port
//! returns [`PlannerError`] with the same messages. Numeric plan fields
//! are `f64` because the original returns `Infinity` for the unbounded
//! (no limits) case — `batch_size`/`max_concurrent` are `f64::INFINITY`
//! there.
use std::fmt;
/// Length of a rate-limit window, in ms.
const WINDOW_MS: f64 = 60_000.0;
/// Default fraction of the limits to target, leaving headroom for retries.
const DEFAULT_SAFETY: f64 = 0.8;
/// Timeline is capped at this many leading batches — enough to act on.
const MAX_TIMELINE: u32 = 10;
/// Provider rate limits; `None` = not limited.
#[derive(Debug, Clone, Copy, Default, PartialEq)]
pub struct RateLimits {
/// Requests per minute.
pub rpm: Option<f64>,
/// Tokens per minute.
pub tpm: Option<f64>,
}
/// The workload to schedule.
#[derive(Debug, Clone, Copy, PartialEq)]
pub struct Workload {
/// Total requests to run.
pub requests: u64,
/// Average tokens per request (prompt + completion).
pub avg_tokens_per_request: f64,
}
/// Optional tweaks; `None` selects the default.
#[derive(Debug, Clone, Copy, Default, PartialEq)]
pub struct PlanOptions {
/// Fraction of the limits to target, leaving headroom for retries.
pub safety_factor: Option<f64>,
}
/// One batch of the schedule: how many requests, when, and their tokens.
#[derive(Debug, Clone, Copy, PartialEq)]
pub struct BatchSlice {
pub batch: u64,
pub at_ms: f64,
pub requests: u64,
pub tokens: f64,
}
/// Which limit binds first.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Bound {
Rpm,
Tpm,
Both,
/// Endpoint treated as unbounded (no limits set).
Unlimited,
}
impl fmt::Display for Bound {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.write_str(match self {
Bound::Rpm => "rpm",
Bound::Tpm => "tpm",
Bound::Both => "both",
Bound::Unlimited => "none",
})
}
}
impl Bound {
/// Uppercase label used by [`describe_plan`] (`toUpperCase()` in TS).
fn upper(self) -> &'static str {
match self {
Bound::Rpm => "RPM",
Bound::Tpm => "TPM",
_ => "",
}
}
}
/// A concrete schedule for the workload.
#[derive(Debug, Clone, PartialEq)]
pub struct RateLimitPlan {
/// Requests to send per 60s window (0 when the workload cannot run).
pub batch_size: f64,
/// Steady-state spacing between individual requests, in ms.
pub interval_ms: f64,
/// Sustainable concurrent in-flight requests under even spacing.
pub max_concurrent: f64,
/// Which limit binds first.
pub bounded_by: Bound,
/// First batches of the schedule (max 10) — enough to act on.
pub timeline: Vec<BatchSlice>,
/// Estimated total wall time, in ms.
pub total_ms: f64,
pub warnings: Vec<String>,
}
/// Error returned instead of the TS lib's `RangeError`.
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct PlannerError(pub String);
impl fmt::Display for PlannerError {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.write_str(&self.0)
}
}
impl std::error::Error for PlannerError {}
/// Plan a schedule under the given limits. `safety_factor` must be in
/// (0, 1] (default 0.8). Errors on a negative `avg_tokens_per_request`
/// or an out-of-range safety factor.
pub fn plan_rate_limit(
limits: RateLimits,
workload: Workload,
opts: PlanOptions,
) -> Result<RateLimitPlan, PlannerError> {
let sf = opts.safety_factor.unwrap_or(DEFAULT_SAFETY);
let mut warnings: Vec<String> = Vec::new();
if workload.avg_tokens_per_request < 0.0 {
// `requests` is u64, so it is always >= 0 — message kept for parity.
return Err(PlannerError(
"requests and avgTokensPerRequest must be >= 0".into(),
));
}
if !(0.0 < sf && sf <= 1.0) {
return Err(PlannerError("safetyFactor must be in (0, 1]".into()));
}
let rpm_eff = limits.rpm.map(|v| v * sf);
let tpm_eff = limits.tpm.map(|v| v * sf);
// Impossible: one request alone exceeds the token budget.
if let Some(tpm) = tpm_eff {
if workload.avg_tokens_per_request > tpm && workload.requests > 0 {
return Ok(RateLimitPlan {
batch_size: 0.0,
interval_ms: 0.0,
max_concurrent: 0.0,
bounded_by: Bound::Tpm,
timeline: Vec::new(),
total_ms: f64::INFINITY,
warnings: vec![format!(
"A single request averages {} tokens but the effective \
token limit is {}/min — no schedule can run this. \
Shrink requests or raise the tier.",
fmt_num(workload.avg_tokens_per_request),
fmt_num(tpm.floor()),
)],
});
}
}
let by_rpm = rpm_eff.unwrap_or(f64::INFINITY);
let by_tokens = match tpm_eff {
None => f64::INFINITY,
Some(_) if workload.avg_tokens_per_request == 0.0 => f64::INFINITY,
Some(t) => t / workload.avg_tokens_per_request,
};
if by_rpm.is_infinite() && by_tokens.is_infinite() {
warnings.push(
"No limits set — the plan assumes an unbounded endpoint. \
Add RPM or TPM for a real schedule."
.into(),
);
}
let steady = 1.0_f64.max(by_rpm.min(by_tokens).floor());
let bounded_by = if by_rpm.is_infinite() && by_tokens.is_infinite() {
Bound::Unlimited
} else if by_rpm.floor() == by_tokens.floor() {
Bound::Both
} else if by_rpm < by_tokens {
Bound::Rpm
} else {
Bound::Tpm
};
// Even pacing inside the window: batch_size requests spread over 60s.
let interval_ms = (WINDOW_MS / steady).round();
// With even spacing and a per-request latency near interval_ms, one
// request is in flight at a time; concurrency >1 only helps
// sub-interval latencies, so the safe published floor is 1 — batch
// bursts raise it to batch_size.
let max_concurrent = if steady == 1.0 {
1.0
} else {
steady.min((steady / 4.0).ceil())
};
let mut timeline: Vec<BatchSlice> = Vec::new();
let mut remaining = workload.requests as f64;
let mut batch: u64 = 0;
while remaining > 0.0 && batch < MAX_TIMELINE as u64 {
let take = steady.min(remaining);
timeline.push(BatchSlice {
batch: batch + 1,
at_ms: batch as f64 * WINDOW_MS,
requests: take as u64,
tokens: take * workload.avg_tokens_per_request,
});
remaining -= take;
batch += 1;
}
let windows_needed = if workload.requests > 0 {
(workload.requests as f64 / steady).ceil()
} else {
0.0
};
let last_window_requests = if windows_needed > 0.0 {
workload.requests as f64 - (windows_needed - 1.0) * steady
} else {
0.0
};
let total_ms = if windows_needed > 0.0 {
(windows_needed - 1.0) * WINDOW_MS + interval_ms * last_window_requests
} else {
0.0
};
if rpm_eff.is_some() && workload.requests > 0 && steady > by_rpm {
warnings.push(
"Rounded up to at least one request per window — even a single \
request per minute keeps the schedule honest."
.into(),
);
}
Ok(RateLimitPlan {
batch_size: steady,
interval_ms,
max_concurrent,
bounded_by,
timeline,
total_ms,
warnings,
})
}
/// Human summary line for the plan (used by the island + docs).
pub fn describe_plan(plan: &RateLimitPlan) -> String {
if plan.batch_size == 0.0 {
return "No viable schedule.".into();
}
if plan.bounded_by == Bound::Unlimited {
return format!(
"{}+ requests per window — endpoint treated as unbounded.",
plain_num(plan.batch_size)
);
}
let limiter = if plan.bounded_by == Bound::Both {
"both limits bind together".to_string()
} else {
format!("the {} limit binds first", plan.bounded_by.upper())
};
format!(
"{} requests per 60s window (one every {}ms) — {}.",
plain_num(plan.batch_size),
plain_num(plan.interval_ms),
limiter
)
}
/// Plain interpolation like TS `${v}` (no separators, no trailing `.0`).
fn plain_num(v: f64) -> String {
if v == f64::INFINITY {
return "Infinity".into();
}
if v == f64::NEG_INFINITY {
return "-Infinity".into();
}
if v == v.trunc() && v.abs() < 9.0e15 {
return format!("{}", v as i64);
}
format!("{v}")
}
/// Format like TS `toLocaleString('en-US')`: thousands separators.
fn fmt_num(v: f64) -> String {
if v == f64::INFINITY {
return "Infinity".into();
}
if v == f64::NEG_INFINITY {
return "-Infinity".into();
}
let neg = v < 0.0;
let a = v.abs();
let int_part = a.trunc();
let digits = format!("{}", int_part as u64);
let bytes = digits.as_bytes();
let mut grouped = String::with_capacity(bytes.len() + bytes.len() / 3);
for (i, b) in bytes.iter().enumerate() {
if i > 0 && (bytes.len() - i) % 3 == 0 {
grouped.push(',');
}
grouped.push(*b as char);
}
let frac = a - int_part;
if frac > 0.0 {
grouped.push_str(&format!("{frac:?}"));
}
if neg {
format!("-{grouped}")
} else {
grouped
}
}
Also available in 13 other languages
Every CosmoDev tool ships its pure logic in TypeScript (web) and Go (CLI), with authored implementations in a dozen-plus languages — the same contract, ported. Compare all languages side by side →