1use crate::sync::{Arc, AtomicI64, AtomicU64, Mutex, Ordering, RwLock};
28use std::collections::HashMap;
29
30#[derive(Clone, Debug, PartialEq)]
32pub enum BudgetReservation {
33 Reserved {
34 remaining_microcents: i64,
35 },
36 Insufficient {
37 remaining_microcents: i64,
38 required_microcents: i64,
39 },
40 MissingTenant,
41 MissingReservation,
42 ExposureLimitExceeded {
43 current_reserved_microcents: i64,
44 max_reserved_microcents: i64,
45 },
46 Overflow {
47 current_reserved_microcents: i64,
48 },
49}
50
51#[derive(Clone, Debug, PartialEq, Eq)]
53pub enum TopUpResult {
54 ToppedUp {
55 added_microcents: i64,
56 new_initial_microcents: i64,
57 remaining_microcents: i64,
58 },
59 MissingTenant,
60 InvalidAmount,
61 Overflow,
62}
63
64#[derive(Clone, Debug, PartialEq, Eq, thiserror::Error)]
66pub enum BudgetConfigurationError {
67 #[error("{field} must be >= 0, got {value}")]
68 NegativeAmount { field: &'static str, value: i64 },
69}
70
71#[derive(Clone, Debug, PartialEq)]
73pub enum BudgetSettlement {
74 Committed {
75 remaining_microcents: i64,
76 actual_microcents: i64,
77 },
78 Released {
79 remaining_microcents: i64,
80 returned_microcents: i64,
81 },
82 Overrun {
83 remaining_microcents: i64,
84 },
85 InvalidAmount,
86 MissingReservation,
87 MissingTenant,
88 Overflow {
89 remaining_microcents: i64,
90 },
91}
92
93#[derive(Debug)]
94struct ReservationRecord {
95 tenant_id: Arc<str>,
96 reserved_microcents: i64,
97}
98
99#[inline]
104fn debit_if_available(budget: &AtomicI64, amount: i64) -> Result<i64, i64> {
105 let mut current = budget.load(Ordering::Acquire);
106 loop {
107 if current < amount {
108 return Err(current);
109 }
110 match budget.compare_exchange_weak(
111 current,
112 current - amount,
113 Ordering::AcqRel,
114 Ordering::Acquire,
115 ) {
116 Ok(_) => return Ok(current - amount),
117 Err(actual) => current = actual,
118 }
119 }
120}
121
122#[inline]
124fn credit_if_no_overflow(budget: &AtomicI64, amount: i64) -> Result<i64, ()> {
125 let mut current = budget.load(Ordering::Acquire);
126 loop {
127 let new = current.checked_add(amount).ok_or(())?;
128 match budget.compare_exchange_weak(current, new, Ordering::AcqRel, Ordering::Acquire) {
129 Ok(_) => return Ok(new),
130 Err(actual) => current = actual,
131 }
132 }
133}
134
135#[derive(Debug, Clone, Copy, PartialEq, Eq)]
136enum ReservedTotalError {
137 Overflow,
138 ExposureExceeded { current: i64 },
139}
140
141fn try_increment_reserved_total(
143 total: &AtomicI64,
144 amount: i64,
145 max: i64,
146) -> Result<i64, ReservedTotalError> {
147 let mut current = total.load(Ordering::Acquire);
148 loop {
149 let new_total = match current.checked_add(amount) {
150 Some(n) => n,
151 None => return Err(ReservedTotalError::Overflow),
152 };
153 if max > 0 && new_total > max {
154 return Err(ReservedTotalError::ExposureExceeded { current });
155 }
156 match total.compare_exchange_weak(current, new_total, Ordering::AcqRel, Ordering::Acquire) {
157 Ok(_) => return Ok(new_total),
158 Err(actual) => current = actual,
159 }
160 }
161}
162
163pub struct BudgetEngine {
170 checkpoint_gate: RwLock<()>,
172 tenant_budgets: Mutex<HashMap<Arc<str>, Arc<AtomicI64>>>,
173 initial_microcents: Mutex<HashMap<Arc<str>, i64>>,
174 committed_microcents: Mutex<HashMap<Arc<str>, i64>>,
175 reservations: Mutex<HashMap<u64, ReservationRecord>>,
176 max_reserved_microcents: Mutex<HashMap<Arc<str>, i64>>,
178 tenant_reserved_totals: Mutex<HashMap<Arc<str>, Arc<AtomicI64>>>,
180 next_id: AtomicU64,
182 snapshot_version: AtomicU64,
184 last_certified_committed_total: Mutex<i64>,
186}
187
188#[derive(Clone, Debug, PartialEq, Eq)]
190#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
191#[cfg_attr(feature = "serde", serde(deny_unknown_fields))]
192pub struct TenantLedger {
193 pub tenant_id: String,
194 pub initial_microcents: i64,
195 pub remaining_microcents: i64,
196 pub reserved_microcents: i64,
197 pub committed_microcents: i64,
199}
200
201#[derive(Clone, Debug, PartialEq, Eq)]
203#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
204#[cfg_attr(feature = "serde", serde(deny_unknown_fields))]
205pub struct BudgetSnapshot {
206 pub version: u64,
215 pub tenants: Vec<TenantLedger>,
216 pub active_reservations: usize,
217 #[cfg_attr(feature = "serde", serde(default))]
220 pub wal_high_watermark: Option<u64>,
221}
222
223#[must_use]
225pub fn conservation_status_for_snapshot(snapshot: &BudgetSnapshot) -> ConservationStatus {
226 for ledger in &snapshot.tenants {
227 let Some(sum) = ledger
228 .remaining_microcents
229 .checked_add(ledger.reserved_microcents)
230 .and_then(|v| v.checked_add(ledger.committed_microcents))
231 else {
232 return ConservationStatus::AggregateOverflow;
233 };
234 if sum != ledger.initial_microcents {
235 let Some(delta) = sum.checked_sub(ledger.initial_microcents) else {
236 return ConservationStatus::AggregateOverflow;
237 };
238 return ConservationStatus::Violation {
239 tenant_id: ledger.tenant_id.clone(),
240 delta_microcents: delta,
241 };
242 }
243 if ledger.remaining_microcents < 0 {
244 return ConservationStatus::Violation {
245 tenant_id: ledger.tenant_id.clone(),
246 delta_microcents: ledger.remaining_microcents,
247 };
248 }
249 }
250 ConservationStatus::Balanced
251}
252
253#[derive(Clone, Debug, PartialEq, Eq)]
255pub enum ConservationStatus {
256 Balanced,
258 Violation {
260 tenant_id: String,
261 delta_microcents: i64,
262 },
263 AggregateOverflow,
265}
266
267impl std::fmt::Display for ConservationStatus {
268 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
269 match self {
270 Self::Balanced => write!(f, "conservation balanced"),
271 Self::Violation {
272 tenant_id,
273 delta_microcents,
274 } => write!(
275 f,
276 "conservation violated for tenant {tenant_id}: delta={delta_microcents} microcents"
277 ),
278 Self::AggregateOverflow => write!(
279 f,
280 "aggregate ledger totals exceed i64::MAX — certificate totals cannot be represented"
281 ),
282 }
283 }
284}
285
286impl std::error::Error for ConservationStatus {}
287
288#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
290pub enum RestoreError {
291 #[error("cannot restore snapshot with {count} active reservations")]
292 ActiveReservations { count: usize },
293 #[error(
294 "tenant {tenant_id} has reserved_microcents={reserved_microcents} with no active reservations"
295 )]
296 GhostReservation {
297 tenant_id: String,
298 reserved_microcents: i64,
299 },
300 #[error("tenant {tenant_id} has negative {field}: {value}")]
301 NegativeLedgerField {
302 tenant_id: String,
303 field: &'static str,
304 value: i64,
305 },
306 #[error("snapshot failed conservation: {status}")]
307 ConservationViolation { status: ConservationStatus },
308 #[error("duplicate tenant_id in snapshot: {tenant_id}")]
309 DuplicateTenant { tenant_id: String },
310}
311
312#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
314pub enum LegacySnapshotMigrationError {
315 #[error("snapshot is already recovery-aware")]
316 AlreadyRecoveryAware,
317 #[error("trusted next reservation ID must be in 1..={max}, found {value}")]
318 InvalidAllocatorFence { value: u64, max: u64 },
319 #[error("legacy snapshot is not recovery-eligible: {0}")]
320 InvalidSnapshot(#[from] RestoreError),
321}
322
323#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
325pub enum RecoverySnapshotError {
326 #[error("legacy snapshot requires explicit migration with a trusted allocator fence")]
327 LegacySnapshotRequiresMigration,
328 #[error("invalid recovery allocator fence: {value}")]
329 InvalidAllocatorFence { value: u64 },
330 #[error("snapshot is not recovery-eligible: {0}")]
331 InvalidSnapshot(#[from] RestoreError),
332}
333
334#[derive(Clone, Copy, Debug, Eq, PartialEq, thiserror::Error)]
336pub enum SnapshotAllocatorError {
337 #[error("reservation allocator exhausted")]
338 AllocatorExhausted,
339}
340
341const RECOVERY_SNAPSHOT_TAG: u64 = 1 << 63;
342const RECOVERY_ALLOCATOR_MASK: u64 = RECOVERY_SNAPSHOT_TAG - 1;
343
344fn allocator_from_snapshot_version(version: u64) -> Result<u64, RestoreError> {
345 let next_id = version & RECOVERY_ALLOCATOR_MASK;
346 if version & RECOVERY_SNAPSHOT_TAG == 0 || next_id == 0 {
347 return Err(RestoreError::NegativeLedgerField {
348 tenant_id: "<snapshot-metadata>".to_owned(),
349 field: "recovery-aware version tag",
350 value: i64::try_from(version).unwrap_or(i64::MAX),
351 });
352 }
353 Ok(next_id)
354}
355
356pub(crate) fn validate_snapshot_for_restore(snap: &BudgetSnapshot) -> Result<(), RestoreError> {
357 allocator_from_snapshot_version(snap.version)?;
358 validate_snapshot_ledger(snap)
359}
360
361fn validate_snapshot_ledger(snap: &BudgetSnapshot) -> Result<(), RestoreError> {
362 if snap.active_reservations > 0 {
363 return Err(RestoreError::ActiveReservations {
364 count: snap.active_reservations,
365 });
366 }
367 let mut seen = std::collections::HashSet::new();
368 for ledger in &snap.tenants {
369 if !seen.insert(ledger.tenant_id.as_str()) {
370 return Err(RestoreError::DuplicateTenant {
371 tenant_id: ledger.tenant_id.clone(),
372 });
373 }
374 if ledger.reserved_microcents != 0 {
375 return Err(RestoreError::GhostReservation {
376 tenant_id: ledger.tenant_id.clone(),
377 reserved_microcents: ledger.reserved_microcents,
378 });
379 }
380 if ledger.remaining_microcents < 0 {
381 return Err(RestoreError::NegativeLedgerField {
382 tenant_id: ledger.tenant_id.clone(),
383 field: "remaining_microcents",
384 value: ledger.remaining_microcents,
385 });
386 }
387 if ledger.initial_microcents < 0 {
388 return Err(RestoreError::NegativeLedgerField {
389 tenant_id: ledger.tenant_id.clone(),
390 field: "initial_microcents",
391 value: ledger.initial_microcents,
392 });
393 }
394 if ledger.committed_microcents < 0 {
395 return Err(RestoreError::NegativeLedgerField {
396 tenant_id: ledger.tenant_id.clone(),
397 field: "committed_microcents",
398 value: ledger.committed_microcents,
399 });
400 }
401 }
402 match conservation_status_for_snapshot(snap) {
403 ConservationStatus::Balanced => Ok(()),
404 status => Err(RestoreError::ConservationViolation { status }),
405 }
406}
407
408pub fn migrate_legacy_snapshot(
415 mut snapshot: BudgetSnapshot,
416 trusted_next_reservation_id: u64,
417) -> Result<BudgetSnapshot, LegacySnapshotMigrationError> {
418 if snapshot.version & RECOVERY_SNAPSHOT_TAG != 0 {
419 return Err(LegacySnapshotMigrationError::AlreadyRecoveryAware);
420 }
421 if trusted_next_reservation_id == 0 || trusted_next_reservation_id > RECOVERY_ALLOCATOR_MASK {
422 return Err(LegacySnapshotMigrationError::InvalidAllocatorFence {
423 value: trusted_next_reservation_id,
424 max: RECOVERY_ALLOCATOR_MASK,
425 });
426 }
427 validate_snapshot_ledger(&snapshot)?;
428 snapshot.version = RECOVERY_SNAPSHOT_TAG | trusted_next_reservation_id;
429 Ok(snapshot)
430}
431
432impl BudgetEngine {
433 pub fn new() -> Self {
434 Self {
435 checkpoint_gate: RwLock::new(()),
436 tenant_budgets: Mutex::new(HashMap::new()),
437 initial_microcents: Mutex::new(HashMap::new()),
438 committed_microcents: Mutex::new(HashMap::new()),
439 reservations: Mutex::new(HashMap::new()),
440 max_reserved_microcents: Mutex::new(HashMap::new()),
441 tenant_reserved_totals: Mutex::new(HashMap::new()),
442 next_id: AtomicU64::new(1),
443 snapshot_version: AtomicU64::new(0),
444 last_certified_committed_total: Mutex::new(0),
445 }
446 }
447
448 pub fn set_max_reserved_microcents(&self, tenant_id: &str, max_microcents: i64) {
454 let _ = self.try_set_max_reserved_microcents(tenant_id, max_microcents);
455 }
456
457 pub fn try_set_max_reserved_microcents(
459 &self,
460 tenant_id: &str,
461 max_microcents: i64,
462 ) -> Result<(), BudgetConfigurationError> {
463 if max_microcents < 0 {
464 return Err(BudgetConfigurationError::NegativeAmount {
465 field: "max_reserved_microcents",
466 value: max_microcents,
467 });
468 }
469 let _mutation = self
470 .checkpoint_gate
471 .write()
472 .unwrap_or_else(|e| e.into_inner());
473 let key: Arc<str> = Arc::from(tenant_id);
474 let mut limits = self
475 .max_reserved_microcents
476 .lock()
477 .unwrap_or_else(|e| e.into_inner());
478 if max_microcents == 0 {
479 limits.remove(&key);
480 } else {
481 limits.insert(key, max_microcents);
482 }
483 Ok(())
484 }
485
486 #[must_use]
488 pub fn snapshot_version(&self) -> u64 {
489 self.snapshot_version.load(Ordering::Acquire)
490 }
491
492 #[must_use]
494 pub fn total_committed_microcents(&self) -> i64 {
495 self.try_total_committed_microcents().unwrap_or(i64::MAX)
496 }
497
498 pub fn try_total_committed_microcents(&self) -> Result<i64, ConservationStatus> {
500 let committed = self
501 .committed_microcents
502 .lock()
503 .unwrap_or_else(|e| e.into_inner());
504 let mut total: i128 = 0;
505 for amount in committed.values() {
506 total = total
507 .checked_add(i128::from(*amount))
508 .ok_or(ConservationStatus::AggregateOverflow)?;
509 }
510 i64::try_from(total).map_err(|_| ConservationStatus::AggregateOverflow)
511 }
512
513 #[must_use]
515 pub fn committed_since_last_certificate(&self) -> i64 {
516 let current = self.total_committed_microcents();
517 let baseline = *self
518 .last_certified_committed_total
519 .lock()
520 .unwrap_or_else(|e| e.into_inner());
521 current.saturating_sub(baseline)
522 }
523
524 pub(crate) fn rotate_certificate_baseline(&self, snapshot_total_committed: i64) -> i64 {
526 let mut baseline = self
527 .last_certified_committed_total
528 .lock()
529 .unwrap_or_else(|e| e.into_inner());
530 if snapshot_total_committed <= *baseline {
531 return 0;
532 }
533 let delta = snapshot_total_committed - *baseline;
534 *baseline = snapshot_total_committed;
535 delta
536 }
537
538 pub fn restore_from_snapshot(&self, snap: BudgetSnapshot) -> Result<(), RestoreError> {
544 validate_snapshot_for_restore(&snap)?;
545 let next_reservation_id = allocator_from_snapshot_version(snap.version)?;
546 let restored_committed_total = snap
547 .tenants
548 .iter()
549 .try_fold(0_i64, |total, tenant| {
550 total.checked_add(tenant.committed_microcents)
551 })
552 .unwrap_or(i64::MAX);
553 let _checkpoint = self
554 .checkpoint_gate
555 .write()
556 .unwrap_or_else(|e| e.into_inner());
557 {
558 let mut reservations = self.reservations.lock().unwrap_or_else(|e| e.into_inner());
559 reservations.clear();
560 }
561 let mut budgets = self
562 .tenant_budgets
563 .lock()
564 .unwrap_or_else(|e| e.into_inner());
565 let mut initials = self
566 .initial_microcents
567 .lock()
568 .unwrap_or_else(|e| e.into_inner());
569 let mut committed = self
570 .committed_microcents
571 .lock()
572 .unwrap_or_else(|e| e.into_inner());
573 let mut reserved_totals = self
574 .tenant_reserved_totals
575 .lock()
576 .unwrap_or_else(|e| e.into_inner());
577 let mut exposure_limits = self
578 .max_reserved_microcents
579 .lock()
580 .unwrap_or_else(|e| e.into_inner());
581 budgets.clear();
582 initials.clear();
583 committed.clear();
584 reserved_totals.clear();
585 exposure_limits.clear();
586 for ledger in snap.tenants {
587 let key: Arc<str> = Arc::from(ledger.tenant_id.as_str());
588 budgets.insert(
589 Arc::clone(&key),
590 Arc::new(AtomicI64::new(ledger.remaining_microcents)),
591 );
592 initials.insert(Arc::clone(&key), ledger.initial_microcents);
593 committed.insert(Arc::clone(&key), ledger.committed_microcents);
594 reserved_totals.insert(key, Arc::new(AtomicI64::new(ledger.reserved_microcents)));
595 }
596 self.next_id.store(next_reservation_id, Ordering::Release);
597 self.snapshot_version.store(snap.version, Ordering::Release);
598 *self
599 .last_certified_committed_total
600 .lock()
601 .unwrap_or_else(|e| e.into_inner()) = restored_committed_total;
602 Ok(())
603 }
604
605 pub fn restore_from_recovery_snapshot(
611 &self,
612 snap: BudgetSnapshot,
613 ) -> Result<(), RecoverySnapshotError> {
614 if snap.version & RECOVERY_SNAPSHOT_TAG == 0 {
615 return Err(RecoverySnapshotError::LegacySnapshotRequiresMigration);
616 }
617 let fence = snap.version & RECOVERY_ALLOCATOR_MASK;
618 if fence == 0 {
619 return Err(RecoverySnapshotError::InvalidAllocatorFence { value: fence });
620 }
621 validate_snapshot_ledger(&snap)?;
622 self.restore_from_snapshot(snap)
623 .map_err(RecoverySnapshotError::InvalidSnapshot)
624 }
625
626 pub fn ensure_tenant(&self, tenant_id: &str, budget_microcents: i64) {
634 let _ = self.try_ensure_tenant(tenant_id, budget_microcents);
635 }
636
637 pub fn try_ensure_tenant(
639 &self,
640 tenant_id: &str,
641 budget_microcents: i64,
642 ) -> Result<(), BudgetConfigurationError> {
643 if budget_microcents < 0 {
644 return Err(BudgetConfigurationError::NegativeAmount {
645 field: "budget_microcents",
646 value: budget_microcents,
647 });
648 }
649 let _mutation = self
650 .checkpoint_gate
651 .read()
652 .unwrap_or_else(|e| e.into_inner());
653 let mut budgets = self
654 .tenant_budgets
655 .lock()
656 .unwrap_or_else(|e| e.into_inner());
657 let mut initials = self
658 .initial_microcents
659 .lock()
660 .unwrap_or_else(|e| e.into_inner());
661 let mut committed = self
662 .committed_microcents
663 .lock()
664 .unwrap_or_else(|e| e.into_inner());
665 let mut reserved_totals = self
666 .tenant_reserved_totals
667 .lock()
668 .unwrap_or_else(|e| e.into_inner());
669 let key: Arc<str> = Arc::from(tenant_id);
670 if !budgets.contains_key(&key) {
671 budgets.insert(
672 Arc::clone(&key),
673 Arc::new(AtomicI64::new(budget_microcents)),
674 );
675 initials.insert(Arc::clone(&key), budget_microcents);
676 committed.insert(Arc::clone(&key), 0);
677 reserved_totals.insert(key, Arc::new(AtomicI64::new(0)));
678 }
679 Ok(())
680 }
681
682 pub fn top_up_tenant(&self, tenant_id: &str, amount_microcents: i64) -> TopUpResult {
687 if amount_microcents <= 0 {
688 return TopUpResult::InvalidAmount;
689 }
690 let _mutation = self
691 .checkpoint_gate
692 .read()
693 .unwrap_or_else(|e| e.into_inner());
694
695 let key: Arc<str> = Arc::from(tenant_id);
696
697 let budget = {
698 let budgets = self
699 .tenant_budgets
700 .lock()
701 .unwrap_or_else(|e| e.into_inner());
702 match budgets.get(&key) {
703 Some(b) => Arc::clone(b),
704 None => return TopUpResult::MissingTenant,
705 }
706 };
707
708 let mut initials = self
709 .initial_microcents
710 .lock()
711 .unwrap_or_else(|e| e.into_inner());
712
713 let Some(current_initial) = initials.get(&key).copied() else {
714 return TopUpResult::MissingTenant;
715 };
716
717 let Some(new_initial) = current_initial.checked_add(amount_microcents) else {
718 return TopUpResult::Overflow;
719 };
720
721 let remaining = match credit_if_no_overflow(&budget, amount_microcents) {
722 Ok(r) => r,
723 Err(()) => return TopUpResult::Overflow,
724 };
725
726 *initials.get_mut(&key).expect("tenant exists") = new_initial;
727
728 TopUpResult::ToppedUp {
729 added_microcents: amount_microcents,
730 new_initial_microcents: new_initial,
731 remaining_microcents: remaining,
732 }
733 }
734
735 #[must_use]
737 pub fn initial_microcents(&self, tenant_id: &str) -> Option<i64> {
738 let initials = self
739 .initial_microcents
740 .lock()
741 .unwrap_or_else(|e| e.into_inner());
742 let key: Arc<str> = Arc::from(tenant_id);
743 initials.get(&key).copied()
744 }
745
746 #[must_use]
750 pub fn committed_microcents(&self, tenant_id: &str) -> Option<i64> {
751 let committed = self
752 .committed_microcents
753 .lock()
754 .unwrap_or_else(|e| e.into_inner());
755 let key: Arc<str> = Arc::from(tenant_id);
756 committed.get(&key).copied()
757 }
758
759 #[must_use]
761 pub fn reserved_microcents(&self, tenant_id: &str) -> i64 {
762 let key: Arc<str> = Arc::from(tenant_id);
763 let totals = self
764 .tenant_reserved_totals
765 .lock()
766 .unwrap_or_else(|e| e.into_inner());
767 totals
768 .get(&key)
769 .map(|t| t.load(Ordering::Acquire))
770 .unwrap_or(0)
771 }
772
773 #[must_use]
785 pub fn snapshot(&self) -> BudgetSnapshot {
786 self.try_snapshot()
787 .expect("reservation allocator exhausted — use try_snapshot() for fallible capture")
788 }
789
790 pub fn try_snapshot(&self) -> Result<BudgetSnapshot, SnapshotAllocatorError> {
792 let _checkpoint = self
793 .checkpoint_gate
794 .write()
795 .unwrap_or_else(|e| e.into_inner());
796 let reservations = self.reservations.lock().unwrap_or_else(|e| e.into_inner());
797 let budgets = self
798 .tenant_budgets
799 .lock()
800 .unwrap_or_else(|e| e.into_inner());
801 let initials = self
802 .initial_microcents
803 .lock()
804 .unwrap_or_else(|e| e.into_inner());
805 let committed = self
806 .committed_microcents
807 .lock()
808 .unwrap_or_else(|e| e.into_inner());
809 let reserved_totals = self
810 .tenant_reserved_totals
811 .lock()
812 .unwrap_or_else(|e| e.into_inner());
813
814 let mut tenants = Vec::with_capacity(budgets.len());
815 for (tenant_id, balance) in budgets.iter() {
816 tenants.push(TenantLedger {
817 tenant_id: tenant_id.to_string(),
818 initial_microcents: initials.get(tenant_id).copied().unwrap_or(0),
819 remaining_microcents: balance.load(Ordering::Acquire),
820 reserved_microcents: reserved_totals
821 .get(tenant_id)
822 .map(|t| t.load(Ordering::Acquire))
823 .unwrap_or(0),
824 committed_microcents: committed.get(tenant_id).copied().unwrap_or(0),
825 });
826 }
827 tenants.sort_by(|a, b| a.tenant_id.cmp(&b.tenant_id));
828 let previous_id = self
834 .next_id
835 .fetch_update(Ordering::AcqRel, Ordering::Acquire, |current| {
836 (current < RECOVERY_ALLOCATOR_MASK).then_some(current + 1)
837 })
838 .map_err(|_| SnapshotAllocatorError::AllocatorExhausted)?;
839 let version = RECOVERY_SNAPSHOT_TAG | (previous_id + 1);
840 self.snapshot_version.store(version, Ordering::Release);
841 Ok(BudgetSnapshot {
842 version,
843 tenants,
844 active_reservations: reservations.len(),
845 wal_high_watermark: None,
846 })
847 }
848
849 #[must_use]
853 pub fn verify_conservation(&self) -> ConservationStatus {
854 conservation_status_for_snapshot(&self.snapshot())
855 }
856
857 #[must_use]
859 pub fn remaining_microcents(&self, tenant_id: &str) -> Option<i64> {
860 let budgets = self
861 .tenant_budgets
862 .lock()
863 .unwrap_or_else(|e| e.into_inner());
864 let key: Arc<str> = Arc::from(tenant_id);
865 budgets.get(&key).map(|b| b.load(Ordering::Acquire))
866 }
867
868 pub fn try_reserve(
874 &self,
875 tenant_id: &str,
876 cost_microcents: i64,
877 ) -> (BudgetReservation, Option<u64>) {
878 if cost_microcents <= 0 {
879 return (
880 BudgetReservation::Insufficient {
881 remaining_microcents: 0,
882 required_microcents: cost_microcents,
883 },
884 None,
885 );
886 }
887 let _mutation = self
888 .checkpoint_gate
889 .read()
890 .unwrap_or_else(|e| e.into_inner());
891
892 let key: Arc<str> = Arc::from(tenant_id);
893 let (budget, reserved_total) = {
894 let budgets = self
895 .tenant_budgets
896 .lock()
897 .unwrap_or_else(|e| e.into_inner());
898 let totals = self
899 .tenant_reserved_totals
900 .lock()
901 .unwrap_or_else(|e| e.into_inner());
902 match (budgets.get(&key), totals.get(&key)) {
903 (Some(b), Some(t)) => (Arc::clone(b), Arc::clone(t)),
904 _ => return (BudgetReservation::MissingTenant, None),
905 }
906 };
907
908 let max_reserved = {
909 let limits = self
910 .max_reserved_microcents
911 .lock()
912 .unwrap_or_else(|e| e.into_inner());
913 limits.get(&key).copied().unwrap_or(0)
914 };
915 match try_increment_reserved_total(&reserved_total, cost_microcents, max_reserved) {
916 Ok(_) => {}
917 Err(ReservedTotalError::Overflow) => {
918 return (
919 BudgetReservation::Overflow {
920 current_reserved_microcents: reserved_total.load(Ordering::Acquire),
921 },
922 None,
923 );
924 }
925 Err(ReservedTotalError::ExposureExceeded { current }) => {
926 return (
927 BudgetReservation::ExposureLimitExceeded {
928 current_reserved_microcents: current,
929 max_reserved_microcents: max_reserved,
930 },
931 None,
932 );
933 }
934 }
935
936 let id = match self
940 .next_id
941 .fetch_update(Ordering::AcqRel, Ordering::Acquire, |current| {
942 (current < RECOVERY_ALLOCATOR_MASK).then_some(current + 1)
943 }) {
944 Ok(id) => id,
945 Err(_) => {
946 reserved_total.fetch_sub(cost_microcents, Ordering::AcqRel);
947 return (
948 BudgetReservation::Overflow {
949 current_reserved_microcents: reserved_total.load(Ordering::Acquire),
950 },
951 None,
952 );
953 }
954 };
955
956 match debit_if_available(&budget, cost_microcents) {
957 Err(current) => {
958 reserved_total.fetch_sub(cost_microcents, Ordering::AcqRel);
959 (
960 BudgetReservation::Insufficient {
961 remaining_microcents: current,
962 required_microcents: cost_microcents,
963 },
964 None,
965 )
966 }
967 Ok(remaining) => {
968 let mut reservations = self.reservations.lock().unwrap_or_else(|e| e.into_inner());
969 reservations.insert(
970 id,
971 ReservationRecord {
972 tenant_id: Arc::clone(&key),
973 reserved_microcents: cost_microcents,
974 },
975 );
976 (
977 BudgetReservation::Reserved {
978 remaining_microcents: remaining,
979 },
980 Some(id),
981 )
982 }
983 }
984 }
985
986 pub fn commit(&self, reservation_id: u64, actual_microcents: i64) -> BudgetSettlement {
996 if actual_microcents < 0 {
997 return BudgetSettlement::InvalidAmount;
998 }
999 let _mutation = self
1000 .checkpoint_gate
1001 .read()
1002 .unwrap_or_else(|e| e.into_inner());
1003
1004 let mut reservations = self.reservations.lock().unwrap_or_else(|e| e.into_inner());
1005 let Some(reservation) = reservations.remove(&reservation_id) else {
1006 return BudgetSettlement::MissingReservation;
1007 };
1008
1009 let budget = {
1010 let budgets = self
1011 .tenant_budgets
1012 .lock()
1013 .unwrap_or_else(|e| e.into_inner());
1014 match budgets.get(&reservation.tenant_id) {
1015 Some(b) => Arc::clone(b),
1016 None => {
1017 reservations.insert(reservation_id, reservation);
1018 return BudgetSettlement::MissingTenant;
1019 }
1020 }
1021 };
1022 let tenant_key = Arc::clone(&reservation.tenant_id);
1023 let mut committed_guard = self
1024 .committed_microcents
1025 .lock()
1026 .unwrap_or_else(|e| e.into_inner());
1027 let current_committed = committed_guard.get(&tenant_key).copied().unwrap_or(0);
1028 let new_committed = match current_committed.checked_add(actual_microcents) {
1029 Some(v) => v,
1030 None => {
1031 drop(committed_guard);
1032 reservations.insert(reservation_id, reservation);
1033 return BudgetSettlement::Overflow {
1034 remaining_microcents: budget.load(Ordering::Acquire),
1035 };
1036 }
1037 };
1038
1039 let delta: i64 = actual_microcents - reservation.reserved_microcents;
1040 match delta.cmp(&0) {
1041 std::cmp::Ordering::Greater => {
1042 if let Err(remaining) = debit_if_available(&budget, delta) {
1043 drop(committed_guard);
1044 reservations.insert(reservation_id, reservation);
1045 return BudgetSettlement::Overrun {
1046 remaining_microcents: remaining,
1047 };
1048 }
1049 }
1050 std::cmp::Ordering::Less => {}
1051 std::cmp::Ordering::Equal => {}
1052 }
1053
1054 committed_guard.insert(tenant_key.clone(), new_committed);
1055 drop(committed_guard);
1056
1057 if delta < 0 {
1058 budget.fetch_add(-delta, Ordering::AcqRel);
1059 }
1060
1061 {
1062 let totals = self
1063 .tenant_reserved_totals
1064 .lock()
1065 .unwrap_or_else(|e| e.into_inner());
1066 if let Some(total) = totals.get(&tenant_key) {
1067 total.fetch_sub(reservation.reserved_microcents, Ordering::AcqRel);
1068 }
1069 }
1070
1071 let remaining = budget.load(Ordering::Acquire);
1072 BudgetSettlement::Committed {
1073 remaining_microcents: remaining,
1074 actual_microcents,
1075 }
1076 }
1077
1078 pub fn release(&self, reservation_id: u64) -> BudgetSettlement {
1080 let _mutation = self
1081 .checkpoint_gate
1082 .read()
1083 .unwrap_or_else(|e| e.into_inner());
1084 let mut reservations = self.reservations.lock().unwrap_or_else(|e| e.into_inner());
1085 let Some((_, reservation)) = reservations.remove_entry(&reservation_id) else {
1086 return BudgetSettlement::MissingReservation;
1087 };
1088
1089 let budget = {
1090 let budgets = self
1091 .tenant_budgets
1092 .lock()
1093 .unwrap_or_else(|e| e.into_inner());
1094 match budgets.get(&reservation.tenant_id) {
1095 Some(b) => Arc::clone(b),
1096 None => {
1097 reservations.insert(reservation_id, reservation);
1098 return BudgetSettlement::MissingTenant;
1099 }
1100 }
1101 };
1102 {
1103 let totals = self
1104 .tenant_reserved_totals
1105 .lock()
1106 .unwrap_or_else(|e| e.into_inner());
1107 if let Some(total) = totals.get(&reservation.tenant_id) {
1108 total.fetch_sub(reservation.reserved_microcents, Ordering::AcqRel);
1109 }
1110 }
1111
1112 let returned = reservation.reserved_microcents;
1113 let remaining = budget.fetch_add(returned, Ordering::AcqRel) + returned;
1114
1115 BudgetSettlement::Released {
1116 remaining_microcents: remaining,
1117 returned_microcents: returned,
1118 }
1119 }
1120
1121 #[must_use]
1123 pub fn tenant_count(&self) -> usize {
1124 self.tenant_budgets
1125 .lock()
1126 .unwrap_or_else(|e| e.into_inner())
1127 .len()
1128 }
1129
1130 #[must_use]
1132 pub fn active_reservations(&self) -> usize {
1133 self.reservations
1134 .lock()
1135 .unwrap_or_else(|e| e.into_inner())
1136 .len()
1137 }
1138}
1139
1140impl Default for BudgetEngine {
1141 fn default() -> Self {
1144 Self::new()
1145 }
1146}
1147
1148#[cfg(test)]
1149mod tests {
1150 use super::*;
1151
1152 #[test]
1153 fn reserve_and_commit() {
1154 let engine = BudgetEngine::default();
1155 engine.ensure_tenant("t1", 100_000_000);
1156 let (_, id) = engine.try_reserve("t1", 25_000_000);
1157 let settlement = engine.commit(id.unwrap(), 20_000_000);
1158 assert!(matches!(settlement, BudgetSettlement::Committed { .. }));
1159 assert_eq!(engine.remaining_microcents("t1"), Some(80_000_000));
1160 }
1161
1162 #[test]
1163 fn reserve_insufficient() {
1164 let engine = BudgetEngine::default();
1165 engine.ensure_tenant("t1", 10_000_000);
1166 let (res, id) = engine.try_reserve("t1", 50_000_000);
1167 assert!(matches!(res, BudgetReservation::Insufficient { .. }));
1168 assert!(id.is_none());
1169 }
1170
1171 #[test]
1172 fn release_returns_full_amount() {
1173 let engine = BudgetEngine::default();
1174 engine.ensure_tenant("t1", 100_000_000);
1175 let (_, id) = engine.try_reserve("t1", 30_000_000);
1176 engine.release(id.unwrap());
1177 assert_eq!(engine.remaining_microcents("t1"), Some(100_000_000));
1178 }
1179
1180 #[test]
1181 fn missing_tenant_rejected() {
1182 let engine = BudgetEngine::default();
1183 let (res, _) = engine.try_reserve("nonexistent", 1000);
1184 assert!(matches!(res, BudgetReservation::MissingTenant));
1185 }
1186
1187 #[test]
1188 fn missing_reservation_rejected() {
1189 let engine = BudgetEngine::default();
1190 assert!(matches!(
1191 engine.commit(999, 1000),
1192 BudgetSettlement::MissingReservation
1193 ));
1194 assert!(matches!(
1195 engine.release(999),
1196 BudgetSettlement::MissingReservation
1197 ));
1198 }
1199
1200 #[test]
1201 fn conservation_invariant() {
1202 let engine = BudgetEngine::default();
1203 engine.ensure_tenant("t1", 100_000_000);
1204 let (_, id1) = engine.try_reserve("t1", 30_000_000);
1205 let (_, id2) = engine.try_reserve("t1", 20_000_000);
1206 engine.commit(id1.unwrap(), 25_000_000);
1207 engine.release(id2.unwrap());
1208 assert_eq!(engine.remaining_microcents("t1"), Some(75_000_000));
1209 assert_eq!(engine.committed_microcents("t1"), Some(25_000_000));
1210 assert_eq!(engine.verify_conservation(), ConservationStatus::Balanced);
1211 }
1212
1213 #[test]
1214 fn top_up_extends_initial_and_remaining() {
1215 let engine = BudgetEngine::default();
1216 engine.ensure_tenant("desk", 100_000_000);
1217 let result = engine.top_up_tenant("desk", 50_000_000);
1218 assert!(matches!(result, TopUpResult::ToppedUp { .. }));
1219 assert_eq!(engine.initial_microcents("desk"), Some(150_000_000));
1220 assert_eq!(engine.remaining_microcents("desk"), Some(150_000_000));
1221 assert_eq!(engine.verify_conservation(), ConservationStatus::Balanced);
1222 }
1223
1224 #[test]
1225 fn snapshot_balances() {
1226 let engine = BudgetEngine::default();
1227 engine.ensure_tenant("desk-a", 50_000_000);
1228 engine.ensure_tenant("desk-b", 80_000_000);
1229 let (_, id) = engine.try_reserve("desk-a", 10_000_000);
1230 engine.commit(id.unwrap(), 9_000_000);
1231 let snap = engine.snapshot();
1232 assert_eq!(snap.tenants.len(), 2);
1233 assert_eq!(engine.verify_conservation(), ConservationStatus::Balanced);
1234 }
1235
1236 #[test]
1237 #[cfg_attr(miri, ignore = "miri: concurrency stress test")]
1238 fn snapshots_are_linearizable_with_reservation_mutations() {
1239 let engine = Arc::new(BudgetEngine::default());
1240 engine.ensure_tenant("desk", 1_000_000);
1241 let running = Arc::new(std::sync::atomic::AtomicBool::new(true));
1242
1243 let writer_engine = Arc::clone(&engine);
1244 let writer_running = Arc::clone(&running);
1245 let writer = std::thread::spawn(move || {
1246 for _ in 0..20_000 {
1247 let (_, id) = writer_engine.try_reserve("desk", 1);
1248 if let Some(id) = id {
1249 let _ = writer_engine.release(id);
1250 }
1251 }
1252 writer_running.store(false, std::sync::atomic::Ordering::Release);
1253 });
1254
1255 while running.load(std::sync::atomic::Ordering::Acquire) {
1256 let snapshot = engine.snapshot();
1257 assert_eq!(
1258 conservation_status_for_snapshot(&snapshot),
1259 ConservationStatus::Balanced
1260 );
1261 let reserved: i64 = snapshot
1262 .tenants
1263 .iter()
1264 .map(|tenant| tenant.reserved_microcents)
1265 .sum();
1266 assert_eq!(snapshot.active_reservations as i64, reserved);
1267 }
1268 writer.join().unwrap();
1269 }
1270
1271 #[test]
1272 #[cfg_attr(miri, ignore = "miri: use Loom for concurrent interleavings")]
1273 fn concurrent_reserve_never_overspends() {
1274 let engine = Arc::new(BudgetEngine::default());
1275 engine.ensure_tenant("t1", 100_000_000);
1276 let handles: Vec<_> = (0..100)
1277 .map(|_| {
1278 let e = Arc::clone(&engine);
1279 std::thread::spawn(move || {
1280 let (res, _) = e.try_reserve("t1", 1_000_000);
1281 matches!(res, BudgetReservation::Reserved { .. })
1282 })
1283 })
1284 .collect();
1285 let success: usize = handles
1286 .into_iter()
1287 .map(|h| if h.join().unwrap() { 1 } else { 0 })
1288 .sum();
1289 assert_eq!(success, 100);
1290 assert_eq!(engine.remaining_microcents("t1"), Some(0));
1291 }
1292
1293 #[test]
1294 fn zero_cost_rejected() {
1295 let engine = BudgetEngine::default();
1296 engine.ensure_tenant("t1", 100_000_000);
1297 let (res, _) = engine.try_reserve("t1", 0);
1298 assert!(matches!(res, BudgetReservation::Insufficient { .. }));
1299 }
1300
1301 #[test]
1302 fn negative_cost_rejected() {
1303 let engine = BudgetEngine::default();
1304 engine.ensure_tenant("t1", 100_000_000);
1305 let (res, _) = engine.try_reserve("t1", -500);
1306 assert!(matches!(res, BudgetReservation::Insufficient { .. }));
1307 }
1308
1309 #[test]
1310 fn negative_commit_rejected() {
1311 let engine = BudgetEngine::default();
1312 engine.ensure_tenant("t1", 100_000_000);
1313 let (_, id) = engine.try_reserve("t1", 10_000_000);
1314 let result = engine.commit(id.unwrap(), -5_000_000);
1315 assert!(matches!(result, BudgetSettlement::InvalidAmount));
1316 }
1317
1318 #[test]
1319 fn failed_overrun_does_not_create_budget() {
1320 let engine = BudgetEngine::default();
1321 engine.ensure_tenant("t1", 10_000_000);
1322
1323 let (_, id) = engine.try_reserve("t1", 7_000_000);
1324 let id = id.unwrap();
1325
1326 let result = engine.commit(id, 11_000_000);
1328 assert!(matches!(result, BudgetSettlement::Overrun { .. }));
1329
1330 assert_eq!(engine.remaining_microcents("t1"), Some(3_000_000));
1332
1333 engine.release(id);
1335 assert_eq!(engine.remaining_microcents("t1"), Some(10_000_000));
1336 }
1337
1338 #[test]
1339 #[cfg_attr(miri, ignore = "miri: use Loom for concurrent interleavings")]
1340 fn no_deadlock_under_contention() {
1341 let engine = Arc::new(BudgetEngine::default());
1342 engine.ensure_tenant("t1", 1_000_000_000);
1343 let handles: Vec<_> = (0..50)
1344 .map(|i| {
1345 let e = Arc::clone(&engine);
1346 std::thread::spawn(move || {
1347 let (_, id) = e.try_reserve("t1", 1_000_000);
1348 if let Some(id) = id {
1349 if i % 3 == 0 {
1350 e.release(id);
1351 } else {
1352 e.commit(id, 800_000);
1353 }
1354 }
1355 })
1356 })
1357 .collect();
1358 for h in handles {
1359 h.join().unwrap();
1360 }
1361 assert!(engine.remaining_microcents("t1").unwrap() > 0);
1362 }
1363
1364 #[test]
1365 #[cfg_attr(miri, ignore = "miri: use Loom for concurrent interleavings")]
1366 fn concurrent_overrun_never_goes_negative() {
1367 let engine = Arc::new(BudgetEngine::default());
1368 engine.ensure_tenant("t1", 50_000_000);
1369
1370 let handles: Vec<_> = (0..20)
1372 .map(|_| {
1373 let e = Arc::clone(&engine);
1374 std::thread::spawn(move || {
1375 let (_, id) = e.try_reserve("t1", 1_000_000);
1376 if let Some(id) = id {
1377 let _ = e.commit(id, 3_000_000);
1379 }
1380 })
1381 })
1382 .collect();
1383
1384 for h in handles {
1385 h.join().unwrap();
1386 }
1387
1388 let remaining = engine.remaining_microcents("t1").unwrap();
1389 assert!(
1390 remaining >= 0,
1391 "budget must never go negative, got {remaining}"
1392 );
1393 }
1394
1395 #[test]
1396 fn double_commit_rejected() {
1397 let engine = BudgetEngine::default();
1398 engine.ensure_tenant("t1", 100_000_000);
1399 let (_, id) = engine.try_reserve("t1", 10_000_000);
1400 let id = id.unwrap();
1401 assert!(matches!(
1402 engine.commit(id, 8_000_000),
1403 BudgetSettlement::Committed { .. }
1404 ));
1405 assert!(matches!(
1406 engine.commit(id, 8_000_000),
1407 BudgetSettlement::MissingReservation
1408 ));
1409 }
1410
1411 #[test]
1412 fn double_release_rejected() {
1413 let engine = BudgetEngine::default();
1414 engine.ensure_tenant("t1", 100_000_000);
1415 let (_, id) = engine.try_reserve("t1", 10_000_000);
1416 let id = id.unwrap();
1417 assert!(matches!(
1418 engine.release(id),
1419 BudgetSettlement::Released { .. }
1420 ));
1421 assert!(matches!(
1422 engine.release(id),
1423 BudgetSettlement::MissingReservation
1424 ));
1425 }
1426
1427 #[test]
1428 fn restore_from_snapshot_roundtrip() {
1429 let engine = BudgetEngine::default();
1430 engine.ensure_tenant("desk", 1_000_000);
1431 let (_, id) = engine.try_reserve("desk", 100_000);
1432 engine.commit(id.unwrap(), 90_000);
1433 let snap = engine.snapshot();
1434 assert!(snap.active_reservations == 0);
1435 let fresh = BudgetEngine::default();
1436 fresh.restore_from_snapshot(snap).unwrap();
1437 assert_eq!(fresh.remaining_microcents("desk"), Some(910_000));
1438 assert_eq!(fresh.committed_microcents("desk"), Some(90_000));
1439 assert_eq!(fresh.verify_conservation(), ConservationStatus::Balanced);
1440 }
1441
1442 #[test]
1443 fn restore_preserves_reservation_id_monotonicity() {
1444 let source = BudgetEngine::default();
1445 source.ensure_tenant("desk", 1_000_000);
1446 let (_, stale_id) = source.try_reserve("desk", 100_000);
1447 let stale_id = stale_id.expect("source reservation");
1448 assert!(matches!(
1449 source.release(stale_id),
1450 BudgetSettlement::Released { .. }
1451 ));
1452
1453 let snapshot = source.snapshot();
1454 let recovered = BudgetEngine::default();
1455 recovered.restore_from_snapshot(snapshot).unwrap();
1456 let (_, recovered_id) = recovered.try_reserve("desk", 100_000);
1457 let recovered_id = recovered_id.expect("recovered reservation");
1458
1459 assert!(recovered_id > stale_id);
1460 assert_eq!(
1461 recovered.release(stale_id),
1462 BudgetSettlement::MissingReservation
1463 );
1464 assert_eq!(recovered.reserved_microcents("desk"), 100_000);
1465 }
1466
1467 #[test]
1468 fn legacy_snapshot_migration_requires_and_binds_a_trusted_allocator_fence() {
1469 let source = BudgetEngine::default();
1470 source.ensure_tenant("desk", 1_000_000);
1471 let mut legacy = source.snapshot();
1472 legacy.version = 7;
1473
1474 assert!(migrate_legacy_snapshot(legacy.clone(), 0).is_err());
1475 let migrated = migrate_legacy_snapshot(legacy, 100).unwrap();
1476 let recovered = BudgetEngine::default();
1477 recovered.restore_from_snapshot(migrated).unwrap();
1478 let (_, reservation_id) = recovered.try_reserve("desk", 100_000);
1479
1480 assert_eq!(reservation_id, Some(100));
1481 }
1482
1483 #[test]
1484 fn snapshot_fails_closed_when_the_reservation_allocator_is_exhausted() {
1485 let engine = BudgetEngine::default();
1486 engine.ensure_tenant("desk", 1_000_000);
1487 engine
1489 .next_id
1490 .store(RECOVERY_ALLOCATOR_MASK, Ordering::Release);
1491
1492 assert_eq!(
1493 engine.try_snapshot(),
1494 Err(SnapshotAllocatorError::AllocatorExhausted)
1495 );
1496 }
1497
1498 #[test]
1499 fn recovery_restore_reports_legacy_format_without_ledger_error_aliasing() {
1500 let legacy = BudgetSnapshot {
1501 version: 7,
1502 tenants: vec![],
1503 active_reservations: 0,
1504 wal_high_watermark: None,
1505 };
1506 let engine = BudgetEngine::default();
1507 assert_eq!(
1508 engine.restore_from_recovery_snapshot(legacy),
1509 Err(RecoverySnapshotError::LegacySnapshotRequiresMigration)
1510 );
1511 }
1512
1513 #[test]
1514 fn restore_replaces_all_recovery_sensitive_runtime_state() {
1515 let source = BudgetEngine::default();
1516 source.ensure_tenant("desk", 1_000_000);
1517 let (_, id) = source.try_reserve("desk", 100_000);
1518 source.commit(id.unwrap(), 80_000);
1519 let snapshot = source.snapshot();
1520
1521 let recovered = BudgetEngine::default();
1522 recovered.ensure_tenant("old", 1_000_000);
1523 recovered.set_max_reserved_microcents("desk", 1);
1524 assert_eq!(recovered.rotate_certificate_baseline(10), 10);
1525 recovered.restore_from_snapshot(snapshot).unwrap();
1526
1527 assert_eq!(recovered.committed_since_last_certificate(), 0);
1528 let (reservation, id) = recovered.try_reserve("desk", 100_000);
1529 assert!(matches!(reservation, BudgetReservation::Reserved { .. }));
1530 recovered.release(id.unwrap());
1531 }
1532
1533 #[test]
1534 fn ensure_tenant_rejects_negative_budget() {
1535 let engine = BudgetEngine::default();
1536 engine.ensure_tenant("desk", -1);
1537 assert_eq!(engine.tenant_count(), 0);
1538 assert!(engine.try_ensure_tenant("desk", -1).is_err());
1539 }
1540
1541 #[test]
1542 fn negative_exposure_cap_does_not_remove_existing_cap() {
1543 let engine = BudgetEngine::default();
1544 engine.ensure_tenant("desk", 1_000);
1545 engine.set_max_reserved_microcents("desk", 100);
1546 assert!(engine.try_set_max_reserved_microcents("desk", -1).is_err());
1547 engine.set_max_reserved_microcents("desk", -1);
1548 let (result, id) = engine.try_reserve("desk", 101);
1549 assert!(matches!(
1550 result,
1551 BudgetReservation::ExposureLimitExceeded {
1552 max_reserved_microcents: 100,
1553 ..
1554 }
1555 ));
1556 assert!(id.is_none());
1557
1558 engine.set_max_reserved_microcents("desk", 0);
1559 let (result, id) = engine.try_reserve("desk", 101);
1560 assert!(matches!(result, BudgetReservation::Reserved { .. }));
1561 engine.release(id.unwrap());
1562 }
1563
1564 #[test]
1565 fn certificate_baseline_is_monotonic() {
1566 let engine = BudgetEngine::default();
1567 assert_eq!(engine.rotate_certificate_baseline(150), 150);
1568 assert_eq!(engine.rotate_certificate_baseline(100), 0);
1569 assert_eq!(engine.rotate_certificate_baseline(200), 50);
1570 }
1571
1572 #[test]
1573 fn total_committed_microcents_reports_aggregate_overflow() {
1574 let engine = BudgetEngine::default();
1575 engine.ensure_tenant("desk-a", 1);
1576 engine.ensure_tenant("desk-b", 1);
1577 {
1578 let mut committed = engine
1579 .committed_microcents
1580 .lock()
1581 .unwrap_or_else(|e| e.into_inner());
1582 committed.insert(Arc::from("desk-a"), i64::MAX);
1583 committed.insert(Arc::from("desk-b"), 1);
1584 }
1585 assert_eq!(
1586 engine.try_total_committed_microcents(),
1587 Err(ConservationStatus::AggregateOverflow)
1588 );
1589 assert_eq!(engine.total_committed_microcents(), i64::MAX);
1590 }
1591
1592 #[test]
1593 fn adversarial_snapshot_per_tenant_sum_overflow() {
1594 let snap = BudgetSnapshot {
1595 version: 1,
1596 tenants: vec![TenantLedger {
1597 tenant_id: "evil".into(),
1598 initial_microcents: 0,
1599 remaining_microcents: i64::MAX,
1600 reserved_microcents: 1,
1601 committed_microcents: 0,
1602 }],
1603 active_reservations: 0,
1604 wal_high_watermark: None,
1605 };
1606 assert_eq!(
1607 conservation_status_for_snapshot(&snap),
1608 ConservationStatus::AggregateOverflow
1609 );
1610 }
1611
1612 #[test]
1613 fn restore_rejects_duplicate_tenant() {
1614 let engine = BudgetEngine::default();
1615 engine.ensure_tenant("desk", 1_000_000);
1616 let mut snap = engine.snapshot();
1617 snap.tenants.push(snap.tenants[0].clone());
1618 assert!(matches!(
1619 BudgetEngine::default().restore_from_snapshot(snap),
1620 Err(RestoreError::DuplicateTenant { .. })
1621 ));
1622 }
1623
1624 #[test]
1625 fn top_up_rejects_overflow() {
1626 let engine = BudgetEngine::default();
1627 engine.ensure_tenant("desk", i64::MAX - 10);
1628 assert!(matches!(
1629 engine.top_up_tenant("desk", 20),
1630 TopUpResult::Overflow
1631 ));
1632 }
1633
1634 #[test]
1635 fn reserve_rejects_reserved_total_overflow() {
1636 let engine = BudgetEngine::default();
1637 engine.ensure_tenant("desk", i64::MAX);
1638 engine.set_max_reserved_microcents("desk", 0);
1639 let (res, _) = engine.try_reserve("desk", i64::MAX);
1640 assert!(matches!(res, BudgetReservation::Reserved { .. }));
1641 let (res2, _) = engine.try_reserve("desk", 1);
1642 assert!(matches!(res2, BudgetReservation::Overflow { .. }));
1643 }
1644
1645 #[test]
1646 fn commit_rejects_lifetime_overflow() {
1647 let engine = BudgetEngine::default();
1648 engine.ensure_tenant("desk", i64::MAX);
1649 let (_, id) = engine.try_reserve("desk", 1);
1650 assert!(matches!(
1651 engine.commit(id.unwrap(), i64::MAX - 1),
1652 BudgetSettlement::Committed { .. }
1653 ));
1654 let (_, id2) = engine.try_reserve("desk", 1);
1655 assert!(matches!(
1656 engine.commit(id2.unwrap(), 2),
1657 BudgetSettlement::Overflow { .. }
1658 ));
1659 }
1660
1661 #[test]
1662 fn restore_rejects_ghost_reserved() {
1663 let engine = BudgetEngine::default();
1664 engine.ensure_tenant("desk", 1_000_000);
1665 let mut snap = engine.snapshot();
1666 snap.tenants[0].reserved_microcents = 100_000;
1667 assert!(matches!(
1668 BudgetEngine::default().restore_from_snapshot(snap),
1669 Err(RestoreError::GhostReservation { .. })
1670 ));
1671 }
1672
1673 #[test]
1674 fn restore_rejects_unbalanced_snapshot() {
1675 let engine = BudgetEngine::default();
1676 engine.ensure_tenant("desk", 1_000_000);
1677 let mut snap = engine.snapshot();
1678 snap.tenants[0].remaining_microcents = 0;
1679 assert!(matches!(
1680 BudgetEngine::default().restore_from_snapshot(snap),
1681 Err(RestoreError::ConservationViolation { .. })
1682 ));
1683 }
1684
1685 #[test]
1686 fn restore_rejects_active_reservations() {
1687 let engine = BudgetEngine::default();
1688 engine.ensure_tenant("desk", 1_000_000);
1689 let (_, id) = engine.try_reserve("desk", 50_000);
1690 id.unwrap();
1691 let mut snap = engine.snapshot();
1692 snap.active_reservations = 1;
1693 let fresh = BudgetEngine::default();
1694 assert!(matches!(
1695 fresh.restore_from_snapshot(snap),
1696 Err(RestoreError::ActiveReservations { count: 1 })
1697 ));
1698 }
1699
1700 #[test]
1701 #[cfg_attr(miri, ignore = "miri: use Loom for concurrent interleavings")]
1702 fn exposure_limit_holds_under_concurrent_reserve() {
1703 let engine = Arc::new(BudgetEngine::default());
1704 engine.ensure_tenant("desk", 10_000_000);
1705 engine.set_max_reserved_microcents("desk", 100_000);
1706 let handles: Vec<_> = (0..32)
1707 .map(|_| {
1708 let e = Arc::clone(&engine);
1709 std::thread::spawn(move || {
1710 let (res, _) = e.try_reserve("desk", 80_000);
1711 matches!(res, BudgetReservation::Reserved { .. })
1712 })
1713 })
1714 .collect();
1715 let successes: usize = handles
1716 .into_iter()
1717 .map(|h| if h.join().unwrap() { 1 } else { 0 })
1718 .sum();
1719 assert_eq!(successes, 1);
1720 assert!(engine.reserved_microcents("desk") <= 100_000);
1721 assert_eq!(engine.verify_conservation(), ConservationStatus::Balanced);
1722 }
1723
1724 #[test]
1725 fn exposure_limit_blocks_reserve() {
1726 let engine = BudgetEngine::default();
1727 engine.ensure_tenant("desk", 1_000_000);
1728 engine.set_max_reserved_microcents("desk", 100_000);
1729 let (_, id1) = engine.try_reserve("desk", 60_000);
1730 assert!(id1.is_some());
1731 let (res, id2) = engine.try_reserve("desk", 50_000);
1732 assert!(matches!(
1733 res,
1734 BudgetReservation::ExposureLimitExceeded { .. }
1735 ));
1736 assert!(id2.is_none());
1737 assert_eq!(engine.verify_conservation(), ConservationStatus::Balanced);
1738 }
1739
1740 use proptest::prelude::*;
1741
1742 fn edge_amounts() -> impl Strategy<Value = i64> {
1743 prop_oneof![
1744 Just(1_i64),
1745 Just(-1),
1746 Just(0),
1747 Just(1_000_000_i64),
1748 1i64..1_000_000,
1749 ]
1750 }
1751
1752 proptest! {
1753 #[test]
1754 fn aggressive_mixed_ops_maintain_conservation(
1755 tenant_count in 1_usize..8,
1756 seed_ops in prop::collection::vec((0u8..6, any::<u8>(), edge_amounts()), 5..80),
1757 ) {
1758 let engine = BudgetEngine::default();
1759 for t in 0..tenant_count {
1760 engine.ensure_tenant(&format!("tenant-{t}"), 2_000_000);
1761 if t % 2 == 0 {
1762 engine.set_max_reserved_microcents(&format!("tenant-{t}"), 500_000);
1763 }
1764 }
1765 let mut open_ids = Vec::new();
1766
1767 for (op, sel, amount) in seed_ops {
1768 let tenant = format!("tenant-{}", sel as usize % tenant_count);
1769 let snap_before = engine.snapshot();
1770 let digest_before = crate::finance::ledger_digest(&snap_before);
1771
1772 match op % 6 {
1773 0 => {
1774 let (_, id) = engine.try_reserve(&tenant, amount);
1775 if let Some(id) = id {
1776 open_ids.push(id);
1777 }
1778 }
1779 1 if !open_ids.is_empty() => {
1780 let idx = sel as usize % open_ids.len();
1781 let id = open_ids[idx];
1782 let commit_amount = if amount <= 0 {
1783 1
1784 } else {
1785 amount.saturating_mul(2)
1786 };
1787 match engine.commit(id, commit_amount) {
1788 BudgetSettlement::Committed { .. } => {
1789 open_ids.remove(idx);
1790 }
1791 BudgetSettlement::Overrun { .. } => {}
1792 _ => {}
1793 }
1794 }
1795 2 if !open_ids.is_empty() => {
1796 let idx = sel as usize % open_ids.len();
1797 let id = open_ids.remove(idx);
1798 let _ = engine.release(id);
1799 }
1800 3 => {
1801 if amount > 0 {
1802 let _ = engine.top_up_tenant(&tenant, amount);
1803 }
1804 }
1805 4 => {
1806 let extra = format!("extra-{}", sel % 4);
1807 if amount > 0 {
1808 engine.ensure_tenant(&extra, amount);
1809 }
1810 }
1811 _ => {
1812 let _ = engine.snapshot();
1813 }
1814 }
1815
1816 prop_assert_eq!(engine.verify_conservation(), ConservationStatus::Balanced);
1817 let snap_after = engine.snapshot();
1818 if open_ids.is_empty() && snap_after.active_reservations == 0 {
1819 let digest_after = crate::finance::ledger_digest(&snap_after);
1820 if snap_after.tenants == snap_before.tenants
1821 && snap_after.version == snap_before.version
1822 {
1823 prop_assert_eq!(digest_before, digest_after);
1824 }
1825 }
1826 }
1827 }
1828
1829 #[test]
1830 fn random_ops_maintain_conservation(
1831 seed_ops in prop::collection::vec((0u8..4, any::<u8>(), 1i64..50_000), 1..40),
1832 ) {
1833 let engine = BudgetEngine::default();
1834 engine.ensure_tenant("t0", 5_000_000);
1835 engine.ensure_tenant("t1", 5_000_000);
1836 let mut open_ids = Vec::new();
1837
1838 for (op, tenant_sel, amount) in seed_ops {
1839 let tenant = if tenant_sel % 2 == 0 { "t0" } else { "t1" };
1840 match op % 4 {
1841 0 => {
1842 let (_, id) = engine.try_reserve(tenant, amount);
1843 if let Some(id) = id {
1844 open_ids.push(id);
1845 }
1846 }
1847 1 if !open_ids.is_empty() => {
1848 let idx = (tenant_sel as usize) % open_ids.len();
1849 let id = open_ids[idx];
1850 match engine.commit(id, amount) {
1851 BudgetSettlement::Committed { .. } => {
1852 open_ids.remove(idx);
1853 }
1854 BudgetSettlement::Overrun { .. } => {}
1855 _ => {}
1856 }
1857 }
1858 2 if !open_ids.is_empty() => {
1859 let idx = (tenant_sel as usize) % open_ids.len();
1860 let id = open_ids.remove(idx);
1861 let _ = engine.release(id);
1862 }
1863 3 => {
1864 let _ = engine.top_up_tenant(tenant, amount);
1865 }
1866 _ => {}
1867 }
1868 prop_assert_eq!(engine.verify_conservation(), ConservationStatus::Balanced);
1869 }
1870 }
1871 }
1872}