use std::num::NonZeroUsize;
use std::sync::Arc;
use std::sync::atomic::{AtomicU32, AtomicUsize, Ordering};
use jiff::{SignedDuration, Timestamp};
use tollgate_admission::{
AdmissionEngine, ArcSwapSnapshotMap, LeaseSlot, NoCapacityPermit, NoGate, ReadyToStart,
SnapshotMap,
};
use tollgate_client::{
Clock, LeaseManager, LeaseManagerConfig, ManualClock, UsagePermit, UsageWriter,
UsageWriterConfig,
};
use tollgate_core::{
AccountId, AccountSnapshot, AccountStatus, CapacityClass, CostTable, CostUnits, DenyReason,
Generation, LeaseGrant, LocalLease, LocalSharding, OpIndex, PermissionBits, PolicyRevision,
Principal, RequestId, ResolvedLimits, UsageEvent, UsageSource,
};
#[path = "../../tollgate-store/tests/support/delegating.rs"]
mod delegating;
use delegating::{DelegatingStore, RejectingStore, rejecting};
use tollgate_store::{
AccountConfig, GrantPolicy, IngestError, IngestReport, LeaseAllocator, MemoryStore,
ReclaimBatch, StoreError, UsageSink,
};
const ACCOUNT: AccountId = AccountId(1);
const PRINCIPAL: Principal = Principal(1);
#[derive(Clone, Copy)]
struct Operation;
impl OpIndex for Operation {
fn index(&self) -> usize {
0
}
}
fn t(secs: i64) -> Timestamp {
Timestamp::from_second(secs).unwrap()
}
fn ready_from_grant(
grant: LeaseGrant,
permit: UsagePermit,
) -> ReadyToStart<UsagePermit, NoCapacityPermit> {
let engine = AdmissionEngine::new(ArcSwapSnapshotMap::new());
let snapshot = Arc::new(
AccountSnapshot::builder(
ACCOUNT,
Generation(1),
AccountStatus::Active,
t(1_000),
PermissionBits::bit(0),
ResolvedLimits::new(64),
Arc::new(
CostTable::builder(CostUnits(50), CostUnits(50))
.weight(&Operation, CostUnits(1))
.build(),
),
)
.build(),
);
let slot = LeaseSlot::for_account(ACCOUNT);
drop(slot.replace(Arc::new(LocalLease::new(grant, CostUnits::ZERO))));
engine.map().install(PRINCIPAL, snapshot, slot).unwrap();
engine
.begin(PRINCIPAL, PermissionBits::bit(0), t(0))
.and_then(|context| context.admit(&[(&Operation, 1)], permit, t(0)))
.expect("the fixture lease funds one request")
.acquire_capacity(&NoGate)
.expect("capacity is disabled in this fixture")
}
fn store(balance: u64) -> Arc<MemoryStore> {
let store = MemoryStore::new(GrantPolicy {
shrink_divisor: 1,
min_grant: CostUnits(1),
max_ttl: SignedDuration::from_secs(3_600),
reclaim_grace: SignedDuration::ZERO,
})
.unwrap();
store.create_account(AccountConfig {
account_id: ACCOUNT,
initial_balance: CostUnits(balance),
status: AccountStatus::Active,
capacity_class: CapacityClass::Assured,
});
store
}
fn manager_config() -> LeaseManagerConfig {
LeaseManagerConfig {
account: ACCOUNT,
target_grant: CostUnits(1_000),
low_water: CostUnits(250),
lease_ttl: SignedDuration::from_secs(60),
expiry_safety_margin: SignedDuration::ZERO,
poll_interval: std::time::Duration::from_millis(5),
store_call_timeout: std::time::Duration::from_secs(5),
shutdown_release_deadline: std::time::Duration::from_secs(10),
}
}
async fn settle() {
tokio::time::sleep(std::time::Duration::from_millis(50)).await;
}
#[tokio::test(start_paused = true)]
async fn refill_installs_lease_on_cold_start() {
let store = store(10_000);
let slot = LeaseSlot::for_account(ACCOUNT);
let clock = Arc::new(ManualClock::new(t(0)));
let manager =
LeaseManager::spawn(store.clone(), Arc::clone(&slot), clock, manager_config()).unwrap();
settle().await;
let lease = slot.load().expect("lease installed");
assert_eq!(lease.grant().units, CostUnits(1_000));
manager.shutdown().await;
}
#[tokio::test(start_paused = true)]
async fn refill_begins_on_the_crossing_debit_not_the_next_tick() {
let store = store(10_000);
let slot = LeaseSlot::for_account(ACCOUNT);
let clock = Arc::new(ManualClock::new(t(0)));
let manager = LeaseManager::spawn(
store.clone(),
Arc::clone(&slot),
clock,
LeaseManagerConfig {
poll_interval: std::time::Duration::from_secs(60),
..manager_config()
},
)
.unwrap();
settle().await;
let first = slot.load().expect("cold start installs a lease");
let first_id = first.grant().lease_id;
first.try_debit(CostUnits(800), t(0)).unwrap();
assert!(first.needs_refill());
drop(first);
settle().await;
let second = slot.load().expect("a replacement must be installed");
assert_ne!(
second.grant().lease_id,
first_id,
"the crossing debit must have started a refill well inside the \
60s poll interval; on the polling-only design this is still the \
original lease"
);
manager.shutdown().await;
}
#[tokio::test(start_paused = true)]
async fn a_burst_across_a_rotation_never_denies_a_funded_account() {
let store = store(1_000_000);
let slot = LeaseSlot::for_account(ACCOUNT);
let clock = Arc::new(ManualClock::new(t(0)));
let manager = LeaseManager::spawn(
store.clone(),
Arc::clone(&slot),
clock,
LeaseManagerConfig {
poll_interval: std::time::Duration::from_secs(60),
..manager_config()
},
)
.unwrap();
settle().await;
let mut denied = 0;
let mut spent = 0u64;
for _ in 0..200 {
match slot.load() {
Some(lease) => match lease.try_debit(CostUnits(75), t(0)) {
Ok(()) => spent += 75,
Err(_) => denied += 1,
},
None => denied += 1,
}
tokio::time::sleep(std::time::Duration::from_millis(1)).await;
}
assert_eq!(
denied, 0,
"a funded account was refused {denied} times across rotations"
);
assert_eq!(spent, 200 * 75);
manager.shutdown().await;
}
#[tokio::test(start_paused = true)]
async fn adaptive_tail_grant_does_not_rotate_while_unspent() {
let store = MemoryStore::new(GrantPolicy::default()).unwrap();
store.create_account(AccountConfig {
account_id: ACCOUNT,
initial_balance: CostUnits(50),
status: AccountStatus::Active,
capacity_class: CapacityClass::Assured,
});
let clock = Arc::new(ManualClock::new(t(0)));
let slot = LeaseSlot::for_account(ACCOUNT);
let manager = LeaseManager::spawn(
store.clone(),
Arc::clone(&slot),
clock,
LeaseManagerConfig {
target_grant: CostUnits(500),
low_water: CostUnits(100),
..manager_config()
},
)
.unwrap();
tokio::time::sleep(std::time::Duration::from_millis(25)).await;
let first = slot.load().unwrap();
assert_eq!(first.grant().units, CostUnits(25));
let lease_id = first.grant().lease_id;
drop(first);
tokio::time::sleep(std::time::Duration::from_millis(100)).await;
assert_eq!(slot.load().unwrap().grant().lease_id, lease_id);
manager.shutdown().await;
}
#[tokio::test(start_paused = true)]
async fn rotation_at_low_water_installs_fresh_lease() {
let store = store(10_000);
let slot = LeaseSlot::with_sharding(ACCOUNT, LocalSharding::new(NonZeroUsize::new(8).unwrap()));
let clock = Arc::new(ManualClock::new(t(0)));
let manager =
LeaseManager::spawn(store.clone(), Arc::clone(&slot), clock, manager_config()).unwrap();
settle().await;
let first = slot.load().unwrap();
first.try_debit(CostUnits(800), t(0)).unwrap();
assert!(first.needs_refill());
settle().await;
let second = slot.load().unwrap();
assert_ne!(first.grant().lease_id, second.grant().lease_id);
assert!(second.grant().fencing_token > first.grant().fencing_token);
assert_eq!(first.remaining(), CostUnits(200));
settle().await;
assert_eq!(store.balance(ACCOUNT), CostUnits(8_000));
drop(first);
settle().await;
assert_eq!(store.balance(ACCOUNT), CostUnits(8_200));
manager.shutdown().await;
}
#[tokio::test(start_paused = true)]
async fn expired_slot_fails_closed_then_recovers() {
let store = store(2_000);
let slot = LeaseSlot::for_account(ACCOUNT);
let clock = Arc::new(ManualClock::new(t(0)));
let manager = LeaseManager::spawn(
store.clone(),
Arc::clone(&slot),
Arc::clone(&clock) as Arc<dyn Clock>,
manager_config(),
)
.unwrap();
settle().await;
assert!(slot.load().is_some());
let drain = store
.acquire(
ACCOUNT,
CostUnits(u64::MAX),
SignedDuration::from_secs(3_600),
t(0),
)
.await
.unwrap()
.grant;
clock.set(t(120)); settle().await;
assert!(slot.load().is_none(), "expired slot must fail closed");
store
.release(drain.lease_id, drain.fencing_token, drain.units, t(121))
.await
.unwrap();
settle().await;
assert!(
slot.load().is_some(),
"manager recovers once balance exists"
);
manager.shutdown().await;
}
#[tokio::test(start_paused = true)]
async fn usability_window_rollover_returns_unspent_capacity() {
let store = MemoryStore::new(GrantPolicy {
shrink_divisor: 1,
min_grant: CostUnits(1),
max_ttl: SignedDuration::from_secs(3_600),
reclaim_grace: SignedDuration::from_secs(30),
})
.unwrap();
store.create_account(AccountConfig {
account_id: ACCOUNT,
initial_balance: CostUnits(1_000),
status: AccountStatus::Active,
capacity_class: CapacityClass::Assured,
});
let slot = LeaseSlot::for_account(ACCOUNT);
let clock = Arc::new(ManualClock::new(t(0)));
let mut config = manager_config();
config.expiry_safety_margin = SignedDuration::from_secs(5);
let manager = LeaseManager::spawn(
store.clone(),
Arc::clone(&slot),
Arc::clone(&clock) as Arc<dyn Clock>,
config,
)
.unwrap();
settle().await;
let first = slot.load().unwrap().grant().lease_id;
clock.set(t(55));
settle().await;
let replacement = slot.load().expect("capacity should rotate during grace");
assert_ne!(replacement.grant().lease_id, first);
assert_eq!(replacement.remaining(), CostUnits(1_000));
manager.shutdown().await;
}
#[test]
fn invalid_lease_manager_durations_are_rejected() {
for ttl in [SignedDuration::ZERO, SignedDuration::from_nanos(-1)] {
let mut config = manager_config();
config.lease_ttl = ttl;
config.expiry_safety_margin = SignedDuration::ZERO;
assert!(config.validate().is_err());
}
let mut config = manager_config();
config.expiry_safety_margin = SignedDuration::from_secs(-1);
assert!(config.validate().is_err());
let mut config = manager_config();
config.poll_interval = std::time::Duration::ZERO;
assert!(config.validate().is_err());
}
#[test]
fn fractional_and_wide_lease_ttls_remain_valid_configuration() {
for ttl in [
SignedDuration::from_nanos(1),
SignedDuration::from_millis(500),
SignedDuration::from_millis(1_500),
SignedDuration::from_secs(i64::from(u32::MAX) + 1),
SignedDuration::MAX,
] {
let mut config = manager_config();
config.lease_ttl = ttl;
config.expiry_safety_margin = SignedDuration::ZERO;
assert_eq!(config.validate(), Ok(()), "TTL {ttl}");
}
}
#[tokio::test(start_paused = true)]
async fn shutdown_releases_unspent_units() {
let store = store(10_000);
let slot = LeaseSlot::with_sharding(ACCOUNT, LocalSharding::new(NonZeroUsize::new(8).unwrap()));
let clock = Arc::new(ManualClock::new(t(0)));
let manager =
LeaseManager::spawn(store.clone(), Arc::clone(&slot), clock, manager_config()).unwrap();
let health = manager.health();
settle().await;
slot.load()
.unwrap()
.try_debit(CostUnits(300), t(0))
.unwrap();
manager.shutdown().await;
assert!(!*health.borrow());
assert!(health.has_changed().is_err());
assert!(slot.load().is_none());
assert_eq!(store.balance(ACCOUNT), CostUnits(9_700));
}
fn event(request: u128, units: u64, lease: &tollgate_core::LeaseGrant) -> UsageEvent {
UsageEvent::new(
RequestId(request),
lease.account_id,
UsageSource::Leased {
lease_id: lease.lease_id,
fencing_token: lease.fencing_token,
},
CostUnits(units),
t(0),
PolicyRevision::UNSTATED,
None,
)
}
fn writer_config(capacity: usize) -> UsageWriterConfig {
UsageWriterConfig {
queue_capacity: capacity,
max_batch: 4,
flush_interval: std::time::Duration::from_millis(10),
retry_backoff: std::time::Duration::from_millis(10),
shutdown_drain_deadline: std::time::Duration::from_secs(60),
ingest_timeout: std::time::Duration::from_secs(5),
}
}
#[tokio::test(start_paused = true)]
async fn refill_counters_attribute_refusals_to_their_reason() {
let store = MemoryStore::new(GrantPolicy::default()).unwrap();
let slot = LeaseSlot::for_account(ACCOUNT);
let clock = Arc::new(ManualClock::new(t(0)));
let manager =
LeaseManager::spawn(store.clone(), Arc::clone(&slot), clock, manager_config()).unwrap();
let counters = manager.counters();
assert_eq!(counters.snapshot().refused(), 0, "nothing tried yet");
settle().await;
let stats = counters.snapshot();
assert_eq!(stats.acquired, 0);
assert_eq!(stats.acquired_units, 0);
assert!(stats.refused() > 0, "the allocator refused every acquire");
let refusals: Vec<_> = stats
.refusals_by_name()
.filter(|(_, count)| *count > 0)
.map(|(name, _)| name)
.collect();
assert_eq!(
refusals,
vec!["unknown_account"],
"one reason, and the others left alone: {:?}",
stats.acquire_refused
);
assert_eq!(stats.acquire_timeouts, 0, "a refusal is not a timeout");
manager.shutdown().await;
}
#[tokio::test(start_paused = true)]
async fn refill_counters_record_grants_and_releases() {
let store = store(10_000);
let slot = LeaseSlot::for_account(ACCOUNT);
let clock = Arc::new(ManualClock::new(t(0)));
let manager =
LeaseManager::spawn(store.clone(), Arc::clone(&slot), clock, manager_config()).unwrap();
let counters = manager.counters();
settle().await;
let stats = counters.snapshot();
assert_eq!(stats.acquired, 1);
assert_eq!(stats.acquired_units, 1_000, "the grant, not the request");
assert_eq!(stats.refused(), 0);
assert_eq!(stats.released, 0, "nothing returned while still serving");
assert_eq!(stats.abandoned, 0);
manager.shutdown().await;
assert_eq!(counters.snapshot().released, 1);
assert_eq!(counters.snapshot().abandoned, 0);
}
#[tokio::test(start_paused = true)]
async fn running_totals_are_readable_and_match_the_final_report() {
let store = store(10_000);
let lease = store
.acquire(
ACCOUNT,
CostUnits(1_000),
SignedDuration::from_secs(60),
t(0),
)
.await
.unwrap()
.grant;
let clock = Arc::new(ManualClock::new(t(100)));
let (recorder, writer) = UsageWriter::spawn(store.clone(), clock, writer_config(64)).unwrap();
let health = recorder.health();
assert_eq!(health.stats, tollgate_client::WriterStats::ZERO);
assert_eq!(health.last_ingest_at, None);
assert_eq!(health.ingest_age(t(100)), None);
assert_eq!(health.queue_capacity, 64);
for i in 0..6u128 {
let request = if i == 5 { 0 } else { i };
recorder
.try_reserve()
.unwrap()
.record(event(request, 10, &lease));
}
settle().await;
let mid_flight = recorder.health();
assert_eq!(mid_flight.stats.accepted, 5);
assert_eq!(mid_flight.stats.duplicate, 1);
assert_eq!(mid_flight.stats.lost, 0);
assert_eq!(
mid_flight.unaccounted, 0,
"every event has a billing outcome once flushed"
);
assert_eq!(
mid_flight.last_ingest_at,
Some(t(100)),
"the sink answered, at the instant the batch was ingested with"
);
assert_eq!(
mid_flight.ingest_age(t(160)),
Some(SignedDuration::from_secs(60))
);
let stats = writer.shutdown().await.unwrap();
assert_eq!(
stats, mid_flight.stats,
"the final report is a read of the same counters, not a second tally"
);
}
#[tokio::test(start_paused = true)]
async fn queue_depth_rises_before_the_shed_and_sheds_are_counted() {
let store = store(10_000);
let clock = Arc::new(ManualClock::new(t(0)));
let (recorder, writer) = UsageWriter::spawn(store.clone(), clock, writer_config(2)).unwrap();
assert_eq!(recorder.health().queue_depth, 0);
let first = recorder.try_reserve().unwrap();
assert_eq!(recorder.health().queue_depth, 1, "a held permit is depth");
assert_eq!(writer.health().queue_depth, 1);
assert_eq!(writer.health().queue_capacity, 2);
let second = recorder.try_reserve().unwrap();
let full = recorder.health();
assert_eq!(full.queue_depth, full.queue_capacity);
assert_eq!(full.shed, 0, "full is not yet shed");
assert_eq!(writer.health().queue_depth, 2);
assert_eq!(
recorder.try_reserve().err(),
Some(DenyReason::AccountingBackpressure)
);
assert_eq!(recorder.health().shed, 1);
assert_eq!(
recorder.try_reserve().err(),
Some(DenyReason::AccountingBackpressure)
);
assert_eq!(recorder.health().shed, 2);
drop(first);
drop(second);
let stats = writer.shutdown().await.unwrap();
assert_eq!(stats, tollgate_client::WriterStats::ZERO);
}
#[tokio::test(start_paused = true)]
async fn a_failing_sink_does_not_advance_the_last_ingest_time() {
let store = store(10_000);
let lease = store
.acquire(
ACCOUNT,
CostUnits(1_000),
SignedDuration::from_secs(60),
t(0),
)
.await
.unwrap()
.grant;
let failures_left = Arc::new(AtomicU32::new(u32::MAX));
let sink = flaky_sink(&store.clone(), &failures_left);
let clock = Arc::new(ManualClock::new(t(500)));
let (recorder, writer) =
UsageWriter::spawn(sink.clone(), clock.clone(), writer_config(64)).unwrap();
recorder.try_reserve().unwrap().record(event(1, 10, &lease));
settle().await;
clock.advance(SignedDuration::from_secs(1_200));
settle().await;
let during = recorder.health();
assert_eq!(
during.last_ingest_at, None,
"a sink that has never answered leaves no timestamp to age"
);
assert_eq!(during.stats.accepted, 0);
assert_eq!(
during.unaccounted, 1,
"the charge is still in the queue with no billing outcome"
);
failures_left.store(0, Ordering::Relaxed);
settle().await;
let after = recorder.health();
assert_eq!(after.stats.accepted, 1);
assert_eq!(after.last_ingest_at, Some(t(1_700)));
assert_eq!(after.unaccounted, 0);
let stats = writer.shutdown().await.unwrap();
assert_eq!(stats.accepted, 1);
assert_eq!(stats.lost, 0);
}
#[tokio::test(start_paused = true)]
async fn rejected_events_are_visible_while_running() {
let store = store(10_000);
let lease = store
.acquire(
ACCOUNT,
CostUnits(1_000),
SignedDuration::from_secs(60),
t(0),
)
.await
.unwrap()
.grant;
let clock = Arc::new(ManualClock::new(t(0)));
let (recorder, writer) = UsageWriter::spawn(store.clone(), clock, writer_config(64)).unwrap();
let mut mismatched = event(1, 10, &lease);
mismatched.source = UsageSource::Leased {
lease_id: lease.lease_id,
fencing_token: tollgate_core::FencingToken(lease.fencing_token.0 + 99),
};
recorder.try_reserve().unwrap().record(mismatched);
settle().await;
let health = recorder.health();
assert_eq!(health.stats.rejected, 1);
assert_eq!(health.stats.accepted, 0);
assert_eq!(
health.unaccounted, 0,
"refused is still an outcome; it is not left unaccounted"
);
let stats = writer.shutdown().await.unwrap();
assert_eq!(stats.rejected, 1);
}
#[tokio::test(start_paused = true)]
async fn writer_flushes_batches_idempotently() {
let store = store(10_000);
let lease = store
.acquire(
ACCOUNT,
CostUnits(1_000),
SignedDuration::from_secs(60),
t(0),
)
.await
.unwrap()
.grant;
let clock = Arc::new(ManualClock::new(t(0)));
let (recorder, writer) = UsageWriter::spawn(store.clone(), clock, writer_config(64)).unwrap();
for i in 0..10u128 {
let request = if i == 9 { 0 } else { i };
recorder
.try_reserve()
.unwrap()
.record(event(request, 10, &lease));
}
settle().await;
assert_eq!(store.usage_recorded(ACCOUNT), CostUnits(90));
let stats = writer.shutdown().await.unwrap();
assert_eq!(stats.accepted, 9);
assert_eq!(stats.duplicate, 1);
assert_eq!(stats.lost, 0);
}
#[tokio::test(start_paused = true)]
async fn full_queue_sheds_before_admission() {
let store = store(10_000);
let clock = Arc::new(ManualClock::new(t(0)));
let (recorder, writer) = UsageWriter::spawn(store.clone(), clock, writer_config(2)).unwrap();
let p1 = recorder.try_reserve().unwrap();
let p2 = recorder.try_reserve().unwrap();
assert_eq!(
recorder.try_reserve().err(),
Some(DenyReason::AccountingBackpressure)
);
drop(p1);
let _p3 = recorder.try_reserve().unwrap();
drop(p2);
writer.shutdown().await.unwrap();
}
#[derive(Default)]
struct BatchCap {
largest: AtomicUsize,
ingested: AtomicUsize,
}
fn batch_capped_sink(cap: usize, seen: &Arc<BatchCap>) -> Arc<DelegatingStore<RejectingStore>> {
let seen = Arc::clone(seen);
Arc::new(
rejecting("a capped-sink fixture answers ingest and nothing else").on_ingest(
move |_, events, _now| {
let seen = Arc::clone(&seen);
async move {
seen.largest.fetch_max(events.len(), Ordering::AcqRel);
if events.len() > cap {
return Err(StoreError("batch exceeds sink limit".into()).into());
}
seen.ingested.fetch_add(events.len(), Ordering::AcqRel);
Ok(IngestReport {
accepted: events.len() as u64,
..IngestReport::default()
})
}
},
),
)
}
#[tokio::test(start_paused = true)]
async fn steady_state_flushes_in_configured_batch_sizes() {
let seen = Arc::new(BatchCap::default());
let sink = batch_capped_sink(2, &seen);
let grant = tollgate_core::LeaseGrant {
lease_id: tollgate_core::LeaseId(1),
account_id: ACCOUNT,
fencing_token: tollgate_core::FencingToken(1),
units: CostUnits(100),
expires_at: t(100),
};
let (recorder, writer) = UsageWriter::spawn(
Arc::clone(&sink) as Arc<dyn UsageSink>,
Arc::new(ManualClock::new(t(0))),
UsageWriterConfig {
queue_capacity: 8,
max_batch: 2,
flush_interval: std::time::Duration::from_secs(60),
retry_backoff: std::time::Duration::from_millis(1),
shutdown_drain_deadline: std::time::Duration::from_secs(60),
ingest_timeout: std::time::Duration::from_secs(5),
},
)
.unwrap();
recorder.try_reserve().unwrap().record(event(0, 1, &grant));
tokio::task::yield_now().await;
for request in 1..4 {
recorder
.try_reserve()
.unwrap()
.record(event(request, 1, &grant));
}
settle().await;
assert_eq!(seen.ingested.load(Ordering::Acquire), 4);
assert_eq!(seen.largest.load(Ordering::Acquire), 2);
let stats = writer.shutdown().await.unwrap();
assert_eq!(stats.accepted, 4);
assert_eq!(stats.lost, 0);
}
#[tokio::test(start_paused = true)]
async fn shutdown_flushes_in_configured_batch_sizes() {
let seen = Arc::new(BatchCap::default());
let sink = batch_capped_sink(2, &seen);
let grant = tollgate_core::LeaseGrant {
lease_id: tollgate_core::LeaseId(1),
account_id: ACCOUNT,
fencing_token: tollgate_core::FencingToken(1),
units: CostUnits(100),
expires_at: t(100),
};
let (recorder, writer) = UsageWriter::spawn(
Arc::clone(&sink) as Arc<dyn UsageSink>,
Arc::new(ManualClock::new(t(0))),
UsageWriterConfig {
queue_capacity: 8,
max_batch: 2,
flush_interval: std::time::Duration::from_secs(60),
retry_backoff: std::time::Duration::from_millis(1),
shutdown_drain_deadline: std::time::Duration::from_secs(60),
ingest_timeout: std::time::Duration::from_secs(5),
},
)
.unwrap();
for request in 0..6 {
recorder
.try_reserve()
.unwrap()
.record(event(request, 1, &grant));
}
let stats = writer.shutdown().await.unwrap();
assert_eq!(stats.accepted, 6);
assert_eq!(stats.lost, 0);
assert_eq!(seen.largest.load(Ordering::Acquire), 2);
}
fn flaky_sink(
store: &Arc<MemoryStore>,
failures_left: &Arc<AtomicU32>,
) -> Arc<DelegatingStore<MemoryStore>> {
let failures_left = Arc::clone(failures_left);
Arc::new(
DelegatingStore::wrapping(Arc::clone(store)).on_ingest(move |inner, events, now| {
let failures_left = Arc::clone(&failures_left);
async move {
if failures_left
.fetch_update(Ordering::AcqRel, Ordering::Acquire, |n| n.checked_sub(1))
.is_ok()
{
return Err(StoreError("injected outage".into()).into());
}
inner.ingest(&events, now).await
}
}),
)
}
#[tokio::test(start_paused = true)]
async fn writer_retries_through_outage_without_losing_events() {
let store = store(10_000);
let lease = store
.acquire(
ACCOUNT,
CostUnits(1_000),
SignedDuration::from_secs(60),
t(0),
)
.await
.unwrap()
.grant;
let clock = Arc::new(ManualClock::new(t(0)));
let failures_left = Arc::new(AtomicU32::new(3));
let sink = flaky_sink(&store.clone(), &failures_left);
let (recorder, writer) = UsageWriter::spawn(sink, clock, writer_config(64)).unwrap();
recorder.try_reserve().unwrap().record(event(1, 25, &lease));
settle().await;
settle().await;
assert_eq!(store.usage_recorded(ACCOUNT), CostUnits(25));
let stats = writer.shutdown().await.unwrap();
assert_eq!(stats.accepted, 1);
assert_eq!(stats.lost, 0);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn panic_after_commit_still_bills() {
let store = store(10_000);
let grant = store
.acquire(
ACCOUNT,
CostUnits(1_000),
SignedDuration::from_secs(60),
t(0),
)
.await
.unwrap()
.grant;
let clock = Arc::new(ManualClock::new(t(0)));
let (recorder, writer) = UsageWriter::spawn(store.clone(), clock, writer_config(8)).unwrap();
let permit = recorder.try_reserve().unwrap();
let worker = tokio::spawn(async move {
let charge = ready_from_grant(grant, permit)
.commit(RequestId(9), t(0))
.expect("a live lease commits");
assert_eq!(charge.units(), CostUnits(51));
panic!("kernel exploded mid-execution");
});
assert!(worker.await.is_err(), "the worker must have panicked");
let stats = writer.shutdown().await.unwrap();
assert_eq!(stats.accepted, 1);
assert_eq!(stats.lost, 0);
assert_eq!(store.usage_recorded(ACCOUNT), CostUnits(51));
}
#[tokio::test]
async fn a_panicking_kernel_under_catch_unwind_still_bills() {
let store = store(10_000);
let grant = store
.acquire(
ACCOUNT,
CostUnits(1_000),
SignedDuration::from_secs(60),
t(0),
)
.await
.unwrap()
.grant;
let clock = Arc::new(ManualClock::new(t(0)));
let (recorder, writer) = UsageWriter::spawn(store.clone(), clock, writer_config(8)).unwrap();
let permit = recorder.try_reserve().unwrap();
let ready = ready_from_grant(grant, permit);
let worker = std::thread::spawn(move || {
std::panic::catch_unwind(std::panic::AssertUnwindSafe(move || {
let charge = ready
.commit(RequestId(9), t(0))
.expect("a live lease commits");
assert_eq!(charge.units(), CostUnits(51));
panic!("kernel exploded mid-execution");
}))
});
let outcome = worker.join().expect("the boundary contains the panic");
assert!(outcome.is_err(), "the kernel panicked");
let stats = writer.shutdown().await.unwrap();
assert_eq!(stats.accepted, 1, "the committed charge was still emitted");
assert_eq!(stats.lost, 0);
assert_eq!(store.usage_recorded(ACCOUNT), CostUnits(51));
}
#[tokio::test(start_paused = true)]
async fn shutdown_during_outage_terminates_and_reports_loss() {
let store = store(10_000);
let lease = store
.acquire(
ACCOUNT,
CostUnits(1_000),
SignedDuration::from_secs(60),
t(0),
)
.await
.unwrap()
.grant;
let clock = Arc::new(ManualClock::new(t(0)));
let failures_left = Arc::new(AtomicU32::new(u32::MAX));
let sink = flaky_sink(&store.clone(), &failures_left);
let (recorder, writer) = UsageWriter::spawn(sink, clock, writer_config(64)).unwrap();
recorder.try_reserve().unwrap().record(event(1, 25, &lease));
recorder.try_reserve().unwrap().record(event(2, 25, &lease));
settle().await;
let stats = tokio::time::timeout(std::time::Duration::from_secs(60), writer.shutdown())
.await
.expect("shutdown must terminate during an outage")
.unwrap();
assert_eq!(stats.lost, 2);
assert_eq!(stats.accepted, 0);
assert_eq!(store.usage_recorded(ACCOUNT), CostUnits::ZERO);
}
#[tokio::test(start_paused = true)]
async fn shutdown_after_recovery_delivers_everything() {
let store = store(10_000);
let lease = store
.acquire(
ACCOUNT,
CostUnits(1_000),
SignedDuration::from_secs(60),
t(0),
)
.await
.unwrap()
.grant;
let clock = Arc::new(ManualClock::new(t(0)));
let failures_left = Arc::new(AtomicU32::new(4));
let sink = flaky_sink(&store.clone(), &failures_left);
let (recorder, writer) = UsageWriter::spawn(sink, clock, writer_config(64)).unwrap();
recorder.try_reserve().unwrap().record(event(1, 25, &lease));
settle().await;
let stats = tokio::time::timeout(std::time::Duration::from_secs(60), writer.shutdown())
.await
.expect("shutdown must terminate")
.unwrap();
assert_eq!(stats.lost, 0);
assert_eq!(stats.accepted, 1);
assert_eq!(store.usage_recorded(ACCOUNT), CostUnits(25));
}
#[tokio::test(start_paused = true)]
async fn a_survivor_flushes_after_reclaim_and_is_billed() {
let store = store(10_000);
let lease = store
.acquire(
ACCOUNT,
CostUnits(1_000),
SignedDuration::from_secs(60),
t(0),
)
.await
.unwrap()
.grant;
store.reclaim_expired(t(120)).await.unwrap();
let clock = Arc::new(ManualClock::new(t(121)));
let (recorder, writer) = UsageWriter::spawn(store.clone(), clock, writer_config(8)).unwrap();
recorder.try_reserve().unwrap().record(event(1, 10, &lease));
recorder
.try_reserve()
.unwrap()
.record(event(2, 991, &lease));
settle().await;
let stats = writer.shutdown().await.unwrap();
assert_eq!(stats.accepted, 1);
assert_eq!(stats.rejected, 1, "991 exceeds the 990 still forfeited");
assert_eq!(store.usage_recorded(ACCOUNT), CostUnits(10));
let c = store.conservation(ACCOUNT).unwrap();
assert_eq!(c.settlement_loss, CostUnits(990));
assert!(c.holds(), "{c:?}");
}
#[tokio::test(start_paused = true)]
async fn reserve_fails_once_shutdown_begins() {
let store = store(10_000);
let lease = store
.acquire(
ACCOUNT,
CostUnits(1_000),
SignedDuration::from_secs(60),
t(0),
)
.await
.unwrap()
.grant;
let clock = Arc::new(ManualClock::new(t(0)));
let (recorder, writer) = UsageWriter::spawn(store.clone(), clock, writer_config(8)).unwrap();
let held = recorder.try_reserve().unwrap();
let shutdown = tokio::spawn(writer.shutdown());
settle().await;
assert_eq!(
recorder.try_reserve().err(),
Some(DenyReason::AccountingBackpressure)
);
assert!(recorder.is_closed(), "readiness must observe the shutdown");
held.record(event(1, 25, &lease));
let stats = shutdown.await.unwrap().unwrap();
assert_eq!(stats.accepted, 1);
assert_eq!(stats.unresolved, 0);
assert_eq!(store.usage_recorded(ACCOUNT), CostUnits(25));
}
#[tokio::test(start_paused = true)]
async fn shutdown_waits_for_outstanding_permit() {
let store = store(10_000);
let lease = store
.acquire(
ACCOUNT,
CostUnits(1_000),
SignedDuration::from_secs(60),
t(0),
)
.await
.unwrap()
.grant;
let clock = Arc::new(ManualClock::new(t(0)));
let (recorder, writer) = UsageWriter::spawn(store.clone(), clock, writer_config(8)).unwrap();
let permit = recorder.try_reserve().unwrap();
let late = event(1, 25, &lease);
let sender = tokio::spawn(async move {
tokio::time::sleep(std::time::Duration::from_millis(20)).await;
permit.record(late);
});
let stats = writer.shutdown().await.unwrap();
sender.await.unwrap();
assert_eq!(stats.accepted, 1);
assert_eq!(stats.lost, 0);
assert_eq!(stats.unresolved, 0);
assert_eq!(store.usage_recorded(ACCOUNT), CostUnits(25));
}
#[tokio::test(start_paused = true)]
async fn shutdown_waits_for_committed_guard() {
let store = store(10_000);
let grant = store
.acquire(
ACCOUNT,
CostUnits(1_000),
SignedDuration::from_secs(60),
t(0),
)
.await
.unwrap()
.grant;
let clock = Arc::new(ManualClock::new(t(0)));
let (recorder, writer) = UsageWriter::spawn(store.clone(), clock, writer_config(8)).unwrap();
let permit = recorder.try_reserve().unwrap();
let holder = tokio::spawn(async move {
let charge = ready_from_grant(grant, permit)
.commit(RequestId(9), t(0))
.expect("a live lease commits");
assert_eq!(charge.units(), CostUnits(51));
tokio::time::sleep(std::time::Duration::from_millis(20)).await;
drop(charge);
});
let stats = writer.shutdown().await.unwrap();
holder.await.unwrap();
assert_eq!(stats.accepted, 1);
assert_eq!(stats.unresolved, 0);
assert_eq!(store.usage_recorded(ACCOUNT), CostUnits(51));
}
#[tokio::test(start_paused = true)]
async fn late_permit_drop_completes_drain() {
let store = store(10_000);
let clock = Arc::new(ManualClock::new(t(0)));
let (recorder, writer) = UsageWriter::spawn(store.clone(), clock, writer_config(8)).unwrap();
let permit = recorder.try_reserve().unwrap();
let dropper = tokio::spawn(async move {
tokio::time::sleep(std::time::Duration::from_millis(20)).await;
drop(permit);
});
let stats = writer.shutdown().await.unwrap();
dropper.await.unwrap();
assert_eq!(stats, tollgate_client::WriterStats::ZERO);
}
#[tokio::test(start_paused = true)]
async fn drain_deadline_expiry_reports_unresolved() {
let store = store(10_000);
let clock = Arc::new(ManualClock::new(t(0)));
let (recorder, writer) = UsageWriter::spawn(store.clone(), clock, writer_config(8)).unwrap();
std::mem::forget(recorder.try_reserve().unwrap());
let stats = tokio::time::timeout(std::time::Duration::from_secs(120), writer.shutdown())
.await
.expect("shutdown must return at the drain deadline, not hang")
.unwrap();
assert_eq!(stats.unresolved, 1);
assert_eq!(stats.lost, 0);
assert_eq!(stats.accepted, 0);
}
#[tokio::test(start_paused = true)]
async fn final_flush_counts_duplicates_and_rejections() {
let store = store(10_000);
let lease = store
.acquire(
ACCOUNT,
CostUnits(1_000),
SignedDuration::from_secs(60),
t(0),
)
.await
.unwrap()
.grant;
let clock = Arc::new(ManualClock::new(t(0)));
let mut config = writer_config(16);
config.flush_interval = std::time::Duration::from_secs(3_600);
config.max_batch = 16;
let (recorder, writer) = UsageWriter::spawn(store.clone(), clock, config).unwrap();
recorder.try_reserve().unwrap().record(event(1, 10, &lease));
recorder.try_reserve().unwrap().record(event(1, 10, &lease));
let mut mismatched = event(2, 10, &lease);
mismatched.source = UsageSource::Leased {
lease_id: lease.lease_id,
fencing_token: tollgate_core::FencingToken(999),
};
recorder.try_reserve().unwrap().record(mismatched);
let stats = writer.shutdown().await.unwrap();
assert_eq!(stats.accepted, 1);
assert_eq!(stats.duplicate, 1);
assert_eq!(stats.rejected, 1);
assert_eq!(stats.lost, 0);
assert_eq!(stats.unresolved, 0);
}
#[tokio::test(start_paused = true)]
async fn final_flush_backs_off_only_between_attempts() {
let store = store(10_000);
let lease = store
.acquire(
ACCOUNT,
CostUnits(1_000),
SignedDuration::from_secs(60),
t(0),
)
.await
.unwrap()
.grant;
let clock = Arc::new(ManualClock::new(t(0)));
let failures_left = Arc::new(AtomicU32::new(u32::MAX));
let sink = flaky_sink(&store.clone(), &failures_left);
let backoff = std::time::Duration::from_millis(100);
let mut config = writer_config(16);
config.flush_interval = std::time::Duration::from_secs(3_600);
config.retry_backoff = backoff;
let (recorder, writer) = UsageWriter::spawn(sink, clock, config).unwrap();
recorder.try_reserve().unwrap().record(event(1, 10, &lease));
let start = tokio::time::Instant::now();
let stats = writer.shutdown().await.unwrap();
assert_eq!(stats.lost, 1);
assert_eq!(start.elapsed(), 2 * backoff, "three attempts, two backoffs");
}
#[tokio::test(start_paused = true)]
async fn drain_deadline_reports_every_unresolved_permit() {
let store = store(10_000);
let clock = Arc::new(ManualClock::new(t(0)));
let (recorder, writer) = UsageWriter::spawn(store.clone(), clock, writer_config(8)).unwrap();
std::mem::forget(recorder.try_reserve().unwrap());
std::mem::forget(recorder.try_reserve().unwrap());
let stats = tokio::time::timeout(std::time::Duration::from_secs(120), writer.shutdown())
.await
.expect("shutdown must return at the drain deadline")
.unwrap();
assert_eq!(stats.unresolved, 2);
}
fn panicking_sink(
store: &Arc<MemoryStore>,
calls_before_panic: u32,
) -> Arc<DelegatingStore<MemoryStore>> {
let budget = Arc::new(AtomicU32::new(calls_before_panic));
Arc::new(
DelegatingStore::wrapping(Arc::clone(store)).on_ingest(move |inner, events, now| {
let budget = Arc::clone(&budget);
async move {
if budget
.fetch_update(Ordering::AcqRel, Ordering::Acquire, |n| n.checked_sub(1))
.is_err()
{
panic!("sink exploded mid-ingest");
}
inner.ingest(&events, now).await
}
}),
)
}
#[tokio::test(start_paused = true)]
async fn panicked_writer_reports_unaccounted_charges() {
let store = store(10_000);
let lease = store
.acquire(
ACCOUNT,
CostUnits(1_000),
SignedDuration::from_secs(60),
t(0),
)
.await
.unwrap()
.grant;
let clock = Arc::new(ManualClock::new(t(0)));
let sink = panicking_sink(&store.clone(), 0);
let mut config = writer_config(16);
config.flush_interval = std::time::Duration::from_secs(3_600);
let (recorder, writer) = UsageWriter::spawn(sink, clock, config).unwrap();
for request in 0..3u128 {
recorder
.try_reserve()
.unwrap()
.record(event(request, 10, &lease));
}
let error = writer
.shutdown()
.await
.expect_err("a panicked writer must not report a clean shutdown");
assert!(error.panicked);
assert_eq!(error.unaccounted, 3);
assert!(
error.to_string().contains("no billing record"),
"operator-facing message must name the consequence: {error}"
);
assert_eq!(store.usage_recorded(ACCOUNT), CostUnits::ZERO);
}
#[tokio::test(start_paused = true)]
async fn panic_after_partial_flush_counts_only_unflushed() {
let store = store(10_000);
let lease = store
.acquire(
ACCOUNT,
CostUnits(1_000),
SignedDuration::from_secs(60),
t(0),
)
.await
.unwrap()
.grant;
let clock = Arc::new(ManualClock::new(t(0)));
let sink = panicking_sink(&store.clone(), 1);
let mut config = writer_config(16);
config.max_batch = 2;
config.flush_interval = std::time::Duration::from_secs(3_600);
let (recorder, writer) = UsageWriter::spawn(sink, clock, config).unwrap();
for request in 0..2u128 {
recorder
.try_reserve()
.unwrap()
.record(event(request, 10, &lease));
}
settle().await; assert_eq!(store.usage_recorded(ACCOUNT), CostUnits(20));
for request in 2..5u128 {
recorder
.try_reserve()
.unwrap()
.record(event(request, 10, &lease));
}
let error = writer.shutdown().await.expect_err("the sink panicked");
assert_eq!(
error.unaccounted, 3,
"the two billed charges must not be counted again"
);
}
fn hanging_sink() -> Arc<DelegatingStore<RejectingStore>> {
Arc::new(
rejecting("a hanging-sink fixture answers ingest and nothing else")
.on_ingest(|_, _events, _now| async { std::future::pending().await }),
)
}
#[tokio::test(start_paused = true)]
async fn hung_ingest_cannot_stall_shutdown() {
let store = store(10_000);
let lease = store
.acquire(
ACCOUNT,
CostUnits(1_000),
SignedDuration::from_secs(60),
t(0),
)
.await
.unwrap()
.grant;
let clock = Arc::new(ManualClock::new(t(0)));
let mut config = writer_config(16);
config.flush_interval = std::time::Duration::from_secs(3_600);
let (recorder, writer) =
UsageWriter::spawn(hanging_sink() as Arc<dyn UsageSink>, clock, config).unwrap();
recorder.try_reserve().unwrap().record(event(1, 25, &lease));
let stats = tokio::time::timeout(std::time::Duration::from_secs(300), writer.shutdown())
.await
.expect("a hung sink must not stall shutdown")
.unwrap();
assert_eq!(stats.lost, 1);
assert_eq!(stats.accepted, 0);
assert_eq!(store.usage_recorded(ACCOUNT), CostUnits::ZERO);
}
#[tokio::test(start_paused = true)]
async fn hung_ingest_times_out_into_the_retry_path() {
let store = store(10_000);
let lease = store
.acquire(
ACCOUNT,
CostUnits(1_000),
SignedDuration::from_secs(60),
t(0),
)
.await
.unwrap()
.grant;
let clock = Arc::new(ManualClock::new(t(0)));
let mut config = writer_config(16);
config.ingest_timeout = std::time::Duration::from_millis(100);
let (recorder, writer) =
UsageWriter::spawn(hanging_sink() as Arc<dyn UsageSink>, clock, config).unwrap();
recorder.try_reserve().unwrap().record(event(1, 25, &lease));
settle().await;
let stats = tokio::time::timeout(std::time::Duration::from_secs(300), writer.shutdown())
.await
.expect("the retry loop must remain interruptible")
.unwrap();
assert_eq!(stats.accepted, 0);
assert_eq!(stats.lost, 1);
}
fn slow_sink(
store: &Arc<MemoryStore>,
delay: std::time::Duration,
) -> Arc<DelegatingStore<MemoryStore>> {
Arc::new(DelegatingStore::wrapping(Arc::clone(store)).on_ingest(
move |inner, events, now| async move {
tokio::time::sleep(delay).await;
inner.ingest(&events, now).await
},
))
}
#[tokio::test(start_paused = true)]
async fn slow_but_healthy_sink_still_delivers_at_shutdown() {
let store = store(10_000);
let lease = store
.acquire(
ACCOUNT,
CostUnits(1_000),
SignedDuration::from_secs(60),
t(0),
)
.await
.unwrap()
.grant;
let clock = Arc::new(ManualClock::new(t(0)));
let sink = slow_sink(&store.clone(), std::time::Duration::from_millis(50));
let mut config = writer_config(16);
config.flush_interval = std::time::Duration::from_secs(3_600);
config.ingest_timeout = std::time::Duration::from_secs(5);
let (recorder, writer) = UsageWriter::spawn(sink as Arc<dyn UsageSink>, clock, config).unwrap();
recorder.try_reserve().unwrap().record(event(1, 25, &lease));
let stats = writer.shutdown().await.unwrap();
assert_eq!(stats.accepted, 1, "a slow call inside its budget delivers");
assert_eq!(stats.lost, 0);
assert_eq!(store.usage_recorded(ACCOUNT), CostUnits(25));
}
fn hanging_acquire() -> Arc<DelegatingStore<RejectingStore>> {
Arc::new(
rejecting("this fixture never consolidates")
.on_acquire(|_, _account, _requested, _ttl, _now| async {
std::future::pending().await
})
.on_release(|_, _lease, _fence, _unspent, _now| async { Ok(()) })
.on_reclaim_expired_batch(|_, _now, limit| async move {
ReclaimBatch::try_new(Vec::new(), limit)
}),
)
}
#[tokio::test(start_paused = true)]
async fn shutdown_during_a_hung_acquire_is_not_delayed_by_it() {
let slot = LeaseSlot::for_account(ACCOUNT);
let clock = Arc::new(ManualClock::new(t(0)));
let config = LeaseManagerConfig {
store_call_timeout: std::time::Duration::from_secs(120),
shutdown_release_deadline: std::time::Duration::from_secs(10),
..manager_config()
};
let manager = LeaseManager::spawn(
hanging_acquire() as Arc<dyn LeaseAllocator>,
Arc::clone(&slot),
clock,
config,
)
.unwrap();
settle().await;
assert!(slot.load().is_none(), "the acquire is still unanswered");
let bound = config.shutdown_release_deadline + std::time::Duration::from_secs(1);
let began = tokio::time::Instant::now();
let report = tokio::time::timeout(bound, manager.shutdown())
.await
.expect("a hung acquire must not hold the shutdown signal");
let took = began.elapsed();
assert!(!report.task_died);
assert_eq!(report.released, 0, "there was never a lease to return");
assert_eq!(report.abandoned, 0);
assert!(
took < config.store_call_timeout,
"shutdown took {took:?}, which is the hung acquire's timeout, not the signal"
);
}
#[tokio::test(start_paused = true)]
async fn a_hung_acquire_is_counted_as_a_timeout_not_a_refusal() {
let slot = LeaseSlot::for_account(ACCOUNT);
let clock = Arc::new(ManualClock::new(t(0)));
let manager = LeaseManager::spawn(
hanging_acquire() as Arc<dyn LeaseAllocator>,
Arc::clone(&slot),
clock,
LeaseManagerConfig {
store_call_timeout: std::time::Duration::from_millis(10),
..manager_config()
},
)
.unwrap();
let counters = manager.counters();
settle().await;
let stats = counters.snapshot();
assert!(
stats.acquire_timeouts > 0,
"the refill loop must not park on a hung allocator"
);
assert_eq!(
stats.refused(),
0,
"a timeout is not a refusal: {:?}",
stats.acquire_refused
);
assert_eq!(stats.acquired, 0);
assert!(slot.load().is_none(), "fail closed while unanswered");
manager.shutdown().await;
}
fn hanging_release(store: &Arc<MemoryStore>) -> Arc<DelegatingStore<MemoryStore>> {
Arc::new(
DelegatingStore::wrapping(Arc::clone(store))
.on_release(|_, _lease, _fence, _unspent, _now| async { std::future::pending().await })
.on_consolidate(|_, _, _, _, _, _, _, _| async {
unreachable!("this fixture never consolidates")
}),
)
}
#[tokio::test(start_paused = true)]
async fn a_refill_does_not_wait_behind_the_release_pass() {
let store = store(100_000);
let slot = LeaseSlot::for_account(ACCOUNT);
let clock = Arc::new(ManualClock::new(t(0)));
let allocator = hanging_release(&store.clone());
let manager = LeaseManager::spawn(
allocator as Arc<dyn LeaseAllocator>,
Arc::clone(&slot),
Arc::clone(&clock) as Arc<dyn tollgate_client::Clock>,
manager_config(),
)
.unwrap();
settle().await;
for tick in 1..=5 {
clock.set(t(tick * 120));
let installed = tokio::time::timeout(std::time::Duration::from_secs(600), async {
loop {
match slot.load() {
Some(lease) if lease.grant().expires_at > clock.now() => break,
_ => tokio::time::sleep(std::time::Duration::from_millis(50)).await,
}
}
})
.await;
installed.expect("the manager must install a lease after each rotation");
}
assert_eq!(manager.counters().snapshot().acquired, 6);
assert_eq!(manager.counters().snapshot().released, 0);
let before = slot.load().expect("a lease is installed");
assert!(
!before.needs_refill(),
"the fixture must start from a lease that has not already crossed"
);
let began = tokio::time::Instant::now();
before.try_debit(CostUnits(800), clock.now()).unwrap();
assert!(before.needs_refill());
tokio::time::timeout(std::time::Duration::from_secs(600), async {
loop {
match slot.load() {
Some(fresh) if fresh.grant().lease_id != before.grant().lease_id => break,
_ => tokio::time::sleep(std::time::Duration::from_millis(10)).await,
}
}
})
.await
.expect("the crossing debit must be followed by a fresh lease");
let waited = began.elapsed();
let ceiling = manager_config().store_call_timeout * 3;
assert!(
waited < ceiling,
"refill waited {waited:?} behind the release pass, over the {ceiling:?} ceiling"
);
manager.shutdown().await;
}
#[tokio::test(start_paused = true)]
async fn shutdown_reports_released_leases() {
let store = store(10_000);
let slot = LeaseSlot::for_account(ACCOUNT);
let clock = Arc::new(ManualClock::new(t(0)));
let manager =
LeaseManager::spawn(store.clone(), Arc::clone(&slot), clock, manager_config()).unwrap();
settle().await;
assert!(slot.load().is_some());
let report = manager.shutdown().await;
assert_eq!(report.released, 1);
assert_eq!(report.abandoned, 0);
assert!(!report.task_died);
assert_eq!(store.balance(ACCOUNT), CostUnits(10_000), "units came back");
}
#[tokio::test(start_paused = true)]
async fn hung_release_cannot_stall_shutdown() {
let store = store(10_000);
let slot = LeaseSlot::for_account(ACCOUNT);
let clock = Arc::new(ManualClock::new(t(0)));
let allocator = hanging_release(&store.clone());
let manager = LeaseManager::spawn(
allocator as Arc<dyn LeaseAllocator>,
Arc::clone(&slot),
Arc::clone(&clock) as Arc<dyn tollgate_client::Clock>,
manager_config(),
)
.unwrap();
let counters = manager.counters();
settle().await;
assert!(slot.load().is_some());
assert_eq!(counters.snapshot().abandoned, 0, "nothing abandoned yet");
let mut readers = Vec::new();
for tick in 1..=5 {
readers.push(
slot.load()
.expect("the preceding rotation installed a grant"),
);
clock.set(t(tick * 120));
settle().await;
}
drop(readers);
let parked = counters.snapshot();
assert_eq!(parked.released, 0, "a hung release settles nothing");
let config = manager_config();
let bound = config.shutdown_release_deadline + std::time::Duration::from_secs(1);
let began = tokio::time::Instant::now();
let report = tokio::time::timeout(bound, manager.shutdown())
.await
.expect("a hung release must not stall shutdown past its stated budget");
let took = began.elapsed();
assert!(
report.abandoned >= 2,
"the test must exercise a multi-lease parked list, got {}",
report.abandoned
);
assert_eq!(report.released, 0);
assert!(!report.task_died);
assert!(
took <= bound,
"shutdown took {took:?} against a stated {:?} budget",
config.shutdown_release_deadline
);
let stats = counters.snapshot();
assert_eq!(stats.abandoned, report.abandoned);
assert_eq!(stats.released, 0);
}
#[tokio::test]
async fn invalid_writer_config_is_rejected() {
let zero = std::time::Duration::ZERO;
let mut zero_capacity = writer_config(8);
zero_capacity.queue_capacity = 0;
let mut zero_batch = writer_config(8);
zero_batch.max_batch = 0;
let mut zero_flush = writer_config(8);
zero_flush.flush_interval = zero;
let mut zero_backoff = writer_config(8);
zero_backoff.retry_backoff = zero;
let mut zero_drain = writer_config(8);
zero_drain.shutdown_drain_deadline = zero;
let mut zero_ingest = writer_config(8);
zero_ingest.ingest_timeout = zero;
for (field, config) in [
("queue_capacity", zero_capacity),
("max_batch", zero_batch),
("flush_interval", zero_flush),
("retry_backoff", zero_backoff),
("shutdown_drain_deadline", zero_drain),
("ingest_timeout", zero_ingest),
] {
assert!(
config.validate().is_err(),
"zero {field} must be rejected, not coerced"
);
let store = store(10_000);
let clock = Arc::new(ManualClock::new(t(0)));
assert!(
UsageWriter::spawn(store, clock, config).is_err(),
"zero {field} must not start a task"
);
}
}
#[test]
fn invalid_lease_manager_timeouts_are_rejected() {
let mut config = manager_config();
config.store_call_timeout = std::time::Duration::ZERO;
assert!(config.validate().is_err());
let mut config = manager_config();
config.shutdown_release_deadline = std::time::Duration::ZERO;
assert!(config.validate().is_err());
}
#[tokio::test(start_paused = true)]
async fn shutdown_abandons_a_lease_an_in_flight_request_still_holds() {
let store = store(10_000);
let slot = LeaseSlot::for_account(ACCOUNT);
let clock = Arc::new(ManualClock::new(t(0)));
let manager = LeaseManager::spawn(
store.clone(),
Arc::clone(&slot),
Arc::clone(&clock) as Arc<dyn Clock>,
manager_config(),
)
.unwrap();
settle().await;
let in_flight = slot.load().expect("a lease is installed");
let balance_before = store.balance(ACCOUNT);
let report = manager.shutdown().await;
assert_eq!(
report.released, 0,
"a lease still reachable by a request must not be released"
);
assert_eq!(
report.abandoned, 1,
"it is abandoned instead, and reported — never silently dropped"
);
assert_eq!(
store.balance(ACCOUNT),
balance_before,
"no units were credited back while a request could still spend them; \
crediting early is what lets committed usage exceed the allocation"
);
assert!(in_flight.try_debit(CostUnits(50), t(0)).is_ok());
}
#[tokio::test(start_paused = true)]
async fn the_final_flush_backoff_cannot_overrun_the_drain_deadline() {
let store = store(10_000);
let lease = store
.acquire(
ACCOUNT,
CostUnits(1_000),
SignedDuration::from_secs(60),
t(0),
)
.await
.unwrap()
.grant;
let clock = Arc::new(ManualClock::new(t(0)));
let failures_left = Arc::new(AtomicU32::new(u32::MAX));
let sink = flaky_sink(&store.clone(), &failures_left);
let mut config = writer_config(16);
config.flush_interval = std::time::Duration::from_secs(3_600);
config.shutdown_drain_deadline = std::time::Duration::from_secs(2);
config.retry_backoff = std::time::Duration::from_secs(5);
config.ingest_timeout = std::time::Duration::from_millis(50);
let (recorder, writer) = UsageWriter::spawn(sink, clock, config).unwrap();
recorder.try_reserve().unwrap().record(event(1, 10, &lease));
let start = tokio::time::Instant::now();
let stats = writer.shutdown().await.unwrap();
let took = start.elapsed();
assert_eq!(
stats.lost, 1,
"an undeliverable event is counted, not dropped"
);
assert!(
took <= config.shutdown_drain_deadline + config.ingest_timeout,
"the drain took {took:?} against a {:?} budget; a backoff that sleeps \
past the deadline spends the margin that keeps a straggler billable",
config.shutdown_drain_deadline
);
}
fn refused_event(request: u128) -> UsageEvent {
UsageEvent::new(
RequestId(request),
ACCOUNT,
UsageSource::Overage,
CostUnits(50),
t(0),
PolicyRevision::UNSTATED,
None,
)
}
fn always_refuses_sink(attempts: &Arc<AtomicUsize>) -> Arc<DelegatingStore<RejectingStore>> {
let attempts = Arc::clone(attempts);
Arc::new(
rejecting("a refusing-sink fixture answers ingest and nothing else").on_ingest(
move |_, _events, _now| {
let attempts = Arc::clone(&attempts);
async move {
attempts.fetch_add(1, Ordering::AcqRel);
Err(IngestError::Refused(StoreError(
"413 batch-too-large: request body exceeds this endpoint's limit".into(),
)))
}
},
),
)
}
#[test]
fn a_max_batch_beyond_the_ingest_limit_is_rejected() {
let mut oversized = writer_config(8);
oversized.max_batch = tollgate_store::MAX_INGEST_BATCH + 1;
assert_eq!(
oversized.validate().unwrap_err().to_string(),
"max_batch exceeds the ingest endpoint's documented limit"
);
let mut at_limit = writer_config(8);
at_limit.max_batch = tollgate_store::MAX_INGEST_BATCH;
assert!(
at_limit.validate().is_ok(),
"the documented limit itself must be usable, or it is not the limit"
);
}
#[tokio::test(start_paused = true)]
async fn a_refused_batch_is_counted_lost_rather_than_retried_forever() {
let attempts = Arc::new(AtomicUsize::new(0));
let sink = always_refuses_sink(&attempts);
let clock = Arc::new(ManualClock::new(t(0)));
let (recorder, writer) = UsageWriter::spawn(sink, clock, writer_config(8)).unwrap();
recorder.try_reserve().unwrap().record(refused_event(1));
settle().await;
assert_eq!(
attempts.load(Ordering::Acquire),
1,
"a refused batch must be attempted once, not retried"
);
assert_eq!(
recorder.health().unaccounted,
0,
"a reported loss is a completed accounting outcome"
);
recorder
.try_reserve()
.expect("the queue drained rather than filling with a wedged batch")
.record(refused_event(2));
settle().await;
let stats = writer.shutdown().await.unwrap();
assert!(
stats.lost >= 2,
"refused events are counted lost, never silently dropped; got {stats:?}"
);
}
fn observed_hung_sink(active: &Arc<AtomicUsize>) -> Arc<DelegatingStore<RejectingStore>> {
let active = Arc::clone(active);
Arc::new(
rejecting("a hung-sink fixture answers ingest and nothing else").on_ingest(
move |_, _events, _now| {
let active = Arc::clone(&active);
async move {
struct Active(Arc<AtomicUsize>);
impl Drop for Active {
fn drop(&mut self) {
self.0.fetch_sub(1, Ordering::AcqRel);
}
}
active.fetch_add(1, Ordering::AcqRel);
let _active = Active(Arc::clone(&active));
std::future::pending().await
}
},
),
)
}
#[tokio::test(start_paused = true)]
async fn cancelling_writer_shutdown_aborts_the_owned_ingest_task() {
let active = Arc::new(AtomicUsize::new(0));
let mut config = writer_config(8);
config.max_batch = 1;
config.ingest_timeout = std::time::Duration::from_secs(120);
config.shutdown_drain_deadline = std::time::Duration::from_secs(60);
let (recorder, writer) = UsageWriter::spawn(
observed_hung_sink(&active),
Arc::new(ManualClock::new(t(0))),
config,
)
.unwrap();
recorder.try_reserve().unwrap().record(refused_event(1));
settle().await;
assert_eq!(active.load(Ordering::Acquire), 1);
assert!(
tokio::time::timeout(std::time::Duration::from_millis(1), writer.shutdown())
.await
.is_err()
);
for _ in 0..10 {
tokio::task::yield_now().await;
}
assert_eq!(
active.load(Ordering::Acquire),
0,
"cancelled shutdown must not detach its writer"
);
}
#[tokio::test(start_paused = true)]
async fn shutdown_interrupts_normal_ingest_before_its_long_timeout() {
let active = Arc::new(AtomicUsize::new(0));
let mut config = writer_config(8);
config.max_batch = 1;
config.ingest_timeout = std::time::Duration::from_secs(120);
config.shutdown_drain_deadline = std::time::Duration::from_millis(10);
let (recorder, writer) = UsageWriter::spawn(
observed_hung_sink(&active),
Arc::new(ManualClock::new(t(0))),
config,
)
.unwrap();
recorder.try_reserve().unwrap().record(refused_event(1));
settle().await;
assert_eq!(active.load(Ordering::Acquire), 1);
let result = tokio::time::timeout(std::time::Duration::from_millis(11), writer.shutdown())
.await
.unwrap()
.unwrap();
assert_eq!(result.lost, 1);
assert_eq!(active.load(Ordering::Acquire), 0);
}
#[tokio::test(start_paused = true)]
async fn cancelling_lease_shutdown_aborts_the_owned_release_task() {
let store = store(10_000);
let slot = LeaseSlot::for_account(ACCOUNT);
let manager = LeaseManager::spawn(
hanging_release(&store),
slot.clone(),
Arc::new(ManualClock::new(t(0))),
manager_config(),
)
.unwrap();
let health = manager.health();
settle().await;
assert!(slot.load().is_some());
assert!(
tokio::time::timeout(std::time::Duration::from_millis(1), manager.shutdown())
.await
.is_err()
);
for _ in 0..10 {
tokio::task::yield_now().await;
}
assert!(
health.has_changed().is_err(),
"cancelled shutdown must not detach its manager"
);
assert!(!*health.borrow(), "cancelled refill task must retain false");
}
#[tokio::test(start_paused = true)]
async fn a_refused_lease_consolidates_rather_than_stranding_the_tail() {
let store = store(160);
let slot = LeaseSlot::for_account(ACCOUNT);
let clock = Arc::new(ManualClock::new(t(0)));
let manager = LeaseManager::spawn(
store.clone(),
Arc::clone(&slot),
clock,
LeaseManagerConfig {
target_grant: CostUnits(100),
low_water: CostUnits(50),
poll_interval: std::time::Duration::from_secs(60),
..manager_config()
},
)
.unwrap();
settle().await;
for spent in 1..=2 {
let lease = slot.load().expect("a lease is installed");
lease
.try_debit(CostUnits(51), t(0))
.unwrap_or_else(|e| panic!("debit {spent} should be funded: {e:?}"));
drop(lease);
settle().await;
}
let tail = slot.load().expect("a lease is installed");
assert_eq!(
tail.remaining(),
CostUnits(49),
"the shrunken tail grant is the state the defect needs"
);
assert!(
!tail.needs_refill(),
"and it sits above its own capped low-water mark, so nothing crosses"
);
let refused = tail.try_debit(CostUnits(51), t(0));
assert!(
matches!(refused, Err(DenyReason::LeaseExhausted { .. })),
"this one debit is still refused: {refused:?}"
);
drop(tail);
settle().await;
let consolidated = slot.load().expect("a lease is installed");
assert_eq!(
consolidated.remaining(),
CostUnits(58),
"the refusal folded the tail's 49 unspent units back in with the \
ledger's 9, so the account's whole remaining balance is reachable"
);
consolidated
.try_debit(CostUnits(51), t(0))
.expect("the request the account could always fund is now admitted");
drop(consolidated);
let stats = manager.counters().snapshot();
assert_eq!(stats.consolidated, 1, "exactly one consolidating rotation");
assert_eq!(
stats.consolidations_deferred, 0,
"nothing was in flight to defer it"
);
manager.shutdown().await;
}
#[tokio::test(start_paused = true)]
async fn consolidation_under_a_shrinking_policy_never_returns_less_than_it_folded() {
let store = Arc::new(
MemoryStore::new(GrantPolicy {
shrink_divisor: 2,
min_grant: CostUnits(1),
max_ttl: SignedDuration::from_secs(3_600),
reclaim_grace: SignedDuration::ZERO,
})
.unwrap(),
);
store.create_account(AccountConfig {
account_id: ACCOUNT,
initial_balance: CostUnits(58),
status: AccountStatus::Active,
capacity_class: CapacityClass::Assured,
});
let grant = store
.acquire(ACCOUNT, CostUnits(49), SignedDuration::from_secs(60), t(0))
.await
.expect("the first grant is funded")
.grant;
assert_eq!(grant.units, CostUnits(29), "the policy shrinks it by half");
let folded = store
.consolidate(
grant.lease_id,
grant.fencing_token,
grant.units,
CostUnits(100),
CostUnits::ZERO,
SignedDuration::from_secs(60),
t(1),
)
.await
.expect("consolidation is funded by the units it returns")
.grant;
assert!(
folded.units >= grant.units,
"a consolidation may grow the holding or leave it alone, never shrink it: \
{} folded in, {} granted",
grant.units.get(),
folded.units.get()
);
assert_eq!(
folded.units,
CostUnits(29),
"the policy still caps growth at balance/2; the floor only forbids the downgrade"
);
}
#[tokio::test(start_paused = true)]
async fn refill_publishes_exhaustion_and_a_topup_clears_it() {
let store = store(0);
let slot = LeaseSlot::for_account(ACCOUNT);
let clock = Arc::new(ManualClock::new(t(0)));
let manager =
LeaseManager::spawn(store.clone(), slot.clone(), clock, manager_config()).unwrap();
settle().await;
assert!(
slot.funding_evidence(t(0))
.is_some_and(|remaining| remaining.is_zero())
);
tollgate_store::AdminStore::deposit(&*store, ACCOUNT, CostUnits(100))
.await
.unwrap();
settle().await;
assert!(
!slot
.funding_evidence(t(0))
.is_some_and(|remaining| remaining.is_zero())
);
assert_eq!(slot.load_observed().unwrap().remaining(), CostUnits(100));
manager.shutdown().await;
}
#[tokio::test(start_paused = true)]
async fn a_shrinking_policy_funds_a_quote_above_half_the_balance() {
let store = MemoryStore::new(GrantPolicy::default()).unwrap();
store.create_account(AccountConfig {
account_id: ACCOUNT,
initial_balance: CostUnits(60),
status: AccountStatus::Active,
capacity_class: CapacityClass::Assured,
});
let slot = LeaseSlot::for_account(ACCOUNT);
let clock = Arc::new(ManualClock::new(t(0)));
let manager = LeaseManager::spawn(
store.clone(),
Arc::clone(&slot),
clock,
LeaseManagerConfig {
poll_interval: std::time::Duration::from_secs(60),
..manager_config()
},
)
.unwrap();
settle().await;
let half = slot.load().expect("a lease is installed");
assert_eq!(half.remaining(), CostUnits(30));
assert!(matches!(
half.try_debit(CostUnits(51), t(0)),
Err(DenyReason::LeaseExhausted { .. })
));
drop(half);
settle().await;
let grown = slot.load().expect("a lease is installed");
assert_eq!(
grown.remaining(),
CostUnits(51),
"grown to the refused quote"
);
grown
.try_debit(CostUnits(51), t(0))
.expect("the quote the account could always fund is admitted");
drop(grown);
assert_eq!(manager.counters().snapshot().consolidated, 1);
manager.shutdown().await;
}
#[tokio::test(start_paused = true)]
async fn refill_publishes_shortfall_and_a_topup_clears_it() {
let store = store(1);
let slot = LeaseSlot::for_account(ACCOUNT);
let clock = Arc::new(ManualClock::new(t(0)));
let manager = LeaseManager::spawn(
store.clone(),
slot.clone(),
clock,
LeaseManagerConfig {
poll_interval: std::time::Duration::from_secs(60),
..manager_config()
},
)
.unwrap();
settle().await;
assert_eq!(slot.load_observed().unwrap().remaining(), CostUnits(1));
assert_eq!(slot.funding_evidence(t(0)), Some(CostUnits(1)));
let refuse_252 = || {
let lease = slot.load().expect("a lease is installed");
assert!(matches!(
lease.try_debit(CostUnits(252), t(0)),
Err(DenyReason::LeaseExhausted { .. })
));
};
refuse_252();
settle().await;
assert_eq!(
slot.funding_evidence(t(0)),
Some(CostUnits(1)),
"the consolidated tail carries the same attested remainder"
);
tollgate_store::AdminStore::deposit(&*store, ACCOUNT, CostUnits(300))
.await
.unwrap();
refuse_252();
settle().await;
assert_eq!(slot.funding_evidence(t(0)), Some(CostUnits(301)));
slot.load()
.expect("a lease is installed")
.try_debit(CostUnits(252), t(0))
.expect("the top-up funds the quote");
manager.shutdown().await;
}
#[tokio::test(start_paused = true)]
async fn a_refill_gap_with_credit_held_elsewhere_is_not_exhaustion() {
let store = store(100);
let held = store
.acquire(ACCOUNT, CostUnits(100), SignedDuration::from_secs(60), t(0))
.await
.unwrap()
.grant;
let slot = LeaseSlot::for_account(ACCOUNT);
let clock = Arc::new(ManualClock::new(t(0)));
let manager =
LeaseManager::spawn(store.clone(), slot.clone(), clock, manager_config()).unwrap();
settle().await;
assert_eq!(slot.funding_evidence(t(0)), Some(CostUnits(100)));
assert!(slot.load_observed().is_none());
store
.release(held.lease_id, held.fencing_token, held.units, t(0))
.await
.unwrap();
settle().await;
assert!(slot.load_observed().is_some());
assert!(
!slot
.funding_evidence(t(0))
.is_some_and(|remaining| remaining.is_zero())
);
manager.shutdown().await;
}