1#![allow(dead_code)]
27
28use std::collections::HashMap;
29
30use chrono::{DateTime, Duration, Utc};
31
32use crate::ir_nodes::IRBudgetQuota;
33
34#[derive(Debug, Clone, Copy, PartialEq, Eq)]
41pub enum BudgetPeriod {
42 Second,
43 Minute,
44 Hour,
45 Day,
46}
47
48impl BudgetPeriod {
49 pub fn parse(s: &str) -> Option<Self> {
52 Some(match s {
53 "second" => BudgetPeriod::Second,
54 "minute" => BudgetPeriod::Minute,
55 "hour" => BudgetPeriod::Hour,
56 "day" => BudgetPeriod::Day,
57 _ => return None,
58 })
59 }
60
61 pub fn as_secs(self) -> f64 {
63 match self {
64 BudgetPeriod::Second => 1.0,
65 BudgetPeriod::Minute => 60.0,
66 BudgetPeriod::Hour => 3600.0,
67 BudgetPeriod::Day => 86400.0,
68 }
69 }
70}
71
72#[derive(Debug, Clone, PartialEq)]
79pub enum AcquireOutcome {
80 Granted,
82 Denied { retry_at: DateTime<Utc> },
85}
86
87impl AcquireOutcome {
88 pub fn is_granted(&self) -> bool {
89 matches!(self, AcquireOutcome::Granted)
90 }
91}
92
93#[derive(Debug, Clone)]
98enum RateState {
99 Bucket { tokens: f64, last_refill: DateTime<Utc> },
102 Window { window_start: DateTime<Utc>, consumed: i64 },
104}
105
106#[derive(Debug, Clone)]
110pub struct RateLease {
111 pub effect: String,
113 pub limit: i64,
115 pub period: BudgetPeriod,
117 state: RateState,
118}
119
120impl RateLease {
121 pub fn rate(effect: impl Into<String>, limit: i64, period: BudgetPeriod, now: DateTime<Utc>) -> Self {
124 RateLease {
125 effect: effect.into(),
126 limit,
127 period,
128 state: RateState::Bucket { tokens: limit.max(0) as f64, last_refill: now },
129 }
130 }
131
132 pub fn max(effect: impl Into<String>, limit: i64, period: BudgetPeriod, now: DateTime<Utc>) -> Self {
134 RateLease {
135 effect: effect.into(),
136 limit,
137 period,
138 state: RateState::Window { window_start: now, consumed: 0 },
139 }
140 }
141
142 pub fn from_quota(q: &IRBudgetQuota, now: DateTime<Utc>) -> Option<Self> {
146 let period = BudgetPeriod::parse(&q.period)?;
147 Some(match q.kind.as_str() {
148 "max" => RateLease::max(q.effect.clone(), q.limit, period, now),
149 _ => RateLease::rate(q.effect.clone(), q.limit, period, now),
151 })
152 }
153
154 fn refill_per_sec(&self) -> f64 {
156 self.limit.max(0) as f64 / self.period.as_secs()
157 }
158
159 pub fn refill(&mut self, now: DateTime<Utc>) {
162 let rate = self.refill_per_sec();
163 let capacity = self.limit.max(0) as f64;
164 let period_secs = self.period.as_secs();
165 match &mut self.state {
166 RateState::Bucket { tokens, last_refill } => {
167 let elapsed = (now - *last_refill).num_milliseconds() as f64 / 1000.0;
168 if elapsed > 0.0 {
169 *tokens = (*tokens + elapsed * rate).min(capacity);
170 *last_refill = now;
171 }
172 }
173 RateState::Window { window_start, consumed } => {
174 let elapsed = (now - *window_start).num_milliseconds() as f64 / 1000.0;
175 if elapsed >= period_secs {
176 *window_start = now;
177 *consumed = 0;
178 }
179 }
180 }
181 }
182
183 pub fn try_acquire(&mut self, now: DateTime<Utc>) -> AcquireOutcome {
187 self.refill(now);
188 let rate = self.refill_per_sec();
189 let period_secs = self.period.as_secs();
190 match &mut self.state {
191 RateState::Bucket { tokens, .. } => {
192 if *tokens >= 1.0 {
193 *tokens -= 1.0;
194 AcquireOutcome::Granted
195 } else {
196 let deficit = 1.0 - *tokens;
198 let wait_secs = if rate > 0.0 { deficit / rate } else { f64::INFINITY };
199 let retry_at = now + secs_to_duration(wait_secs);
200 AcquireOutcome::Denied { retry_at }
201 }
202 }
203 RateState::Window { window_start, consumed } => {
204 if *consumed < self.limit {
205 *consumed += 1;
206 AcquireOutcome::Granted
207 } else {
208 let retry_at = *window_start + secs_to_duration(period_secs);
209 AcquireOutcome::Denied { retry_at }
210 }
211 }
212 }
213 }
214
215 pub fn available(&self, now: DateTime<Utc>) -> f64 {
218 let mut probe = self.clone();
219 probe.refill(now);
220 match probe.state {
221 RateState::Bucket { tokens, .. } => tokens,
222 RateState::Window { consumed, .. } => (self.limit - consumed).max(0) as f64,
223 }
224 }
225
226 pub fn peek(&self, now: DateTime<Utc>) -> AcquireOutcome {
230 let mut probe = self.clone();
231 probe.try_acquire(now)
232 }
233
234 pub fn snapshot(&self) -> RateLeaseSnapshot {
240 match &self.state {
241 RateState::Bucket { tokens, last_refill } => RateLeaseSnapshot {
242 kind: "rate".to_string(),
243 tokens: *tokens,
244 last_refill_ms: last_refill.timestamp_millis(),
245 window_start_ms: 0,
246 consumed: 0,
247 },
248 RateState::Window { window_start, consumed } => RateLeaseSnapshot {
249 kind: "max".to_string(),
250 tokens: 0.0,
251 last_refill_ms: 0,
252 window_start_ms: window_start.timestamp_millis(),
253 consumed: *consumed,
254 },
255 }
256 }
257
258 pub fn restore(&mut self, snap: &RateLeaseSnapshot) {
264 match (&mut self.state, snap.kind.as_str()) {
265 (RateState::Bucket { tokens, last_refill }, "rate") => {
266 *tokens = snap.tokens.min(self.limit.max(0) as f64);
267 if let Some(t) = DateTime::from_timestamp_millis(snap.last_refill_ms) {
268 *last_refill = t;
269 }
270 }
271 (RateState::Window { window_start, consumed }, "max") => {
272 *consumed = snap.consumed.clamp(0, self.limit.max(0));
273 if let Some(t) = DateTime::from_timestamp_millis(snap.window_start_ms) {
274 *window_start = t;
275 }
276 }
277 _ => { }
278 }
279 }
280}
281
282#[derive(Debug, Clone, PartialEq, serde::Serialize, serde::Deserialize)]
287pub struct RateLeaseSnapshot {
288 pub kind: String,
289 pub tokens: f64,
290 pub last_refill_ms: i64,
291 pub window_start_ms: i64,
292 pub consumed: i64,
293}
294
295fn secs_to_duration(secs: f64) -> Duration {
298 if !secs.is_finite() || secs <= 0.0 {
299 return Duration::zero();
300 }
301 let capped = secs.min(8_640_000.0);
303 Duration::milliseconds((capped * 1000.0) as i64)
304}
305
306#[derive(Default)]
316pub struct RateLeaseKernel {
317 leases: HashMap<String, RateLease>,
318}
319
320impl RateLeaseKernel {
321 pub fn new() -> Self {
322 Self::default()
323 }
324
325 pub fn register(&mut self, key: impl Into<String>, lease: RateLease) {
327 self.leases.insert(key.into(), lease);
328 }
329
330 pub fn contains(&self, key: &str) -> bool {
332 self.leases.contains_key(key)
333 }
334
335 pub fn try_acquire(&mut self, key: &str, now: DateTime<Utc>) -> AcquireOutcome {
339 match self.leases.get_mut(key) {
340 Some(lease) => lease.try_acquire(now),
341 None => AcquireOutcome::Granted,
342 }
343 }
344
345 pub fn try_acquire_all(&mut self, keys: &[String], now: DateTime<Utc>) -> AcquireOutcome {
353 let mut latest_retry: Option<DateTime<Utc>> = None;
355 for key in keys {
356 if let Some(lease) = self.leases.get(key) {
357 if let AcquireOutcome::Denied { retry_at } = lease.peek(now) {
358 latest_retry = Some(match latest_retry {
359 Some(prev) if prev >= retry_at => prev,
360 _ => retry_at,
361 });
362 }
363 }
364 }
365 if let Some(retry_at) = latest_retry {
366 return AcquireOutcome::Denied { retry_at };
367 }
368 for key in keys {
370 if let Some(lease) = self.leases.get_mut(key) {
371 let _ = lease.try_acquire(now);
372 }
373 }
374 AcquireOutcome::Granted
375 }
376
377 pub fn available(&self, key: &str, now: DateTime<Utc>) -> Option<f64> {
379 self.leases.get(key).map(|l| l.available(now))
380 }
381
382 pub fn tick(&mut self, now: DateTime<Utc>) {
386 for lease in self.leases.values_mut() {
387 lease.refill(now);
388 }
389 }
390
391 pub fn snapshot(&self) -> Vec<(String, RateLeaseSnapshot)> {
394 let mut out: Vec<(String, RateLeaseSnapshot)> =
395 self.leases.iter().map(|(k, l)| (k.clone(), l.snapshot())).collect();
396 out.sort_by(|a, b| a.0.cmp(&b.0));
397 out
398 }
399
400 pub fn restore(&mut self, snaps: &[(String, RateLeaseSnapshot)]) {
403 for (key, snap) in snaps {
404 if let Some(lease) = self.leases.get_mut(key) {
405 lease.restore(snap);
406 }
407 }
408 }
409}
410
411#[derive(Debug, Clone, PartialEq)]
417pub enum GateDecision {
418 Allow,
420 Deny {
424 retry_at: DateTime<Utc>,
425 on_exhausted: String,
426 },
427}
428
429pub struct BudgetGate {
436 kernel: RateLeaseKernel,
437 on_exhausted: String,
438 by_effect: HashMap<String, Vec<String>>,
440}
441
442impl BudgetGate {
443 pub fn from_ir(budget: &crate::ir_nodes::IRBudget, scope: &str, now: DateTime<Utc>) -> Self {
447 let mut kernel = RateLeaseKernel::new();
448 let mut by_effect: HashMap<String, Vec<String>> = HashMap::new();
449 for (i, quota) in budget.quotas.iter().enumerate() {
450 let Some(lease) = RateLease::from_quota(quota, now) else {
451 continue;
452 };
453 let key = format!("{scope}:Tool({}):{}:{i}", quota.effect, quota.kind);
454 kernel.register(key.clone(), lease);
455 by_effect.entry(quota.effect.clone()).or_default().push(key);
456 }
457 BudgetGate {
458 kernel,
459 on_exhausted: if budget.on_exhausted.is_empty() {
460 "block".to_string()
461 } else {
462 budget.on_exhausted.clone()
463 },
464 by_effect,
465 }
466 }
467
468 pub fn merged_with(mut self, other: BudgetGate) -> BudgetGate {
484 for (key, lease) in other.kernel.leases {
485 self.kernel.leases.insert(key, lease);
486 }
487 for (effect, keys) in other.by_effect {
488 self.by_effect.entry(effect).or_default().extend(keys);
489 }
490 let strictness = |p: &str| match p {
492 "block" => 2,
493 "defer" => 1,
494 _ => 0, };
496 if strictness(&other.on_exhausted) > strictness(&self.on_exhausted) {
497 self.on_exhausted = other.on_exhausted;
498 }
499 self
500 }
501
502 pub fn gate(&mut self, effect: &str, now: DateTime<Utc>) -> GateDecision {
507 let Some(keys) = self.by_effect.get(effect) else {
508 return GateDecision::Allow;
509 };
510 let keys = keys.clone();
511 match self.kernel.try_acquire_all(&keys, now) {
512 AcquireOutcome::Granted => GateDecision::Allow,
513 AcquireOutcome::Denied { retry_at } => GateDecision::Deny {
514 retry_at,
515 on_exhausted: self.on_exhausted.clone(),
516 },
517 }
518 }
519
520 pub fn on_exhausted(&self) -> &str {
522 &self.on_exhausted
523 }
524
525 pub fn governs(&self, effect: &str) -> bool {
527 self.by_effect.contains_key(effect)
528 }
529
530 pub fn snapshot(&self) -> Vec<(String, RateLeaseSnapshot)> {
534 self.kernel.snapshot()
535 }
536
537 pub fn restore(&mut self, snaps: &[(String, RateLeaseSnapshot)]) {
539 self.kernel.restore(snaps);
540 }
541}
542
543#[cfg(test)]
544mod tests {
545 use super::*;
546
547 fn t0() -> DateTime<Utc> {
548 "2026-06-29T00:00:00Z".parse().unwrap()
549 }
550
551 fn quota(kind: &str, limit: i64, period: &str, effect: &str) -> IRBudgetQuota {
552 IRBudgetQuota {
553 kind: kind.into(),
554 limit,
555 period: period.into(),
556 effect: effect.into(),
557 }
558 }
559
560 fn ir_budget(quotas: Vec<IRBudgetQuota>, on_exhausted: &str) -> crate::ir_nodes::IRBudget {
561 crate::ir_nodes::IRBudget {
562 node_type: "budget",
563 source_line: 1,
564 source_column: 1,
565 name: String::new(),
566 quotas,
567 on_exhausted: on_exhausted.into(),
568 }
569 }
570
571 #[test]
574 fn period_parses_and_maps_to_seconds() {
575 assert_eq!(BudgetPeriod::parse("hour"), Some(BudgetPeriod::Hour));
576 assert_eq!(BudgetPeriod::parse("fortnight"), None);
577 assert_eq!(BudgetPeriod::Second.as_secs(), 1.0);
578 assert_eq!(BudgetPeriod::Minute.as_secs(), 60.0);
579 assert_eq!(BudgetPeriod::Hour.as_secs(), 3600.0);
580 assert_eq!(BudgetPeriod::Day.as_secs(), 86400.0);
581 }
582
583 #[test]
586 fn bucket_starts_full_and_grants_up_to_capacity() {
587 let now = t0();
588 let mut l = RateLease::rate("Telnyx", 3, BudgetPeriod::Hour, now);
589 assert!(l.try_acquire(now).is_granted());
591 assert!(l.try_acquire(now).is_granted());
592 assert!(l.try_acquire(now).is_granted());
593 match l.try_acquire(now) {
595 AcquireOutcome::Denied { retry_at } => {
596 assert_eq!(retry_at, now + Duration::seconds(1200));
598 }
599 other => panic!("expected Denied, got {other:?}"),
600 }
601 }
602
603 #[test]
604 fn bucket_refills_over_time() {
605 let now = t0();
606 let mut l = RateLease::rate("Telnyx", 2, BudgetPeriod::Hour, now);
607 assert!(l.try_acquire(now).is_granted());
609 assert!(l.try_acquire(now).is_granted());
610 assert!(!l.try_acquire(now).is_granted());
611 let later = now + Duration::seconds(1800);
613 assert!(l.try_acquire(later).is_granted());
614 assert!(!l.try_acquire(later).is_granted());
615 }
616
617 #[test]
618 fn bucket_refill_is_capped_at_capacity() {
619 let now = t0();
620 let mut l = RateLease::rate("Telnyx", 5, BudgetPeriod::Minute, now);
621 assert!(l.try_acquire(now).is_granted());
623 let way_later = now + Duration::days(1);
624 for _ in 0..5 {
626 assert!(l.try_acquire(way_later).is_granted());
627 }
628 assert!(!l.try_acquire(way_later).is_granted(), "capped at capacity");
629 }
630
631 #[test]
634 fn window_caps_at_limit_then_rolls() {
635 let now = t0();
636 let mut l = RateLease::max("Telnyx", 50, BudgetPeriod::Day, now);
637 for _ in 0..50 {
639 assert!(l.try_acquire(now).is_granted());
640 }
641 match l.try_acquire(now) {
643 AcquireOutcome::Denied { retry_at } => {
644 assert_eq!(retry_at, now + Duration::seconds(86400));
645 }
646 other => panic!("expected Denied, got {other:?}"),
647 }
648 assert!(!l.try_acquire(now + Duration::seconds(86399)).is_granted());
650 let next_day = now + Duration::seconds(86400);
652 assert!(l.try_acquire(next_day).is_granted());
653 assert_eq!(l.available(next_day), 49.0);
654 }
655
656 #[test]
657 fn window_has_no_intra_window_refill() {
658 let now = t0();
659 let mut l = RateLease::max("Telnyx", 2, BudgetPeriod::Hour, now);
660 assert!(l.try_acquire(now).is_granted());
661 assert!(l.try_acquire(now).is_granted());
662 assert!(!l.try_acquire(now + Duration::seconds(1800)).is_granted());
664 }
665
666 #[test]
669 fn from_quota_builds_the_right_kind() {
670 let now = t0();
671 let rate_q = IRBudgetQuota {
672 kind: "rate".into(),
673 limit: 8,
674 period: "hour".into(),
675 effect: "Telnyx".into(),
676 };
677 let max_q = IRBudgetQuota {
678 kind: "max".into(),
679 limit: 50,
680 period: "day".into(),
681 effect: "Telnyx".into(),
682 };
683 let rate = RateLease::from_quota(&rate_q, now).unwrap();
684 assert_eq!(rate.available(now), 8.0, "a rate bucket starts full");
685 let maxl = RateLease::from_quota(&max_q, now).unwrap();
686 assert_eq!(maxl.available(now), 50.0, "a max window starts with the full allowance");
687 let bad = IRBudgetQuota { period: "fortnight".into(), ..rate_q };
689 assert!(RateLease::from_quota(&bad, now).is_none());
690 }
691
692 #[test]
695 fn kernel_unregistered_key_is_unbudgeted() {
696 let mut k = RateLeaseKernel::new();
697 assert!(k.try_acquire("daemon:X:Tool(Y):rate", t0()).is_granted());
699 assert_eq!(k.available("daemon:X:Tool(Y):rate", t0()), None);
700 }
701
702 #[test]
703 fn kernel_enforces_a_registered_lease() {
704 let now = t0();
705 let mut k = RateLeaseKernel::new();
706 k.register("d:Out:Tool(Telnyx):rate", RateLease::rate("Telnyx", 1, BudgetPeriod::Hour, now));
707 assert!(k.try_acquire("d:Out:Tool(Telnyx):rate", now).is_granted());
708 assert!(!k.try_acquire("d:Out:Tool(Telnyx):rate", now).is_granted());
709 let later = now + Duration::seconds(3600);
711 assert!(k.try_acquire("d:Out:Tool(Telnyx):rate", later).is_granted());
712 }
713
714 #[test]
715 fn kernel_tick_refreshes_available_without_consuming() {
716 let now = t0();
717 let mut k = RateLeaseKernel::new();
718 k.register("k", RateLease::rate("E", 4, BudgetPeriod::Minute, now));
719 for _ in 0..4 {
721 assert!(k.try_acquire("k", now).is_granted());
722 }
723 assert_eq!(k.available("k", now), Some(0.0));
724 let later = now + Duration::seconds(30);
726 k.tick(later);
727 assert_eq!(k.available("k", later), Some(2.0));
728 }
729
730 #[test]
733 fn acquire_all_is_all_or_none() {
734 let now = t0();
735 let mut k = RateLeaseKernel::new();
736 k.register("r", RateLease::rate("E", 5, BudgetPeriod::Hour, now));
738 k.register("m", RateLease::max("E", 1, BudgetPeriod::Day, now));
739 let keys = vec!["r".to_string(), "m".to_string()];
740 assert!(k.try_acquire_all(&keys, now).is_granted());
742 match k.try_acquire_all(&keys, now) {
745 AcquireOutcome::Denied { retry_at } => {
746 assert_eq!(retry_at, now + Duration::seconds(86400), "binding = the daily max");
747 }
748 other => panic!("expected Denied, got {other:?}"),
749 }
750 assert_eq!(k.available("r", now), Some(4.0), "rate token not consumed on denial");
751 }
752
753 #[test]
754 fn acquire_all_empty_keys_is_granted() {
755 let mut k = RateLeaseKernel::new();
756 assert!(k.try_acquire_all(&[], t0()).is_granted());
757 }
758
759 #[test]
762 fn gate_allows_unbudgeted_effects() {
763 let now = t0();
764 let b = ir_budget(vec![quota("rate", 1, "hour", "Telnyx")], "block");
765 let mut gate = BudgetGate::from_ir(&b, "daemon:Out", now);
766 assert_eq!(gate.gate("SomeOtherTool", now), GateDecision::Allow);
768 assert!(!gate.governs("SomeOtherTool"));
769 assert!(gate.governs("Telnyx"));
770 }
771
772 #[test]
773 fn gate_enforces_then_denies_with_policy() {
774 let now = t0();
775 let b = ir_budget(
776 vec![
777 quota("rate", 2, "hour", "Telnyx"),
778 quota("max", 3, "day", "Telnyx"),
779 ],
780 "defer",
781 );
782 let mut gate = BudgetGate::from_ir(&b, "daemon:Out", now);
783 assert_eq!(gate.gate("Telnyx", now), GateDecision::Allow);
785 assert_eq!(gate.gate("Telnyx", now), GateDecision::Allow);
786 match gate.gate("Telnyx", now) {
787 GateDecision::Deny { on_exhausted, retry_at } => {
788 assert_eq!(on_exhausted, "defer");
789 assert_eq!(retry_at, now + Duration::seconds(1800));
791 }
792 other => panic!("expected Deny, got {other:?}"),
793 }
794 }
795
796 #[test]
797 fn gate_omitted_policy_is_block() {
798 let now = t0();
799 let b = ir_budget(vec![quota("rate", 1, "hour", "E")], "");
800 let gate = BudgetGate::from_ir(&b, "d", now);
801 assert_eq!(gate.on_exhausted(), "block");
802 }
803
804 #[test]
807 fn snapshot_restore_carries_max_window_across_ticks() {
808 let now = t0();
809 let b = ir_budget(vec![quota("max", 3, "day", "Telnyx")], "block");
810 let mut g1 = BudgetGate::from_ir(&b, "d", now);
812 assert_eq!(g1.gate("Telnyx", now), GateDecision::Allow);
813 assert_eq!(g1.gate("Telnyx", now), GateDecision::Allow);
814 let snap = g1.snapshot();
815
816 let mut g2 = BudgetGate::from_ir(&b, "d", now + Duration::minutes(5));
819 g2.restore(&snap);
820 assert_eq!(g2.gate("Telnyx", now + Duration::minutes(5)), GateDecision::Allow);
821 match g2.gate("Telnyx", now + Duration::minutes(5)) {
824 GateDecision::Deny { .. } => {}
825 other => panic!("expected the daily cap to hold across ticks, got {other:?}"),
826 }
827 }
828
829 #[test]
830 fn snapshot_round_trips_a_bucket() {
831 let now = t0();
832 let mut l = RateLease::rate("E", 8, BudgetPeriod::Hour, now);
833 l.try_acquire(now); let snap = l.snapshot();
835 assert_eq!(snap.kind, "rate");
836 let mut l2 = RateLease::rate("E", 8, BudgetPeriod::Hour, now);
837 l2.restore(&snap);
838 assert_eq!(l2.available(now), 7.0, "restored bucket carries the consumed token");
839 }
840}