use std::sync::Arc;
use std::sync::atomic::{AtomicBool, AtomicU8, Ordering};
use jiff::Timestamp;
use crate::deny::DenyReason;
use crate::ids::{AccountId, PolicyRevision, RequestId};
use crate::lease::{AccountOverage, LeaseDebit, LocalLease};
use crate::sharding::Locality;
use crate::snapshot::EnforcementMode;
use crate::units::CostUnits;
use crate::usage::{UsageEvent, UsageSource};
const PENDING_LEASE: u8 = 0;
const PENDING_OVERAGE: u8 = 1;
const COMMITTED_LEASE: u8 = 2;
const COMMITTED_OVERAGE: u8 = 3;
const RELEASED: u8 = 4;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum CancelOutcome {
ZeroCharged,
AlreadyCommitted { units: CostUnits },
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum CommitError {
Cancelled,
Denied(DenyReason),
AlreadyReleased,
AlreadyCommitted,
}
#[derive(Debug, Clone, Copy)]
pub enum CommitFunding<'a> {
LeaseOnly,
OverageFallback {
overage: &'a AccountOverage,
cap: CostUnits,
},
}
impl<'a> CommitFunding<'a> {
#[must_use]
#[inline]
pub fn from_mode(mode: EnforcementMode, overage: &'a AccountOverage) -> Self {
match mode.overage_cap() {
None => Self::LeaseOnly,
Some(cap) => Self::OverageFallback { overage, cap },
}
}
}
#[derive(Debug)]
enum ChargeSource {
Lease {
lease: Arc<LocalLease>,
debit: LeaseDebit,
},
Overage(Arc<AccountOverage>),
}
impl ChargeSource {
const fn pending_phase(&self) -> u8 {
match self {
Self::Lease { .. } => PENDING_LEASE,
Self::Overage(_) => PENDING_OVERAGE,
}
}
const fn committed_phase(&self) -> u8 {
match self {
Self::Lease { .. } => COMMITTED_LEASE,
Self::Overage(_) => COMMITTED_OVERAGE,
}
}
}
#[derive(Debug)]
pub struct Reservation {
source: ChargeSource,
units: CostUnits,
phase: AtomicU8,
}
impl Reservation {
#[inline]
pub fn reserve(
lease: &Arc<LocalLease>,
units: CostUnits,
now: Timestamp,
) -> Result<Reservation, DenyReason> {
Self::reserve_at_locality(Arc::clone(lease), units, now, Locality::current())
}
#[doc(hidden)]
#[inline]
pub fn reserve_at_locality(
lease: Arc<LocalLease>,
units: CostUnits,
now: Timestamp,
locality: Locality,
) -> Result<Reservation, DenyReason> {
let debit = lease.try_reserve_at(units, now, locality)?;
Ok(Reservation {
source: ChargeSource::Lease { lease, debit },
units,
phase: AtomicU8::new(PENDING_LEASE),
})
}
#[inline]
pub fn reserve_overage(
overage: &Arc<AccountOverage>,
units: CostUnits,
cap: CostUnits,
) -> Result<Reservation, DenyReason> {
overage.try_debit(units, cap)?;
Ok(Reservation {
source: ChargeSource::Overage(Arc::clone(overage)),
units,
phase: AtomicU8::new(PENDING_OVERAGE),
})
}
#[must_use]
pub fn units(&self) -> CostUnits {
self.units
}
#[must_use]
pub fn admitted_as_overage(&self) -> bool {
matches!(self.source, ChargeSource::Overage(_))
}
fn account_id(&self) -> AccountId {
match &self.source {
ChargeSource::Lease { lease, .. } => lease.grant().account_id,
ChargeSource::Overage(overage) => overage.account_id(),
}
}
fn window_lapsed(&self, now: Timestamp) -> bool {
match &self.source {
ChargeSource::Lease { lease, .. } => now >= lease.usable_until(),
ChargeSource::Overage(_) => false,
}
}
fn refund(&self) {
match &self.source {
ChargeSource::Lease { lease, debit } => lease.credit(debit),
ChargeSource::Overage(overage) => overage.credit(self.units),
}
}
#[cold]
fn release_for_zero(&self, reason: DenyReason) -> Result<CostUnits, CommitError> {
match self.phase.compare_exchange(
self.source.pending_phase(),
RELEASED,
Ordering::AcqRel,
Ordering::Acquire,
) {
Ok(_) => {
self.refund();
Err(CommitError::Denied(reason))
}
Err(RELEASED) => Err(CommitError::AlreadyReleased),
Err(_) => Err(CommitError::AlreadyCommitted),
}
}
#[inline]
pub fn commit_at_execution_start(
&self,
now: Timestamp,
funding: CommitFunding<'_>,
) -> Result<CostUnits, CommitError> {
if self.window_lapsed(now) {
return self.commit_after_lapse(funding);
}
let claim = || {
self.phase.compare_exchange(
self.source.pending_phase(),
self.source.committed_phase(),
Ordering::AcqRel,
Ordering::Acquire,
)
};
let transition = match &self.source {
ChargeSource::Overage(overage) => overage.publish_claim(self.units, claim),
ChargeSource::Lease { .. } => claim(),
};
match transition {
Ok(_) => Ok(self.units),
Err(RELEASED) => Err(CommitError::AlreadyReleased),
Err(_) => Err(CommitError::AlreadyCommitted),
}
}
#[cold]
fn commit_after_lapse(&self, funding: CommitFunding<'_>) -> Result<CostUnits, CommitError> {
let (CommitFunding::OverageFallback { overage, cap }, ChargeSource::Lease { lease, debit }) =
(funding, &self.source)
else {
return self.release_for_zero(DenyReason::FundingExpiredAtStart);
};
debug_assert_eq!(
overage.account_id(),
self.account_id(),
"commit-time overage fallback must use the reservation's own account counter"
);
if self.phase.load(Ordering::Acquire) != PENDING_LEASE {
return self.release_for_zero(DenyReason::FundingExpiredAtStart);
}
let tentative = match overage.debit_tentatively(self.units, cap) {
Ok(tentative) => tentative,
Err(refused) => return self.release_for_zero(refused),
};
match tentative.publish_commit(|| {
self.phase.compare_exchange(
PENDING_LEASE,
COMMITTED_OVERAGE,
Ordering::AcqRel,
Ordering::Acquire,
)
}) {
Ok(_) => {
lease.credit(debit);
Ok(self.units)
}
Err(RELEASED) => Err(CommitError::AlreadyReleased),
Err(_) => Err(CommitError::AlreadyCommitted),
}
}
#[inline]
pub fn cancel(&self) -> CancelOutcome {
match self.phase.compare_exchange(
self.source.pending_phase(),
RELEASED,
Ordering::AcqRel,
Ordering::Acquire,
) {
Ok(_) => {
self.refund();
CancelOutcome::ZeroCharged
}
Err(COMMITTED_LEASE | COMMITTED_OVERAGE) => {
CancelOutcome::AlreadyCommitted { units: self.units }
}
Err(_) => CancelOutcome::ZeroCharged,
}
}
#[must_use]
pub fn usage_event(
&self,
request_id: RequestId,
now: Timestamp,
policy_revision: PolicyRevision,
key_id: Option<crate::KeyId>,
) -> Option<UsageEvent> {
let source = match self.phase.load(Ordering::Acquire) {
COMMITTED_OVERAGE => UsageSource::Overage,
COMMITTED_LEASE => match &self.source {
ChargeSource::Lease { lease, .. } => {
let grant = lease.grant();
UsageSource::Leased {
lease_id: grant.lease_id,
fencing_token: grant.fencing_token,
}
}
ChargeSource::Overage(_) => {
debug_assert!(false, "a committed-lease phase requires a lease receipt");
UsageSource::Overage
}
},
_ => return None,
};
Some(UsageEvent::new(
request_id,
self.account_id(),
source,
self.units,
now,
policy_revision,
key_id,
))
}
}
impl Drop for Reservation {
fn drop(&mut self) {
if self
.phase
.compare_exchange(
self.source.pending_phase(),
RELEASED,
Ordering::AcqRel,
Ordering::Acquire,
)
.is_ok()
{
self.refund();
}
}
}
#[derive(Debug)]
pub struct SharedCharge {
reservation: Reservation,
cancel_requested: AtomicBool,
}
impl SharedCharge {
#[inline]
#[must_use]
pub fn reservation(&self) -> &Reservation {
&self.reservation
}
#[inline]
#[must_use]
pub fn cancel_requested(&self) -> bool {
self.cancel_requested.load(Ordering::Acquire)
}
#[must_use]
pub fn cancel_handle(self: &Arc<Self>) -> CancelHandle {
CancelHandle(Arc::clone(self))
}
}
#[derive(Debug, Clone)]
pub struct CancelHandle(Arc<SharedCharge>);
impl CancelHandle {
#[inline]
pub fn cancel(&self) -> CancelOutcome {
self.0.cancel_requested.store(true, Ordering::Release);
self.0.reservation.cancel()
}
#[inline]
#[must_use]
pub fn is_cancelled(&self) -> bool {
self.0.cancel_requested()
}
}
impl Reservation {
#[must_use]
pub fn split(self) -> (Arc<SharedCharge>, CancelHandle) {
let shared = Arc::new(SharedCharge {
reservation: self,
cancel_requested: AtomicBool::new(false),
});
let handle = CancelHandle(Arc::clone(&shared));
(shared, handle)
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::ids::{AccountId, FencingToken, LeaseId};
use crate::lease::LeaseGrant;
use crate::usage::UsageSource;
fn t(secs: i64) -> Timestamp {
Timestamp::from_second(secs).unwrap()
}
fn overage() -> Arc<AccountOverage> {
Arc::new(AccountOverage::new(AccountId(1)))
}
#[test]
fn a_committed_overage_bills_against_no_lease() {
let o = overage();
let r = Reservation::reserve_overage(&o, CostUnits(30), CostUnits(100)).unwrap();
assert!(r.admitted_as_overage());
assert_eq!(o.spent(), CostUnits(30));
assert_eq!(
r.commit_at_execution_start(t(1), CommitFunding::LeaseOnly)
.unwrap(),
CostUnits(30)
);
let event = r
.usage_event(RequestId(7), t(1), PolicyRevision::UNSTATED, None)
.unwrap();
assert_eq!(event.account_id, AccountId(1));
assert_eq!(event.units, CostUnits(30));
assert_eq!(event.source, UsageSource::Overage);
assert_eq!(event.source.lease_id(), None);
assert_eq!(event.source.fencing_token(), None);
}
#[test]
fn cancelling_an_overage_returns_the_credit() {
let o = overage();
let r = Reservation::reserve_overage(&o, CostUnits(30), CostUnits(100)).unwrap();
assert_eq!(r.cancel(), CancelOutcome::ZeroCharged);
assert_eq!(o.spent(), CostUnits::ZERO);
assert_eq!(
r.usage_event(RequestId(7), t(1), PolicyRevision::UNSTATED, None),
None
);
}
#[test]
fn dropping_a_pending_overage_returns_the_credit() {
let o = overage();
drop(Reservation::reserve_overage(&o, CostUnits(30), CostUnits(100)).unwrap());
assert_eq!(o.spent(), CostUnits::ZERO);
}
#[test]
fn overage_has_no_usability_window_to_lapse() {
let leased = Reservation::reserve(&lease(100), CostUnits(30), t(0)).unwrap();
assert_eq!(
leased.commit_at_execution_start(t(1_000), CommitFunding::LeaseOnly),
Err(CommitError::Denied(DenyReason::FundingExpiredAtStart))
);
let o = overage();
let r = Reservation::reserve_overage(&o, CostUnits(30), CostUnits(100)).unwrap();
assert_eq!(
r.commit_at_execution_start(t(1_000_000), CommitFunding::LeaseOnly)
.unwrap(),
CostUnits(30)
);
assert_eq!(o.spent(), CostUnits(30), "a commit keeps the credit spent");
}
#[test]
fn an_overage_beyond_the_cap_is_refused_and_claims_nothing() {
let o = overage();
assert_eq!(
Reservation::reserve_overage(&o, CostUnits(101), CostUnits(100)).unwrap_err(),
DenyReason::OverageCapExhausted {
spent: CostUnits::ZERO,
overage_cap: CostUnits(100),
}
);
assert_eq!(o.spent(), CostUnits::ZERO);
}
fn lease(units: u64) -> Arc<LocalLease> {
Arc::new(LocalLease::new(
LeaseGrant {
lease_id: LeaseId(7),
account_id: AccountId(1),
fencing_token: FencingToken(3),
units: CostUnits(units),
expires_at: t(1_000),
},
CostUnits::ZERO,
))
}
#[test]
fn commit_charges_and_keeps_units_spent() {
let l = lease(100);
let r = Reservation::reserve(&l, CostUnits(30), t(0)).unwrap();
assert_eq!(
r.commit_at_execution_start(t(0), CommitFunding::LeaseOnly),
Ok(CostUnits(30))
);
assert_eq!(l.remaining(), CostUnits(70));
assert!(
r.usage_event(RequestId(9), t(1), PolicyRevision::UNSTATED, None)
.is_some()
);
drop(r);
assert_eq!(l.remaining(), CostUnits(70));
}
#[test]
fn a_reservation_reports_the_units_it_will_charge() {
let l = lease(100);
let r = Reservation::reserve(&l, CostUnits(30), t(0)).unwrap();
assert_eq!(r.units(), CostUnits(30));
assert_eq!(
l.remaining(),
CostUnits(70),
"the pending debit is the advertised amount"
);
assert_eq!(
r.commit_at_execution_start(t(0), CommitFunding::LeaseOnly),
Ok(r.units())
);
assert_eq!(
r.usage_event(RequestId(1), t(1), PolicyRevision::UNSTATED, None)
.unwrap()
.units,
r.units(),
"and the billing event carries it too"
);
}
#[test]
fn cancel_charges_zero_and_refunds() {
let l = lease(100);
let r = Reservation::reserve(&l, CostUnits(30), t(0)).unwrap();
assert_eq!(r.cancel(), CancelOutcome::ZeroCharged);
assert_eq!(l.remaining(), CostUnits(100));
assert_eq!(
r.usage_event(RequestId(9), t(1), PolicyRevision::UNSTATED, None),
None
);
assert_eq!(r.cancel(), CancelOutcome::ZeroCharged);
assert_eq!(l.remaining(), CostUnits(100));
}
#[test]
fn drop_releases_pending() {
let l = lease(100);
let r = Reservation::reserve(&l, CostUnits(30), t(0)).unwrap();
assert_eq!(l.remaining(), CostUnits(70));
drop(r);
assert_eq!(l.remaining(), CostUnits(100));
}
#[test]
fn cancel_after_commit_reports_full_charge() {
let l = lease(100);
let r = Reservation::reserve(&l, CostUnits(30), t(0)).unwrap();
r.commit_at_execution_start(t(0), CommitFunding::LeaseOnly)
.unwrap();
assert_eq!(
r.cancel(),
CancelOutcome::AlreadyCommitted {
units: CostUnits(30)
}
);
assert_eq!(l.remaining(), CostUnits(70));
}
#[test]
fn commit_after_cancel_is_refused() {
let l = lease(100);
let r = Reservation::reserve(&l, CostUnits(30), t(0)).unwrap();
assert_eq!(r.cancel(), CancelOutcome::ZeroCharged);
assert_eq!(
r.commit_at_execution_start(t(0), CommitFunding::LeaseOnly),
Err(CommitError::AlreadyReleased)
);
assert_eq!(l.remaining(), CostUnits(100));
}
#[test]
fn double_commit_is_a_surfaced_error() {
let l = lease(100);
let r = Reservation::reserve(&l, CostUnits(30), t(0)).unwrap();
r.commit_at_execution_start(t(0), CommitFunding::LeaseOnly)
.unwrap();
assert_eq!(
r.commit_at_execution_start(t(0), CommitFunding::LeaseOnly),
Err(CommitError::AlreadyCommitted)
);
}
#[test]
fn commit_after_window_closes_releases_for_zero() {
let l = lease(100); let r = Reservation::reserve(&l, CostUnits(30), t(999)).unwrap();
assert_eq!(l.remaining(), CostUnits(70));
assert_eq!(
r.commit_at_execution_start(t(1_000), CommitFunding::LeaseOnly),
Err(CommitError::Denied(DenyReason::FundingExpiredAtStart))
);
assert_eq!(l.remaining(), CostUnits(100));
assert_eq!(
r.usage_event(RequestId(1), t(1_001), PolicyRevision::UNSTATED, None),
None
);
assert_eq!(
r.commit_at_execution_start(t(999), CommitFunding::LeaseOnly),
Err(CommitError::AlreadyReleased)
);
}
#[test]
fn an_elastic_lapse_at_execution_start_bills_as_overage_with_no_lease_capability() {
let l = lease(100);
let o = overage();
let r = Reservation::reserve(&l, CostUnits(30), t(999)).unwrap();
assert!(
!r.admitted_as_overage(),
"admission was funded by the lease"
);
assert_eq!(l.remaining(), CostUnits(70));
let funding = CommitFunding::OverageFallback {
overage: &o,
cap: CostUnits(100),
};
assert_eq!(
r.commit_at_execution_start(t(1_000), funding),
Ok(CostUnits(30))
);
assert_eq!(l.remaining(), CostUnits(100));
assert_eq!(o.spent(), CostUnits(30));
let event = r
.usage_event(RequestId(7), t(1_000), PolicyRevision::UNSTATED, None)
.unwrap();
assert_eq!(event.units, CostUnits(30));
assert_eq!(event.account_id, AccountId(1));
assert_eq!(event.source, UsageSource::Overage);
assert_eq!(event.source.lease_id(), None);
assert_eq!(event.source.fencing_token(), None);
assert!(!r.admitted_as_overage());
}
#[test]
fn a_strict_lapse_at_execution_start_releases_for_zero_and_yields_no_event() {
let l = lease(100);
let o = overage();
let r = Reservation::reserve(&l, CostUnits(30), t(999)).unwrap();
assert_eq!(
r.commit_at_execution_start(t(1_000), CommitFunding::LeaseOnly),
Err(CommitError::Denied(DenyReason::FundingExpiredAtStart))
);
assert_eq!(l.remaining(), CostUnits(100));
assert_eq!(
r.usage_event(RequestId(7), t(1_000), PolicyRevision::UNSTATED, None),
None
);
assert_eq!(o.spent(), CostUnits::ZERO);
}
#[test]
fn an_elastic_lapse_with_no_committed_headroom_releases_the_lease_for_zero() {
let l = lease(100);
let o = overage();
let held = Reservation::reserve_overage(&o, CostUnits(100), CostUnits(100)).unwrap();
held.commit_at_execution_start(t(0), CommitFunding::LeaseOnly)
.unwrap();
let r = Reservation::reserve(&l, CostUnits(30), t(999)).unwrap();
let funding = CommitFunding::OverageFallback {
overage: &o,
cap: CostUnits(100),
};
assert_eq!(
r.commit_at_execution_start(t(1_000), funding),
Err(CommitError::Denied(DenyReason::OverageCapExhausted {
spent: CostUnits(100),
overage_cap: CostUnits(100),
}))
);
assert_eq!(l.remaining(), CostUnits(100));
assert_eq!(o.spent(), CostUnits(100));
assert_eq!(
r.usage_event(RequestId(7), t(1_000), PolicyRevision::UNSTATED, None),
None
);
}
#[test]
fn an_elastic_lapse_blocked_by_refundable_occupancy_is_temporarily_exhausted() {
let l = lease(100);
let o = overage();
let _sibling = Reservation::reserve_overage(&o, CostUnits(100), CostUnits(100)).unwrap();
let r = Reservation::reserve(&l, CostUnits(30), t(999)).unwrap();
let funding = CommitFunding::OverageFallback {
overage: &o,
cap: CostUnits(100),
};
let denied = r.commit_at_execution_start(t(1_000), funding).unwrap_err();
assert_eq!(
denied,
CommitError::Denied(DenyReason::OverageCapTemporarilyExhausted {
spent: CostUnits(100),
overage_cap: CostUnits(100),
})
);
let CommitError::Denied(reason) = denied else {
unreachable!("the fallback reports a classified denial")
};
assert_eq!(reason.retry(), crate::deny::Retry::Transient);
assert_eq!(l.remaining(), CostUnits(100));
}
#[test]
fn an_elastic_lapse_overlapping_a_sibling_publication_reports_an_in_flight_commit() {
let l = lease(100);
let o = overage();
let sibling = Reservation::reserve_overage(&o, CostUnits(100), CostUnits(100)).unwrap();
let r = Reservation::reserve(&l, CostUnits(30), t(999)).unwrap();
o.publish_claim(sibling.units, || {
let claimed = sibling.phase.compare_exchange(
PENDING_OVERAGE,
COMMITTED_OVERAGE,
Ordering::AcqRel,
Ordering::Acquire,
);
assert!(claimed.is_ok(), "the sibling wins its own phase claim");
let funding = CommitFunding::OverageFallback {
overage: &o,
cap: CostUnits(100),
};
let denied = r.commit_at_execution_start(t(1_000), funding).unwrap_err();
let CommitError::Denied(reason) = denied else {
unreachable!("the fallback reports a classified denial")
};
assert!(
matches!(reason, DenyReason::OverageCommitInProgress { .. }),
"expected an in-flight publication refusal, got {reason:?}"
);
assert_eq!(reason.retry(), crate::deny::Retry::AfterInFlight);
Ok::<_, ()>(())
})
.unwrap();
assert_eq!(l.remaining(), CostUnits(100));
assert_eq!(o.spent(), CostUnits(100));
}
#[test]
fn a_fallback_commit_refunds_its_lease_exactly_once() {
let l = lease(100);
let o = overage();
let r = Reservation::reserve(&l, CostUnits(30), t(999)).unwrap();
let funding = CommitFunding::OverageFallback {
overage: &o,
cap: CostUnits(100),
};
r.commit_at_execution_start(t(1_000), funding).unwrap();
assert_eq!(l.remaining(), CostUnits(100));
assert_eq!(
r.cancel(),
CancelOutcome::AlreadyCommitted {
units: CostUnits(30)
}
);
drop(r);
assert_eq!(l.remaining(), CostUnits(100));
assert_eq!(o.spent(), CostUnits(30));
}
#[test]
fn a_fallback_that_loses_to_cancel_strands_no_overage_capacity() {
let l = lease(100);
let o = overage();
let r = Reservation::reserve(&l, CostUnits(30), t(999)).unwrap();
assert_eq!(r.cancel(), CancelOutcome::ZeroCharged);
assert_eq!(l.remaining(), CostUnits(100));
let funding = CommitFunding::OverageFallback {
overage: &o,
cap: CostUnits(100),
};
assert_eq!(
r.commit_at_execution_start(t(1_000), funding),
Err(CommitError::AlreadyReleased)
);
assert_eq!(o.spent(), CostUnits::ZERO);
let fresh = Reservation::reserve_overage(&o, CostUnits(100), CostUnits(100));
assert!(fresh.is_ok(), "the full cap must remain reusable");
assert_eq!(l.remaining(), CostUnits(100));
}
#[test]
fn a_native_overage_reservation_ignores_a_commit_time_fallback() {
let o = overage();
let r = Reservation::reserve_overage(&o, CostUnits(30), CostUnits(100)).unwrap();
assert_eq!(o.spent(), CostUnits(30));
let funding = CommitFunding::OverageFallback {
overage: &o,
cap: CostUnits(100),
};
assert_eq!(
r.commit_at_execution_start(t(9_999), funding),
Ok(CostUnits(30))
);
assert_eq!(o.spent(), CostUnits(30));
assert_eq!(
r.usage_event(RequestId(1), t(9_999), PolicyRevision::UNSTATED, None)
.unwrap()
.source,
UsageSource::Overage
);
}
#[test]
fn a_second_commit_after_a_fallback_is_a_surfaced_error() {
let l = lease(100);
let o = overage();
let r = Reservation::reserve(&l, CostUnits(30), t(999)).unwrap();
let funding = CommitFunding::OverageFallback {
overage: &o,
cap: CostUnits(100),
};
r.commit_at_execution_start(t(1_000), funding).unwrap();
assert_eq!(
r.commit_at_execution_start(t(1_000), funding),
Err(CommitError::AlreadyCommitted)
);
assert_eq!(o.spent(), CostUnits(30));
assert_eq!(l.remaining(), CostUnits(100));
}
#[test]
fn safety_margin_closes_window_before_expiry() {
let l = Arc::new(LocalLease::with_safety_margin(
LeaseGrant {
lease_id: LeaseId(7),
account_id: AccountId(1),
fencing_token: FencingToken(3),
units: CostUnits(100),
expires_at: t(1_000),
},
CostUnits::ZERO,
jiff::SignedDuration::from_secs(10),
));
assert_eq!(l.usable_until(), t(990));
assert_eq!(
Reservation::reserve(&l, CostUnits(1), t(990)).unwrap_err(),
DenyReason::LeaseExpired
);
let r = Reservation::reserve(&l, CostUnits(1), t(989)).unwrap();
assert_eq!(
r.commit_at_execution_start(t(990), CommitFunding::LeaseOnly),
Err(CommitError::Denied(DenyReason::FundingExpiredAtStart))
);
}
#[test]
fn commit_cancel_race_one_winner() {
for _ in 0..500 {
let l = lease(100);
let r = Arc::new(Reservation::reserve(&l, CostUnits(10), t(0)).unwrap());
let rc = Arc::clone(&r);
let committer = std::thread::spawn(move || {
rc.commit_at_execution_start(t(0), CommitFunding::LeaseOnly)
});
let canceller = std::thread::spawn({
let rc = Arc::clone(&r);
move || rc.cancel()
});
let commit = committer.join().unwrap();
let cancel = canceller.join().unwrap();
match (commit, cancel) {
(Ok(units), CancelOutcome::AlreadyCommitted { units: seen }) => {
assert_eq!(units, seen);
assert_eq!(l.remaining(), CostUnits(90));
}
(Err(CommitError::AlreadyReleased), CancelOutcome::ZeroCharged) => {
assert_eq!(l.remaining(), CostUnits(100));
}
other => panic!("impossible race outcome: {other:?}"),
}
}
}
#[test]
fn fallback_commit_and_cancel_leave_exactly_one_funding_term() {
for _ in 0..500 {
let l = lease(100);
let o = overage();
let r = Arc::new(Reservation::reserve(&l, CostUnits(10), t(999)).unwrap());
let committer = std::thread::spawn({
let rc = Arc::clone(&r);
let oc = Arc::clone(&o);
move || {
rc.commit_at_execution_start(
t(1_000),
CommitFunding::OverageFallback {
overage: &oc,
cap: CostUnits(100),
},
)
}
});
let canceller = std::thread::spawn({
let rc = Arc::clone(&r);
move || rc.cancel()
});
let commit = committer.join().unwrap();
let cancel = canceller.join().unwrap();
match (commit, cancel) {
(Ok(units), CancelOutcome::AlreadyCommitted { units: seen }) => {
assert_eq!(units, seen);
assert_eq!(l.remaining(), CostUnits(100));
assert_eq!(o.spent(), CostUnits(10));
assert_eq!(
r.usage_event(RequestId(1), t(1_000), PolicyRevision::UNSTATED, None)
.unwrap()
.source,
UsageSource::Overage
);
}
(Err(CommitError::AlreadyReleased), CancelOutcome::ZeroCharged) => {
assert_eq!(l.remaining(), CostUnits(100));
assert_eq!(o.spent(), CostUnits::ZERO);
assert_eq!(
r.usage_event(RequestId(1), t(1_000), PolicyRevision::UNSTATED, None),
None
);
}
other => panic!("impossible race outcome: {other:?}"),
}
}
}
#[test]
fn concurrent_fallbacks_never_exceed_the_overage_cap() {
let o = overage();
let cap = CostUnits(60);
let reservations: Vec<_> = (0..10)
.map(|_| {
let l = lease(100);
let r = Reservation::reserve(&l, CostUnits(10), t(999)).unwrap();
(l, r)
})
.collect();
let committed = std::thread::scope(|scope| {
let handles: Vec<_> = reservations
.iter()
.map(|(_, r)| {
let oc = Arc::clone(&o);
scope.spawn(move || {
r.commit_at_execution_start(
t(1_000),
CommitFunding::OverageFallback { overage: &oc, cap },
)
.is_ok()
})
})
.collect();
handles
.into_iter()
.map(|handle| handle.join().expect("no worker panics"))
.filter(|won| *won)
.count()
});
assert!(o.spent() <= cap, "overage spend exceeded its cap");
assert_eq!(o.spent(), CostUnits(committed as u64 * 10));
assert!(committed <= 6, "at most six ten-unit charges fit in sixty");
for (l, _) in &reservations {
assert_eq!(l.remaining(), CostUnits(100));
}
}
#[test]
fn split_commit_and_handle_cancel_have_exactly_one_winner() {
for _ in 0..500 {
let l = lease(100);
let (shared, handle) = Reservation::reserve(&l, CostUnits(10), t(0))
.unwrap()
.split();
let worker = std::thread::spawn({
let shared = Arc::clone(&shared);
move || {
shared
.reservation()
.commit_at_execution_start(t(0), CommitFunding::LeaseOnly)
}
});
let waiter = std::thread::spawn(move || handle.cancel());
let commit = worker.join().unwrap();
let cancel = waiter.join().unwrap();
match (commit, cancel) {
(Ok(units), CancelOutcome::AlreadyCommitted { units: seen }) => {
assert_eq!(units, seen);
assert_eq!(l.remaining(), CostUnits(90));
}
(Err(CommitError::AlreadyReleased), CancelOutcome::ZeroCharged) => {
assert_eq!(l.remaining(), CostUnits(100));
}
other => panic!("impossible race outcome: {other:?}"),
}
assert!(shared.cancel_requested());
}
}
#[test]
fn a_cancel_handle_before_worker_start_charges_zero() {
let l = lease(100);
let (shared, handle) = Reservation::reserve(&l, CostUnits(30), t(0))
.unwrap()
.split();
assert!(!handle.is_cancelled());
assert_eq!(l.remaining(), CostUnits(70));
assert_eq!(handle.cancel(), CancelOutcome::ZeroCharged);
assert!(handle.is_cancelled());
assert_eq!(l.remaining(), CostUnits(100));
assert_eq!(
shared
.reservation()
.commit_at_execution_start(t(0), CommitFunding::LeaseOnly),
Err(CommitError::AlreadyReleased),
"the worker must not execute after a winning cancellation"
);
assert_eq!(
shared
.reservation()
.usage_event(RequestId(1), t(0), PolicyRevision::UNSTATED, None),
None
);
}
#[test]
fn a_late_cancel_reports_the_full_charge_and_still_records_the_request() {
let l = lease(100);
let (shared, handle) = Reservation::reserve(&l, CostUnits(30), t(0))
.unwrap()
.split();
shared
.reservation()
.commit_at_execution_start(t(0), CommitFunding::LeaseOnly)
.unwrap();
assert_eq!(
handle.cancel(),
CancelOutcome::AlreadyCommitted {
units: CostUnits(30)
}
);
assert_eq!(l.remaining(), CostUnits(70), "the charge stands in full");
assert!(
shared.cancel_requested(),
"a kernel may still learn its caller has gone"
);
assert!(
shared
.reservation()
.usage_event(RequestId(1), t(0), PolicyRevision::UNSTATED, None)
.is_some()
);
}
#[test]
fn a_second_cancel_handle_resolves_the_same_charge() {
let l = lease(100);
let (shared, first) = Reservation::reserve(&l, CostUnits(30), t(0))
.unwrap()
.split();
let second = shared.cancel_handle();
assert_eq!(first.cancel(), CancelOutcome::ZeroCharged);
assert_eq!(second.cancel(), CancelOutcome::ZeroCharged);
assert!(second.is_cancelled());
assert_eq!(l.remaining(), CostUnits(100));
}
#[test]
fn a_shared_fallback_and_a_handle_cancel_leave_one_funding_term() {
for _ in 0..200 {
let l = lease(100);
let o = overage();
let (shared, handle) = Reservation::reserve(&l, CostUnits(10), t(999))
.unwrap()
.split();
let worker = std::thread::spawn({
let shared = Arc::clone(&shared);
let o = Arc::clone(&o);
move || {
shared.reservation().commit_at_execution_start(
t(1_000),
CommitFunding::OverageFallback {
overage: &o,
cap: CostUnits(100),
},
)
}
});
let waiter = std::thread::spawn(move || handle.cancel());
let commit = worker.join().unwrap();
let cancel = waiter.join().unwrap();
match (commit, cancel) {
(Ok(_), CancelOutcome::AlreadyCommitted { .. }) => {
assert_eq!(l.remaining(), CostUnits(100));
assert_eq!(o.spent(), CostUnits(10));
}
(Err(CommitError::AlreadyReleased), CancelOutcome::ZeroCharged) => {
assert_eq!(l.remaining(), CostUnits(100));
assert_eq!(o.spent(), CostUnits::ZERO);
}
other => panic!("impossible race outcome: {other:?}"),
}
}
}
#[test]
fn overage_commit_cancel_race_preserves_retry_classification() {
for _ in 0..500 {
let o = overage();
let r =
Arc::new(Reservation::reserve_overage(&o, CostUnits(10), CostUnits(10)).unwrap());
let committer = std::thread::spawn({
let r = Arc::clone(&r);
move || r.commit_at_execution_start(t(0), CommitFunding::LeaseOnly)
});
let canceller = std::thread::spawn({
let r = Arc::clone(&r);
move || r.cancel()
});
match (committer.join().unwrap(), canceller.join().unwrap()) {
(Ok(units), CancelOutcome::AlreadyCommitted { units: seen }) => {
assert_eq!(units, seen);
assert_eq!(o.spent(), CostUnits(10));
let denied =
Reservation::reserve_overage(&o, CostUnits(1), CostUnits(10)).unwrap_err();
assert!(matches!(denied, DenyReason::OverageCapExhausted { .. }));
assert_eq!(denied.retry(), crate::deny::Retry::Transient);
}
(Err(CommitError::AlreadyReleased), CancelOutcome::ZeroCharged) => {
assert_eq!(o.spent(), CostUnits::ZERO);
drop(
Reservation::reserve_overage(&o, CostUnits(10), CostUnits(10))
.expect("cancelled credit is immediately reusable"),
);
}
other => panic!("impossible overage race outcome: {other:?}"),
}
}
}
#[test]
fn overage_commit_publication_never_looks_refundable() {
let overage = overage();
let reservation =
Reservation::reserve_overage(&overage, CostUnits(100), CostUnits(100)).unwrap();
overage
.publish_claim(reservation.units, || {
let claimed = reservation.phase.compare_exchange(
PENDING_OVERAGE,
COMMITTED_OVERAGE,
Ordering::AcqRel,
Ordering::Acquire,
);
assert!(claimed.is_ok(), "this commit wins the phase claim");
assert_eq!(
{
let denied =
Reservation::reserve_overage(&overage, CostUnits(1), CostUnits(100))
.unwrap_err();
assert_eq!(denied.retry(), crate::deny::Retry::AfterInFlight);
denied
},
DenyReason::OverageCommitInProgress {
spent: CostUnits(100),
overage_cap: CostUnits(100),
},
"an irrevocable commit is identified as an in-flight publication"
);
claimed
})
.unwrap();
assert!(matches!(
Reservation::reserve_overage(&overage, CostUnits(1), CostUnits(100)).unwrap_err(),
DenyReason::OverageCapExhausted { .. }
));
}
}