1use std::collections::HashMap;
27use std::sync::Mutex;
28use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
29
30use anyhow::{Context, Result};
31use tracing::warn;
32
33use crate::config::{BudgetCfg, LlmCfg};
34
35#[derive(Debug, Clone, Copy, PartialEq, Eq)]
37pub enum BudgetUnit {
38 Tokens,
40 UsdMicros,
42}
43
44impl BudgetUnit {
45 fn parse(s: &str) -> Result<BudgetUnit> {
46 match s.trim().to_ascii_lowercase().as_str() {
47 "tokens" | "token" | "" => Ok(BudgetUnit::Tokens),
48 "usd" | "usd_micros" | "cost" => Ok(BudgetUnit::UsdMicros),
49 other => anyhow::bail!("invalid llm budget unit {other:?} (expected tokens|usd)"),
50 }
51 }
52}
53
54#[derive(Debug, Clone, Copy, PartialEq, Eq)]
56pub enum BudgetScope {
57 Global,
59 PerKey,
61 PerModel,
63 PerTeam,
66}
67
68impl BudgetScope {
69 fn parse(s: &str) -> Result<BudgetScope> {
70 match s.trim().to_ascii_lowercase().as_str() {
71 "global" | "" => Ok(BudgetScope::Global),
72 "key" | "per_key" | "per-key" => Ok(BudgetScope::PerKey),
73 "model" | "per_model" | "per-model" => Ok(BudgetScope::PerModel),
74 "team" | "per_team" | "per-team" | "tag" => Ok(BudgetScope::PerTeam),
75 other => {
76 anyhow::bail!("invalid llm budget scope {other:?} (expected global|key|model|team)")
77 }
78 }
79 }
80
81 pub fn label(&self) -> &'static str {
84 match self {
85 BudgetScope::Global => "global",
86 BudgetScope::PerKey => "key",
87 BudgetScope::PerModel => "model",
88 BudgetScope::PerTeam => "team",
89 }
90 }
91}
92
93#[derive(Debug, Clone)]
95struct Budget {
96 name: String,
97 scope: BudgetScope,
98 unit: BudgetUnit,
99 limit: u64,
100 window_secs: u64,
101}
102
103impl Budget {
104 fn build(cfg: &BudgetCfg) -> Result<Budget> {
107 anyhow::ensure!(
108 !cfg.name.trim().is_empty(),
109 "llm budget name must not be empty"
110 );
111 let unit = BudgetUnit::parse(&cfg.unit)?;
112 let scope = BudgetScope::parse(&cfg.scope)?;
113 let window_secs = crate::config::parse_duration(&cfg.window)
114 .with_context(|| format!("llm budget {:?} window", cfg.name))?
115 .as_secs();
116 anyhow::ensure!(
117 window_secs > 0,
118 "llm budget {:?} window must be > 0",
119 cfg.name
120 );
121 anyhow::ensure!(
122 cfg.limit.is_finite() && cfg.limit > 0.0,
123 "llm budget {:?} limit must be > 0",
124 cfg.name
125 );
126 let limit = match unit {
127 BudgetUnit::Tokens => cfg.limit.round() as u64,
128 BudgetUnit::UsdMicros => (cfg.limit * 1_000_000.0).round() as u64,
129 };
130 anyhow::ensure!(
131 limit > 0,
132 "llm budget {:?} limit rounds to zero; use a larger value",
133 cfg.name
134 );
135 Ok(Budget {
136 name: cfg.name.clone(),
137 scope,
138 unit,
139 limit,
140 window_secs,
141 })
142 }
143
144 fn key(&self, prefix: &str, dims: &Dims, now_secs: u64) -> String {
148 let dim = match self.scope {
149 BudgetScope::Global => "_",
150 BudgetScope::PerKey => dims.principal.unwrap_or("_anon"),
151 BudgetScope::PerModel => dims.model,
152 BudgetScope::PerTeam => dims.team.unwrap_or("_none"),
153 };
154 let window = now_secs / self.window_secs;
155 format!("{prefix}:budget:{}:{dim}:{window}", self.name)
156 }
157
158 fn amount(&self, tokens: u64, cost_micros: u64) -> u64 {
160 match self.unit {
161 BudgetUnit::Tokens => tokens,
162 BudgetUnit::UsdMicros => cost_micros,
163 }
164 }
165
166 fn ttl_ms(&self) -> u64 {
169 self.window_secs.saturating_mul(2_000).max(1)
170 }
171}
172
173#[derive(Debug, Clone, Copy, Default)]
175pub struct Dims<'a> {
176 pub principal: Option<&'a str>,
178 pub model: &'a str,
180 pub team: Option<&'a str>,
182}
183
184fn would_reserve(used: u64, amount: u64, limit: u64) -> bool {
187 used.saturating_add(amount) <= limit
188}
189
190#[derive(Debug, Clone, Copy)]
194struct ReserveOutcome {
195 admitted: bool,
196 used_after: u64,
197}
198
199enum Store {
201 Memory(MemoryStore),
202 Redis(Box<RedisStore>),
203}
204
205impl Store {
206 async fn reserve(
209 &self,
210 key: &str,
211 amount: u64,
212 limit: u64,
213 ttl_ms: u64,
214 ) -> Result<ReserveOutcome> {
215 match self {
216 Store::Memory(s) => Ok(s.reserve(key, amount, limit, ttl_ms)),
217 Store::Redis(s) => s.reserve(key, amount, limit, ttl_ms).await,
218 }
219 }
220
221 async fn reconcile(&self, key: &str, delta: i64, ttl_ms: u64, marker: &str) -> bool {
226 match self {
227 Store::Memory(s) => {
228 s.reconcile(key, delta, ttl_ms, marker);
229 true
230 }
231 Store::Redis(s) => match s.reconcile(key, delta, ttl_ms, marker).await {
232 Ok(()) => true,
233 Err(e) => {
234 warn!(error = %format!("{e:#}"), "llm budget reconcile failed after retries — counter has drifted (alert on edgeguard_llm_budget_reconcile_failures_total)");
235 false
236 }
237 },
238 }
239 }
240}
241
242#[derive(Default)]
247struct MemoryStore {
248 used: Mutex<HashMap<String, (u64, Instant)>>,
249 settled: Mutex<HashMap<String, Instant>>,
252}
253
254impl MemoryStore {
255 fn reserve(&self, key: &str, amount: u64, limit: u64, ttl_ms: u64) -> ReserveOutcome {
256 let now = Instant::now();
257 let expires = now + Duration::from_millis(ttl_ms);
258 let mut map = self.used.lock().expect("budget store mutex poisoned");
259 map.retain(|_, (_, exp)| *exp > now);
260 let used = map.get(key).map(|(v, _)| *v).unwrap_or(0);
261 if would_reserve(used, amount, limit) {
262 let used_after = used.saturating_add(amount);
263 map.insert(key.to_string(), (used_after, expires));
264 ReserveOutcome {
265 admitted: true,
266 used_after,
267 }
268 } else {
269 ReserveOutcome {
270 admitted: false,
271 used_after: used,
272 }
273 }
274 }
275
276 fn reconcile(&self, key: &str, delta: i64, ttl_ms: u64, marker: &str) {
277 let now = Instant::now();
278 {
281 let mut seen = self.settled.lock().expect("budget settled mutex poisoned");
282 seen.retain(|_, exp| *exp > now);
283 if seen.contains_key(marker) {
284 return;
285 }
286 seen.insert(
287 marker.to_string(),
288 now + Duration::from_millis(MARKER_TTL_MS),
289 );
290 }
291 let expires = now + Duration::from_millis(ttl_ms);
292 let mut map = self.used.lock().expect("budget store mutex poisoned");
293 map.retain(|_, (_, exp)| *exp > now);
294 let used = map.get(key).map(|(v, _)| *v).unwrap_or(0) as i64;
295 let next = (used + delta).max(0) as u64;
296 map.insert(key.to_string(), (next, expires));
297 }
298}
299
300const RESERVE_LUA: &str = r#"
304local used = tonumber(redis.call('GET', KEYS[1]) or '0')
305local amount = tonumber(ARGV[1])
306local limit = tonumber(ARGV[2])
307local ttl = tonumber(ARGV[3])
308if used + amount > limit then
309 return {0, used}
310end
311local newv = redis.call('INCRBY', KEYS[1], amount)
312redis.call('PEXPIRE', KEYS[1], ttl)
313return {1, newv}
314"#;
315
316const RECONCILE_LUA: &str = r#"
324if redis.call('SET', KEYS[2], '1', 'NX', 'PX', tonumber(ARGV[3])) == false then
325 return tonumber(redis.call('GET', KEYS[1]) or '0')
326end
327local new = redis.call('INCRBY', KEYS[1], tonumber(ARGV[1]))
328if new < 0 then
329 redis.call('SET', KEYS[1], 0)
330 new = 0
331end
332redis.call('PEXPIRE', KEYS[1], tonumber(ARGV[2]))
333return new
334"#;
335
336const MARKER_TTL_MS: u64 = 300_000;
339const RECONCILE_ATTEMPTS: u32 = 3;
342const RECONCILE_BACKOFF: Duration = Duration::from_millis(25);
343
344async fn retry_async<F, Fut, T>(attempts: u32, backoff: Duration, mut op: F) -> Result<T>
347where
348 F: FnMut() -> Fut,
349 Fut: std::future::Future<Output = Result<T>>,
350{
351 let attempts = attempts.max(1);
352 let mut last: Option<anyhow::Error> = None;
353 for i in 0..attempts {
354 match op().await {
355 Ok(v) => return Ok(v),
356 Err(e) => {
357 last = Some(e);
358 if i + 1 < attempts {
359 tokio::time::sleep(backoff).await;
360 }
361 }
362 }
363 }
364 Err(last.expect("retry_async ran at least one attempt"))
365}
366
367struct RedisStore {
370 client: redis::Client,
371 conn: tokio::sync::OnceCell<redis::aio::ConnectionManager>,
372 reserve: redis::Script,
373 reconcile: redis::Script,
374}
375
376impl RedisStore {
377 fn new(url: &str) -> Result<RedisStore> {
378 anyhow::ensure!(
379 !url.trim().is_empty(),
380 "llm.redis_url is required when llm.store = \"redis\""
381 );
382 let client = redis::Client::open(url)
383 .with_context(|| format!("opening redis client for {url:?} (llm.redis_url)"))?;
384 Ok(RedisStore {
385 client,
386 conn: tokio::sync::OnceCell::new(),
387 reserve: redis::Script::new(RESERVE_LUA),
388 reconcile: redis::Script::new(RECONCILE_LUA),
389 })
390 }
391
392 async fn manager(&self) -> Result<redis::aio::ConnectionManager> {
393 self.conn
394 .get_or_try_init(|| redis::aio::ConnectionManager::new(self.client.clone()))
395 .await
396 .context("connecting to redis llm-budget store")
397 .cloned()
398 }
399
400 async fn reserve(
401 &self,
402 key: &str,
403 amount: u64,
404 limit: u64,
405 ttl_ms: u64,
406 ) -> Result<ReserveOutcome> {
407 let mut conn = self.manager().await?;
408 let (admitted, used_after): (i64, i64) = self
409 .reserve
410 .key(key)
411 .arg(amount)
412 .arg(limit)
413 .arg(ttl_ms)
414 .invoke_async(&mut conn)
415 .await
416 .context("evaluating redis budget reserve script")?;
417 Ok(ReserveOutcome {
418 admitted: admitted == 1,
419 used_after: used_after.max(0) as u64,
420 })
421 }
422
423 async fn reconcile(&self, key: &str, delta: i64, ttl_ms: u64, marker: &str) -> Result<()> {
424 retry_async(RECONCILE_ATTEMPTS, RECONCILE_BACKOFF, || async {
427 let mut conn = self.manager().await?;
428 let _: i64 = self
429 .reconcile
430 .key(key)
431 .key(marker)
432 .arg(delta)
433 .arg(ttl_ms)
434 .arg(MARKER_TTL_MS)
435 .invoke_async(&mut conn)
436 .await
437 .context("evaluating redis budget reconcile script")?;
438 Ok(())
439 })
440 .await
441 }
442}
443
444#[derive(Debug, Clone, Copy, Default)]
446pub struct Spend {
447 pub tokens: u64,
448 pub cost_micros: u64,
449}
450
451#[derive(Debug, Clone)]
453pub struct Observation {
454 pub name: String,
456 pub consumed_ratio: f64,
458}
459
460#[derive(Debug, Clone)]
463pub struct Denial {
464 pub name: String,
465 pub scope: BudgetScope,
466 pub unit: BudgetUnit,
467 pub rollback_failures: usize,
471}
472
473#[derive(Debug, Default)]
477pub struct Reservation {
478 id: String,
481 held: Vec<(String, u64, usize)>,
483 observations: Vec<Observation>,
485 rollback_failures: usize,
489}
490
491impl Reservation {
492 pub fn is_empty(&self) -> bool {
493 self.held.is_empty()
494 }
495
496 pub fn observations(&self) -> &[Observation] {
498 &self.observations
499 }
500
501 pub fn rollback_failures(&self) -> usize {
504 self.rollback_failures
505 }
506}
507
508#[derive(Debug)]
510pub enum Reserved {
511 Ok(Reservation),
513 Denied(Denial),
516 Error { rollback_failures: usize },
520}
521
522pub struct BudgetEngine {
525 budgets: Vec<Budget>,
526 store: Store,
527 prefix: String,
528 fail_open: bool,
529}
530
531impl BudgetEngine {
532 pub fn build(cfg: &LlmCfg) -> Result<Option<BudgetEngine>> {
535 if cfg.budgets.is_empty() {
536 return Ok(None);
537 }
538 let store = match crate::limiter::StoreMode::parse(&cfg.store)? {
539 crate::limiter::StoreMode::Local | crate::limiter::StoreMode::Memory => {
542 Store::Memory(MemoryStore::default())
543 }
544 crate::limiter::StoreMode::Redis => {
545 Store::Redis(Box::new(RedisStore::new(&cfg.redis_url)?))
546 }
547 };
548 let budgets = cfg
549 .budgets
550 .iter()
551 .map(Budget::build)
552 .collect::<Result<Vec<_>>>()?;
553 Ok(Some(BudgetEngine {
554 budgets,
555 store,
556 prefix: if cfg.redis_prefix.trim().is_empty() {
557 "edgeguard".to_string()
558 } else {
559 cfg.redis_prefix.clone()
560 },
561 fail_open: cfg.fail_open,
562 }))
563 }
564
565 pub async fn reserve(&self, dims: Dims<'_>, estimate: Spend) -> Reserved {
570 let now = now_secs();
571 let mut held = Vec::new();
572 let mut observations = Vec::new();
573 for (idx, budget) in self.budgets.iter().enumerate() {
574 let amount = budget.amount(estimate.tokens, estimate.cost_micros);
575 if amount == 0 {
578 continue;
579 }
580 let key = budget.key(&self.prefix, &dims, now);
581 match self
582 .store
583 .reserve(&key, amount, budget.limit, budget.ttl_ms())
584 .await
585 {
586 Ok(outcome) if outcome.admitted => {
587 held.push((key, amount, idx));
588 observations.push(Observation {
589 name: budget.name.clone(),
590 consumed_ratio: ratio(outcome.used_after, budget.limit),
591 });
592 }
593 Ok(_) => {
594 let rollback_failures = self.rollback(&held).await;
595 return Reserved::Denied(Denial {
596 name: budget.name.clone(),
597 scope: budget.scope,
598 unit: budget.unit,
599 rollback_failures,
600 });
601 }
602 Err(e) => {
603 let rollback_failures = self.rollback(&held).await;
604 if self.fail_open {
605 warn!(error = %format!("{e:#}"), budget = %budget.name, "llm budget store error; failing open (allowing request)");
606 return Reserved::Ok(Reservation {
607 rollback_failures,
608 ..Reservation::default()
609 });
610 }
611 warn!(error = %format!("{e:#}"), budget = %budget.name, "llm budget store error; failing closed (503)");
612 return Reserved::Error { rollback_failures };
613 }
614 }
615 }
616 Reserved::Ok(Reservation {
617 id: uuid::Uuid::new_v4().to_string(),
618 held,
619 observations,
620 rollback_failures: 0,
621 })
622 }
623
624 pub async fn reconcile(&self, reservation: &Reservation, actual: Spend) -> usize {
630 let mut failures = 0usize;
631 for (key, reserved, idx) in &reservation.held {
632 let budget = &self.budgets[*idx];
633 let actual_amount = budget.amount(actual.tokens, actual.cost_micros);
634 let delta = actual_amount as i64 - *reserved as i64;
635 if delta != 0 {
636 let marker = format!("{key}:s:{}", reservation.id);
638 if !self
639 .store
640 .reconcile(key, delta, budget.ttl_ms(), &marker)
641 .await
642 {
643 failures += 1;
644 }
645 }
646 }
647 failures
648 }
649
650 pub async fn release(&self, reservation: &Reservation) -> usize {
653 self.reconcile(reservation, Spend::default()).await
654 }
655
656 async fn rollback(&self, held: &[(String, u64, usize)]) -> usize {
661 let mut failures = 0usize;
662 for (key, amount, idx) in held {
663 let budget = &self.budgets[*idx];
664 let marker = format!("{key}:rb:{}", uuid::Uuid::new_v4());
665 if !self
666 .store
667 .reconcile(key, -(*amount as i64), budget.ttl_ms(), &marker)
668 .await
669 {
670 failures += 1;
671 }
672 }
673 failures
674 }
675}
676
677fn ratio(used: u64, limit: u64) -> f64 {
680 if limit == 0 {
681 return 0.0;
682 }
683 used as f64 / limit as f64
684}
685
686fn now_secs() -> u64 {
688 SystemTime::now()
689 .duration_since(UNIX_EPOCH)
690 .map(|d| d.as_secs())
691 .unwrap_or(0)
692}
693
694#[cfg(test)]
695mod tests {
696 use super::*;
697 use crate::config::BudgetCfg;
698
699 fn token_budget(limit: f64, window: &str) -> BudgetCfg {
700 BudgetCfg {
701 name: "test".into(),
702 scope: "key".into(),
703 unit: "tokens".into(),
704 limit,
705 window: window.into(),
706 }
707 }
708
709 fn engine(budgets: Vec<BudgetCfg>) -> BudgetEngine {
710 BudgetEngine::build(&LlmCfg {
711 enabled: true,
712 budgets,
713 store: "memory".into(),
714 ..Default::default()
715 })
716 .unwrap()
717 .expect("budgets configured")
718 }
719
720 fn dims<'a>(principal: Option<&'a str>, model: &'a str) -> Dims<'a> {
722 Dims {
723 principal,
724 model,
725 team: None,
726 }
727 }
728
729 #[test]
730 fn would_reserve_caps_at_limit() {
731 assert!(would_reserve(0, 100, 100));
732 assert!(would_reserve(90, 10, 100));
733 assert!(!would_reserve(90, 11, 100));
734 assert!(!would_reserve(u64::MAX, 1, 100));
736 }
737
738 #[test]
739 fn unit_and_scope_parse() {
740 assert_eq!(BudgetUnit::parse("tokens").unwrap(), BudgetUnit::Tokens);
741 assert_eq!(BudgetUnit::parse("USD").unwrap(), BudgetUnit::UsdMicros);
742 assert!(BudgetUnit::parse("bananas").is_err());
743 assert_eq!(BudgetScope::parse("global").unwrap(), BudgetScope::Global);
744 assert_eq!(BudgetScope::parse("per-key").unwrap(), BudgetScope::PerKey);
745 assert!(BudgetScope::parse("galaxy").is_err());
746 }
747
748 #[test]
749 fn build_is_none_without_budgets() {
750 let none = BudgetEngine::build(&LlmCfg::default()).unwrap();
751 assert!(none.is_none());
752 }
753
754 #[test]
755 fn limit_rounding_to_zero_is_rejected() {
756 assert!(Budget::build(&BudgetCfg {
758 name: "tiny".into(),
759 scope: "global".into(),
760 unit: "tokens".into(),
761 limit: 0.3,
762 window: "1h".into(),
763 })
764 .is_err());
765 }
766
767 #[test]
768 fn usd_budget_compiles_to_micros() {
769 let b = Budget::build(&BudgetCfg {
770 name: "spend".into(),
771 scope: "global".into(),
772 unit: "usd".into(),
773 limit: 2.50,
774 window: "24h".into(),
775 })
776 .unwrap();
777 assert_eq!(b.limit, 2_500_000); assert_eq!(b.unit, BudgetUnit::UsdMicros);
779 }
780
781 #[tokio::test]
782 async fn retry_async_succeeds_after_transient_failures() {
783 use std::sync::atomic::{AtomicU32, Ordering};
784 let calls = AtomicU32::new(0);
785 let r: Result<u32> = retry_async(3, Duration::from_millis(0), || {
786 let n = calls.fetch_add(1, Ordering::SeqCst);
787 async move {
788 if n < 2 {
789 anyhow::bail!("transient blip")
790 } else {
791 Ok(n)
792 }
793 }
794 })
795 .await;
796 assert_eq!(r.unwrap(), 2);
797 assert_eq!(calls.load(Ordering::SeqCst), 3);
798 }
799
800 #[tokio::test]
801 async fn retry_async_gives_up_after_exhausting_attempts() {
802 use std::sync::atomic::{AtomicU32, Ordering};
803 let calls = AtomicU32::new(0);
804 let r: Result<()> = retry_async(2, Duration::from_millis(0), || {
805 calls.fetch_add(1, Ordering::SeqCst);
806 async move { anyhow::bail!("always fails") }
807 })
808 .await;
809 assert!(r.is_err());
810 assert_eq!(calls.load(Ordering::SeqCst), 2);
811 }
812
813 #[test]
814 fn memory_reconcile_is_idempotent_per_marker() {
815 let store = MemoryStore::default();
818 assert!(store.reserve("k", 100, 1000, 60_000).admitted); store.reconcile("k", -40, 60_000, "settle-1"); store.reconcile("k", -40, 60_000, "settle-1"); assert_eq!(store.reserve("k", 0, 1000, 60_000).used_after, 60);
822 store.reconcile("k", -10, 60_000, "settle-2");
824 assert_eq!(store.reserve("k", 0, 1000, 60_000).used_after, 50);
825 }
826
827 #[tokio::test]
828 async fn reserve_admits_until_limit_then_denies() {
829 let eng = engine(vec![token_budget(100.0, "1h")]);
830 let est = Spend {
831 tokens: 60,
832 cost_micros: 0,
833 };
834 assert!(matches!(
836 eng.reserve(dims(Some("alice"), "gpt-4o"), est).await,
837 Reserved::Ok(_)
838 ));
839 assert!(matches!(
840 eng.reserve(dims(Some("alice"), "gpt-4o"), est).await,
841 Reserved::Denied(d) if d.name == "test"
842 ));
843 assert!(matches!(
845 eng.reserve(dims(Some("bob"), "gpt-4o"), est).await,
846 Reserved::Ok(_)
847 ));
848 }
849
850 #[tokio::test]
851 async fn reconcile_releases_overestimate() {
852 let eng = engine(vec![token_budget(100.0, "1h")]);
853 let reserved = match eng
855 .reserve(
856 dims(Some("c"), "m"),
857 Spend {
858 tokens: 80,
859 cost_micros: 0,
860 },
861 )
862 .await
863 {
864 Reserved::Ok(r) => r,
865 other => panic!("expected Ok, got {other:?}"),
866 };
867 eng.reconcile(
868 &reserved,
869 Spend {
870 tokens: 30,
871 cost_micros: 0,
872 },
873 )
874 .await;
875 assert!(matches!(
877 eng.reserve(
878 dims(Some("c"), "m"),
879 Spend {
880 tokens: 70,
881 cost_micros: 0
882 }
883 )
884 .await,
885 Reserved::Ok(_)
886 ));
887 }
888
889 #[tokio::test]
890 async fn release_returns_full_reservation() {
891 let eng = engine(vec![token_budget(100.0, "1h")]);
892 let reserved = match eng
893 .reserve(
894 dims(Some("d"), "m"),
895 Spend {
896 tokens: 100,
897 cost_micros: 0,
898 },
899 )
900 .await
901 {
902 Reserved::Ok(r) => r,
903 other => panic!("expected Ok, got {other:?}"),
904 };
905 assert!(matches!(
907 eng.reserve(
908 dims(Some("d"), "m"),
909 Spend {
910 tokens: 1,
911 cost_micros: 0
912 }
913 )
914 .await,
915 Reserved::Denied(_)
916 ));
917 eng.release(&reserved).await;
919 assert!(matches!(
920 eng.reserve(
921 dims(Some("d"), "m"),
922 Spend {
923 tokens: 100,
924 cost_micros: 0
925 }
926 )
927 .await,
928 Reserved::Ok(_)
929 ));
930 }
931
932 #[tokio::test]
933 async fn multi_budget_denial_rolls_back_prior_reserve() {
934 let eng = engine(vec![
937 BudgetCfg {
938 name: "tok".into(),
939 scope: "global".into(),
940 unit: "tokens".into(),
941 limit: 1000.0,
942 window: "1h".into(),
943 },
944 BudgetCfg {
945 name: "cost".into(),
946 scope: "global".into(),
947 unit: "usd".into(),
948 limit: 0.000010, window: "1h".into(),
950 },
951 ]);
952 assert!(matches!(
954 eng.reserve(dims(None, "m"), Spend { tokens: 100, cost_micros: 20 }).await,
955 Reserved::Denied(d) if d.name == "cost" && d.unit == BudgetUnit::UsdMicros
956 ));
957 assert!(matches!(
959 eng.reserve(
960 dims(None, "m"),
961 Spend {
962 tokens: 1000,
963 cost_micros: 0
964 }
965 )
966 .await,
967 Reserved::Ok(_)
968 ));
969 }
970
971 #[tokio::test]
972 async fn per_team_scope_is_keyed_by_team() {
973 let eng = engine(vec![BudgetCfg {
974 name: "team-cap".into(),
975 scope: "team".into(),
976 unit: "tokens".into(),
977 limit: 100.0,
978 window: "1h".into(),
979 }]);
980 let est = Spend {
981 tokens: 60,
982 cost_micros: 0,
983 };
984 let team_a = Dims {
985 principal: Some("alice"),
986 model: "gpt-4o",
987 team: Some("team-a"),
988 };
989 assert!(matches!(eng.reserve(team_a, est).await, Reserved::Ok(_)));
992 let team_a_bob = Dims {
993 principal: Some("bob"),
994 model: "gpt-4o",
995 team: Some("team-a"),
996 };
997 assert!(matches!(
998 eng.reserve(team_a_bob, est).await,
999 Reserved::Denied(d) if d.scope == BudgetScope::PerTeam
1000 ));
1001 let team_b = Dims {
1003 principal: Some("alice"),
1004 model: "gpt-4o",
1005 team: Some("team-b"),
1006 };
1007 assert!(matches!(eng.reserve(team_b, est).await, Reserved::Ok(_)));
1008 }
1009
1010 #[tokio::test]
1011 async fn reserve_reports_consumed_ratio() {
1012 let eng = engine(vec![token_budget(100.0, "1h")]);
1013 let r = match eng
1014 .reserve(
1015 dims(Some("alice"), "m"),
1016 Spend {
1017 tokens: 75,
1018 cost_micros: 0,
1019 },
1020 )
1021 .await
1022 {
1023 Reserved::Ok(r) => r,
1024 other => panic!("expected Ok, got {other:?}"),
1025 };
1026 let obs = r.observations();
1027 assert_eq!(obs.len(), 1);
1028 assert_eq!(obs[0].name, "test");
1029 assert!(
1030 (obs[0].consumed_ratio - 0.75).abs() < 1e-9,
1031 "{}",
1032 obs[0].consumed_ratio
1033 );
1034 }
1035
1036 fn redis_url() -> String {
1042 std::env::var("EDGEGUARD_TEST_REDIS_URL")
1043 .unwrap_or_else(|_| "redis://127.0.0.1:6379".into())
1044 }
1045
1046 #[tokio::test]
1047 #[ignore = "requires a live Redis (EDGEGUARD_TEST_REDIS_URL, default redis://127.0.0.1:6379)"]
1048 async fn redis_budget_reserve_and_reconcile_live() {
1049 let eng = BudgetEngine::build(&LlmCfg {
1050 enabled: true,
1051 store: "redis".into(),
1052 redis_url: redis_url(),
1053 redis_prefix: format!("egtest:budget:{}:{}", std::process::id(), now_secs()),
1054 budgets: vec![token_budget(100.0, "1h")],
1055 ..Default::default()
1056 })
1057 .unwrap()
1058 .expect("budgets configured");
1059
1060 let est = Spend {
1061 tokens: 60,
1062 cost_micros: 0,
1063 };
1064 let reserved = match eng.reserve(dims(Some("alice"), "m"), est).await {
1066 Reserved::Ok(r) => r,
1067 Reserved::Error { .. } => {
1068 eprintln!("skipping redis_budget_reserve_and_reconcile_live: Redis unreachable");
1069 return;
1070 }
1071 other => panic!("unexpected first reserve: {other:?}"),
1072 };
1073 assert!(matches!(
1075 eng.reserve(dims(Some("alice"), "m"), est).await,
1076 Reserved::Denied(_)
1077 ));
1078 eng.reconcile(
1080 &reserved,
1081 Spend {
1082 tokens: 10,
1083 cost_micros: 0,
1084 },
1085 )
1086 .await;
1087 assert!(matches!(
1088 eng.reserve(dims(Some("alice"), "m"), est).await,
1089 Reserved::Ok(_)
1090 ));
1091 }
1092}