Skip to content

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 →