#![allow(dead_code)]
use std::collections::HashMap;
use chrono::{DateTime, Duration, Utc};
use crate::ir_nodes::IRBudgetQuota;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum BudgetPeriod {
Second,
Minute,
Hour,
Day,
}
impl BudgetPeriod {
pub fn parse(s: &str) -> Option<Self> {
Some(match s {
"second" => BudgetPeriod::Second,
"minute" => BudgetPeriod::Minute,
"hour" => BudgetPeriod::Hour,
"day" => BudgetPeriod::Day,
_ => return None,
})
}
pub fn as_secs(self) -> f64 {
match self {
BudgetPeriod::Second => 1.0,
BudgetPeriod::Minute => 60.0,
BudgetPeriod::Hour => 3600.0,
BudgetPeriod::Day => 86400.0,
}
}
}
#[derive(Debug, Clone, PartialEq)]
pub enum AcquireOutcome {
Granted,
Denied { retry_at: DateTime<Utc> },
}
impl AcquireOutcome {
pub fn is_granted(&self) -> bool {
matches!(self, AcquireOutcome::Granted)
}
}
#[derive(Debug, Clone)]
enum RateState {
Bucket { tokens: f64, last_refill: DateTime<Utc> },
Window { window_start: DateTime<Utc>, consumed: i64 },
}
#[derive(Debug, Clone)]
pub struct RateLease {
pub effect: String,
pub limit: i64,
pub period: BudgetPeriod,
state: RateState,
}
impl RateLease {
pub fn rate(effect: impl Into<String>, limit: i64, period: BudgetPeriod, now: DateTime<Utc>) -> Self {
RateLease {
effect: effect.into(),
limit,
period,
state: RateState::Bucket { tokens: limit.max(0) as f64, last_refill: now },
}
}
pub fn max(effect: impl Into<String>, limit: i64, period: BudgetPeriod, now: DateTime<Utc>) -> Self {
RateLease {
effect: effect.into(),
limit,
period,
state: RateState::Window { window_start: now, consumed: 0 },
}
}
pub fn from_quota(q: &IRBudgetQuota, now: DateTime<Utc>) -> Option<Self> {
let period = BudgetPeriod::parse(&q.period)?;
Some(match q.kind.as_str() {
"max" => RateLease::max(q.effect.clone(), q.limit, period, now),
_ => RateLease::rate(q.effect.clone(), q.limit, period, now),
})
}
fn refill_per_sec(&self) -> f64 {
self.limit.max(0) as f64 / self.period.as_secs()
}
pub fn refill(&mut self, now: DateTime<Utc>) {
let rate = self.refill_per_sec();
let capacity = self.limit.max(0) as f64;
let period_secs = self.period.as_secs();
match &mut self.state {
RateState::Bucket { tokens, last_refill } => {
let elapsed = (now - *last_refill).num_milliseconds() as f64 / 1000.0;
if elapsed > 0.0 {
*tokens = (*tokens + elapsed * rate).min(capacity);
*last_refill = now;
}
}
RateState::Window { window_start, consumed } => {
let elapsed = (now - *window_start).num_milliseconds() as f64 / 1000.0;
if elapsed >= period_secs {
*window_start = now;
*consumed = 0;
}
}
}
}
pub fn try_acquire(&mut self, now: DateTime<Utc>) -> AcquireOutcome {
self.refill(now);
let rate = self.refill_per_sec();
let period_secs = self.period.as_secs();
match &mut self.state {
RateState::Bucket { tokens, .. } => {
if *tokens >= 1.0 {
*tokens -= 1.0;
AcquireOutcome::Granted
} else {
let deficit = 1.0 - *tokens;
let wait_secs = if rate > 0.0 { deficit / rate } else { f64::INFINITY };
let retry_at = now + secs_to_duration(wait_secs);
AcquireOutcome::Denied { retry_at }
}
}
RateState::Window { window_start, consumed } => {
if *consumed < self.limit {
*consumed += 1;
AcquireOutcome::Granted
} else {
let retry_at = *window_start + secs_to_duration(period_secs);
AcquireOutcome::Denied { retry_at }
}
}
}
}
pub fn available(&self, now: DateTime<Utc>) -> f64 {
let mut probe = self.clone();
probe.refill(now);
match probe.state {
RateState::Bucket { tokens, .. } => tokens,
RateState::Window { consumed, .. } => (self.limit - consumed).max(0) as f64,
}
}
pub fn peek(&self, now: DateTime<Utc>) -> AcquireOutcome {
let mut probe = self.clone();
probe.try_acquire(now)
}
pub fn snapshot(&self) -> RateLeaseSnapshot {
match &self.state {
RateState::Bucket { tokens, last_refill } => RateLeaseSnapshot {
kind: "rate".to_string(),
tokens: *tokens,
last_refill_ms: last_refill.timestamp_millis(),
window_start_ms: 0,
consumed: 0,
},
RateState::Window { window_start, consumed } => RateLeaseSnapshot {
kind: "max".to_string(),
tokens: 0.0,
last_refill_ms: 0,
window_start_ms: window_start.timestamp_millis(),
consumed: *consumed,
},
}
}
pub fn restore(&mut self, snap: &RateLeaseSnapshot) {
match (&mut self.state, snap.kind.as_str()) {
(RateState::Bucket { tokens, last_refill }, "rate") => {
*tokens = snap.tokens.min(self.limit.max(0) as f64);
if let Some(t) = DateTime::from_timestamp_millis(snap.last_refill_ms) {
*last_refill = t;
}
}
(RateState::Window { window_start, consumed }, "max") => {
*consumed = snap.consumed.clamp(0, self.limit.max(0));
if let Some(t) = DateTime::from_timestamp_millis(snap.window_start_ms) {
*window_start = t;
}
}
_ => { }
}
}
}
#[derive(Debug, Clone, PartialEq, serde::Serialize, serde::Deserialize)]
pub struct RateLeaseSnapshot {
pub kind: String,
pub tokens: f64,
pub last_refill_ms: i64,
pub window_start_ms: i64,
pub consumed: i64,
}
fn secs_to_duration(secs: f64) -> Duration {
if !secs.is_finite() || secs <= 0.0 {
return Duration::zero();
}
let capped = secs.min(8_640_000.0);
Duration::milliseconds((capped * 1000.0) as i64)
}
#[derive(Default)]
pub struct RateLeaseKernel {
leases: HashMap<String, RateLease>,
}
impl RateLeaseKernel {
pub fn new() -> Self {
Self::default()
}
pub fn register(&mut self, key: impl Into<String>, lease: RateLease) {
self.leases.insert(key.into(), lease);
}
pub fn contains(&self, key: &str) -> bool {
self.leases.contains_key(key)
}
pub fn try_acquire(&mut self, key: &str, now: DateTime<Utc>) -> AcquireOutcome {
match self.leases.get_mut(key) {
Some(lease) => lease.try_acquire(now),
None => AcquireOutcome::Granted,
}
}
pub fn try_acquire_all(&mut self, keys: &[String], now: DateTime<Utc>) -> AcquireOutcome {
let mut latest_retry: Option<DateTime<Utc>> = None;
for key in keys {
if let Some(lease) = self.leases.get(key) {
if let AcquireOutcome::Denied { retry_at } = lease.peek(now) {
latest_retry = Some(match latest_retry {
Some(prev) if prev >= retry_at => prev,
_ => retry_at,
});
}
}
}
if let Some(retry_at) = latest_retry {
return AcquireOutcome::Denied { retry_at };
}
for key in keys {
if let Some(lease) = self.leases.get_mut(key) {
let _ = lease.try_acquire(now);
}
}
AcquireOutcome::Granted
}
pub fn available(&self, key: &str, now: DateTime<Utc>) -> Option<f64> {
self.leases.get(key).map(|l| l.available(now))
}
pub fn tick(&mut self, now: DateTime<Utc>) {
for lease in self.leases.values_mut() {
lease.refill(now);
}
}
pub fn snapshot(&self) -> Vec<(String, RateLeaseSnapshot)> {
let mut out: Vec<(String, RateLeaseSnapshot)> =
self.leases.iter().map(|(k, l)| (k.clone(), l.snapshot())).collect();
out.sort_by(|a, b| a.0.cmp(&b.0));
out
}
pub fn restore(&mut self, snaps: &[(String, RateLeaseSnapshot)]) {
for (key, snap) in snaps {
if let Some(lease) = self.leases.get_mut(key) {
lease.restore(snap);
}
}
}
}
#[derive(Debug, Clone, PartialEq)]
pub enum GateDecision {
Allow,
Deny {
retry_at: DateTime<Utc>,
on_exhausted: String,
},
}
pub struct BudgetGate {
kernel: RateLeaseKernel,
on_exhausted: String,
by_effect: HashMap<String, Vec<String>>,
}
impl BudgetGate {
pub fn from_ir(budget: &crate::ir_nodes::IRBudget, scope: &str, now: DateTime<Utc>) -> Self {
let mut kernel = RateLeaseKernel::new();
let mut by_effect: HashMap<String, Vec<String>> = HashMap::new();
for (i, quota) in budget.quotas.iter().enumerate() {
let Some(lease) = RateLease::from_quota(quota, now) else {
continue;
};
let key = format!("{scope}:Tool({}):{}:{i}", quota.effect, quota.kind);
kernel.register(key.clone(), lease);
by_effect.entry(quota.effect.clone()).or_default().push(key);
}
BudgetGate {
kernel,
on_exhausted: if budget.on_exhausted.is_empty() {
"block".to_string()
} else {
budget.on_exhausted.clone()
},
by_effect,
}
}
pub fn gate(&mut self, effect: &str, now: DateTime<Utc>) -> GateDecision {
let Some(keys) = self.by_effect.get(effect) else {
return GateDecision::Allow;
};
let keys = keys.clone();
match self.kernel.try_acquire_all(&keys, now) {
AcquireOutcome::Granted => GateDecision::Allow,
AcquireOutcome::Denied { retry_at } => GateDecision::Deny {
retry_at,
on_exhausted: self.on_exhausted.clone(),
},
}
}
pub fn on_exhausted(&self) -> &str {
&self.on_exhausted
}
pub fn governs(&self, effect: &str) -> bool {
self.by_effect.contains_key(effect)
}
pub fn snapshot(&self) -> Vec<(String, RateLeaseSnapshot)> {
self.kernel.snapshot()
}
pub fn restore(&mut self, snaps: &[(String, RateLeaseSnapshot)]) {
self.kernel.restore(snaps);
}
}
#[cfg(test)]
mod tests {
use super::*;
fn t0() -> DateTime<Utc> {
"2026-06-29T00:00:00Z".parse().unwrap()
}
fn quota(kind: &str, limit: i64, period: &str, effect: &str) -> IRBudgetQuota {
IRBudgetQuota {
kind: kind.into(),
limit,
period: period.into(),
effect: effect.into(),
}
}
fn ir_budget(quotas: Vec<IRBudgetQuota>, on_exhausted: &str) -> crate::ir_nodes::IRBudget {
crate::ir_nodes::IRBudget {
node_type: "budget",
source_line: 1,
source_column: 1,
quotas,
on_exhausted: on_exhausted.into(),
}
}
#[test]
fn period_parses_and_maps_to_seconds() {
assert_eq!(BudgetPeriod::parse("hour"), Some(BudgetPeriod::Hour));
assert_eq!(BudgetPeriod::parse("fortnight"), None);
assert_eq!(BudgetPeriod::Second.as_secs(), 1.0);
assert_eq!(BudgetPeriod::Minute.as_secs(), 60.0);
assert_eq!(BudgetPeriod::Hour.as_secs(), 3600.0);
assert_eq!(BudgetPeriod::Day.as_secs(), 86400.0);
}
#[test]
fn bucket_starts_full_and_grants_up_to_capacity() {
let now = t0();
let mut l = RateLease::rate("Telnyx", 3, BudgetPeriod::Hour, now);
assert!(l.try_acquire(now).is_granted());
assert!(l.try_acquire(now).is_granted());
assert!(l.try_acquire(now).is_granted());
match l.try_acquire(now) {
AcquireOutcome::Denied { retry_at } => {
assert_eq!(retry_at, now + Duration::seconds(1200));
}
other => panic!("expected Denied, got {other:?}"),
}
}
#[test]
fn bucket_refills_over_time() {
let now = t0();
let mut l = RateLease::rate("Telnyx", 2, BudgetPeriod::Hour, now);
assert!(l.try_acquire(now).is_granted());
assert!(l.try_acquire(now).is_granted());
assert!(!l.try_acquire(now).is_granted());
let later = now + Duration::seconds(1800);
assert!(l.try_acquire(later).is_granted());
assert!(!l.try_acquire(later).is_granted());
}
#[test]
fn bucket_refill_is_capped_at_capacity() {
let now = t0();
let mut l = RateLease::rate("Telnyx", 5, BudgetPeriod::Minute, now);
assert!(l.try_acquire(now).is_granted());
let way_later = now + Duration::days(1);
for _ in 0..5 {
assert!(l.try_acquire(way_later).is_granted());
}
assert!(!l.try_acquire(way_later).is_granted(), "capped at capacity");
}
#[test]
fn window_caps_at_limit_then_rolls() {
let now = t0();
let mut l = RateLease::max("Telnyx", 50, BudgetPeriod::Day, now);
for _ in 0..50 {
assert!(l.try_acquire(now).is_granted());
}
match l.try_acquire(now) {
AcquireOutcome::Denied { retry_at } => {
assert_eq!(retry_at, now + Duration::seconds(86400));
}
other => panic!("expected Denied, got {other:?}"),
}
assert!(!l.try_acquire(now + Duration::seconds(86399)).is_granted());
let next_day = now + Duration::seconds(86400);
assert!(l.try_acquire(next_day).is_granted());
assert_eq!(l.available(next_day), 49.0);
}
#[test]
fn window_has_no_intra_window_refill() {
let now = t0();
let mut l = RateLease::max("Telnyx", 2, BudgetPeriod::Hour, now);
assert!(l.try_acquire(now).is_granted());
assert!(l.try_acquire(now).is_granted());
assert!(!l.try_acquire(now + Duration::seconds(1800)).is_granted());
}
#[test]
fn from_quota_builds_the_right_kind() {
let now = t0();
let rate_q = IRBudgetQuota {
kind: "rate".into(),
limit: 8,
period: "hour".into(),
effect: "Telnyx".into(),
};
let max_q = IRBudgetQuota {
kind: "max".into(),
limit: 50,
period: "day".into(),
effect: "Telnyx".into(),
};
let rate = RateLease::from_quota(&rate_q, now).unwrap();
assert_eq!(rate.available(now), 8.0, "a rate bucket starts full");
let maxl = RateLease::from_quota(&max_q, now).unwrap();
assert_eq!(maxl.available(now), 50.0, "a max window starts with the full allowance");
let bad = IRBudgetQuota { period: "fortnight".into(), ..rate_q };
assert!(RateLease::from_quota(&bad, now).is_none());
}
#[test]
fn kernel_unregistered_key_is_unbudgeted() {
let mut k = RateLeaseKernel::new();
assert!(k.try_acquire("daemon:X:Tool(Y):rate", t0()).is_granted());
assert_eq!(k.available("daemon:X:Tool(Y):rate", t0()), None);
}
#[test]
fn kernel_enforces_a_registered_lease() {
let now = t0();
let mut k = RateLeaseKernel::new();
k.register("d:Out:Tool(Telnyx):rate", RateLease::rate("Telnyx", 1, BudgetPeriod::Hour, now));
assert!(k.try_acquire("d:Out:Tool(Telnyx):rate", now).is_granted());
assert!(!k.try_acquire("d:Out:Tool(Telnyx):rate", now).is_granted());
let later = now + Duration::seconds(3600);
assert!(k.try_acquire("d:Out:Tool(Telnyx):rate", later).is_granted());
}
#[test]
fn kernel_tick_refreshes_available_without_consuming() {
let now = t0();
let mut k = RateLeaseKernel::new();
k.register("k", RateLease::rate("E", 4, BudgetPeriod::Minute, now));
for _ in 0..4 {
assert!(k.try_acquire("k", now).is_granted());
}
assert_eq!(k.available("k", now), Some(0.0));
let later = now + Duration::seconds(30);
k.tick(later);
assert_eq!(k.available("k", later), Some(2.0));
}
#[test]
fn acquire_all_is_all_or_none() {
let now = t0();
let mut k = RateLeaseKernel::new();
k.register("r", RateLease::rate("E", 5, BudgetPeriod::Hour, now));
k.register("m", RateLease::max("E", 1, BudgetPeriod::Day, now));
let keys = vec!["r".to_string(), "m".to_string()];
assert!(k.try_acquire_all(&keys, now).is_granted());
match k.try_acquire_all(&keys, now) {
AcquireOutcome::Denied { retry_at } => {
assert_eq!(retry_at, now + Duration::seconds(86400), "binding = the daily max");
}
other => panic!("expected Denied, got {other:?}"),
}
assert_eq!(k.available("r", now), Some(4.0), "rate token not consumed on denial");
}
#[test]
fn acquire_all_empty_keys_is_granted() {
let mut k = RateLeaseKernel::new();
assert!(k.try_acquire_all(&[], t0()).is_granted());
}
#[test]
fn gate_allows_unbudgeted_effects() {
let now = t0();
let b = ir_budget(vec![quota("rate", 1, "hour", "Telnyx")], "block");
let mut gate = BudgetGate::from_ir(&b, "daemon:Out", now);
assert_eq!(gate.gate("SomeOtherTool", now), GateDecision::Allow);
assert!(!gate.governs("SomeOtherTool"));
assert!(gate.governs("Telnyx"));
}
#[test]
fn gate_enforces_then_denies_with_policy() {
let now = t0();
let b = ir_budget(
vec![
quota("rate", 2, "hour", "Telnyx"),
quota("max", 3, "day", "Telnyx"),
],
"defer",
);
let mut gate = BudgetGate::from_ir(&b, "daemon:Out", now);
assert_eq!(gate.gate("Telnyx", now), GateDecision::Allow);
assert_eq!(gate.gate("Telnyx", now), GateDecision::Allow);
match gate.gate("Telnyx", now) {
GateDecision::Deny { on_exhausted, retry_at } => {
assert_eq!(on_exhausted, "defer");
assert_eq!(retry_at, now + Duration::seconds(1800));
}
other => panic!("expected Deny, got {other:?}"),
}
}
#[test]
fn gate_omitted_policy_is_block() {
let now = t0();
let b = ir_budget(vec![quota("rate", 1, "hour", "E")], "");
let gate = BudgetGate::from_ir(&b, "d", now);
assert_eq!(gate.on_exhausted(), "block");
}
#[test]
fn snapshot_restore_carries_max_window_across_ticks() {
let now = t0();
let b = ir_budget(vec![quota("max", 3, "day", "Telnyx")], "block");
let mut g1 = BudgetGate::from_ir(&b, "d", now);
assert_eq!(g1.gate("Telnyx", now), GateDecision::Allow);
assert_eq!(g1.gate("Telnyx", now), GateDecision::Allow);
let snap = g1.snapshot();
let mut g2 = BudgetGate::from_ir(&b, "d", now + Duration::minutes(5));
g2.restore(&snap);
assert_eq!(g2.gate("Telnyx", now + Duration::minutes(5)), GateDecision::Allow);
match g2.gate("Telnyx", now + Duration::minutes(5)) {
GateDecision::Deny { .. } => {}
other => panic!("expected the daily cap to hold across ticks, got {other:?}"),
}
}
#[test]
fn snapshot_round_trips_a_bucket() {
let now = t0();
let mut l = RateLease::rate("E", 8, BudgetPeriod::Hour, now);
l.try_acquire(now); let snap = l.snapshot();
assert_eq!(snap.kind, "rate");
let mut l2 = RateLease::rate("E", 8, BudgetPeriod::Hour, now);
l2.restore(&snap);
assert_eq!(l2.available(now), 7.0, "restored bucket carries the consumed token");
}
}