use std::num::{NonZeroU32, NonZeroUsize};
use std::sync::Arc;
use jiff::{SignedDuration, Timestamp};
use tokio::sync::Mutex;
use tollgate_core::{
AccountId, AccountSnapshot, AccountStatus, BudgetSchedule, BudgetView, CapacityClass,
CostTable, CostUnits, FencingToken, Generation, KeyId, LeaseId, PermissionBits, PolicyRevision,
Principal, PublishableSnapshot, RequestId, ResolvedLimits, UsageEvent, UsageSource,
};
use tollgate_store::{
AccountConfig, AdminAuthority, AdminState, AdminStore, AllocateError, BudgetError,
Conservation, CreateAccountError, DEFAULT_RECLAIM_BATCH_LIMIT, GrantPolicy, KeyDirectory,
KeyError, KeyRecord, LeaseAllocator, PublishSnapshotError, Revocation, RolledAccount,
SetStatusError, SnapshotResolution, SnapshotSource, UsageSink,
};
use tollgate_store_postgres::PostgresStore;
static DB_LOCK: Mutex<()> = Mutex::const_new(());
fn t(secs: i64) -> Timestamp {
Timestamp::from_second(secs).unwrap()
}
const TTL: SignedDuration = SignedDuration::from_secs(60);
const ACCOUNT: AccountId = AccountId(1);
fn publishable(snapshot: Arc<AccountSnapshot>) -> PublishableSnapshot {
PublishableSnapshot::try_new(snapshot).expect("test snapshot limits are valid")
}
fn full_grant_policy() -> GrantPolicy {
GrantPolicy {
shrink_divisor: 1,
min_grant: CostUnits(1),
max_ttl: SignedDuration::from_secs(300),
reclaim_grace: SignedDuration::ZERO,
}
}
#[tokio::test]
async fn invalid_pool_config_is_rejected_before_connecting() {
use tollgate_store_postgres::PoolConfig;
for config in [
PoolConfig {
max_connections: 0,
..PoolConfig::default()
},
PoolConfig {
acquire_timeout: std::time::Duration::ZERO,
..PoolConfig::default()
},
] {
assert!(config.validate().is_err());
assert!(
PostgresStore::connect_with("postgres://invalid", full_grant_policy(), config)
.await
.is_err()
);
}
}
#[tokio::test]
async fn invalid_grant_policy_is_rejected_before_connecting() {
let policy = GrantPolicy {
reclaim_grace: SignedDuration::from_secs(-1),
..GrantPolicy::default()
};
assert!(
PostgresStore::connect("postgres://invalid", policy)
.await
.is_err()
);
}
fn redact_url(url: &str) -> String {
match url.split_once('@') {
Some((_, host)) => format!("postgres://<redacted>@{host}"),
None => url.to_owned(),
}
}
async fn empty_store(policy: GrantPolicy) -> Option<Arc<PostgresStore>> {
let Ok(url) = std::env::var("TOLLGATE_PG_URL") else {
assert!(
std::env::var_os("TOLLGATE_REQUIRE_PG").is_none(),
"TOLLGATE_PG_URL is unset but TOLLGATE_REQUIRE_PG is set; \
the Postgres gate must not pass without a database"
);
eprintln!("SKIPPED: TOLLGATE_PG_URL not set (see docker-compose.yml)");
return None;
};
let store = PostgresStore::connect(&url, policy)
.await
.unwrap_or_else(|e| panic!("postgres unreachable at {}: {e}", redact_url(&url)));
tollgate_store_postgres::test_support::truncate_all(&store)
.await
.unwrap();
Some(store)
}
async fn store_with_balance(policy: GrantPolicy, balance: u64) -> Option<Arc<PostgresStore>> {
let Ok(url) = std::env::var("TOLLGATE_PG_URL") else {
assert!(
std::env::var_os("TOLLGATE_REQUIRE_PG").is_none(),
"TOLLGATE_PG_URL is unset but TOLLGATE_REQUIRE_PG is set; \
the Postgres gate must not pass without a database"
);
eprintln!("SKIPPED: TOLLGATE_PG_URL not set (see docker-compose.yml)");
return None;
};
let store = PostgresStore::connect(&url, policy)
.await
.unwrap_or_else(|e| panic!("postgres unreachable at {}: {e}", redact_url(&url)));
tollgate_store_postgres::test_support::truncate_all(&store)
.await
.unwrap();
AdminStore::create_account(
&*store,
AccountConfig {
account_id: ACCOUNT,
initial_balance: CostUnits(balance),
status: AccountStatus::Active,
capacity_class: CapacityClass::Assured,
},
)
.await
.unwrap();
Some(store)
}
#[tokio::test]
async fn nonpositive_lease_ttl_is_rejected_without_debiting() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 100).await else {
return;
};
assert_eq!(
store
.acquire(ACCOUNT, CostUnits(10), SignedDuration::ZERO, t(0))
.await
.unwrap_err(),
AllocateError::InvalidTtl
);
assert_eq!(store.balance(ACCOUNT).await.unwrap(), CostUnits(100));
assert!(
store
.acquire(ACCOUNT, CostUnits(10), TTL, Timestamp::MAX)
.await
.is_err()
);
assert_eq!(store.balance(ACCOUNT).await.unwrap(), CostUnits(100));
}
fn usage(lease: &tollgate_core::LeaseGrant, request: u128, units: u64, at: i64) -> UsageEvent {
UsageEvent::new(
RequestId(request),
lease.account_id,
UsageSource::Leased {
lease_id: lease.lease_id,
fencing_token: lease.fencing_token,
},
CostUnits(units),
t(at),
PolicyRevision::UNSTATED,
None,
)
}
fn overage_usage(account: AccountId, request: u128, units: u64, at: i64) -> UsageEvent {
UsageEvent::new(
RequestId(request),
account,
UsageSource::Overage,
CostUnits(units),
t(at),
PolicyRevision::UNSTATED,
None,
)
}
async fn assert_conserved(store: &PostgresStore) {
let conservation = store.conservation(ACCOUNT).await.unwrap().unwrap();
assert!(
conservation.holds(),
"conservation violated: {conservation:?}"
);
}
#[tokio::test]
async fn adaptive_grant_shrinks_near_exhaustion() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(GrantPolicy::default(), 1_000).await else {
return;
};
let grant = store
.acquire(ACCOUNT, CostUnits(600), TTL, t(0))
.await
.unwrap()
.grant;
assert_eq!(grant.units, CostUnits(500));
let grant = store
.acquire(ACCOUNT, CostUnits(600), TTL, t(0))
.await
.unwrap()
.grant;
assert_eq!(grant.units, CostUnits(250));
let mut drained = 0u64;
loop {
match store.acquire(ACCOUNT, CostUnits(600), TTL, t(0)).await {
Ok(g) => drained += g.grant.units.get(),
Err(AllocateError::BalanceInsufficient(evidence)) => {
assert_eq!(evidence.remaining, CostUnits(1_000));
break;
}
Err(other) => panic!("unexpected: {other}"),
}
}
assert_eq!(drained, 250);
assert_eq!(store.balance(ACCOUNT).await.unwrap(), CostUnits::ZERO);
assert_conserved(&store).await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 8)]
async fn no_double_spend_across_instances() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(GrantPolicy::default(), 100_000).await else {
return;
};
let mut tasks = Vec::new();
for _ in 0..16 {
let store = Arc::clone(&store);
tasks.push(tokio::spawn(async move {
let mut granted = 0u64;
loop {
match store.acquire(ACCOUNT, CostUnits(1_000), TTL, t(0)).await {
Ok(g) => granted += g.grant.units.get(),
Err(AllocateError::BalanceInsufficient(evidence)) => {
assert_eq!(evidence.remaining, CostUnits(100_000));
return granted;
}
Err(other) => panic!("unexpected: {other}"),
}
}
}));
}
let mut total = 0u64;
for task in tasks {
total += task.await.unwrap();
}
assert_eq!(total, 100_000);
assert_eq!(store.balance(ACCOUNT).await.unwrap(), CostUnits::ZERO);
assert_conserved(&store).await;
}
#[tokio::test]
async fn newer_lease_does_not_invalidate_older_active_capability() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 1_000).await else {
return;
};
let older = store
.acquire(ACCOUNT, CostUnits(400), TTL, t(0))
.await
.unwrap()
.grant;
let newer = store
.acquire(ACCOUNT, CostUnits(400), TTL, t(0))
.await
.unwrap()
.grant;
assert!(newer.fencing_token > older.fencing_token);
let report = store
.ingest(&[usage(&older, 1, 50, 1)], t(1))
.await
.unwrap();
assert_eq!((report.accepted, report.rejected), (1, 0));
store
.release(older.lease_id, older.fencing_token, CostUnits(350), t(2))
.await
.unwrap();
let report = store
.ingest(&[usage(&newer, 2, 25, 2)], t(2))
.await
.unwrap();
assert_eq!((report.accepted, report.rejected), (1, 0));
store
.release(newer.lease_id, newer.fencing_token, CostUnits(375), t(3))
.await
.unwrap();
assert_eq!(store.balance(ACCOUNT).await.unwrap(), CostUnits(925));
assert_eq!(store.usage_recorded(ACCOUNT).await.unwrap(), CostUnits(75));
assert_conserved(&store).await;
}
#[tokio::test]
async fn wrong_token_release_leaves_lease_reclaimable() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 1_000).await else {
return;
};
let lease = store
.acquire(ACCOUNT, CostUnits(400), TTL, t(0))
.await
.unwrap()
.grant;
assert_eq!(
store
.release(lease.lease_id, FencingToken(999), CostUnits(400), t(1))
.await
.unwrap_err(),
AllocateError::Fenced
);
let reclaimed = store.reclaim_expired(t(61)).await.unwrap();
assert_eq!(reclaimed.len(), 1);
assert_eq!(reclaimed[0].lease_id, lease.lease_id);
assert_eq!(
store.usage_recorded(ACCOUNT).await.unwrap(),
CostUnits::ZERO
);
assert_conserved(&store).await;
}
#[tokio::test]
async fn usage_rejects_mismatched_lease_capability() {
const OTHER: AccountId = AccountId(2);
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 1_000).await else {
return;
};
AdminStore::create_account(
&*store,
AccountConfig {
account_id: OTHER,
initial_balance: CostUnits(100),
status: AccountStatus::Active,
capacity_class: CapacityClass::Assured,
},
)
.await
.unwrap();
let lease = store
.acquire(ACCOUNT, CostUnits(400), TTL, t(0))
.await
.unwrap()
.grant;
let mut wrong_token = usage(&lease, 1, 10, 1);
wrong_token.source = UsageSource::Leased {
lease_id: lease.lease_id,
fencing_token: FencingToken(lease.fencing_token.0.checked_add(1).unwrap()),
};
let report = store.ingest(&[wrong_token], t(1)).await.unwrap();
assert_eq!(
(report.accepted, report.duplicate, report.rejected),
(0, 0, 1)
);
let mut wrong_account = usage(&lease, 2, 10, 1);
wrong_account.account_id = OTHER;
let report = store.ingest(&[wrong_account], t(1)).await.unwrap();
assert_eq!(
(report.accepted, report.duplicate, report.rejected),
(0, 0, 1)
);
assert_eq!(
store.usage_recorded(ACCOUNT).await.unwrap(),
CostUnits::ZERO
);
assert_eq!(store.usage_recorded(OTHER).await.unwrap(), CostUnits::ZERO);
assert_conserved(&store).await;
assert!(store.conservation(OTHER).await.unwrap().unwrap().holds());
}
#[tokio::test]
async fn reclaim_forfeits_an_unreleased_remainder_as_provisional_loss() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 1_000).await else {
return;
};
let lease = store
.acquire(ACCOUNT, CostUnits(500), TTL, t(0))
.await
.unwrap()
.grant;
assert_eq!(store.balance(ACCOUNT).await.unwrap(), CostUnits(500));
let report = store
.ingest(&[usage(&lease, 1, 70, 10), usage(&lease, 2, 50, 20)], t(20))
.await
.unwrap();
assert_eq!(report.accepted, 2);
assert!(store.reclaim_expired(t(59)).await.unwrap().is_empty());
let reclaimed = store.reclaim_expired(t(60)).await.unwrap();
assert_eq!(reclaimed[0].forfeited, CostUnits(380));
assert_eq!(
store.balance(ACCOUNT).await.unwrap(),
CostUnits(500),
"nothing returns"
);
assert_eq!(store.usage_recorded(ACCOUNT).await.unwrap(), CostUnits(120));
assert_eq!(
store
.conservation(ACCOUNT)
.await
.unwrap()
.unwrap()
.settlement_loss,
CostUnits(380)
);
assert_conserved(&store).await;
}
#[tokio::test]
async fn straggler_usage_after_reclaim_is_billed_against_the_forfeit() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 1_000).await else {
return;
};
let lease = store
.acquire(ACCOUNT, CostUnits(500), TTL, t(0))
.await
.unwrap()
.grant;
store
.ingest(&[usage(&lease, 1, 120, 10)], t(10))
.await
.unwrap();
store.reclaim_expired(t(60)).await.unwrap();
let report = store
.ingest(&[usage(&lease, 2, 300, 50)], t(61))
.await
.unwrap();
assert_eq!(report.accepted, 1, "billed, not dropped");
assert_eq!(store.usage_recorded(ACCOUNT).await.unwrap(), CostUnits(420));
let over = store
.ingest(&[usage(&lease, 3, 81, 55)], t(62))
.await
.unwrap();
assert_eq!(over.rejected, 1, "80 forfeited units remain, not 81");
let exact = store
.ingest(&[usage(&lease, 4, 80, 55)], t(62))
.await
.unwrap();
assert_eq!(exact.accepted, 1);
let c = store.conservation(ACCOUNT).await.unwrap().unwrap();
assert_eq!(c.settlement_loss, CostUnits::ZERO);
assert_eq!(c.settled_usage, CostUnits(500));
assert_eq!(store.balance(ACCOUNT).await.unwrap(), CostUnits(500));
assert_conserved(&store).await;
}
#[tokio::test]
async fn expired_backlog_is_reclaimed_in_bounded_batches() {
const OTHER: AccountId = AccountId(2);
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 200).await else {
return;
};
AdminStore::create_account(
&*store,
AccountConfig {
account_id: OTHER,
initial_balance: CostUnits(400),
status: AccountStatus::Active,
capacity_class: CapacityClass::Assured,
},
)
.await
.unwrap();
let leases = [
store
.acquire(ACCOUNT, CostUnits(200), TTL, t(0))
.await
.unwrap()
.grant,
store
.acquire(OTHER, CostUnits(200), TTL, t(0))
.await
.unwrap()
.grant,
store
.acquire(OTHER, CostUnits(200), TTL, t(0))
.await
.unwrap()
.grant,
];
store
.ingest(
&[usage(&leases[0], 10, 25, 10), usage(&leases[2], 11, 50, 10)],
t(10),
)
.await
.unwrap();
let limit = NonZeroUsize::new(2).unwrap();
let first = store.reclaim_expired_batch(t(60), limit).await.unwrap();
assert_eq!(first.len(), 2);
assert!(first.is_saturated());
let second = store.reclaim_expired_batch(t(60), limit).await.unwrap();
assert_eq!(second.len(), 1);
assert!(!second.is_saturated());
let mut expected_ids: Vec<_> = leases.iter().map(|lease| lease.lease_id).collect();
expected_ids.sort_by_key(|lease_id| lease_id.0);
let mut reclaimed_ids: Vec<_> = first
.reclaimed()
.iter()
.chain(second.reclaimed())
.map(|lease| lease.lease_id)
.collect();
reclaimed_ids.sort_by_key(|lease_id| lease_id.0);
assert_eq!(reclaimed_ids, expected_ids);
assert_eq!(store.balance(ACCOUNT).await.unwrap(), CostUnits::ZERO);
assert_eq!(store.balance(OTHER).await.unwrap(), CostUnits::ZERO);
assert_eq!(
store
.conservation(ACCOUNT)
.await
.unwrap()
.unwrap()
.settlement_loss,
CostUnits(175)
);
assert_eq!(
store
.conservation(OTHER)
.await
.unwrap()
.unwrap()
.settlement_loss,
CostUnits(350)
);
assert_eq!(store.usage_recorded(ACCOUNT).await.unwrap(), CostUnits(25));
assert_eq!(store.usage_recorded(OTHER).await.unwrap(), CostUnits(50));
assert_conserved(&store).await;
let other_conservation = store.conservation(OTHER).await.unwrap().unwrap();
assert!(
other_conservation.holds(),
"conservation violated: {other_conservation:?}"
);
}
#[tokio::test]
async fn a_bounded_reclaim_page_settles_the_oldest_due_leases_first() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 100).await else {
return;
};
for id in 2..=4u128 {
AdminStore::create_account(
&*store,
AccountConfig {
account_id: AccountId(id),
initial_balance: CostUnits(100),
status: AccountStatus::Active,
capacity_class: CapacityClass::Assured,
},
)
.await
.unwrap();
}
let mut leases = Vec::new();
for (id, acquired_at) in [(4u128, 0i64), (3, 10), (2, 20), (1, 600)] {
leases.push(
store
.acquire(AccountId(id), CostUnits(100), TTL, t(acquired_at))
.await
.unwrap()
.grant,
);
}
let limit = NonZeroUsize::new(2).unwrap();
let first = store.reclaim_expired_batch(t(100), limit).await.unwrap();
assert!(
first.is_saturated(),
"the page is full, so ordering decides it"
);
assert_eq!(
first
.reclaimed()
.iter()
.map(|lease| lease.account_id.0)
.collect::<Vec<_>>(),
vec![4, 3],
"the two oldest expiries settle first; account order would give [2, 3]"
);
let second = store.reclaim_expired_batch(t(100), limit).await.unwrap();
assert!(
!second.is_saturated(),
"the due prefix ends at the third lease"
);
assert_eq!(
second
.reclaimed()
.iter()
.map(|lease| lease.lease_id)
.collect::<Vec<_>>(),
vec![leases[2].lease_id],
"the walk stops at the lease that is not yet due"
);
assert_eq!(store.balance(AccountId(1)).await.unwrap(), CostUnits::ZERO);
assert_eq!(
store
.conservation(AccountId(1))
.await
.unwrap()
.unwrap()
.settlement_loss,
CostUnits::ZERO
);
let last = store.reclaim_expired_batch(t(700), limit).await.unwrap();
assert_eq!(
last.reclaimed()
.iter()
.map(|lease| lease.lease_id)
.collect::<Vec<_>>(),
vec![leases[3].lease_id]
);
assert_eq!(store.balance(AccountId(1)).await.unwrap(), CostUnits::ZERO);
assert_eq!(
store
.conservation(AccountId(1))
.await
.unwrap()
.unwrap()
.settlement_loss,
CostUnits(100)
);
assert_conserved(&store).await;
}
#[tokio::test]
async fn usage_replay_is_idempotent() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 1_000).await else {
return;
};
let lease = store
.acquire(ACCOUNT, CostUnits(500), TTL, t(0))
.await
.unwrap()
.grant;
let batch = [usage(&lease, 1, 70, 10), usage(&lease, 2, 50, 10)];
let first = store.ingest(&batch, t(10)).await.unwrap();
assert_eq!((first.accepted, first.duplicate), (2, 0));
let replay = store.ingest(&batch, t(11)).await.unwrap();
assert_eq!((replay.accepted, replay.duplicate), (0, 2));
assert_eq!(store.usage_recorded(ACCOUNT).await.unwrap(), CostUnits(120));
let with_dup = [usage(&lease, 3, 10, 12), usage(&lease, 3, 10, 12)];
let report = store.ingest(&with_dup, t(12)).await.unwrap();
assert_eq!((report.accepted, report.duplicate), (1, 1));
assert_eq!(store.usage_recorded(ACCOUNT).await.unwrap(), CostUnits(130));
assert_conserved(&store).await;
}
#[tokio::test]
async fn mixed_usage_batch_preserves_partial_acceptance() {
const OTHER: AccountId = AccountId(2);
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 1_000).await else {
return;
};
AdminStore::create_account(
&*store,
AccountConfig {
account_id: OTHER,
initial_balance: CostUnits(400),
status: AccountStatus::Active,
capacity_class: CapacityClass::Assured,
},
)
.await
.unwrap();
let active = store
.acquire(ACCOUNT, CostUnits(300), TTL, t(0))
.await
.unwrap()
.grant;
let settled = store
.acquire(ACCOUNT, CostUnits(200), TTL, t(0))
.await
.unwrap()
.grant;
let other = store
.acquire(OTHER, CostUnits(300), TTL, t(0))
.await
.unwrap()
.grant;
store
.release(
settled.lease_id,
settled.fencing_token,
CostUnits(170),
t(5),
)
.await
.unwrap();
let prior = usage(&active, 1, 10, 2);
assert_eq!(store.ingest(&[prior], t(2)).await.unwrap().accepted, 1);
let mut replay_with_unrepresentable_units = prior;
replay_with_unrepresentable_units.units = CostUnits(u64::MAX);
let repeated = usage(&other, 4, 40, 6);
let mut wrong_fence = usage(&other, 6, 10, 6);
wrong_fence.source = UsageSource::Leased {
lease_id: other.lease_id,
fencing_token: FencingToken(other.fencing_token.0.checked_add(1).unwrap()),
};
let mut unknown_lease = usage(&other, 7, 10, 6);
unknown_lease.source = UsageSource::Leased {
lease_id: LeaseId(u128::MAX),
fencing_token: other.fencing_token,
};
let batch = [
usage(&active, 2, 20, 6),
usage(&active, 10, i64::MAX as u64 + 1, 6),
usage(&active, 11, u64::MAX, 6),
replay_with_unrepresentable_units,
wrong_fence,
unknown_lease,
usage(&settled, 3, 30, 4),
repeated,
repeated,
usage(&active, 5, 15, 6),
overage_usage(ACCOUNT, 8, 25, 6),
overage_usage(AccountId(u128::MAX), 9, 25, 6),
];
let report = store.ingest(&batch, t(6)).await.unwrap();
assert_eq!(
(report.accepted, report.duplicate, report.rejected),
(5, 2, 5)
);
assert_eq!(store.usage_recorded(ACCOUNT).await.unwrap(), CostUnits(100));
assert_eq!(store.usage_recorded(OTHER).await.unwrap(), CostUnits(40));
assert_eq!(
store
.conservation(ACCOUNT)
.await
.unwrap()
.unwrap()
.overage_recorded,
CostUnits(25),
"the overage event funds the units it billed"
);
assert_eq!(
store
.conservation(ACCOUNT)
.await
.unwrap()
.unwrap()
.settlement_loss,
CostUnits::ZERO
);
assert_conserved(&store).await;
let other_conservation = store.conservation(OTHER).await.unwrap().unwrap();
assert!(
other_conservation.holds(),
"conservation violated: {other_conservation:?}"
);
}
#[tokio::test]
async fn unrepresentable_usage_rejects_only_that_event_and_does_not_claim_its_id() {
let _guard = DB_LOCK.lock().await;
for leased in [false, true] {
for invalid in [i64::MAX as u64 + 1, u64::MAX] {
let Some(store) = store_with_balance(full_grant_policy(), 100).await else {
return;
};
let lease = store
.acquire(ACCOUNT, CostUnits(100), TTL, t(0))
.await
.unwrap()
.grant;
let event = |id, units| {
if leased {
usage(&lease, id, units, 0)
} else {
overage_usage(ACCOUNT, id, units, 0)
}
};
let batch = [
event(1, 10),
event(2, invalid),
event(2, 20),
event(2, invalid),
event(3, 30),
];
let report = store.ingest(&batch, t(1)).await.unwrap();
report.validate(batch.len()).unwrap();
assert_eq!(
(report.accepted, report.duplicate, report.rejected),
(3, 1, 1)
);
assert_eq!(store.usage_recorded(ACCOUNT).await.unwrap(), CostUnits(60));
assert_conserved(&store).await;
let replay = store.ingest(&batch, t(2)).await.unwrap();
assert_eq!(
(replay.accepted, replay.duplicate, replay.rejected),
(0, 5, 0)
);
let pool = corruption_pool().await;
let stored: Vec<i64> =
sqlx::query_scalar("SELECT units FROM tollgate_usage_events ORDER BY request_id")
.fetch_all(&pool)
.await
.unwrap();
assert_eq!(stored, [10, 20, 30]);
}
}
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn concurrent_multi_account_batches_use_stable_lock_order() {
const OTHER: AccountId = AccountId(2);
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 1_000).await else {
return;
};
AdminStore::create_account(
&*store,
AccountConfig {
account_id: OTHER,
initial_balance: CostUnits(1_000),
status: AccountStatus::Active,
capacity_class: CapacityClass::Assured,
},
)
.await
.unwrap();
let a1 = store
.acquire(ACCOUNT, CostUnits(400), TTL, t(0))
.await
.unwrap()
.grant;
let a2 = store
.acquire(ACCOUNT, CostUnits(400), TTL, t(0))
.await
.unwrap()
.grant;
let b1 = store
.acquire(OTHER, CostUnits(400), TTL, t(0))
.await
.unwrap()
.grant;
let b2 = store
.acquire(OTHER, CostUnits(400), TTL, t(0))
.await
.unwrap()
.grant;
let barrier = Arc::new(tokio::sync::Barrier::new(2));
let left = tokio::spawn({
let store = Arc::clone(&store);
let barrier = Arc::clone(&barrier);
async move {
for i in 0..64u128 {
barrier.wait().await;
store
.ingest(
&[usage(&a1, 10_000 + i, 1, 1), usage(&b1, 20_000 + i, 1, 1)],
t(1),
)
.await?;
}
Ok::<(), tollgate_store::IngestError>(())
}
});
let right = tokio::spawn({
let store = Arc::clone(&store);
async move {
for i in 0..64u128 {
barrier.wait().await;
store
.ingest(
&[usage(&b2, 30_000 + i, 1, 1), usage(&a2, 40_000 + i, 1, 1)],
t(1),
)
.await?;
}
Ok::<(), tollgate_store::IngestError>(())
}
});
left.await.unwrap().unwrap();
right.await.unwrap().unwrap();
assert_eq!(store.usage_recorded(ACCOUNT).await.unwrap(), CostUnits(128));
assert_eq!(store.usage_recorded(OTHER).await.unwrap(), CostUnits(128));
}
#[tokio::test]
async fn graceful_release_returns_unspent() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 1_000).await else {
return;
};
let lease = store
.acquire(ACCOUNT, CostUnits(500), TTL, t(0))
.await
.unwrap()
.grant;
store
.ingest(&[usage(&lease, 1, 100, 5)], t(5))
.await
.unwrap();
assert_eq!(
store
.release(lease.lease_id, lease.fencing_token, CostUnits(450), t(10))
.await
.unwrap_err(),
AllocateError::InvalidRelease
);
store
.release(lease.lease_id, lease.fencing_token, CostUnits(400), t(10))
.await
.unwrap();
assert_eq!(store.balance(ACCOUNT).await.unwrap(), CostUnits(900));
assert_eq!(store.usage_recorded(ACCOUNT).await.unwrap(), CostUnits(100));
assert_conserved(&store).await;
assert_eq!(
store
.release(lease.lease_id, lease.fencing_token, CostUnits(0), t(11))
.await
.unwrap_err(),
AllocateError::LeaseNotActive
);
}
#[tokio::test]
async fn consolidation_never_grants_less_than_it_folded_in() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(GrantPolicy::default(), 58).await else {
return;
};
let first = store
.acquire(ACCOUNT, CostUnits(100), TTL, t(0))
.await
.unwrap()
.grant;
assert_eq!(first.units, CostUnits(29));
assert_eq!(store.balance(ACCOUNT).await.unwrap(), CostUnits(29));
let folded = store
.consolidate(
first.lease_id,
first.fencing_token,
first.units,
CostUnits(100),
CostUnits::ZERO,
TTL,
t(1),
)
.await
.unwrap()
.grant;
assert!(
folded.units >= first.units,
"{} folded in, {} granted",
first.units.get(),
folded.units.get()
);
assert_ne!(folded.lease_id, first.lease_id);
assert!(folded.fencing_token > first.fencing_token);
assert_conserved(&store).await;
assert_eq!(
store
.release(first.lease_id, first.fencing_token, CostUnits(0), t(2))
.await
.unwrap_err(),
AllocateError::LeaseNotActive
);
}
#[tokio::test]
async fn consolidation_folds_the_tail_grant_and_the_ledger_into_one_lease() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 58).await else {
return;
};
let tail = store
.acquire(ACCOUNT, CostUnits(49), TTL, t(0))
.await
.unwrap()
.grant;
assert_eq!(tail.units, CostUnits(49));
assert_eq!(store.balance(ACCOUNT).await.unwrap(), CostUnits(9));
let folded = store
.consolidate(
tail.lease_id,
tail.fencing_token,
tail.units,
CostUnits(100),
CostUnits::ZERO,
TTL,
t(1),
)
.await
.unwrap()
.grant;
assert_eq!(folded.units, CostUnits(58));
assert_eq!(store.balance(ACCOUNT).await.unwrap(), CostUnits::ZERO);
assert_conserved(&store).await;
}
#[tokio::test]
async fn consolidation_grows_to_a_proven_quote_under_a_shrinking_policy() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(GrantPolicy::default(), 60).await else {
return;
};
let held = store
.acquire(ACCOUNT, CostUnits(1_000), TTL, t(0))
.await
.unwrap()
.grant;
assert_eq!(held.units, CostUnits(30));
let same = store
.consolidate(
held.lease_id,
held.fencing_token,
held.units,
CostUnits(1_000),
CostUnits::ZERO,
TTL,
t(1),
)
.await
.unwrap()
.grant;
assert_eq!(
same.units,
CostUnits(30),
"without demand, the GL-109 floor"
);
let grown = store
.consolidate(
same.lease_id,
same.fencing_token,
same.units,
CostUnits(1_000),
CostUnits(51),
TTL,
t(2),
)
.await
.unwrap()
.grant;
assert_eq!(grown.units, CostUnits(51));
assert_eq!(store.balance(ACCOUNT).await.unwrap(), CostUnits(9));
assert_conserved(&store).await;
}
#[tokio::test]
async fn consolidation_never_grows_past_the_restored_balance() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(GrantPolicy::default(), 60).await else {
return;
};
let held = store
.acquire(ACCOUNT, CostUnits(1_000), TTL, t(0))
.await
.unwrap()
.grant;
let kept = store
.consolidate(
held.lease_id,
held.fencing_token,
held.units,
CostUnits(1_000),
CostUnits(61),
TTL,
t(1),
)
.await
.unwrap()
.grant;
assert_eq!(kept.units, CostUnits(30));
assert_eq!(store.balance(ACCOUNT).await.unwrap(), CostUnits(30));
assert_conserved(&store).await;
}
#[tokio::test]
async fn growth_leaves_the_rest_for_another_instance() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(GrantPolicy::default(), 100).await else {
return;
};
let held = store
.acquire(ACCOUNT, CostUnits(1_000), TTL, t(0))
.await
.unwrap()
.grant;
assert_eq!(held.units, CostUnits(50));
let grown = store
.consolidate(
held.lease_id,
held.fencing_token,
held.units,
CostUnits(1_000),
CostUnits(70),
TTL,
t(1),
)
.await
.unwrap()
.grant;
assert_eq!(grown.units, CostUnits(70));
let other = store
.acquire(ACCOUNT, CostUnits(1_000), TTL, t(1))
.await
.unwrap()
.grant;
assert_eq!(other.units, CostUnits(15), "30 left, halved by the policy");
assert_conserved(&store).await;
}
#[tokio::test]
async fn a_refused_consolidation_leaves_the_original_lease_spendable() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 100).await else {
return;
};
let lease = store
.acquire(ACCOUNT, CostUnits(100), TTL, t(0))
.await
.unwrap()
.grant;
assert_eq!(store.balance(ACCOUNT).await.unwrap(), CostUnits::ZERO);
assert_eq!(
store
.consolidate(
lease.lease_id,
lease.fencing_token,
CostUnits::ZERO,
CostUnits(100),
CostUnits::ZERO,
TTL,
t(1)
)
.await
.unwrap_err(),
AllocateError::InsufficientBalance
);
store
.release(lease.lease_id, lease.fencing_token, CostUnits(100), t(2))
.await
.unwrap();
assert_eq!(store.balance(ACCOUNT).await.unwrap(), CostUnits(100));
assert_conserved(&store).await;
}
#[tokio::test]
async fn consolidating_without_the_lease_capability_moves_no_units() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 1_000).await else {
return;
};
let lease = store
.acquire(ACCOUNT, CostUnits(400), TTL, t(0))
.await
.unwrap()
.grant;
let before = store.balance(ACCOUNT).await.unwrap();
assert_eq!(
store
.consolidate(
lease.lease_id,
FencingToken(lease.fencing_token.0 + 7),
CostUnits(400),
CostUnits(400),
CostUnits::ZERO,
TTL,
t(1)
)
.await
.unwrap_err(),
AllocateError::Fenced
);
assert_eq!(
store
.consolidate(
LeaseId(9_999),
lease.fencing_token,
CostUnits(400),
CostUnits(400),
CostUnits::ZERO,
TTL,
t(1)
)
.await
.unwrap_err(),
AllocateError::UnknownLease
);
assert_eq!(
store
.consolidate(
lease.lease_id,
lease.fencing_token,
CostUnits(401),
CostUnits(400),
CostUnits::ZERO,
TTL,
t(1)
)
.await
.unwrap_err(),
AllocateError::InvalidRelease
);
assert_eq!(store.balance(ACCOUNT).await.unwrap(), before);
assert_conserved(&store).await;
}
#[tokio::test]
async fn a_consolidation_with_an_invalid_ttl_settles_nothing() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 1_000).await else {
return;
};
let lease = store
.acquire(ACCOUNT, CostUnits(400), TTL, t(0))
.await
.unwrap()
.grant;
assert_eq!(
store
.consolidate(
lease.lease_id,
lease.fencing_token,
CostUnits(400),
CostUnits(400),
CostUnits::ZERO,
SignedDuration::ZERO,
t(1)
)
.await
.unwrap_err(),
AllocateError::InvalidTtl
);
store
.release(lease.lease_id, lease.fencing_token, CostUnits(400), t(2))
.await
.expect("the lease was never settled");
assert_conserved(&store).await;
}
#[tokio::test]
async fn reclaim_waits_for_grace_and_release_works_within_it() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(
GrantPolicy {
shrink_divisor: 1,
min_grant: CostUnits(1),
max_ttl: SignedDuration::from_secs(300),
reclaim_grace: SignedDuration::from_secs(30),
},
1_000,
)
.await
else {
return;
};
let lease = store
.acquire(ACCOUNT, CostUnits(500), TTL, t(0))
.await
.unwrap()
.grant;
assert!(store.reclaim_expired(t(60)).await.unwrap().is_empty());
assert!(store.reclaim_expired(t(89)).await.unwrap().is_empty());
let report = store
.ingest(&[usage(&lease, 1, 120, 59)], t(65))
.await
.unwrap();
assert_eq!(report.accepted, 1);
store
.release(lease.lease_id, lease.fencing_token, CostUnits(380), t(70))
.await
.unwrap();
assert_eq!(store.balance(ACCOUNT).await.unwrap(), CostUnits(880));
assert_conserved(&store).await;
let lease2 = store
.acquire(ACCOUNT, CostUnits(400), TTL, t(70))
.await
.unwrap()
.grant; assert!(store.reclaim_expired(t(159)).await.unwrap().is_empty());
let reclaimed = store.reclaim_expired(t(160)).await.unwrap();
assert_eq!(reclaimed[0].lease_id, lease2.lease_id);
assert_eq!(reclaimed[0].forfeited, CostUnits(400));
assert_conserved(&store).await;
}
#[tokio::test]
async fn straggler_usage_after_release_is_billed() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 1_000).await else {
return;
};
let lease = store
.acquire(ACCOUNT, CostUnits(500), TTL, t(0))
.await
.unwrap()
.grant;
store
.release(lease.lease_id, lease.fencing_token, CostUnits(470), t(10))
.await
.unwrap();
let c = store.conservation(ACCOUNT).await.unwrap().unwrap();
assert_eq!(c.settlement_loss, CostUnits(30));
let report = store
.ingest(&[usage(&lease, 1, 30, 5)], t(11))
.await
.unwrap();
assert_eq!(report.accepted, 1);
let c = store.conservation(ACCOUNT).await.unwrap().unwrap();
assert_eq!(c.settlement_loss, CostUnits::ZERO);
assert_eq!(store.usage_recorded(ACCOUNT).await.unwrap(), CostUnits(30));
assert!(c.holds());
let report = store
.ingest(&[usage(&lease, 2, 1, 12)], t(12))
.await
.unwrap();
assert_eq!(report.rejected, 1);
}
#[tokio::test]
async fn publish_pushes_to_subscribers_only_when_the_row_changes() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 100).await else {
return;
};
let principal = Principal(4242);
let mut updates = SnapshotSource::subscribe(&*store);
let snapshot = |generation: u64| {
publishable(Arc::new(
AccountSnapshot::builder(
ACCOUNT,
Generation(generation),
AccountStatus::Active,
t(10_000),
PermissionBits::ALL,
ResolvedLimits::new(64).with_weighted_rate(1_000, 1_000),
Arc::new(CostTable::builder(CostUnits(1), CostUnits(1)).build()),
)
.build(),
))
};
store
.publish_snapshot(principal, snapshot(5))
.await
.unwrap();
let pushed = updates.recv().await.expect("a new generation is pushed");
assert_eq!(pushed.principal, principal);
store
.publish_snapshot(principal, snapshot(4))
.await
.unwrap();
store.remove_snapshot(principal).await.unwrap();
let pushed = updates.recv().await.expect("the revocation is pushed");
assert!(
matches!(pushed.resolution, SnapshotResolution::Revoked { .. }),
"the stale publish must not have produced a push of its own"
);
}
#[tokio::test]
async fn recreate_account_is_refused_and_nondestructive() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 1_000).await else {
return;
};
let lease = store
.acquire(ACCOUNT, CostUnits(400), TTL, t(0))
.await
.unwrap()
.grant;
assert_eq!(
AdminStore::create_account(
&*store,
AccountConfig {
account_id: ACCOUNT,
initial_balance: CostUnits(5),
status: AccountStatus::Active,
capacity_class: CapacityClass::Assured,
},
)
.await
.unwrap_err(),
CreateAccountError::AlreadyExists
);
assert_eq!(store.balance(ACCOUNT).await.unwrap(), CostUnits(600));
let replacement = store
.acquire(ACCOUNT, CostUnits(100), TTL, t(1))
.await
.unwrap()
.grant;
assert!(replacement.fencing_token > lease.fencing_token);
assert_conserved(&store).await;
}
#[tokio::test]
async fn a_non_active_account_refuses_leases() {
let _guard = DB_LOCK.lock().await;
for status in [AccountStatus::Suspended, AccountStatus::Closed] {
let Some(store) = store_with_balance(GrantPolicy::default(), 1_000).await else {
return;
};
AdminStore::set_account_status(&*store, ACCOUNT, status)
.await
.unwrap();
assert_eq!(
store
.acquire(ACCOUNT, CostUnits(100), TTL, t(0))
.await
.unwrap_err(),
AllocateError::AccountInactive,
"{status:?} must refuse leases"
);
}
}
#[tokio::test]
async fn unknown_account_status_update_is_refused() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(GrantPolicy::default(), 1_000).await else {
return;
};
let unknown = AccountId(999);
for status in [
AccountStatus::Active,
AccountStatus::Suspended,
AccountStatus::Closed,
] {
assert_eq!(
AdminStore::set_account_status(&*store, unknown, status)
.await
.unwrap_err(),
SetStatusError::UnknownAccount,
"an unknown account cannot be set to {status:?}"
);
}
store
.acquire(ACCOUNT, CostUnits(100), TTL, t(0))
.await
.expect("refusing the unknown account leaves existing accounts active");
}
#[tokio::test]
async fn creating_a_suspended_account_denies_from_birth() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(GrantPolicy::default(), 1_000).await else {
return;
};
let suspended = AccountId(4_242);
AdminStore::create_account(
&*store,
AccountConfig {
account_id: suspended,
initial_balance: CostUnits(1_000),
status: AccountStatus::Suspended,
capacity_class: CapacityClass::Assured,
},
)
.await
.expect("creation succeeds");
assert_eq!(
store
.acquire(suspended, CostUnits(100), TTL, t(0))
.await
.unwrap_err(),
AllocateError::AccountInactive,
"a suspended account refuses from creation, not only after a transition"
);
}
#[tokio::test]
async fn enumerating_principals_includes_revoked_ones() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 1_000).await else {
return;
};
assert_eq!(
store.principals().await.unwrap(),
Some(Vec::new()),
"an empty catalogue is Some(empty), never None"
);
let snapshot = || {
publishable(Arc::new(
AccountSnapshot::builder(
ACCOUNT,
Generation(1),
AccountStatus::Active,
t(10_000),
PermissionBits::ALL,
ResolvedLimits::new(64).with_weighted_rate(1_000, 1_000),
Arc::new(CostTable::builder(CostUnits(1), CostUnits(1)).build()),
)
.build(),
))
};
let live = Principal(1);
let revoked = Principal(2);
AdminStore::publish_snapshot(&*store, live, snapshot())
.await
.unwrap();
AdminStore::publish_snapshot(&*store, revoked, snapshot())
.await
.unwrap();
AdminStore::remove_snapshot(&*store, revoked).await.unwrap();
let mut listed = store.principals().await.unwrap().expect("enumerable");
listed.sort_by_key(|principal| principal.0);
assert_eq!(listed, vec![live, revoked]);
}
#[tokio::test]
async fn a_ttl_beyond_the_policy_maximum_is_clamped() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 1_000).await else {
return;
};
let lease = store
.acquire(
ACCOUNT,
CostUnits(100),
SignedDuration::from_secs(3_600),
t(0),
)
.await
.unwrap()
.grant;
assert_eq!(
lease.expires_at,
t(300),
"the policy caps the lease at max_ttl, not at what the caller asked for"
);
let exact = store
.acquire(
ACCOUNT,
CostUnits(100),
SignedDuration::from_secs(300),
t(0),
)
.await
.unwrap()
.grant;
assert_eq!(exact.expires_at, t(300));
let shorter = store
.acquire(ACCOUNT, CostUnits(100), SignedDuration::from_secs(30), t(0))
.await
.unwrap()
.grant;
assert_eq!(shorter.expires_at, t(30));
}
#[tokio::test]
async fn depositing_funds_the_account_and_the_ledger_agrees() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 100).await else {
return;
};
assert_eq!(store.balance(ACCOUNT).await.unwrap(), CostUnits(100));
AdminStore::deposit(&*store, ACCOUNT, CostUnits(400))
.await
.unwrap();
assert_eq!(store.balance(ACCOUNT).await.unwrap(), CostUnits(500));
let conservation = store.conservation(ACCOUNT).await.unwrap().unwrap();
assert_eq!(conservation.deposited, CostUnits(500));
assert!(conservation.holds(), "conservation: {conservation:?}");
let lease = store
.acquire(ACCOUNT, CostUnits(500), TTL, t(0))
.await
.unwrap()
.grant;
assert_eq!(lease.units, CostUnits(500));
assert_conserved(&store).await;
assert_eq!(
AdminStore::deposit(&*store, AccountId(999), CostUnits(1))
.await
.unwrap_err(),
AllocateError::UnknownAccount,
"an unknown account is refused, not silently created"
);
}
#[tokio::test]
async fn snapshot_publish_fetch_and_generation_monotonicity() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(GrantPolicy::default(), 1_000).await else {
return;
};
let principal = Principal(42);
let snapshot = |generation: u64| {
publishable(Arc::new(
AccountSnapshot::builder(
ACCOUNT,
Generation(generation),
AccountStatus::Active,
t(10_000),
PermissionBits::ALL,
ResolvedLimits::new(64).with_weighted_rate(1_000, 1_000),
Arc::new(CostTable::builder(CostUnits(1), CostUnits(1)).build()),
)
.build(),
))
};
assert!(matches!(
store.snapshot(principal).await.unwrap(),
SnapshotResolution::Unknown
));
store
.publish_snapshot(principal, snapshot(3))
.await
.unwrap();
let SnapshotResolution::Present(fetched) = store.snapshot(principal).await.unwrap() else {
panic!("published snapshot must be present");
};
assert_eq!(fetched.generation, Generation(3));
store
.publish_snapshot(principal, snapshot(2))
.await
.unwrap();
let SnapshotResolution::Present(fetched) = store.snapshot(principal).await.unwrap() else {
panic!("newest snapshot must remain present");
};
assert_eq!(fetched.generation, Generation(3));
store.remove_snapshot(principal).await.unwrap();
assert!(matches!(
store.snapshot(principal).await.unwrap(),
SnapshotResolution::Revoked {
generation: Generation(3)
}
));
store
.publish_snapshot(principal, snapshot(2))
.await
.unwrap();
assert!(matches!(
store.snapshot(principal).await.unwrap(),
SnapshotResolution::Revoked {
generation: Generation(3)
}
));
store
.publish_snapshot(principal, snapshot(4))
.await
.unwrap();
let SnapshotResolution::Present(fetched) = store.snapshot(principal).await.unwrap() else {
panic!("newer snapshot must supersede revocation");
};
assert_eq!(fetched.generation, Generation(4));
assert!(
store
.publish_snapshot(principal, snapshot(i64::MAX as u64 + 1))
.await
.is_err()
);
}
#[tokio::test]
async fn staged_limits_round_trip_through_postgres() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(GrantPolicy::default(), 1_000).await else {
return;
};
let principal = Principal(43);
let limits = ResolvedLimits::new(64)
.with_weighted_rate_compatibility_fallback(1_000, 2_000)
.with_request_rate(NonZeroU32::new(10).unwrap(), NonZeroU32::new(20).unwrap())
.with_concurrency(
NonZeroU32::new(4).unwrap(),
Some(NonZeroU32::new(2).unwrap()),
)
.unwrap();
let snapshot = AccountSnapshot::builder(
ACCOUNT,
Generation(1),
AccountStatus::Active,
t(10_000),
PermissionBits::ALL,
limits,
Arc::new(CostTable::builder(CostUnits(1), CostUnits(1)).build()),
)
.build();
store
.publish_snapshot(principal, publishable(Arc::new(snapshot)))
.await
.unwrap();
let SnapshotResolution::Present(fetched) = store.snapshot(principal).await.unwrap() else {
panic!("published snapshot must be present");
};
assert_eq!(fetched.limits, limits);
}
#[test]
fn legacy_limit_wire_defaults_new_dimensions_without_changing_weighted_rate() {
let limits: ResolvedLimits = serde_json::from_value(serde_json::json!({
"max_items_per_request": 64,
"rate_units_per_second": 1_000,
"rate_burst_units": 2_000
}))
.unwrap();
assert!(limits.weighted_rate().is_some());
assert_eq!(limits.request_rate(), None);
assert_eq!(limits.max_concurrent_requests(), None);
assert_eq!(limits.principal_max_concurrent_requests(), None);
}
#[tokio::test]
async fn snapshot_json_preserves_legacy_numbers_and_encodes_high_ids_exactly() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(GrantPolicy::default(), 1_000).await else {
return;
};
let key = KeyId((1u128 << 127) | 2);
let principal = Principal((1u128 << 127) | 3);
let mut digest = [0x77; 32];
digest[..16].copy_from_slice(&principal.0.to_be_bytes());
store
.insert_key(KeyRecord {
key_id: key,
account_id: ACCOUNT,
principal,
digest,
not_after: None,
})
.await
.unwrap();
let snapshot = publishable(Arc::new(
AccountSnapshot::builder(
ACCOUNT,
Generation(1),
AccountStatus::Active,
t(10_000),
PermissionBits::ALL,
ResolvedLimits::new(64).with_weighted_rate(1_000, 1_000),
Arc::new(CostTable::builder(CostUnits(1), CostUnits(1)).build()),
)
.key_id(key)
.build(),
));
store.publish_snapshot(principal, snapshot).await.unwrap();
let pool = corruption_pool().await;
let row = sqlx::query("SELECT snapshot FROM tollgate_snapshots WHERE principal = $1")
.bind(principal.0.to_be_bytes().to_vec())
.fetch_one(&pool)
.await
.unwrap();
let stored: serde_json::Value = sqlx::Row::get(&row, 0);
assert_eq!(stored["account_id"], serde_json::json!(ACCOUNT.0));
assert_eq!(stored["key_id"], serde_json::json!(key.to_string()));
assert!(
stored.get("generation").is_none(),
"the column is the only stored copy of the generation: {stored}"
);
let SnapshotResolution::Present(fetched) = store.snapshot(principal).await.unwrap() else {
panic!("the storage-local representation must decode");
};
assert_eq!(fetched.account_id, ACCOUNT);
assert_eq!(fetched.key_id, Some(key));
}
#[tokio::test]
async fn legacy_invalid_snapshot_is_rejected_on_read() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(GrantPolicy::default(), 1_000).await else {
return;
};
let principal = Principal(47);
let pool = corruption_pool().await;
let snapshot = serde_json::json!({
"account_id": 1,
"key_id": null,
"generation": 1,
"status": "Active",
"valid_until": "2100-01-01T00:00:00Z",
"permissions": 1,
"limits": {
"max_items_per_request": 64,
"rate_units_per_second": 1000,
"rate_burst_units": 113
},
"cost_table": {
"fixed_request": 50,
"minimum_charge": 50,
"weights": [1]
}
});
sqlx::query(
"INSERT INTO tollgate_snapshots (principal, generation, snapshot, deleted)
VALUES ($1, 1, $2, FALSE)",
)
.bind(principal.0.to_be_bytes().to_vec())
.bind(snapshot)
.execute(&pool)
.await
.unwrap();
let error = store.snapshot(principal).await.unwrap_err();
assert!(
error.0.contains("invalid stored snapshot") && error.0.contains("exceeding the burst"),
"unexpected error: {error}"
);
}
#[tokio::test]
async fn the_account_filter_is_answered_by_an_index_not_by_discarding_rows() {
const OTHER_ACCOUNTS: u128 = 40;
const LIVE_PER_ACCOUNT: usize = 25;
const SETTLED: usize = 500;
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 1_000_000).await else {
return;
};
for id in 2..=(OTHER_ACCOUNTS + 1) {
AdminStore::create_account(
&*store,
AccountConfig {
account_id: AccountId(id),
initial_balance: CostUnits(1_000_000),
status: AccountStatus::Active,
capacity_class: CapacityClass::Assured,
},
)
.await
.unwrap();
}
let long = SignedDuration::from_secs(300);
let short = SignedDuration::from_secs(60);
for _ in 0..SETTLED {
store
.acquire(ACCOUNT, CostUnits(1), short, t(0))
.await
.unwrap();
}
for account in std::iter::once(ACCOUNT).chain((2..=(OTHER_ACCOUNTS + 1)).map(AccountId)) {
for _ in 0..LIVE_PER_ACCOUNT {
store
.acquire(account, CostUnits(1), long, t(0))
.await
.unwrap();
}
}
let settled = store.reclaim_expired(t(120)).await.unwrap();
assert_eq!(settled.len(), SETTLED, "only the short-TTL leases are due");
let conservation = store.conservation(ACCOUNT).await.unwrap().unwrap();
assert!(
conservation.holds(),
"conservation violated: {conservation:?}"
);
assert_eq!(
conservation.active_lease_grants,
CostUnits(LIVE_PER_ACCOUNT as u64),
"the measured account holds its own live leases and no one else's"
);
let plan = tollgate_store_postgres::test_support::explain_active_lease_sum(&store, ACCOUNT)
.await
.unwrap();
assert!(
plan.contains("tollgate_leases_account_active"),
"the reconciliation query must reach its index; plan was:\n{plan}"
);
assert!(
!plan.contains("Filter: (account_id"),
"the account predicate is still being applied by discarding rows other \
accounts own, which is the cost GL-12 exists to remove; plan was:\n{plan}"
);
}
#[tokio::test]
async fn the_expiry_sweep_stops_at_its_batch_instead_of_sorting_the_backlog() {
const BACKLOG: usize = 1_200;
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 1_000_000).await else {
return;
};
for i in 0..BACKLOG {
store
.acquire(ACCOUNT, CostUnits(1), TTL, t(i as i64))
.await
.unwrap();
}
let limit = i64::try_from(DEFAULT_RECLAIM_BATCH_LIMIT.get()).unwrap();
let plan = tollgate_store_postgres::test_support::explain_reclaim_due_leases(
&store,
t(BACKLOG as i64 + 1_000),
limit,
)
.await
.unwrap();
assert!(
plan.contains("tollgate_leases_expiry"),
"the sweep must reach its partial expiry index; plan was:\n{plan}"
);
assert!(
!plan.contains("Sort"),
"the batch is still being taken from a sort of the whole backlog, which \
is the cost GL-65 exists to remove; plan was:\n{plan}"
);
}
#[tokio::test]
async fn the_rollover_sweep_reaches_its_index_instead_of_sorting_the_due_set() {
const ACCOUNTS: u128 = 400;
const DUE: u128 = 5;
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 0).await else {
return;
};
for id in 2..=ACCOUNTS {
AdminStore::create_account(
&*store,
AccountConfig {
account_id: AccountId(id),
initial_balance: CostUnits::ZERO,
status: AccountStatus::Active,
capacity_class: CapacityClass::Assured,
},
)
.await
.unwrap();
}
for id in 1..=ACCOUNTS {
AdminStore::set_budget_schedule(&*store, AccountId(id), Some(monthly(100)))
.await
.unwrap();
}
let pool = corruption_pool().await;
sqlx::query("UPDATE tollgate_accounts SET period_start_us = $1")
.bind(FEB * 1_000_000)
.execute(&pool)
.await
.unwrap();
for id in 1..=DUE {
sqlx::query("UPDATE tollgate_accounts SET period_start_us = $2 WHERE account_id = $1")
.bind(id.to_be_bytes().to_vec())
.bind(JAN * 1_000_000)
.execute(&pool)
.await
.unwrap();
}
let plan = tollgate_store_postgres::test_support::explain_due_periods(
&store,
"monthly",
FEB * 1_000_000,
64,
)
.await
.unwrap();
assert!(
plan.contains("tollgate_accounts_due_rollover"),
"the rollover sweep must reach its partial index; plan was:\n{plan}"
);
assert!(
!plan.contains("Sort"),
"the bounded page is still taken from a sort of every due account, \
which is the cost GL-65 exists to remove; plan was:\n{plan}"
);
}
#[tokio::test]
async fn a_bounded_rollover_page_crosses_the_oldest_boundaries_first() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 0).await else {
return;
};
for id in 2..=4u128 {
AdminStore::create_account(
&*store,
AccountConfig {
account_id: AccountId(id),
initial_balance: CostUnits::ZERO,
status: AccountStatus::Active,
capacity_class: CapacityClass::Assured,
},
)
.await
.unwrap();
}
for id in 1..=4u128 {
AdminStore::set_budget_schedule(&*store, AccountId(id), Some(monthly(100)))
.await
.unwrap();
}
let pool = corruption_pool().await;
for (id, months_behind) in [(1u128, 1i64), (2, 2), (3, 3), (4, 4)] {
sqlx::query("UPDATE tollgate_accounts SET period_start_us = $2 WHERE account_id = $1")
.bind(id.to_be_bytes().to_vec())
.bind((JAN - months_behind * 86_400 * 31) * 1_000_000)
.execute(&pool)
.await
.unwrap();
}
let batch = AdminStore::roll_due_periods(&*store, t(FEB), NonZeroUsize::new(2).unwrap())
.await
.unwrap();
assert!(
batch.is_saturated(),
"the page is full, so ordering decides it"
);
let rolled: Vec<u128> = batch
.rolled()
.iter()
.map(|account| account.account_id.0)
.collect();
assert_eq!(
rolled,
vec![4, 3],
"the two oldest boundaries cross first; account order would have given [1, 2]"
);
}
async fn corruption_pool() -> sqlx::PgPool {
let url = std::env::var("TOLLGATE_PG_URL").expect("caller already gated on TOLLGATE_PG_URL");
sqlx::PgPool::connect(&url)
.await
.unwrap_or_else(|e| panic!("postgres unreachable at {}: {e}", redact_url(&url)))
}
fn account_bytes() -> Vec<u8> {
ACCOUNT.0.to_be_bytes().to_vec()
}
async fn suspend_checks(
pool: &sqlx::PgPool,
table: &str,
column: &str,
keep: Option<&str>,
) -> Vec<(String, String)> {
let rows: Vec<(String, String)> = sqlx::query_as(
"SELECT conname, pg_get_constraintdef(oid)
FROM pg_constraint
WHERE conrelid = $1::regclass AND contype = 'c'
AND pg_get_constraintdef(oid) LIKE '%' || $2 || '%'",
)
.bind(table)
.bind(column)
.fetch_all(pool)
.await
.unwrap();
let dropped: Vec<_> = rows
.into_iter()
.filter(|(name, _)| Some(name.as_str()) != keep)
.collect();
for (name, _) in &dropped {
sqlx::raw_sql(&format!("ALTER TABLE {table} DROP CONSTRAINT {name}"))
.execute(pool)
.await
.unwrap();
}
dropped
}
async fn restore_checks(pool: &sqlx::PgPool, table: &str, dropped: Vec<(String, String)>) {
for (name, definition) in dropped {
sqlx::raw_sql(&format!(
"ALTER TABLE {table} ADD CONSTRAINT {name} {definition} NOT VALID"
))
.execute(pool)
.await
.unwrap();
}
}
async fn plant_value(
pool: &sqlx::PgPool,
table: &str,
column: &str,
id_column: &str,
id: Vec<u8>,
value: i64,
) {
let dropped = suspend_checks(pool, table, column, None).await;
sqlx::query(&format!(
"UPDATE {table} SET {column} = $1 WHERE {id_column} = $2"
))
.bind(value)
.bind(id)
.execute(pool)
.await
.unwrap();
restore_checks(pool, table, dropped).await;
}
async fn set_account_column(pool: &sqlx::PgPool, column: &str, value: i64) {
plant_value(
pool,
"tollgate_accounts",
column,
"account_id",
account_bytes(),
value,
)
.await;
}
async fn set_lease_column(pool: &sqlx::PgPool, lease: u128, column: &str, value: i64) {
plant_value(
pool,
"tollgate_leases",
column,
"lease_id",
lease.to_be_bytes().to_vec(),
value,
)
.await;
}
#[tokio::test]
async fn negative_account_column_fails_conservation_read() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 100).await else {
return;
};
let pool = corruption_pool().await;
for (column, restore) in [
("deposited", 100),
("balance", 100),
("usage_recorded", 0),
("settlement_loss", 0),
("overage_recorded", 0),
] {
set_account_column(&pool, column, -1).await;
let err = store.conservation(ACCOUNT).await.unwrap_err();
assert!(
err.0.contains(column) && err.0.contains("negative"),
"conservation over negative {column} must name it, got: {err}"
);
set_account_column(&pool, column, restore).await;
}
assert_conserved(&store).await;
}
#[tokio::test]
async fn negative_account_column_fails_balance_and_usage_reads() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 100).await else {
return;
};
let pool = corruption_pool().await;
set_account_column(&pool, "balance", -1).await;
let err = store.balance(ACCOUNT).await.unwrap_err();
assert!(err.0.contains("account balance"), "got: {err}");
set_account_column(&pool, "balance", 100).await;
set_account_column(&pool, "usage_recorded", -1).await;
let err = store.usage_recorded(ACCOUNT).await.unwrap_err();
assert!(err.0.contains("account usage_recorded"), "got: {err}");
}
#[tokio::test]
async fn negative_lease_sum_fails_conservation_read() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 1_000).await else {
return;
};
let lease = store
.acquire(ACCOUNT, CostUnits(500), TTL, t(0))
.await
.unwrap()
.grant;
let pool = corruption_pool().await;
set_lease_column(&pool, lease.lease_id.0, "granted", -1).await;
let err = store.conservation(ACCOUNT).await.unwrap_err();
assert!(err.0.contains("active lease grants"), "got: {err}");
set_lease_column(&pool, lease.lease_id.0, "granted", 500).await;
set_lease_column(&pool, lease.lease_id.0, "used", -1).await;
let err = store.conservation(ACCOUNT).await.unwrap_err();
assert!(err.0.contains("active lease usage"), "got: {err}");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn a_reconciliation_read_never_observes_a_torn_ledger() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 100_000).await else {
return;
};
let lease = store
.acquire(ACCOUNT, CostUnits(50_000), TTL, t(0))
.await
.unwrap()
.grant;
const WRITES: u32 = 200;
let written = Arc::new(std::sync::atomic::AtomicU32::new(0));
let writers = tokio::spawn({
let store = Arc::clone(&store);
let written = Arc::clone(&written);
async move {
for request in 0..u128::from(WRITES) {
UsageSink::ingest(&*store, &[usage(&lease, request, 1, 0)], t(1))
.await
.unwrap();
written.fetch_add(1, std::sync::atomic::Ordering::Release);
}
}
});
let reader = tokio::spawn({
let store = Arc::clone(&store);
let written = Arc::clone(&written);
async move {
let mut overlapping = 0u32;
let mut reads = 0u32;
loop {
let before = written.load(std::sync::atomic::Ordering::Acquire);
let observed = store
.conservation(ACCOUNT)
.await
.expect("a reconciliation read must not fail on a healthy ledger")
.expect("the account exists");
assert!(
observed.holds(),
"reconciliation reported corruption on a correct ledger: {observed:?}"
);
reads += 1;
let after = written.load(std::sync::atomic::Ordering::Acquire);
if before > 0 && after < WRITES {
overlapping += 1;
}
if after >= WRITES || reads > 10_000 {
break;
}
}
(reads, overlapping)
}
});
writers.await.unwrap();
let (reads, overlapping) = reader.await.unwrap();
assert!(reads > 0, "the reader must have run at all");
assert!(
overlapping > 0,
"every read landed outside the write sequence, so nothing was raced: \
{reads} reads, {overlapping} overlapping"
);
assert_conserved(&store).await;
}
#[tokio::test]
async fn usage_below_the_live_lease_sum_fails_the_conservation_read() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 1_000).await else {
return;
};
let lease = store
.acquire(ACCOUNT, CostUnits(500), TTL, t(0))
.await
.unwrap()
.grant;
let pool = corruption_pool().await;
set_lease_column(&pool, lease.lease_id.0, "used", 100).await;
let err = store
.conservation(ACCOUNT)
.await
.expect_err("a read that cannot subtract must report, not panic");
assert!(
err.0.contains("exceeds recorded usage"),
"the error must name the discrepancy it found; got: {err}"
);
}
#[tokio::test]
async fn acquire_surfaces_negative_stored_balance() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 100).await else {
return;
};
let pool = corruption_pool().await;
set_account_column(&pool, "balance", -5).await;
let err = store
.acquire(ACCOUNT, CostUnits(10), TTL, t(0))
.await
.unwrap_err();
assert!(
matches!(&err, AllocateError::Storage(e) if e.0.contains("account balance")),
"corrupt balance must surface as storage corruption, got: {err}"
);
}
#[tokio::test]
async fn reclaim_refuses_a_negative_remainder() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 1_000).await else {
return;
};
let lease = store
.acquire(ACCOUNT, CostUnits(500), TTL, t(0))
.await
.unwrap()
.grant;
let pool = corruption_pool().await;
sqlx::query("UPDATE tollgate_leases SET used = 600 WHERE lease_id = $1")
.bind(lease.lease_id.0.to_be_bytes().to_vec())
.execute(&pool)
.await
.unwrap();
let err = store.reclaim_expired(t(120)).await.unwrap_err();
assert!(err.0.contains("reclaim remainder"), "got: {err}");
assert_eq!(store.balance(ACCOUNT).await.unwrap(), CostUnits(500));
sqlx::query("UPDATE tollgate_leases SET used = 0 WHERE lease_id = $1")
.bind(lease.lease_id.0.to_be_bytes().to_vec())
.execute(&pool)
.await
.unwrap();
assert_eq!(
store
.conservation(ACCOUNT)
.await
.unwrap()
.unwrap()
.settlement_loss,
CostUnits::ZERO
);
}
#[tokio::test]
async fn straggler_exceeding_recorded_loss_fails_ingest() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 1_000).await else {
return;
};
let lease = store
.acquire(ACCOUNT, CostUnits(500), TTL, t(0))
.await
.unwrap()
.grant;
store
.release(lease.lease_id, lease.fencing_token, CostUnits(470), t(10))
.await
.unwrap();
let pool = corruption_pool().await;
set_account_column(&pool, "settlement_loss", 10).await;
let err = store
.ingest(&[usage(&lease, 1, 30, 5)], t(11))
.await
.unwrap_err();
assert!(
err.to_string().contains("settlement_loss underflow"),
"got: {err}"
);
assert_eq!(
store.usage_recorded(ACCOUNT).await.unwrap(),
CostUnits::ZERO
);
}
#[tokio::test]
async fn acquire_surfaces_nonpositive_stored_fence() {
for invalid in [-1, 0] {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 100).await else {
return;
};
let pool = corruption_pool().await;
set_account_column(&pool, "next_fence", invalid).await;
let err = store
.acquire(ACCOUNT, CostUnits(10), TTL, t(0))
.await
.unwrap_err();
assert!(
matches!(&err, AllocateError::Storage(e) if e.0.contains("fencing token")),
"corrupt fence must surface as storage corruption, got: {err}"
);
}
}
#[tokio::test]
async fn release_and_ingest_surface_nonpositive_stored_fence() {
for invalid in [-1, 0] {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 1_000).await else {
return;
};
let lease = store
.acquire(ACCOUNT, CostUnits(500), TTL, t(0))
.await
.unwrap()
.grant;
let pool = corruption_pool().await;
set_lease_column(&pool, lease.lease_id.0, "fencing_token", invalid).await;
let err = store
.release(lease.lease_id, lease.fencing_token, CostUnits(500), t(1))
.await
.unwrap_err();
assert!(
matches!(&err, AllocateError::Storage(e) if e.0.contains("fencing token")),
"got: {err}"
);
let err = store
.ingest(&[usage(&lease, 1, 10, 2)], t(2))
.await
.unwrap_err();
assert!(err.to_string().contains("fencing token"), "got: {err}");
}
}
#[tokio::test]
async fn checked_ledger_columns_reject_negative_writes() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 100).await else {
return;
};
let lease = store
.acquire(ACCOUNT, CostUnits(10), TTL, t(0))
.await
.unwrap()
.grant;
let event = usage(&lease, 1, 0, 0);
store.ingest(&[event], t(0)).await.unwrap();
let pool = corruption_pool().await;
let account_columns = [
"deposited",
"balance",
"usage_recorded",
"settlement_loss",
"overage_recorded",
"next_fence",
"allowance_balance",
"expired",
"budget_allowance",
]
.map(|c| ("tollgate_accounts", c, "account_id", account_bytes()));
let lease_columns = [
"fencing_token",
"granted",
"used",
"credited",
"from_allowance",
]
.map(|c| {
(
"tollgate_leases",
c,
"lease_id",
lease.lease_id.0.to_be_bytes().to_vec(),
)
});
let usage_columns = ["units", "fencing_token"].map(|column| {
(
"tollgate_usage_events",
column,
"request_id",
event.request_id.0.to_be_bytes().to_vec(),
)
});
for (table, column, id_column, id) in account_columns
.into_iter()
.chain(lease_columns)
.chain(usage_columns)
{
let suffix = if matches!(column, "next_fence" | "fencing_token") {
"positive"
} else {
"nonneg"
};
let nonneg = format!("{table}_{column}_{suffix}");
let dropped = suspend_checks(&pool, table, column, Some(&nonneg)).await;
let result = sqlx::query(&format!(
"UPDATE {table} SET {column} = -1 WHERE {id_column} = $1"
))
.bind(id)
.execute(&pool)
.await;
restore_checks(&pool, table, dropped).await;
let err = result.expect_err("a negative write must be refused");
assert!(
err.to_string().contains(&nonneg),
"negative {table}.{column} write must violate its CHECK constraint, got: {err}"
);
}
}
#[tokio::test]
async fn zero_fences_are_refused_in_every_persisted_capability() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 100).await else {
return;
};
let lease = store
.acquire(ACCOUNT, CostUnits(10), TTL, t(0))
.await
.unwrap()
.grant;
assert_eq!(lease.fencing_token, FencingToken(1));
let event = usage(&lease, 1, 0, 0);
store
.ingest(&[event, overage_usage(ACCOUNT, 2, 0, 0)], t(0))
.await
.unwrap();
let pool = corruption_pool().await;
for (table, column, id_column, id) in [
(
"tollgate_accounts",
"next_fence",
"account_id",
account_bytes(),
),
(
"tollgate_leases",
"fencing_token",
"lease_id",
lease.lease_id.0.to_be_bytes().to_vec(),
),
(
"tollgate_usage_events",
"fencing_token",
"request_id",
event.request_id.0.to_be_bytes().to_vec(),
),
] {
let error = sqlx::query(&format!(
"UPDATE {table} SET {column} = 0 WHERE {id_column} = $1"
))
.bind(id)
.execute(&pool)
.await
.unwrap_err();
assert_eq!(
error.as_database_error().unwrap().constraint(),
Some(format!("{table}_{column}_positive").as_str())
);
}
let second = store
.acquire(ACCOUNT, CostUnits(10), TTL, t(1))
.await
.unwrap()
.grant;
assert_eq!(second.fencing_token, FencingToken(2));
assert_conserved(&store).await;
}
#[tokio::test]
async fn the_allowance_split_cannot_exceed_what_it_is_part_of() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 100).await else {
return;
};
let lease = store
.acquire(ACCOUNT, CostUnits(10), TTL, t(0))
.await
.unwrap()
.grant;
let pool = corruption_pool().await;
for (table, column, id_column, id, value, constraint) in [
(
"tollgate_accounts",
"allowance_balance",
"account_id",
account_bytes(),
1_000,
"tollgate_accounts_allowance_within_balance",
),
(
"tollgate_leases",
"from_allowance",
"lease_id",
lease.lease_id.0.to_be_bytes().to_vec(),
1_000,
"tollgate_leases_from_allowance_within_grant",
),
] {
let err = sqlx::query(&format!(
"UPDATE {table} SET {column} = $1 WHERE {id_column} = $2"
))
.bind(value)
.bind(id)
.execute(&pool)
.await
.unwrap_err();
assert!(
err.to_string().contains(constraint),
"{table}.{column} beyond its whole must violate {constraint}, got: {err}"
);
}
}
fn account_snapshot(account: AccountId, generation: u64, status: AccountStatus) -> AccountSnapshot {
AccountSnapshot::builder(
account,
Generation(generation),
status,
t(10_000),
PermissionBits::ALL,
ResolvedLimits::new(64).with_weighted_rate(1_000, 1_000),
Arc::new(CostTable::builder(CostUnits(1), CostUnits(1)).build()),
)
.build()
}
async fn status_of(store: &PostgresStore, principal: Principal) -> (AccountStatus, Generation) {
let SnapshotResolution::Present(snapshot) = store.snapshot(principal).await.unwrap() else {
panic!("principal {principal} must be present");
};
(snapshot.status, snapshot.generation)
}
async fn publish(store: &PostgresStore, principal: Principal, snapshot: AccountSnapshot) {
AdminStore::publish_snapshot(store, principal, publishable(Arc::new(snapshot)))
.await
.unwrap();
}
#[tokio::test]
async fn suspending_an_account_stops_leases_and_republishes_its_snapshots() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(GrantPolicy::default(), 1_000).await else {
return;
};
let first = Principal(10);
let second = Principal(11);
for principal in [first, second] {
publish(
&store,
principal,
account_snapshot(ACCOUNT, 3, AccountStatus::Active),
)
.await;
}
let change = AdminStore::set_account_status(&*store, ACCOUNT, AccountStatus::Suspended)
.await
.unwrap();
assert_eq!(
(change.outcome.republished, change.outcome.unreadable),
(2, 0),
"the reported blast radius is both live principals"
);
assert_eq!(
store
.acquire(ACCOUNT, CostUnits(100), TTL, t(0))
.await
.unwrap_err(),
AllocateError::AccountInactive,
"the ledger half"
);
for principal in [first, second] {
assert_eq!(
status_of(&store, principal).await,
(AccountStatus::Suspended, Generation(4)),
"the snapshot half, for {principal}"
);
}
}
#[tokio::test]
async fn suspension_republishes_only_the_suspended_accounts_snapshots() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(GrantPolicy::default(), 1_000).await else {
return;
};
let other_account = AccountId(2);
AdminStore::create_account(
&*store,
AccountConfig {
account_id: other_account,
initial_balance: CostUnits(1_000),
status: AccountStatus::Active,
capacity_class: CapacityClass::Assured,
},
)
.await
.unwrap();
let mine = Principal(10);
let theirs = Principal(20);
publish(
&store,
mine,
account_snapshot(ACCOUNT, 3, AccountStatus::Active),
)
.await;
publish(
&store,
theirs,
account_snapshot(other_account, 7, AccountStatus::Active),
)
.await;
AdminStore::set_account_status(&*store, ACCOUNT, AccountStatus::Suspended)
.await
.unwrap();
assert_eq!(
status_of(&store, mine).await,
(AccountStatus::Suspended, Generation(4))
);
assert_eq!(
status_of(&store, theirs).await,
(AccountStatus::Active, Generation(7)),
"another account's snapshot is untouched, generation included"
);
store
.acquire(other_account, CostUnits(100), TTL, t(0))
.await
.expect("and it can still lease");
}
#[tokio::test]
async fn suspending_an_account_does_not_resurrect_revoked_principals() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(GrantPolicy::default(), 1_000).await else {
return;
};
let live = Principal(10);
let revoked = Principal(11);
for principal in [live, revoked] {
publish(
&store,
principal,
account_snapshot(ACCOUNT, 3, AccountStatus::Active),
)
.await;
}
AdminStore::remove_snapshot(&*store, revoked).await.unwrap();
AdminStore::set_account_status(&*store, ACCOUNT, AccountStatus::Suspended)
.await
.unwrap();
assert_eq!(
status_of(&store, live).await,
(AccountStatus::Suspended, Generation(4))
);
assert!(
matches!(
store.snapshot(revoked).await.unwrap(),
SnapshotResolution::Revoked {
generation: Generation(3)
}
),
"a tombstone stays revoked, at its own generation"
);
}
#[tokio::test]
async fn reactivating_an_account_restores_admission_and_bumps_generations() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(GrantPolicy::default(), 1_000).await else {
return;
};
let principal = Principal(10);
publish(
&store,
principal,
account_snapshot(ACCOUNT, 3, AccountStatus::Active),
)
.await;
AdminStore::set_account_status(&*store, ACCOUNT, AccountStatus::Suspended)
.await
.unwrap();
AdminStore::set_account_status(&*store, ACCOUNT, AccountStatus::Active)
.await
.unwrap();
assert_eq!(
status_of(&store, principal).await,
(AccountStatus::Active, Generation(5)),
"each transition is its own generation; they never move backward"
);
store
.acquire(ACCOUNT, CostUnits(100), TTL, t(0))
.await
.expect("reactivation restores leasing");
}
#[tokio::test]
async fn a_closed_account_cannot_be_reactivated() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(GrantPolicy::default(), 1_000).await else {
return;
};
let principal = Principal(10);
publish(
&store,
principal,
account_snapshot(ACCOUNT, 3, AccountStatus::Active),
)
.await;
AdminStore::set_account_status(&*store, ACCOUNT, AccountStatus::Closed)
.await
.unwrap();
let after_close = status_of(&store, principal).await;
for status in [AccountStatus::Active, AccountStatus::Suspended] {
assert_eq!(
AdminStore::set_account_status(&*store, ACCOUNT, status)
.await
.unwrap_err(),
SetStatusError::AccountClosed,
"a closed account cannot become {status:?}"
);
assert_eq!(
status_of(&store, principal).await,
after_close,
"and the refusal moved nothing"
);
assert_eq!(
store
.acquire(ACCOUNT, CostUnits(100), TTL, t(0))
.await
.unwrap_err(),
AllocateError::AccountInactive,
"including the ledger"
);
}
AdminStore::set_account_status(&*store, ACCOUNT, AccountStatus::Closed)
.await
.expect("Closed -> Closed is a no-op, not a refusal");
}
#[tokio::test]
async fn a_capacity_class_change_republishes_every_live_snapshot() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 1_000).await else {
return;
};
let first = Principal(10);
let second = Principal(11);
for principal in [first, second] {
AdminStore::publish_snapshot(
&*store,
principal,
publishable(Arc::new(account_snapshot(
ACCOUNT,
3,
AccountStatus::Active,
))),
)
.await
.unwrap();
}
let change = AdminStore::set_capacity_class(&*store, ACCOUNT, CapacityClass::BestEffort)
.await
.unwrap();
assert_eq!(
(change.outcome.republished, change.outcome.unreadable),
(2, 0)
);
for principal in [first, second] {
let SnapshotResolution::Present(snapshot) = store.snapshot(principal).await.unwrap() else {
panic!("the snapshot must be present");
};
assert_eq!(snapshot.capacity_class, CapacityClass::BestEffort);
assert_eq!(snapshot.generation, Generation(4));
assert_eq!(snapshot.status, AccountStatus::Active);
}
let pool = corruption_pool().await;
let stored: serde_json::Value =
sqlx::query_scalar("SELECT snapshot FROM tollgate_snapshots WHERE principal = $1")
.bind(first.0.to_be_bytes().to_vec())
.fetch_one(&pool)
.await
.unwrap();
assert_eq!(
stored["capacity_class"],
serde_json::Value::String("BestEffort".to_owned())
);
}
#[tokio::test]
async fn repeating_a_capacity_class_change_publishes_nothing_new() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 1_000).await else {
return;
};
let principal = Principal(10);
AdminStore::publish_snapshot(
&*store,
principal,
publishable(Arc::new(account_snapshot(
ACCOUNT,
3,
AccountStatus::Active,
))),
)
.await
.unwrap();
AdminStore::set_capacity_class(&*store, ACCOUNT, CapacityClass::BestEffort)
.await
.unwrap();
let change = AdminStore::set_capacity_class(&*store, ACCOUNT, CapacityClass::BestEffort)
.await
.unwrap();
assert_eq!(
(change.outcome.republished, change.outcome.unreadable),
(0, 0)
);
let SnapshotResolution::Present(snapshot) = store.snapshot(principal).await.unwrap() else {
panic!("the snapshot must be present");
};
assert_eq!(snapshot.generation, Generation(4));
}
#[tokio::test]
async fn a_capacity_class_change_does_not_resurrect_revoked_principals() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 1_000).await else {
return;
};
let live = Principal(10);
let revoked = Principal(11);
for principal in [live, revoked] {
AdminStore::publish_snapshot(
&*store,
principal,
publishable(Arc::new(account_snapshot(
ACCOUNT,
3,
AccountStatus::Active,
))),
)
.await
.unwrap();
}
AdminStore::remove_snapshot(&*store, revoked).await.unwrap();
let change = AdminStore::set_capacity_class(&*store, ACCOUNT, CapacityClass::BestEffort)
.await
.unwrap();
assert_eq!(change.outcome.republished, 1);
assert!(matches!(
store.snapshot(revoked).await.unwrap(),
SnapshotResolution::Revoked { .. }
));
}
#[tokio::test]
async fn a_publish_contradicting_the_ledger_class_is_refused() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 1_000).await else {
return;
};
let mut snapshot = account_snapshot(ACCOUNT, 3, AccountStatus::Active);
snapshot.capacity_class = CapacityClass::BestEffort;
let error =
AdminStore::publish_snapshot(&*store, Principal(10), publishable(Arc::new(snapshot)))
.await
.expect_err("a publish may not change the account's class");
assert!(matches!(
error,
PublishSnapshotError::CapacityClassMismatch {
ledger: CapacityClass::Assured,
submitted: CapacityClass::BestEffort,
}
));
}
#[tokio::test]
async fn a_publish_matching_a_reclassified_ledger_is_accepted() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 1_000).await else {
return;
};
let principal = Principal(10);
AdminStore::set_capacity_class(&*store, ACCOUNT, CapacityClass::BestEffort)
.await
.unwrap();
let mut matching = account_snapshot(ACCOUNT, 3, AccountStatus::Active);
matching.capacity_class = CapacityClass::BestEffort;
AdminStore::publish_snapshot(&*store, principal, publishable(Arc::new(matching)))
.await
.expect("a publish carrying the ledger's current class is accepted");
let error = AdminStore::publish_snapshot(
&*store,
principal,
publishable(Arc::new(account_snapshot(
ACCOUNT,
4,
AccountStatus::Active,
))),
)
.await
.expect_err("`Assured` now contradicts the ledger");
assert!(matches!(
error,
PublishSnapshotError::CapacityClassMismatch {
ledger: CapacityClass::BestEffort,
submitted: CapacityClass::Assured,
}
));
}
#[tokio::test]
async fn an_account_can_be_born_best_effort() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(GrantPolicy::default(), 1_000).await else {
return;
};
let born = AccountId(77);
AdminStore::create_account(
&*store,
AccountConfig {
account_id: born,
initial_balance: CostUnits(1_000),
status: AccountStatus::Active,
capacity_class: CapacityClass::BestEffort,
},
)
.await
.unwrap();
let mut snapshot = account_snapshot(born, 3, AccountStatus::Active);
snapshot.capacity_class = CapacityClass::BestEffort;
AdminStore::publish_snapshot(&*store, Principal(20), publishable(Arc::new(snapshot)))
.await
.expect("the account was created best-effort, so a matching publish is accepted");
let error = AdminStore::publish_snapshot(
&*store,
Principal(21),
publishable(Arc::new(account_snapshot(born, 3, AccountStatus::Active))),
)
.await
.expect_err("`Assured` contradicts the class the account was born with");
assert!(matches!(
error,
PublishSnapshotError::CapacityClassMismatch {
ledger: CapacityClass::BestEffort,
submitted: CapacityClass::Assured,
}
));
}
#[tokio::test]
async fn a_capacity_class_change_that_cannot_republish_moves_neither_record() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(GrantPolicy::default(), 1_000).await else {
return;
};
let principal = Principal(10);
let ceiling = i64::MAX as u64;
publish(
&store,
principal,
account_snapshot(ACCOUNT, ceiling, AccountStatus::Active),
)
.await;
let error = AdminStore::set_capacity_class(&*store, ACCOUNT, CapacityClass::BestEffort)
.await
.unwrap_err();
assert!(
matches!(error, SetStatusError::Storage(_)),
"an unrepublishable snapshot surfaces, rather than being skipped: {error:?}"
);
let SnapshotResolution::Present(snapshot) = store.snapshot(principal).await.unwrap() else {
panic!("the snapshot must still be present");
};
assert_eq!(
(snapshot.capacity_class, snapshot.generation),
(CapacityClass::Assured, Generation(ceiling)),
"the snapshot half did not move"
);
publish(
&store,
Principal(11),
account_snapshot(ACCOUNT, 3, AccountStatus::Active),
)
.await;
}
#[tokio::test]
async fn a_capacity_class_change_reports_rows_it_changed_but_could_not_push() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(GrantPolicy::default(), 1_000).await else {
return;
};
let readable = Principal(62);
let undecodable = Principal(63);
publish(
&store,
readable,
account_snapshot(ACCOUNT, 3, AccountStatus::Active),
)
.await;
let pool = corruption_pool().await;
sqlx::query(
"INSERT INTO tollgate_snapshots (principal, generation, snapshot, deleted)
VALUES ($1, 3, $2, FALSE)",
)
.bind(undecodable.0.to_be_bytes().to_vec())
.bind(serde_json::json!({
"account_id": ACCOUNT.0,
"key_id": null,
"generation": 3,
"status": "Active",
"valid_until": "2100-01-01T00:00:00Z",
"permissions": 1,
"limits": {
"max_items_per_request": 64,
"rate_units_per_second": 1000,
"rate_burst_units": 113
},
"cost_table": { "fixed_request": 50, "minimum_charge": 50, "weights": [1] }
}))
.execute(&pool)
.await
.unwrap();
let mut updates = store.subscribe();
let change = AdminStore::set_capacity_class(&*store, ACCOUNT, CapacityClass::BestEffort)
.await
.unwrap();
assert_eq!(
(change.outcome.republished, change.outcome.unreadable),
(2, 1),
"both rows changed durably; one of them could not be pushed"
);
let mut pushed = Vec::new();
while let Ok(push) = updates.try_recv() {
pushed.push(push.principal);
}
assert_eq!(
pushed,
vec![readable],
"only the decodable principal is pushed; the other waits for a refresh"
);
let SnapshotResolution::Present(snapshot) = store.snapshot(readable).await.unwrap() else {
panic!("the readable snapshot must still be present");
};
assert_eq!(snapshot.capacity_class, CapacityClass::BestEffort);
assert_eq!(snapshot.generation, Generation(4));
assert!(
store.snapshot(undecodable).await.is_err(),
"the corrupt row is still refused on read, as it was before the change"
);
}
#[tokio::test]
async fn reclassifying_an_unknown_or_closed_account_is_refused() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 1_000).await else {
return;
};
assert!(matches!(
AdminStore::set_capacity_class(&*store, AccountId(u128::MAX), CapacityClass::BestEffort)
.await,
Err(SetStatusError::UnknownAccount)
));
AdminStore::set_account_status(&*store, ACCOUNT, AccountStatus::Closed)
.await
.unwrap();
assert!(matches!(
AdminStore::set_capacity_class(&*store, ACCOUNT, CapacityClass::BestEffort).await,
Err(SetStatusError::AccountClosed)
));
}
#[tokio::test]
async fn an_unrecognized_stored_capacity_class_is_refused_on_publish() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 1_000).await else {
return;
};
let pool = corruption_pool().await;
let dropped = suspend_checks(&pool, "tollgate_accounts", "capacity_class", None).await;
sqlx::query("UPDATE tollgate_accounts SET capacity_class = $2 WHERE account_id = $1")
.bind(account_bytes())
.bind("Preferred")
.execute(&pool)
.await
.unwrap();
let error = AdminStore::publish_snapshot(
&*store,
Principal(10),
publishable(Arc::new(account_snapshot(
ACCOUNT,
3,
AccountStatus::Active,
))),
)
.await
.expect_err("an unrecognized stored class must not decode");
assert!(
format!("{error}").contains("unrecognized capacity class"),
"expected a decode refusal, got {error}"
);
restore_checks(&pool, "tollgate_accounts", dropped).await;
}
#[tokio::test]
async fn repeating_a_status_change_publishes_nothing_new() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(GrantPolicy::default(), 1_000).await else {
return;
};
let principal = Principal(10);
publish(
&store,
principal,
account_snapshot(ACCOUNT, 3, AccountStatus::Active),
)
.await;
AdminStore::set_account_status(&*store, ACCOUNT, AccountStatus::Suspended)
.await
.unwrap();
let mut updates = store.subscribe();
let change = AdminStore::set_account_status(&*store, ACCOUNT, AccountStatus::Suspended)
.await
.unwrap();
assert_eq!(
change.outcome.republished, 0,
"a repeat reports zero, which is how an operator sees it changed nothing"
);
assert_eq!(
status_of(&store, principal).await,
(AccountStatus::Suspended, Generation(4)),
"already at the target status, so not rewritten"
);
assert!(
updates.try_recv().is_err(),
"and nothing was pushed for a change that did not happen"
);
}
#[tokio::test]
async fn a_status_change_pushes_every_republished_principal() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(GrantPolicy::default(), 1_000).await else {
return;
};
let first = Principal(10);
let second = Principal(11);
let revoked = Principal(12);
for principal in [first, second, revoked] {
publish(
&store,
principal,
account_snapshot(ACCOUNT, 3, AccountStatus::Active),
)
.await;
}
AdminStore::remove_snapshot(&*store, revoked).await.unwrap();
let mut updates = store.subscribe();
let change = AdminStore::set_account_status(&*store, ACCOUNT, AccountStatus::Suspended)
.await
.unwrap();
assert_eq!(
change.outcome.republished, 2,
"the tombstone is not republished, so it is not counted either"
);
let mut pushed = Vec::new();
while let Ok(push) = updates.try_recv() {
let SnapshotResolution::Present(snapshot) = push.resolution else {
panic!("a status change republishes; it never revokes");
};
assert_eq!(snapshot.status, AccountStatus::Suspended);
pushed.push(push.principal);
}
assert_eq!(
pushed,
vec![first, second],
"one push per live principal, ordered, and none for the tombstone"
);
}
#[tokio::test]
async fn publishing_a_snapshot_that_contradicts_the_ledger_is_refused() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(GrantPolicy::default(), 1_000).await else {
return;
};
let principal = Principal(10);
AdminStore::set_account_status(&*store, ACCOUNT, AccountStatus::Suspended)
.await
.unwrap();
assert_eq!(
AdminStore::publish_snapshot(
&*store,
principal,
publishable(Arc::new(account_snapshot(
ACCOUNT,
3,
AccountStatus::Active
))),
)
.await
.unwrap_err(),
PublishSnapshotError::StatusMismatch {
ledger: AccountStatus::Suspended,
submitted: AccountStatus::Active,
},
);
assert!(
matches!(
store.snapshot(principal).await.unwrap(),
SnapshotResolution::Unknown
),
"and the refusal wrote nothing"
);
publish(
&store,
principal,
account_snapshot(ACCOUNT, 3, AccountStatus::Suspended),
)
.await;
publish(
&store,
Principal(99),
account_snapshot(AccountId(4_242), 1, AccountStatus::Active),
)
.await;
}
#[tokio::test]
async fn the_account_column_is_derived_for_both_stored_id_spellings() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(GrantPolicy::default(), 1_000).await else {
return;
};
let high_account = AccountId((1u128 << 127) | 5);
let legacy_account = AccountId(u128::from(u64::MAX));
for account in [high_account, legacy_account] {
AdminStore::create_account(
&*store,
AccountConfig {
account_id: account,
initial_balance: CostUnits(1_000),
status: AccountStatus::Active,
capacity_class: CapacityClass::Assured,
},
)
.await
.unwrap();
}
let high_principal = Principal(30);
let legacy_principal = Principal(31);
publish(
&store,
high_principal,
account_snapshot(high_account, 3, AccountStatus::Active),
)
.await;
publish(
&store,
legacy_principal,
account_snapshot(legacy_account, 3, AccountStatus::Active),
)
.await;
for (account, principal) in [
(high_account, high_principal),
(legacy_account, legacy_principal),
] {
AdminStore::set_account_status(&*store, account, AccountStatus::Suspended)
.await
.unwrap();
assert_eq!(
status_of(&store, principal).await,
(AccountStatus::Suspended, Generation(4)),
"the status change must reach {account}, whichever way its id is spelled"
);
}
}
#[tokio::test]
async fn every_account_status_variant_round_trips_the_status_column() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(GrantPolicy::default(), 1_000).await else {
return;
};
let principal = Principal(10);
publish(
&store,
principal,
account_snapshot(ACCOUNT, 3, AccountStatus::Active),
)
.await;
for (step, status) in [
AccountStatus::Suspended,
AccountStatus::Active,
AccountStatus::Closed,
]
.into_iter()
.enumerate()
{
AdminStore::set_account_status(&*store, ACCOUNT, status)
.await
.unwrap_or_else(|e| panic!("{status:?} must be storable: {e}"));
let (stored, _) = status_of(&store, principal).await;
assert_eq!(
stored, status,
"step {step}: {status:?} must round-trip the ledger column and the JSONB"
);
}
}
#[tokio::test]
async fn an_unrecognized_status_column_is_a_storage_error() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(GrantPolicy::default(), 1_000).await else {
return;
};
let pool = corruption_pool().await;
sqlx::raw_sql(
"ALTER TABLE tollgate_accounts DROP CONSTRAINT IF EXISTS tollgate_accounts_status_known",
)
.execute(&pool)
.await
.unwrap();
sqlx::query("UPDATE tollgate_accounts SET status = 'Bogus' WHERE account_id = $1")
.bind(account_bytes())
.execute(&pool)
.await
.unwrap();
sqlx::raw_sql(
"ALTER TABLE tollgate_accounts ADD CONSTRAINT tollgate_accounts_status_known \
CHECK (status IN ('Active', 'Suspended', 'Closed')) NOT VALID",
)
.execute(&pool)
.await
.unwrap();
let error = store
.acquire(ACCOUNT, CostUnits(100), TTL, t(0))
.await
.unwrap_err();
assert!(
matches!(error, AllocateError::Storage(_)),
"an unrecognized status must surface, not be read as Active: {error:?}"
);
}
#[tokio::test]
async fn a_status_change_that_cannot_republish_moves_neither_record() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(GrantPolicy::default(), 1_000).await else {
return;
};
let principal = Principal(10);
let ceiling = i64::MAX as u64;
publish(
&store,
principal,
account_snapshot(ACCOUNT, ceiling, AccountStatus::Active),
)
.await;
let error = AdminStore::set_account_status(&*store, ACCOUNT, AccountStatus::Suspended)
.await
.unwrap_err();
assert!(
matches!(error, SetStatusError::Storage(_)),
"an unrepublishable snapshot surfaces, rather than being skipped: {error:?}"
);
assert_eq!(
status_of(&store, principal).await,
(AccountStatus::Active, Generation(ceiling)),
"the snapshot half did not move"
);
store
.acquire(ACCOUNT, CostUnits(100), TTL, t(0))
.await
.expect("and neither did the ledger half");
}
#[tokio::test]
async fn a_status_change_reports_rows_it_changed_but_could_not_push() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(GrantPolicy::default(), 1_000).await else {
return;
};
let readable = Principal(60);
let undecodable = Principal(61);
publish(
&store,
readable,
account_snapshot(ACCOUNT, 3, AccountStatus::Active),
)
.await;
let pool = corruption_pool().await;
sqlx::query(
"INSERT INTO tollgate_snapshots (principal, generation, snapshot, deleted)
VALUES ($1, 3, $2, FALSE)",
)
.bind(undecodable.0.to_be_bytes().to_vec())
.bind(serde_json::json!({
"account_id": ACCOUNT.0,
"key_id": null,
"generation": 3,
"status": "Active",
"valid_until": "2100-01-01T00:00:00Z",
"permissions": 1,
"limits": {
"max_items_per_request": 64,
"rate_units_per_second": 1000,
"rate_burst_units": 113
},
"cost_table": { "fixed_request": 50, "minimum_charge": 50, "weights": [1] }
}))
.execute(&pool)
.await
.unwrap();
let mut updates = store.subscribe();
let change = AdminStore::set_account_status(&*store, ACCOUNT, AccountStatus::Suspended)
.await
.unwrap();
assert_eq!(
(change.outcome.republished, change.outcome.unreadable),
(2, 1),
"both rows changed durably; one of them could not be pushed"
);
let mut pushed = Vec::new();
while let Ok(push) = updates.try_recv() {
pushed.push(push.principal);
}
assert_eq!(
pushed,
vec![readable],
"only the decodable principal is pushed; the other waits for a refresh"
);
assert_eq!(
status_of(&store, readable).await,
(AccountStatus::Suspended, Generation(4))
);
assert!(
store.snapshot(undecodable).await.is_err(),
"the corrupt row is still refused on read, as it was before the change"
);
}
#[tokio::test]
async fn a_vestigial_jsonb_generation_is_ignored_in_favour_of_the_column() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(GrantPolicy::default(), 1_000).await else {
return;
};
let principal = Principal(70);
let pool = corruption_pool().await;
sqlx::query(
"INSERT INTO tollgate_snapshots (principal, generation, snapshot, deleted)
VALUES ($1, 7, $2, FALSE)",
)
.bind(principal.0.to_be_bytes().to_vec())
.bind(serde_json::json!({
"account_id": ACCOUNT.0,
"key_id": null,
"generation": 1,
"status": "Active",
"valid_until": "2100-01-01T00:00:00Z",
"permissions": 1,
"limits": {
"max_items_per_request": 64,
"rate_units_per_second": 1000,
"rate_burst_units": 1000
},
"cost_table": { "fixed_request": 50, "minimum_charge": 50, "weights": [1] }
}))
.execute(&pool)
.await
.unwrap();
let SnapshotResolution::Present(live) = store.snapshot(principal).await.unwrap() else {
panic!("the planted row is live");
};
assert_eq!(
live.generation,
Generation(7),
"the live read resolves from the column, not the vestigial JSONB key"
);
AdminStore::remove_snapshot(&*store, principal)
.await
.unwrap();
assert!(
matches!(
store.snapshot(principal).await.unwrap(),
SnapshotResolution::Revoked {
generation: Generation(7)
}
),
"both branches report the same generation for the same row"
);
}
#[tokio::test]
async fn overage_usage_is_billed_and_funds_itself() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 1_000).await else {
return;
};
let before = store.conservation(ACCOUNT).await.unwrap().unwrap();
assert_eq!(before.overage_recorded, CostUnits::ZERO);
let report = store
.ingest(&[overage_usage(ACCOUNT, 1, 40, 1)], t(1))
.await
.unwrap();
assert_eq!(
(report.accepted, report.duplicate, report.rejected),
(1, 0, 0)
);
let after = store.conservation(ACCOUNT).await.unwrap().unwrap();
assert_eq!(after.overage_recorded, CostUnits(40));
assert_eq!(after.settled_usage, CostUnits(40));
assert_eq!(after.deposited, before.deposited, "no deposit was made");
assert_eq!(after.balance, before.balance, "and no balance was spent");
assert!(after.holds(), "conservation violated: {after:?}");
assert!(
!Conservation {
overage_recorded: CostUnits::ZERO,
..after
}
.holds(),
"the funding term must be what closes the equation"
);
assert_eq!(store.usage_recorded(ACCOUNT).await.unwrap(), CostUnits(40));
}
#[tokio::test]
async fn a_commit_time_fallback_is_ingested_as_overage_after_its_lease_settles() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 1_000).await else {
return;
};
let released = store
.acquire(ACCOUNT, CostUnits(500), TTL, t(0))
.await
.unwrap()
.grant;
store
.release(
released.lease_id,
released.fencing_token,
released.units,
t(59),
)
.await
.unwrap();
let receipt_sourced = store
.ingest(&[usage(&released, 1, 30, 61)], t(61))
.await
.unwrap();
assert_eq!(
(receipt_sourced.accepted, receipt_sourced.rejected),
(0, 1),
"a leased event naming a fully credited lease has nowhere to fit"
);
assert_eq!(
store.usage_recorded(ACCOUNT).await.unwrap(),
CostUnits::ZERO
);
let swept = store
.acquire(ACCOUNT, CostUnits(500), TTL, t(61))
.await
.unwrap()
.grant;
store.reclaim_expired(t(121)).await.unwrap();
for (request, at) in [(2, 61), (3, 122)] {
let phase_sourced = store
.ingest(&[overage_usage(ACCOUNT, request, 30, at)], t(at))
.await
.unwrap();
assert_eq!((phase_sourced.accepted, phase_sourced.rejected), (1, 0));
}
let after = store.conservation(ACCOUNT).await.unwrap().unwrap();
assert_eq!(after.settled_usage, CostUnits(60), "the work is billed");
assert_eq!(after.overage_recorded, CostUnits(60), "and it is funded");
assert_eq!(
after.settlement_loss, swept.units,
"the swept lease's forfeit is untouched by the fallback's bill"
);
assert!(after.holds(), "conservation violated: {after:?}");
assert_conserved(&store).await;
}
#[tokio::test]
async fn overage_usage_for_an_unknown_account_is_rejected() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 1_000).await else {
return;
};
let report = store
.ingest(&[overage_usage(AccountId(u128::MAX), 1, 40, 1)], t(1))
.await
.unwrap();
assert_eq!(
(report.accepted, report.duplicate, report.rejected),
(0, 0, 1)
);
assert_conserved(&store).await;
}
#[tokio::test]
async fn an_overage_replay_is_idempotent() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 1_000).await else {
return;
};
let event = overage_usage(ACCOUNT, 1, 40, 1);
assert_eq!(store.ingest(&[event], t(1)).await.unwrap().accepted, 1);
let report = store.ingest(&[event, event], t(2)).await.unwrap();
assert_eq!(
(report.accepted, report.duplicate, report.rejected),
(0, 2, 0)
);
assert_eq!(
store
.conservation(ACCOUNT)
.await
.unwrap()
.unwrap()
.overage_recorded,
CostUnits(40),
"a replay bills once, so it funds once"
);
assert_conserved(&store).await;
}
#[tokio::test]
async fn overage_accounting_overflow_is_surfaced() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 1_000).await else {
return;
};
let ceiling = u64::try_from(i64::MAX).unwrap();
assert_eq!(
store
.ingest(&[overage_usage(ACCOUNT, 1, ceiling, 1)], t(1))
.await
.unwrap()
.accepted,
1
);
let before = store.conservation(ACCOUNT).await.unwrap().unwrap();
assert_eq!(before.overage_recorded, CostUnits(ceiling));
let error = store
.ingest(&[overage_usage(ACCOUNT, 2, 1, 2)], t(2))
.await
.expect_err("a total that cannot be represented must be surfaced");
assert!(
error.to_string().contains("overflow"),
"unexpected error: {error}"
);
assert!(
!error.is_retryable(),
"a monotonic overflow cannot recover on retry"
);
let after = store.conservation(ACCOUNT).await.unwrap().unwrap();
assert_eq!(
(after.overage_recorded, after.settled_usage),
(before.overage_recorded, before.settled_usage),
"a rolled-back batch moves neither column"
);
}
#[tokio::test]
async fn settlement_is_unaffected_by_an_account_carrying_overage() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 1_000).await else {
return;
};
let released = store
.acquire(ACCOUNT, CostUnits(300), TTL, t(0))
.await
.unwrap()
.grant;
let expired = store
.acquire(ACCOUNT, CostUnits(200), TTL, t(0))
.await
.unwrap()
.grant;
assert_eq!(
store
.ingest(
&[
overage_usage(ACCOUNT, 1, 90, 1),
usage(&released, 2, 100, 1),
],
t(1)
)
.await
.unwrap()
.accepted,
2
);
store
.release(
released.lease_id,
released.fencing_token,
CostUnits(200),
t(2),
)
.await
.unwrap();
let reclaimed = store
.reclaim_expired(
expired
.expires_at
.checked_add(SignedDuration::from_secs(3_600))
.unwrap(),
)
.await
.unwrap();
assert_eq!(reclaimed.len(), 1);
let conservation = store.conservation(ACCOUNT).await.unwrap().unwrap();
assert_eq!(conservation.overage_recorded, CostUnits(90));
assert!(
conservation.holds(),
"conservation violated: {conservation:?}"
);
}
#[tokio::test]
async fn a_usage_row_cannot_carry_half_a_capability() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 1_000).await else {
return;
};
let pool = corruption_pool().await;
for (request, lease_id, fencing_token) in [
(901u128, Some(7u128.to_be_bytes().to_vec()), None::<i64>),
(902u128, None::<Vec<u8>>, Some(3i64)),
] {
let error = sqlx::query(
"INSERT INTO tollgate_usage_events
(request_id, account_id, lease_id, fencing_token, units, occurred_at_us)
VALUES ($1, $2, $3, $4, 1, 0)",
)
.bind(request.to_be_bytes().to_vec())
.bind(account_bytes())
.bind(&lease_id)
.bind(fencing_token)
.execute(&pool)
.await
.expect_err("half a capability must be refused");
assert!(
error.to_string().contains("lease_all_or_nothing"),
"unexpected error: {error}"
);
}
assert_conserved(&store).await;
}
fn key(id: u128, principal: u128, digest_byte: u8, not_after: Option<Timestamp>) -> KeyRecord {
KeyRecord {
key_id: KeyId(id),
account_id: ACCOUNT,
principal: Principal(principal),
digest: [digest_byte; 32],
not_after,
}
}
#[tokio::test]
async fn a_recorded_credential_is_active_until_it_is_revoked() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 100).await else {
return;
};
store.insert_key(key(1, 11, 0xa1, None)).await.unwrap();
let active = store.active_keys(t(0)).await.unwrap();
assert_eq!(active.len(), 1);
assert_eq!(active[0].key_id, KeyId(1));
assert_eq!(active[0].principal, Principal(11));
assert_eq!(active[0].digest, [0xa1; 32]);
assert_eq!(
store.revoke_key(KeyId(1), t(10)).await,
Ok(Revocation::Retired)
);
assert!(
store.active_keys(t(11)).await.unwrap().is_empty(),
"a retired credential is never active again"
);
}
#[tokio::test]
async fn revoking_reports_whether_anything_was_retired() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 100).await else {
return;
};
store.insert_key(key(1, 11, 0xa1, None)).await.unwrap();
assert_eq!(
store.revoke_key(KeyId(1), t(10)).await,
Ok(Revocation::Retired)
);
assert_eq!(
store.revoke_key(KeyId(1), t(20)).await,
Ok(Revocation::AlreadyRetired),
"the second call changed nothing and says so"
);
assert_eq!(
store.revoke_key(KeyId(404), t(20)).await,
Err(KeyError::UnknownKey),
"revoking a key that never existed is a mistake, not a no-op"
);
}
#[tokio::test]
async fn a_credential_expires_out_of_the_active_set_without_being_revoked() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 100).await else {
return;
};
store
.insert_key(key(1, 11, 0xa1, Some(t(100))))
.await
.unwrap();
assert_eq!(store.active_keys(t(99)).await.unwrap().len(), 1);
assert!(
store.active_keys(t(100)).await.unwrap().is_empty(),
"expiry is exclusive, and needs no operator action"
);
assert_eq!(
store.revoke_key(KeyId(1), t(200)).await,
Ok(Revocation::Retired),
"an expired credential is still revocable: expiry and retirement are different facts"
);
}
#[tokio::test]
async fn issuance_is_never_destructive() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 100).await else {
return;
};
store.insert_key(key(1, 11, 0xa1, None)).await.unwrap();
assert_eq!(
store.insert_key(key(1, 22, 0xb2, None)).await,
Err(KeyError::AlreadyExists),
"an overwrite would retire a live credential whose digest cannot be recovered"
);
let active = store.active_keys(t(0)).await.unwrap();
assert_eq!(active.len(), 1);
assert_eq!(
active[0].digest, [0xa1; 32],
"the refused insert changed nothing"
);
}
#[tokio::test]
async fn a_credential_cannot_belong_to_an_account_that_does_not_exist() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 100).await else {
return;
};
let orphan = KeyRecord {
account_id: AccountId(404),
..key(1, 11, 0xa1, None)
};
assert_eq!(
store.insert_key(orphan).await,
Err(KeyError::UnknownAccount),
"a credential nothing can authenticate as is refused at issuance, not at every admission"
);
assert!(store.active_keys(t(0)).await.unwrap().is_empty());
}
#[tokio::test]
async fn the_active_set_is_ordered_so_two_instances_project_alike() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 100).await else {
return;
};
for id in [3u128, 1, 2] {
store
.insert_key(key(id, 100 + id, id as u8, None))
.await
.unwrap();
}
let ids: Vec<_> = store
.active_keys(t(0))
.await
.unwrap()
.into_iter()
.map(|record| record.key_id)
.collect();
assert_eq!(ids, [KeyId(1), KeyId(2), KeyId(3)]);
}
#[tokio::test]
async fn a_half_written_credential_expiry_is_refused_not_read_as_absent() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 100).await else {
return;
};
store
.insert_key(key(77, 77, 0xc3, None))
.await
.expect("a credential with no expiry is the honest form of this row");
let pool = corruption_pool().await;
let dropped = suspend_checks(
&pool,
"tollgate_credential_keys",
"not_after_is_lower_bound",
None,
)
.await;
assert!(
!dropped.is_empty(),
"0018's expiry-domain check must exist, or this fixture is proving nothing"
);
sqlx::query(
"UPDATE tollgate_credential_keys SET not_after_is_lower_bound = TRUE WHERE key_id = $1",
)
.bind(77u128.to_be_bytes().to_vec())
.execute(&pool)
.await
.unwrap();
restore_checks(&pool, "tollgate_credential_keys", dropped).await;
let listed = store
.account_keys(ACCOUNT, None, NonZeroUsize::new(10).unwrap())
.await;
match listed {
Err(error) => assert!(
error
.to_string()
.contains("incomplete stored credential expiry"),
"the refusal must name the corruption it found: {error}"
),
Ok(summaries) => panic!(
"a lower-bound flag with no expiry was absorbed as absent: {:?}",
summaries
.iter()
.map(|summary| (summary.key_id, summary.not_after))
.collect::<Vec<_>>()
),
}
}
#[tokio::test]
async fn an_unclassified_database_refusal_is_storage_not_a_domain_answer() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 100).await else {
return;
};
let pool = corruption_pool().await;
sqlx::raw_sql(
"ALTER TABLE tollgate_credential_keys
ADD CONSTRAINT tollgate_credential_keys_test_refuse CHECK (false) NOT VALID",
)
.execute(&pool)
.await
.unwrap();
let refused = store.insert_key(key(1, 11, 0xa1, None)).await;
sqlx::raw_sql(
"ALTER TABLE tollgate_credential_keys
DROP CONSTRAINT tollgate_credential_keys_test_refuse",
)
.execute(&pool)
.await
.unwrap();
match refused {
Err(KeyError::Storage(_)) => {}
other => panic!(
"an unclassifiable database refusal must be Storage, not a domain answer: {other:?}"
),
}
}
#[tokio::test]
async fn two_credentials_cannot_share_a_principal() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 100).await else {
return;
};
store.insert_key(key(1, 11, 0xa1, None)).await.unwrap();
assert_eq!(
store.insert_key(key(2, 11, 0xb2, None)).await,
Err(KeyError::AlreadyExists)
);
assert_eq!(store.active_keys(t(0)).await.unwrap().len(), 1);
}
#[tokio::test]
async fn a_refused_deposit_moves_neither_column() {
let _guard = DB_LOCK.lock().await;
let ceiling = u64::try_from(i64::MAX).unwrap();
let Some(store) = store_with_balance(full_grant_policy(), ceiling).await else {
return;
};
store
.acquire(ACCOUNT, CostUnits(1_000), TTL, t(0))
.await
.unwrap();
let before = store.conservation(ACCOUNT).await.unwrap().unwrap();
assert!(before.holds(), "fixture must start conserved: {before:?}");
let refused = AdminStore::deposit(&*store, ACCOUNT, CostUnits(500)).await;
assert_eq!(refused.unwrap_err(), AllocateError::BalanceOverflow);
let after = store.conservation(ACCOUNT).await.unwrap().unwrap();
assert_eq!(
after.balance, before.balance,
"a refused deposit credited balance anyway"
);
assert_eq!(after.deposited, before.deposited);
assert!(
after.holds(),
"the ledger no longer conserves after a refused deposit: {after:?}"
);
}
#[tokio::test]
async fn a_failed_ingest_batch_leaves_the_ledger_untouched() {
for overage_first in [false, true] {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 10_000).await else {
return;
};
let lease = store
.acquire(ACCOUNT, CostUnits(1_000), TTL, t(0))
.await
.unwrap()
.grant;
let before = store.conservation(ACCOUNT).await.unwrap().unwrap();
let ceiling = u64::try_from(i64::MAX).unwrap();
let mut events = [
usage(&lease, 1, 100, 0),
overage_usage(ACCOUNT, 2, ceiling, 0),
];
if overage_first {
events.reverse();
}
let failed = UsageSink::ingest(&*store, &events, t(1)).await;
assert!(
matches!(failed, Err(tollgate_store::IngestError::Refused(_))),
"the accounting total cannot be represented: {failed:?}"
);
assert_eq!(
store.usage_recorded(ACCOUNT).await.unwrap(),
CostUnits::ZERO,
"the first event of a failed batch was applied anyway"
);
let after = store.conservation(ACCOUNT).await.unwrap().unwrap();
assert_eq!(after.settled_usage, before.settled_usage);
assert!(
after.holds(),
"conservation after a failed batch: {after:?}"
);
let replay = UsageSink::ingest(&*store, &[usage(&lease, 1, 100, 0)], t(2))
.await
.unwrap();
assert_eq!(
(replay.accepted, replay.duplicate),
(1, 0),
"an event from a failed batch was left indexed, so the replay saw a duplicate"
);
}
}
const JAN: i64 = 1_767_225_600;
const FEB: i64 = 1_769_904_000;
const MAR: i64 = 1_772_323_200;
fn monthly(allowance: u64) -> BudgetSchedule {
BudgetSchedule::monthly(CostUnits(allowance))
}
async fn roll(store: &PostgresStore, now: i64) -> Option<RolledAccount> {
let batch = AdminStore::roll_due_periods(store, t(now), NonZeroUsize::new(64).unwrap())
.await
.expect("the rollover pass succeeds");
batch
.rolled()
.iter()
.find(|rolled| rolled.account_id == ACCOUNT)
.copied()
}
#[tokio::test]
async fn setting_a_schedule_deposits_nothing() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 100).await else {
return;
};
AdminStore::set_budget_schedule(&*store, ACCOUNT, Some(monthly(500)))
.await
.unwrap();
assert_eq!(store.balance(ACCOUNT).await.unwrap(), CostUnits(100));
assert_conserved(&store).await;
}
#[tokio::test]
async fn setting_a_schedule_on_an_unknown_account_is_refused() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 100).await else {
return;
};
assert_eq!(
AdminStore::set_budget_schedule(&*store, AccountId(999), Some(monthly(500)))
.await
.unwrap_err(),
BudgetError::UnknownAccount
);
}
#[tokio::test]
async fn racing_passes_cross_a_boundary_exactly_once() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 0).await else {
return;
};
AdminStore::set_budget_schedule(&*store, ACCOUNT, Some(monthly(500)))
.await
.unwrap();
let (first, second) = tokio::join!(roll(&store, FEB), roll(&store, FEB));
let rolled = [first, second];
let winners: Vec<_> = rolled.iter().flatten().collect();
assert_eq!(winners.len(), 1, "exactly one pass crosses the boundary");
assert_eq!(winners[0].deposited, CostUnits(500));
assert_eq!(winners[0].expired, CostUnits::ZERO);
assert_eq!(
store.balance(ACCOUNT).await.unwrap(),
CostUnits(500),
"the loser must not deposit a second allowance"
);
assert_conserved(&store).await;
}
#[tokio::test]
async fn a_pass_inside_the_current_period_changes_nothing() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 0).await else {
return;
};
AdminStore::set_budget_schedule(&*store, ACCOUNT, Some(monthly(500)))
.await
.unwrap();
roll(&store, FEB).await;
assert_eq!(roll(&store, FEB + 10 * 86_400).await, None);
assert_eq!(store.balance(ACCOUNT).await.unwrap(), CostUnits(500));
}
#[tokio::test]
async fn an_account_without_a_schedule_is_untouched() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 100).await else {
return;
};
assert_eq!(roll(&store, FEB).await, None);
assert_eq!(store.balance(ACCOUNT).await.unwrap(), CostUnits(100));
assert_conserved(&store).await;
}
#[tokio::test]
async fn a_top_up_survives_rollover_but_the_allowance_does_not() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 0).await else {
return;
};
AdminStore::set_budget_schedule(&*store, ACCOUNT, Some(monthly(500)))
.await
.unwrap();
roll(&store, JAN + 1).await;
AdminStore::deposit(&*store, ACCOUNT, CostUnits(70))
.await
.unwrap();
assert_eq!(store.balance(ACCOUNT).await.unwrap(), CostUnits(570));
assert_eq!(
roll(&store, FEB).await,
Some(RolledAccount {
account_id: ACCOUNT,
deposited: CostUnits(500),
expired: CostUnits(500),
}),
"the whole unspent allowance expires; the top-up is not touched"
);
assert_eq!(
store.balance(ACCOUNT).await.unwrap(),
CostUnits(570),
"one fresh allowance plus the surviving top-up"
);
assert_conserved(&store).await;
}
#[tokio::test]
async fn a_missed_period_does_not_accrue_a_backlog() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 0).await else {
return;
};
AdminStore::set_budget_schedule(&*store, ACCOUNT, Some(monthly(500)))
.await
.unwrap();
roll(&store, JAN + 1).await;
assert_eq!(
roll(&store, MAR + 86_400).await,
Some(RolledAccount {
account_id: ACCOUNT,
deposited: CostUnits(500),
expired: CostUnits(500),
})
);
assert_eq!(store.balance(ACCOUNT).await.unwrap(), CostUnits(500));
assert_conserved(&store).await;
}
#[tokio::test]
async fn a_lease_from_the_closed_period_expires_its_unspent_allowance() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 0).await else {
return;
};
AdminStore::set_budget_schedule(&*store, ACCOUNT, Some(monthly(500)))
.await
.unwrap();
roll(&store, JAN + 1).await;
let lease = store
.acquire(ACCOUNT, CostUnits(200), TTL, t(FEB - 30))
.await
.unwrap()
.grant;
roll(&store, FEB).await;
store
.ingest(&[usage(&lease, 1, 50, FEB + 5)], t(FEB + 5))
.await
.unwrap();
store
.release(
lease.lease_id,
lease.fencing_token,
CostUnits(150),
t(FEB + 10),
)
.await
.unwrap();
assert_eq!(
store.balance(ACCOUNT).await.unwrap(),
CostUnits(500),
"the released units belonged to the closed period; only the new allowance remains"
);
let conservation = store.conservation(ACCOUNT).await.unwrap().unwrap();
assert_eq!(
conservation.expired,
CostUnits(450),
"300 unspent at the boundary plus the lease's 150"
);
assert_eq!(
conservation.settled_usage,
CostUnits(50),
"the straggler bills against the period the lease was granted in"
);
assert!(conservation.holds(), "{conservation:?}");
}
#[tokio::test]
async fn a_lease_funded_by_a_top_up_is_unaffected_by_a_boundary() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 300).await else {
return;
};
AdminStore::set_budget_schedule(&*store, ACCOUNT, Some(monthly(500)))
.await
.unwrap();
let lease = store
.acquire(ACCOUNT, CostUnits(200), TTL, t(FEB - 30))
.await
.unwrap()
.grant;
roll(&store, FEB).await;
store
.release(
lease.lease_id,
lease.fencing_token,
CostUnits(200),
t(FEB + 10),
)
.await
.unwrap();
assert_eq!(
store.balance(ACCOUNT).await.unwrap(),
CostUnits(800),
"the 300 top-up is intact and the new allowance sits beside it"
);
assert_eq!(
store.conservation(ACCOUNT).await.unwrap().unwrap().expired,
CostUnits::ZERO
);
assert_conserved(&store).await;
}
#[tokio::test]
async fn a_split_funded_lease_charges_the_allowance_half_first() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 100).await else {
return;
};
AdminStore::set_budget_schedule(&*store, ACCOUNT, Some(monthly(60)))
.await
.unwrap();
roll(&store, JAN + 1).await;
let lease = store
.acquire(ACCOUNT, CostUnits(160), TTL, t(FEB - 30))
.await
.unwrap()
.grant;
roll(&store, FEB).await;
store
.ingest(&[usage(&lease, 1, 10, FEB + 1)], t(FEB + 1))
.await
.unwrap();
store
.release(
lease.lease_id,
lease.fencing_token,
CostUnits(150),
t(FEB + 2),
)
.await
.unwrap();
let conservation = store.conservation(ACCOUNT).await.unwrap().unwrap();
assert_eq!(
conservation.expired,
CostUnits(50),
"the 10 units spent came out of the 60-unit allowance half, leaving 50 to expire"
);
assert_eq!(
store.balance(ACCOUNT).await.unwrap(),
CostUnits(160),
"the whole top-up survives, beside the new allowance"
);
assert!(conservation.holds(), "{conservation:?}");
}
#[tokio::test]
async fn reclaim_never_resurrects_a_closed_period_lease_it_sweeps() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 0).await else {
return;
};
AdminStore::set_budget_schedule(&*store, ACCOUNT, Some(monthly(500)))
.await
.unwrap();
roll(&store, JAN + 1).await;
let lease = store
.acquire(ACCOUNT, CostUnits(200), TTL, t(FEB - 30))
.await
.unwrap()
.grant;
roll(&store, FEB).await;
let batch = store
.reclaim_expired_batch(t(FEB + 3_600), NonZeroUsize::new(8).unwrap())
.await
.unwrap();
assert_eq!(batch.len(), 1);
assert_eq!(batch.reclaimed()[0].lease_id, lease.lease_id);
assert_eq!(
store.balance(ACCOUNT).await.unwrap(),
CostUnits(500),
"only the new period's allowance is spendable"
);
let c = store.conservation(ACCOUNT).await.unwrap().unwrap();
assert_eq!(
c.expired,
CostUnits(300),
"the rollover expired what it held"
);
assert_eq!(
c.settlement_loss,
CostUnits(200),
"the sweep forfeited the rest"
);
assert_conserved(&store).await;
}
#[tokio::test]
async fn clearing_a_schedule_leaves_the_balance_alone() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 0).await else {
return;
};
AdminStore::set_budget_schedule(&*store, ACCOUNT, Some(monthly(500)))
.await
.unwrap();
roll(&store, JAN + 1).await;
AdminStore::set_budget_schedule(&*store, ACCOUNT, None)
.await
.unwrap();
assert_eq!(roll(&store, FEB).await, None);
assert_eq!(store.balance(ACCOUNT).await.unwrap(), CostUnits(500));
assert_conserved(&store).await;
}
#[tokio::test]
async fn the_rollover_pass_is_bounded_and_saturation_says_there_is_more() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 0).await else {
return;
};
for id in 2..=5u128 {
AdminStore::create_account(
&*store,
AccountConfig {
account_id: AccountId(id),
initial_balance: CostUnits::ZERO,
status: AccountStatus::Active,
capacity_class: CapacityClass::Assured,
},
)
.await
.unwrap();
}
for id in 1..=5u128 {
AdminStore::set_budget_schedule(&*store, AccountId(id), Some(monthly(100)))
.await
.unwrap();
}
let limit = NonZeroUsize::new(2).unwrap();
let mut seen = Vec::new();
let mut batches = 0;
loop {
let batch = AdminStore::roll_due_periods(&*store, t(FEB), limit)
.await
.unwrap();
batches += 1;
assert!(batch.len() <= limit.get(), "the backend honoured the limit");
seen.extend(batch.rolled().iter().map(|rolled| rolled.account_id));
if !batch.is_saturated() {
break;
}
assert!(batches < 10, "the drain must terminate");
}
assert_eq!(seen.len(), 5, "every due account is rolled exactly once");
seen.sort_unstable_by_key(|account| account.0);
seen.dedup();
assert_eq!(seen.len(), 5, "no account is rolled twice across batches");
assert_eq!(
batches, 3,
"two saturated batches of two, then the partial one that ends the drain"
);
}
#[tokio::test]
async fn an_unscheduled_account_is_never_selected() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 10).await else {
return;
};
for id in 2..=3u128 {
AdminStore::create_account(
&*store,
AccountConfig {
account_id: AccountId(id),
initial_balance: CostUnits(10),
status: AccountStatus::Active,
capacity_class: CapacityClass::Assured,
},
)
.await
.unwrap();
}
AdminStore::set_budget_schedule(&*store, AccountId(2), Some(monthly(100)))
.await
.unwrap();
let batch = AdminStore::roll_due_periods(&*store, t(FEB), NonZeroUsize::new(64).unwrap())
.await
.unwrap();
assert_eq!(batch.len(), 1);
assert_eq!(batch.rolled()[0].account_id, AccountId(2));
assert_eq!(store.balance(ACCOUNT).await.unwrap(), CostUnits(10));
assert_eq!(store.balance(AccountId(3)).await.unwrap(), CostUnits(10));
}
#[tokio::test]
async fn a_grant_spends_the_expiring_allowance_before_a_top_up() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 100).await else {
return;
};
AdminStore::set_budget_schedule(&*store, ACCOUNT, Some(monthly(60)))
.await
.unwrap();
roll(&store, JAN + 1).await;
let lease = store
.acquire(ACCOUNT, CostUnits(40), TTL, t(JAN + 2))
.await
.unwrap()
.grant;
store
.ingest(&[usage(&lease, 1, 40, JAN + 3)], t(JAN + 3))
.await
.unwrap();
store
.release(
lease.lease_id,
lease.fencing_token,
CostUnits::ZERO,
t(JAN + 4),
)
.await
.unwrap();
assert_eq!(
store.balance(ACCOUNT).await.unwrap(),
CostUnits(120),
"60 - 40 allowance, plus the untouched top-up"
);
assert_eq!(
roll(&store, FEB).await.map(|rolled| rolled.expired),
Some(CostUnits(20)),
"the 40 spent came out of the allowance, so only 20 of it was left to expire"
);
assert_eq!(
store.balance(ACCOUNT).await.unwrap(),
CostUnits(160),
"the whole 100 top-up survives beside the new allowance"
);
assert_conserved(&store).await;
}
#[tokio::test]
async fn consolidating_inside_a_lease_own_period_expires_nothing() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 0).await else {
return;
};
AdminStore::set_budget_schedule(&*store, ACCOUNT, Some(monthly(500)))
.await
.unwrap();
roll(&store, JAN + 1).await;
let lease = store
.acquire(ACCOUNT, CostUnits(200), TTL, t(JAN + 2))
.await
.unwrap()
.grant;
assert_eq!(lease.units, CostUnits(200), "drawn from the allowance");
let folded = store
.consolidate(
lease.lease_id,
lease.fencing_token,
lease.units,
CostUnits(500),
CostUnits::ZERO,
TTL,
t(JAN + 3),
)
.await
.unwrap()
.grant;
assert_eq!(
folded.units,
CostUnits(500),
"the returned allowance is spendable again inside its own period"
);
assert_eq!(
store.conservation(ACCOUNT).await.unwrap().unwrap().expired,
CostUnits::ZERO,
"nothing lapsed"
);
assert_conserved(&store).await;
}
#[tokio::test]
async fn consolidating_across_a_boundary_regrants_only_what_the_credit_restores() {
let _guard = DB_LOCK.lock().await;
const LONG: SignedDuration = SignedDuration::from_secs(60 * 24 * 3_600);
let Some(store) = store_with_balance(
GrantPolicy {
max_ttl: LONG,
..full_grant_policy()
},
0,
)
.await
else {
return;
};
AdminStore::set_budget_schedule(&*store, ACCOUNT, Some(monthly(500)))
.await
.unwrap();
roll(&store, JAN + 1).await;
AdminStore::deposit(&*store, ACCOUNT, CostUnits(100))
.await
.unwrap();
let january = store
.acquire(ACCOUNT, CostUnits(600), LONG, t(JAN + 2))
.await
.unwrap()
.grant;
assert_eq!(january.units, CostUnits(600));
roll(&store, FEB + 1).await;
let february = store
.consolidate(
january.lease_id,
january.fencing_token,
CostUnits(600),
CostUnits(1_000),
CostUnits::ZERO,
LONG,
t(FEB + 2),
)
.await
.unwrap()
.grant;
assert_eq!(february.units, CostUnits(600));
assert_eq!(
store.conservation(ACCOUNT).await.unwrap().unwrap().expired,
CostUnits(500),
"the lapsed allowance is written off, not re-leased"
);
assert_conserved(&store).await;
}
#[tokio::test]
async fn consolidation_after_a_budget_reduction_uses_only_restored_credit_as_floor() {
const LONG: SignedDuration = SignedDuration::from_secs(60 * 24 * 3_600);
let policy = GrantPolicy {
shrink_divisor: 2,
max_ttl: LONG,
..full_grant_policy()
};
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(policy, 0).await else {
return;
};
AdminStore::set_budget_schedule(&*store, ACCOUNT, Some(monthly(500)))
.await
.unwrap();
roll(&store, JAN + 1).await;
let old = store
.acquire(ACCOUNT, CostUnits(500), LONG, t(JAN + 2))
.await
.unwrap()
.grant;
assert_eq!(old.units, CostUnits(250));
AdminStore::set_budget_schedule(&*store, ACCOUNT, Some(monthly(100)))
.await
.unwrap();
roll(&store, FEB + 1).await;
let fresh = store
.consolidate(
old.lease_id,
old.fencing_token,
old.units,
CostUnits(1_000),
CostUnits::ZERO,
LONG,
t(FEB + 2),
)
.await
.unwrap()
.grant;
assert_eq!(
fresh.units,
CostUnits(50),
"expired allowance cannot increase the policy floor"
);
assert_conserved(&store).await;
}
#[tokio::test]
async fn a_lease_released_inside_its_own_period_expires_nothing() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 0).await else {
return;
};
AdminStore::set_budget_schedule(&*store, ACCOUNT, Some(monthly(500)))
.await
.unwrap();
roll(&store, JAN + 1).await;
let lease = store
.acquire(ACCOUNT, CostUnits(200), TTL, t(JAN + 2))
.await
.unwrap()
.grant;
store
.release(
lease.lease_id,
lease.fencing_token,
CostUnits(200),
t(JAN + 3),
)
.await
.unwrap();
assert_eq!(
store.conservation(ACCOUNT).await.unwrap().unwrap().expired,
CostUnits::ZERO
);
assert_eq!(
store.balance(ACCOUNT).await.unwrap(),
CostUnits(500),
"the allowance is whole again"
);
assert_conserved(&store).await;
}
#[tokio::test]
async fn a_lease_reclaimed_inside_its_own_period_expires_nothing() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 0).await else {
return;
};
AdminStore::set_budget_schedule(&*store, ACCOUNT, Some(monthly(500)))
.await
.unwrap();
roll(&store, JAN + 1).await;
store
.acquire(ACCOUNT, CostUnits(200), TTL, t(JAN + 2))
.await
.unwrap();
let batch = store
.reclaim_expired_batch(t(JAN + 3_600), NonZeroUsize::new(8).unwrap())
.await
.unwrap();
assert_eq!(batch.len(), 1);
let c = store.conservation(ACCOUNT).await.unwrap().unwrap();
assert_eq!(c.expired, CostUnits::ZERO);
assert_eq!(c.settlement_loss, CostUnits(200));
assert_eq!(store.balance(ACCOUNT).await.unwrap(), CostUnits(300));
assert_conserved(&store).await;
}
#[tokio::test]
async fn a_snapshot_revision_round_trips_through_postgres() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 1_000).await else {
return;
};
let principal = Principal(0x94);
let revision = PolicyRevision([0x7e; 32]);
let submitted = AccountSnapshot::builder(
ACCOUNT,
Generation(1),
AccountStatus::Active,
t(10_000),
PermissionBits::ALL,
ResolvedLimits::new(64).with_weighted_rate(1_000, 1_000),
Arc::new(CostTable::builder(CostUnits(1), CostUnits(1)).build()),
)
.policy_revision(revision)
.build();
AdminStore::publish_snapshot(&*store, principal, publishable(Arc::new(submitted)))
.await
.unwrap();
let SnapshotResolution::Present(fetched) = store.snapshot(principal).await.unwrap() else {
panic!("the snapshot must be present");
};
assert_eq!(fetched.policy_revision, revision);
let pool = corruption_pool().await;
let stored: serde_json::Value =
sqlx::query_scalar("SELECT snapshot FROM tollgate_snapshots WHERE principal = $1")
.bind(principal.0.to_be_bytes().to_vec())
.fetch_one(&pool)
.await
.unwrap();
assert_eq!(
stored["policy_revision"],
serde_json::Value::String("7e".repeat(32))
);
}
#[tokio::test]
async fn a_legacy_snapshot_document_defaults_the_revision_to_unstated() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 1_000).await else {
return;
};
let principal = Principal(0x96);
let submitted = account_snapshot(ACCOUNT, 1, AccountStatus::Active);
AdminStore::publish_snapshot(&*store, principal, publishable(Arc::new(submitted)))
.await
.unwrap();
let pool = corruption_pool().await;
sqlx::query("UPDATE tollgate_snapshots SET snapshot = snapshot - 'policy_revision' WHERE principal = $1")
.bind(principal.0.to_be_bytes().to_vec())
.execute(&pool)
.await
.unwrap();
let stored: serde_json::Value =
sqlx::query_scalar("SELECT snapshot FROM tollgate_snapshots WHERE principal = $1")
.bind(principal.0.to_be_bytes().to_vec())
.fetch_one(&pool)
.await
.unwrap();
assert!(
stored.get("policy_revision").is_none(),
"the fixture must actually be a pre-GL-94 document"
);
let SnapshotResolution::Present(fetched) = store.snapshot(principal).await.unwrap() else {
panic!("a legacy document must still decode");
};
assert!(fetched.policy_revision.is_unstated());
}
#[tokio::test]
async fn a_usage_row_carries_its_policy_revision() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 1_000).await else {
return;
};
let lease = store
.acquire(ACCOUNT, CostUnits(500), TTL, t(0))
.await
.unwrap()
.grant;
let revision = PolicyRevision([0x31; 32]);
let event = UsageEvent::new(
RequestId(1),
lease.account_id,
UsageSource::Leased {
lease_id: lease.lease_id,
fencing_token: lease.fencing_token,
},
CostUnits(40),
t(1),
revision,
None,
);
let report = store.ingest(&[event], t(1)).await.unwrap();
assert_eq!((report.accepted, report.rejected), (1, 0));
let pool = corruption_pool().await;
let stored: Vec<u8> = sqlx::query_scalar(
"SELECT policy_revision FROM tollgate_usage_events WHERE request_id = $1",
)
.bind(1u128.to_be_bytes().to_vec())
.fetch_one(&pool)
.await
.unwrap();
assert_eq!(stored, revision.as_bytes().to_vec());
let plain = UsageEvent::new(
RequestId(2),
lease.account_id,
UsageSource::Leased {
lease_id: lease.lease_id,
fencing_token: lease.fencing_token,
},
CostUnits(10),
t(2),
PolicyRevision::UNSTATED,
None,
);
store.ingest(&[plain], t(2)).await.unwrap();
let unstated: Vec<u8> = sqlx::query_scalar(
"SELECT policy_revision FROM tollgate_usage_events WHERE request_id = $1",
)
.bind(2u128.to_be_bytes().to_vec())
.fetch_one(&pool)
.await
.unwrap();
assert_eq!(unstated, vec![0u8; 32]);
assert_conserved(&store).await;
}
#[tokio::test]
async fn a_usage_row_cannot_carry_a_wrong_width_revision() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 1_000).await else {
return;
};
let lease = store
.acquire(ACCOUNT, CostUnits(500), TTL, t(0))
.await
.unwrap()
.grant;
let pool = corruption_pool().await;
for wrong in [vec![0u8; 31], vec![0u8; 33], Vec::new()] {
let error = sqlx::query(
"INSERT INTO tollgate_usage_events
(request_id, account_id, lease_id, fencing_token, units, occurred_at_us, policy_revision)
VALUES ($1, $2, $3, $4, $5, $6, $7)",
)
.bind(vec![9u8; 16])
.bind(lease.account_id.0.to_be_bytes().to_vec())
.bind(lease.lease_id.0.to_be_bytes().to_vec())
.bind(i64::try_from(lease.fencing_token.0).unwrap())
.bind(1i64)
.bind(1i64)
.bind(&wrong)
.execute(&pool)
.await
.expect_err("a revision that is not 32 bytes must be refused");
assert!(
error.to_string().contains("policy_revision_len"),
"expected the length constraint, got {error}"
);
}
}
#[tokio::test]
async fn a_published_snapshot_carries_the_ledgers_budget() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 0).await else {
return;
};
AdminStore::set_budget_schedule(&*store, ACCOUNT, Some(monthly(500)))
.await
.unwrap();
roll(&store, JAN + 1).await;
let principal = Principal(0x51);
let submitted = account_snapshot(ACCOUNT, 1, AccountStatus::Active);
assert_eq!(
submitted.budget, None,
"a builder cannot describe a balance; only the store can"
);
AdminStore::publish_snapshot(&*store, principal, publishable(Arc::new(submitted)))
.await
.unwrap();
let SnapshotResolution::Present(fetched) = store.snapshot(principal).await.unwrap() else {
panic!("the snapshot must be present");
};
let budget = fetched.budget.expect("the store stamps a budget view");
assert_eq!(budget.balance_at_publish, CostUnits(500));
assert_eq!(
budget.period_end,
Some(t(FEB)),
"the current period ends at the next boundary"
);
}
#[tokio::test]
async fn the_budget_view_counts_units_out_on_lease() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 1_000).await else {
return;
};
let principal = Principal(0x52);
let lease = store
.acquire(ACCOUNT, CostUnits(400), TTL, t(0))
.await
.unwrap()
.grant;
store
.ingest(&[usage(&lease, 1, 150, 1)], t(1))
.await
.unwrap();
AdminStore::publish_snapshot(
&*store,
principal,
publishable(Arc::new(account_snapshot(
ACCOUNT,
1,
AccountStatus::Active,
))),
)
.await
.unwrap();
let SnapshotResolution::Present(fetched) = store.snapshot(principal).await.unwrap() else {
panic!("the snapshot must be present");
};
assert_eq!(
fetched.budget.unwrap().balance_at_publish,
CostUnits(850),
"600 in the balance plus the lease's unspent 250; only the 150 spent is gone"
);
assert_eq!(
fetched.budget.unwrap().period_end,
None,
"an account with no schedule has no period that ends"
);
}
#[tokio::test]
async fn a_supplied_budget_never_survives_publication() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 700).await else {
return;
};
let fabricated = BudgetView {
balance_at_publish: CostUnits(999_999),
period_end: Some(t(MAR)),
};
let known = publishable(Arc::new(account_snapshot(
ACCOUNT,
1,
AccountStatus::Active,
)))
.with_budget(Some(fabricated));
AdminStore::publish_snapshot(&*store, Principal(0x53), known)
.await
.unwrap();
let SnapshotResolution::Present(fetched) = store.snapshot(Principal(0x53)).await.unwrap()
else {
panic!("the snapshot must be present");
};
assert_eq!(
fetched.budget.unwrap().balance_at_publish,
CostUnits(700),
"the ledger's number, not the publisher's"
);
let stranger = AccountId(0xdead);
let unknown = publishable(Arc::new(account_snapshot(
stranger,
1,
AccountStatus::Active,
)))
.with_budget(Some(fabricated));
AdminStore::publish_snapshot(&*store, Principal(0x54), unknown)
.await
.unwrap();
let SnapshotResolution::Present(fetched) = store.snapshot(Principal(0x54)).await.unwrap()
else {
panic!("the snapshot must be present");
};
assert_eq!(
fetched.budget, None,
"an account the ledger does not hold reports no budget, not a fabricated one"
);
}
#[tokio::test]
async fn the_budget_view_does_not_report_expired_units_as_spendable() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 100).await else {
return;
};
AdminStore::set_budget_schedule(&*store, ACCOUNT, Some(monthly(500)))
.await
.unwrap();
roll(&store, JAN + 1).await;
let principal = Principal(0x55);
roll(&store, FEB).await;
AdminStore::publish_snapshot(
&*store,
principal,
publishable(Arc::new(account_snapshot(
ACCOUNT,
1,
AccountStatus::Active,
))),
)
.await
.unwrap();
let SnapshotResolution::Present(fetched) = store.snapshot(principal).await.unwrap() else {
panic!("the snapshot must be present");
};
assert_eq!(
fetched.budget.unwrap().balance_at_publish,
CostUnits(600),
"the new allowance plus the surviving top-up, and not a unit of the expired 500"
);
assert_eq!(fetched.budget.unwrap().period_end, Some(t(MAR)));
}
#[tokio::test]
async fn the_budget_view_writes_off_settlement_loss() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 1_000).await else {
return;
};
let principal = Principal(0x56);
let lease = store
.acquire(ACCOUNT, CostUnits(400), TTL, t(0))
.await
.unwrap()
.grant;
store
.ingest(&[usage(&lease, 1, 150, 1)], t(1))
.await
.unwrap();
store
.release(lease.lease_id, lease.fencing_token, CostUnits(200), t(2))
.await
.unwrap();
AdminStore::publish_snapshot(
&*store,
principal,
publishable(Arc::new(account_snapshot(
ACCOUNT,
1,
AccountStatus::Active,
))),
)
.await
.unwrap();
let SnapshotResolution::Present(fetched) = store.snapshot(principal).await.unwrap() else {
panic!("the snapshot must be present");
};
assert_eq!(
fetched.budget.unwrap().balance_at_publish,
CostUnits(800),
"1000 funded, less the 150 billed and the 50 written off"
);
assert_conserved(&store).await;
}
#[tokio::test]
async fn a_partially_populated_schedule_is_reported_not_interpreted() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 100).await else {
return;
};
let pool = corruption_pool().await;
let dropped = suspend_checks(&pool, "tollgate_accounts", "budget_allowance", None).await;
sqlx::query(
"UPDATE tollgate_accounts SET budget_period = 'utc_calendar_month' WHERE account_id = $1",
)
.bind(account_bytes())
.execute(&pool)
.await
.unwrap();
restore_checks(&pool, "tollgate_accounts", dropped).await;
let error = AdminStore::publish_snapshot(
&*store,
Principal(0x57),
publishable(Arc::new(account_snapshot(
ACCOUNT,
1,
AccountStatus::Active,
))),
)
.await
.unwrap_err();
assert!(
error.to_string().contains("partially populated"),
"the publish must say what it found: {error}"
);
}
#[tokio::test]
async fn an_unrecognized_stored_schedule_name_is_refused() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 100).await else {
return;
};
let pool = corruption_pool().await;
for (column, planted, principal) in [
("budget_period", "lunar_month", Principal(0x58)),
("budget_rollover", "carry_over", Principal(0x59)),
] {
sqlx::query(
"UPDATE tollgate_accounts
SET budget_allowance = 500, budget_period = 'utc_calendar_month',
budget_rollover = 'none'
WHERE account_id = $1",
)
.bind(account_bytes())
.execute(&pool)
.await
.unwrap();
sqlx::query(&format!(
"UPDATE tollgate_accounts SET {column} = $1 WHERE account_id = $2"
))
.bind(planted)
.bind(account_bytes())
.execute(&pool)
.await
.unwrap();
let error = AdminStore::publish_snapshot(
&*store,
principal,
publishable(Arc::new(account_snapshot(
ACCOUNT,
1,
AccountStatus::Active,
))),
)
.await
.unwrap_err();
assert!(
error.to_string().contains(planted),
"the refusal must name what it could not read in {column}: {error}"
);
}
sqlx::query("UPDATE tollgate_accounts SET budget_period = 'lunar_month' WHERE account_id = $1")
.bind(account_bytes())
.execute(&pool)
.await
.unwrap();
let batch = AdminStore::roll_due_periods(&*store, t(FEB), NonZeroUsize::new(8).unwrap())
.await
.unwrap();
assert!(batch.is_empty());
assert_eq!(store.balance(ACCOUNT).await.unwrap(), CostUnits(100));
}
#[tokio::test]
async fn admin_receipts_identify_the_state_each_operation_replaced() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 100).await else {
return;
};
use tollgate_store::AdminState;
let created = AdminStore::create_account(
&*store,
AccountConfig {
account_id: AccountId(2),
initial_balance: CostUnits(10),
status: AccountStatus::Active,
capacity_class: CapacityClass::Assured,
},
)
.await
.unwrap();
assert_eq!(created.before, AdminState::Absent);
assert_eq!(
created.after,
AdminState::AccountCreated {
initial_balance: CostUnits(10),
status: AccountStatus::Active,
capacity_class: CapacityClass::Assured,
origin: AdminAuthority::Operator,
}
);
let deposit = AdminStore::deposit(&*store, ACCOUNT, CostUnits(25))
.await
.unwrap();
assert_eq!(
deposit.before,
AdminState::Funding {
topup: CostUnits(100),
deposited: CostUnits(100)
}
);
assert_eq!(
deposit.after,
AdminState::Funding {
topup: CostUnits(125),
deposited: CostUnits(125)
}
);
let status = AdminStore::set_account_status(&*store, ACCOUNT, AccountStatus::Suspended)
.await
.unwrap();
assert_eq!(
status.before,
AdminState::Status {
status: AccountStatus::Active,
set_by: AdminAuthority::Operator,
}
);
assert_eq!(
status.after,
AdminState::Status {
status: AccountStatus::Suspended,
set_by: AdminAuthority::Operator,
}
);
let repeated = AdminStore::set_account_status(&*store, ACCOUNT, AccountStatus::Suspended)
.await
.unwrap();
assert_eq!(repeated.before, repeated.after);
let class = AdminStore::set_capacity_class(&*store, ACCOUNT, CapacityClass::BestEffort)
.await
.unwrap();
assert_eq!(
class.before,
AdminState::CapacityClass {
capacity_class: CapacityClass::Assured
}
);
assert_eq!(
class.after,
AdminState::CapacityClass {
capacity_class: CapacityClass::BestEffort
}
);
let unknown = AdminStore::remove_snapshot(&*store, Principal(99))
.await
.unwrap();
assert_eq!(unknown.before, AdminState::Absent);
assert_eq!(unknown.after, AdminState::Absent);
}
#[tokio::test]
async fn concurrent_deposit_receipts_form_one_exact_funding_history() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 100).await else {
return;
};
use tollgate_store::AdminState;
let mut tasks = tokio::task::JoinSet::new();
for _ in 0..32 {
let store = Arc::clone(&store);
tasks.spawn(async move {
AdminStore::deposit(&*store, ACCOUNT, CostUnits(7))
.await
.unwrap()
});
}
let mut receipts = Vec::new();
while let Some(receipt) = tasks.join_next().await {
receipts.push(receipt.unwrap());
}
receipts.sort_by_key(|receipt| match receipt.before {
AdminState::Funding { deposited, .. } => deposited.get(),
_ => panic!("wrong receipt state"),
});
for (index, receipt) in receipts.iter().enumerate() {
let before = 100 + index as u64 * 7;
assert_eq!(
receipt.before,
AdminState::Funding {
topup: CostUnits(before),
deposited: CostUnits(before)
}
);
assert_eq!(
receipt.after,
AdminState::Funding {
topup: CostUnits(before + 7),
deposited: CostUnits(before + 7)
}
);
}
}
#[tokio::test]
async fn racing_publication_and_revocation_receipts_name_the_actual_predecessor() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 100).await else {
return;
};
use tollgate_store::AdminState;
let mut tasks = tokio::task::JoinSet::new();
for generation in 1..=24 {
let publisher = Arc::clone(&store);
tasks.spawn(async move {
let snapshot = publishable(Arc::new(
AccountSnapshot::builder(
ACCOUNT,
Generation(generation),
AccountStatus::Active,
t(1000),
PermissionBits(0),
ResolvedLimits::new(1),
Arc::new(CostTable::builder(CostUnits(1), CostUnits(1)).build()),
)
.build(),
));
AdminStore::publish_snapshot(&*publisher, Principal(99), snapshot)
.await
.unwrap()
});
if generation % 3 == 0 {
let store = Arc::clone(&store);
tasks.spawn(async move {
AdminStore::remove_snapshot(&*store, Principal(99))
.await
.unwrap()
});
}
}
let mut changed = Vec::new();
while let Some(receipt) = tasks.join_next().await {
let receipt = receipt.unwrap();
if receipt.before != receipt.after {
changed.push(receipt);
}
}
changed.sort_by_key(|receipt| match receipt.after {
AdminState::Snapshot {
generation,
revoked,
} => (generation.0, revoked),
_ => panic!("wrong publication receipt"),
});
let mut previous = AdminState::Absent;
for receipt in changed {
assert_eq!(
receipt.before, previous,
"receipts must describe one serialized history"
);
previous = receipt.after;
}
let final_state = match SnapshotSource::snapshot(&*store, Principal(99))
.await
.unwrap()
{
SnapshotResolution::Present(snapshot) => AdminState::Snapshot {
generation: snapshot.generation,
revoked: false,
},
SnapshotResolution::Revoked { generation } => AdminState::Snapshot {
generation,
revoked: true,
},
SnapshotResolution::Unknown => AdminState::Absent,
};
assert_eq!(previous, final_state);
}
#[tokio::test]
async fn credential_projection_preserves_lifecycle_truth_and_refuses_corrupt_identity() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 100).await else {
return;
};
use tollgate_store::KeySource;
for (id, expiry) in [(1u128, None), (2, Some(t(100))), (3, None)] {
let mut digest = [0u8; 32];
digest[..16].copy_from_slice(&id.to_be_bytes());
store
.insert_key(KeyRecord {
key_id: KeyId(id),
account_id: ACCOUNT,
principal: Principal(id),
digest,
not_after: expiry,
})
.await
.unwrap();
}
store.revoke_key(KeyId(3), t(99)).await.unwrap();
let before = store
.active_keys_page(t(99), None, tollgate_store::DEFAULT_KEY_PAGE_LIMIT)
.await
.unwrap();
assert_eq!(before.records().len(), 2);
assert_eq!(before.records()[0].principal, Principal(1));
assert_eq!(before.records()[1].not_after, Some(t(100)));
let at_expiry = store
.active_keys_page(t(100), None, tollgate_store::DEFAULT_KEY_PAGE_LIMIT)
.await
.unwrap();
assert_eq!(at_expiry.records(), &before.records()[..1]);
store.revoke_key(KeyId(1), t(100)).await.unwrap();
assert!(
store
.active_keys_page(t(100), None, tollgate_store::DEFAULT_KEY_PAGE_LIMIT)
.await
.unwrap()
.records()
.is_empty()
);
store
.insert_key(KeyRecord {
key_id: KeyId(4),
account_id: ACCOUNT,
principal: Principal(4),
digest: [0; 32],
not_after: None,
})
.await
.unwrap();
assert!(
store
.active_keys_page(t(100), None, tollgate_store::DEFAULT_KEY_PAGE_LIMIT)
.await
.is_err()
);
}
#[tokio::test]
async fn credential_pages_order_bound_skip_retired_and_expose_every_mutation() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 100).await else {
return;
};
use tollgate_store::KeySource;
let limit = NonZeroUsize::new(2).unwrap();
let initial = store.active_keys_page(t(100), None, limit).await.unwrap();
let record = |id: u128, expiry| {
let mut digest = [0; 32];
digest[..16].copy_from_slice(&id.to_be_bytes());
KeyRecord {
key_id: KeyId(id),
account_id: ACCOUNT,
principal: Principal(id),
digest,
not_after: expiry,
}
};
for (id, expiry) in [
(8, None),
(0, None),
(3, Some(t(100))),
(4, None),
(2, None),
(6, None),
] {
store.insert_key(record(id, expiry)).await.unwrap();
}
store.revoke_key(KeyId(4), t(99)).await.unwrap();
let first = store.active_keys_page(t(100), None, limit).await.unwrap();
assert!(first.revision() > initial.revision());
assert_eq!(first.as_of(), t(100));
assert_eq!(
first.records().iter().map(|k| k.key_id).collect::<Vec<_>>(),
vec![KeyId(0), KeyId(2)]
);
assert_eq!(first.next_after(), Some(KeyId(2)));
let last = store
.active_keys_page(t(100), first.next_after(), limit)
.await
.unwrap();
assert_eq!(last.revision(), first.revision());
assert_eq!(
last.records().iter().map(|k| k.key_id).collect::<Vec<_>>(),
vec![KeyId(6), KeyId(8)]
);
assert_eq!(
last.next_after(),
None,
"lookahead makes an exactly full final page terminal"
);
let empty = store
.active_keys_page(t(100), Some(KeyId(8)), limit)
.await
.unwrap();
assert!(empty.records().is_empty());
assert_eq!(empty.next_after(), None);
assert_eq!(empty.revision(), first.revision());
let skipped = store
.active_keys_page(t(100), Some(KeyId(3)), limit)
.await
.unwrap();
assert_eq!(skipped.records(), last.records());
store.revoke_key(KeyId(0), t(101)).await.unwrap();
store.insert_key(record(1, None)).await.unwrap(); let changed = store.active_keys_page(t(101), None, limit).await.unwrap();
assert!(changed.revision() > first.revision());
assert_eq!(
changed
.records()
.iter()
.map(|k| k.key_id)
.collect::<Vec<_>>(),
vec![KeyId(1), KeyId(2)]
);
let mut duplicate = record(0, None);
duplicate.key_id = KeyId(99);
assert!(
matches!(
store.insert_key(duplicate).await,
Err(KeyError::AlreadyExists)
),
"retirement must not free the principal uniqueness index"
);
assert!(
store
.active_keys_page(
t(100),
None,
NonZeroUsize::new(tollgate_store::MAX_KEY_PAGE_LIMIT + 1).unwrap()
)
.await
.is_err()
);
}
#[tokio::test]
async fn credential_revision_covers_direct_writes_rollback_and_overflow() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 100).await else {
return;
};
use tollgate_store::{DEFAULT_KEY_PAGE_LIMIT, KeySource};
let pool = corruption_pool().await;
let revision = || store.active_keys_page(t(100), None, DEFAULT_KEY_PAGE_LIMIT);
let initial = revision().await.unwrap().revision();
let id = 1u128.to_be_bytes();
let mut digest = [0; 32];
digest[..16].copy_from_slice(&id);
sqlx::query("INSERT INTO tollgate_credential_keys (key_id, account_id, principal, digest, not_after_is_lower_bound) VALUES ($1,$2,$1,$3,FALSE)")
.bind(id.as_slice()).bind(ACCOUNT.0.to_be_bytes().as_slice()).bind(digest.as_slice())
.execute(&pool).await.unwrap();
let inserted = revision().await.unwrap();
assert!(inserted.revision() > initial);
assert_eq!(inserted.records().len(), 1);
let mut tx = pool.begin().await.unwrap();
sqlx::query("UPDATE tollgate_credential_keys SET revoked_at_us = 100000000")
.execute(&mut *tx)
.await
.unwrap();
let before_commit = revision().await.unwrap();
assert_eq!(before_commit.revision(), inserted.revision());
assert_eq!(before_commit.records(), inserted.records());
tx.rollback().await.unwrap();
assert_eq!(revision().await.unwrap(), inserted);
sqlx::query("UPDATE tollgate_credential_keys SET revoked_at_us = 100000000")
.execute(&pool)
.await
.unwrap();
let revoked = revision().await.unwrap();
assert!(revoked.revision() > inserted.revision());
assert!(revoked.records().is_empty());
sqlx::query("UPDATE tollgate_credential_revision SET revision = $1")
.bind(i64::MAX)
.execute(&pool)
.await
.unwrap();
let repeated_retirement = store.revoke_key(KeyId(1), t(100)).await;
let unknown_retirement = store.revoke_key(KeyId(404), t(100)).await;
let duplicate = store
.insert_key(KeyRecord {
key_id: KeyId(1),
account_id: ACCOUNT,
principal: Principal(1),
digest,
not_after: None,
})
.await;
let mut second_digest = [0; 32];
second_digest[..16].copy_from_slice(&2u128.to_be_bytes());
let failed = store
.insert_key(KeyRecord {
key_id: KeyId(2),
account_id: ACCOUNT,
principal: Principal(2),
digest: second_digest,
not_after: None,
})
.await;
let overflow_page = revision().await;
sqlx::query("UPDATE tollgate_credential_revision SET revision = $1")
.bind(i64::try_from(revoked.revision() + 2).unwrap())
.execute(&pool)
.await
.unwrap();
assert_eq!(repeated_retirement, Ok(Revocation::AlreadyRetired));
assert_eq!(unknown_retirement, Err(KeyError::UnknownKey));
assert_eq!(duplicate, Err(KeyError::AlreadyExists));
assert!(matches!(failed, Err(KeyError::Storage(_))));
let overflow_page = overflow_page.unwrap();
assert_eq!(overflow_page.revision(), i64::MAX as u64);
assert!(overflow_page.records().is_empty());
sqlx::query("DELETE FROM tollgate_credential_keys")
.execute(&pool)
.await
.unwrap();
assert!(revision().await.unwrap().revision() > revoked.revision());
}
#[tokio::test]
async fn usage_accepts_zero_and_the_backends_unit_ceiling() {
let _guard = DB_LOCK.lock().await;
for leased in [false, true] {
let ceiling = i64::MAX as u64;
let Some(store) =
store_with_balance(full_grant_policy(), if leased { ceiling } else { 0 }).await
else {
return;
};
let lease = if leased {
Some(
store
.acquire(ACCOUNT, CostUnits(ceiling), TTL, t(0))
.await
.unwrap()
.grant,
)
} else {
None
};
let event = |id, units| {
if let Some(lease) = &lease {
usage(lease, id, units, 0)
} else {
overage_usage(ACCOUNT, id, units, 0)
}
};
let batch = [
event(1, 0),
event(2, ceiling),
event(2, u64::MAX),
event(3, 0),
];
let report = UsageSink::ingest(&*store, &batch, t(1)).await.unwrap();
report.validate(batch.len()).unwrap();
assert_eq!(
(report.accepted, report.duplicate, report.rejected),
(3, 1, 0)
);
assert_eq!(
store.usage_recorded(ACCOUNT).await.unwrap(),
CostUnits(ceiling)
);
assert_conserved(&store).await;
}
}
#[tokio::test]
async fn nanosecond_lease_boundaries_preserve_release_reclaim_and_consolidation() {
let _guard = DB_LOCK.lock().await;
let ns = SignedDuration::from_nanos;
for (now, ttl, grace, max_ttl) in [
(t(100), ns(1), SignedDuration::from_secs(30), TTL),
(t(100) + ns(1), ns(999), ns(1), TTL),
(t(-100) - ns(2), ns(1), ns(1), TTL),
(t(0) - ns(2), ns(1), ns(0), TTL),
(t(0), ns(2_001), ns(999), ns(1_001)),
(Timestamp::MIN, ns(1), ns(1), TTL),
(Timestamp::MAX - ns(1), ns(1), ns(0), TTL),
(Timestamp::MAX - ns(1_001), ns(1), ns(999), TTL),
] {
let policy = GrantPolicy {
max_ttl,
reclaim_grace: grace,
..full_grant_policy()
};
let Some(store) = store_with_balance(policy, 100).await else {
return;
};
let lease = store
.acquire(ACCOUNT, CostUnits(10), ttl, now)
.await
.unwrap()
.grant;
let expiry = now.checked_add(ttl.min(max_ttl)).unwrap();
assert_eq!(lease.expires_at, expiry);
let lease = store
.consolidate(
lease.lease_id,
lease.fencing_token,
CostUnits(10),
CostUnits(10),
CostUnits::ZERO,
ttl,
now,
)
.await
.unwrap()
.grant;
assert_eq!(lease.expires_at, expiry);
let released = store
.acquire(ACCOUNT, CostUnits(10), ttl, now)
.await
.unwrap()
.grant;
let deadline = expiry.checked_add(grace).unwrap();
let before = deadline - ns(1);
assert!(
store.reclaim_expired(before).await.unwrap().is_empty(),
"early reclaim: {now}, {ttl}, {grace}"
);
store
.release(
released.lease_id,
released.fencing_token,
CostUnits(10),
before,
)
.await
.expect("release remains valid until the exact deadline");
let c = store.conservation(ACCOUNT).await.unwrap().unwrap();
assert!(c.holds());
assert_eq!(c.balance, CostUnits(90));
assert_eq!(c.active_lease_grants, CostUnits(10));
assert_eq!(
store
.release(lease.lease_id, lease.fencing_token, CostUnits(10), deadline)
.await
.unwrap_err(),
AllocateError::LeaseNotActive
);
let reclaimed = store.reclaim_expired(deadline).await.unwrap();
assert_eq!(reclaimed.len(), 1);
assert_eq!(reclaimed[0].lease_id, lease.lease_id);
assert_eq!(reclaimed[0].forfeited, CostUnits(10));
assert!(store.reclaim_expired(deadline).await.unwrap().is_empty());
let c = store.conservation(ACCOUNT).await.unwrap().unwrap();
assert!(c.holds());
assert_eq!(c.balance, CostUnits(90));
assert_eq!(c.active_lease_grants, CostUnits::ZERO);
assert_eq!(c.settlement_loss, CostUnits(10));
}
}
#[tokio::test]
async fn a_grace_deadline_beyond_timestamp_max_never_reclaims_early() {
let _guard = DB_LOCK.lock().await;
for grace in [SignedDuration::from_nanos(2), SignedDuration::MAX] {
let policy = GrantPolicy {
reclaim_grace: grace,
..full_grant_policy()
};
let Some(store) = store_with_balance(policy, 100).await else {
return;
};
let now = Timestamp::MAX - SignedDuration::from_nanos(2);
let lease = store
.acquire(ACCOUNT, CostUnits(10), SignedDuration::from_nanos(1), now)
.await
.unwrap()
.grant;
assert!(
store
.reclaim_expired(Timestamp::MIN)
.await
.unwrap()
.is_empty()
);
assert!(
store
.reclaim_expired(Timestamp::MAX)
.await
.unwrap()
.is_empty()
);
store
.release(
lease.lease_id,
lease.fencing_token,
CostUnits(10),
Timestamp::MAX,
)
.await
.expect("a deadline beyond the timestamp domain has not elapsed");
let c = store.conservation(ACCOUNT).await.unwrap().unwrap();
assert!(c.holds());
assert_eq!(c.balance, CostUnits(100));
assert_eq!(c.active_lease_grants, CostUnits::ZERO);
}
}
#[tokio::test]
async fn an_unrepresentable_replacement_expiry_leaves_the_original_grant_untouched() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 100).await else {
return;
};
let lease = store
.acquire(ACCOUNT, CostUnits(10), TTL, t(0))
.await
.unwrap()
.grant;
for ttl in [SignedDuration::from_nanos(1), SignedDuration::MAX] {
assert!(matches!(
store
.acquire(ACCOUNT, CostUnits(10), ttl, Timestamp::MAX)
.await,
Err(AllocateError::Storage(_))
));
assert!(matches!(
store
.consolidate(
lease.lease_id,
lease.fencing_token,
CostUnits(10),
CostUnits(10),
CostUnits::ZERO,
ttl,
Timestamp::MAX
)
.await,
Err(AllocateError::Storage(_))
));
let c = store.conservation(ACCOUNT).await.unwrap().unwrap();
assert!(c.holds());
assert_eq!(c.balance, CostUnits(90));
assert_eq!(c.active_lease_grants, CostUnits(10));
}
store
.release(lease.lease_id, lease.fencing_token, CostUnits(10), t(1))
.await
.unwrap();
let c = store.conservation(ACCOUNT).await.unwrap().unwrap();
assert!(c.holds());
assert_eq!(c.balance, CostUnits(100));
assert_eq!(c.active_lease_grants, CostUnits::ZERO);
}
#[path = "../../tollgate-store/tests/support/credential_expiry.rs"]
mod credential_expiry;
#[tokio::test]
async fn credential_expiry_is_exact_in_directory_and_every_page() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 100).await else {
return;
};
credential_expiry::exact_expiry(&*store, ACCOUNT).await;
}
#[path = "../../tollgate-store/tests/support/account_keys.rs"]
mod account_keys;
#[tokio::test]
async fn account_key_listing_is_scoped_ordered_and_paged() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(GrantPolicy::default(), 1_000).await else {
return;
};
account_keys::listing_is_scoped_ordered_and_paged(&*store).await;
}
#[tokio::test]
async fn account_key_listing_separates_expiry_from_revocation() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(GrantPolicy::default(), 1_000).await else {
return;
};
account_keys::listing_separates_expiry_from_revocation(&*store).await;
}
#[tokio::test]
async fn the_active_key_bound_counts_only_live_credentials() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(GrantPolicy::default(), 1_000).await else {
return;
};
account_keys::the_active_bound_counts_only_live_credentials(&*store).await;
}
#[tokio::test]
async fn the_active_key_bound_is_per_account_and_preserves_issuance_rules() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(GrantPolicy::default(), 1_000).await else {
return;
};
account_keys::the_bound_is_per_account_and_preserves_issuance_rules(&*store).await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn concurrent_issuers_cannot_exceed_the_active_key_bound() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(GrantPolicy::default(), 1_000).await else {
return;
};
account_keys::concurrent_issuers_cannot_exceed_the_bound(store).await;
}
#[tokio::test]
async fn without_the_account_lock_two_issuers_both_see_room_under_the_bound() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(GrantPolicy::default(), 1_000).await else {
return;
};
let pool = tollgate_store_postgres::test_support::pool(&store);
let count_sql = tollgate_store_postgres::test_support::issuance_sql::LIVE_KEY_COUNT;
let bytes = |value: u128| value.to_be_bytes().to_vec();
let account = bytes(ACCOUNT.0);
let mut first = pool.begin().await.unwrap();
let mut second = pool.begin().await.unwrap();
let seen_by_first: i64 = sqlx::query_scalar(count_sql)
.bind(account.clone())
.bind(0i64)
.bind(0i16)
.fetch_one(&mut *first)
.await
.unwrap();
sqlx::query(
"INSERT INTO tollgate_credential_keys
(key_id, account_id, principal, digest, not_after_floor_us,
not_after_submicro_ns, not_after_is_lower_bound, revoked_at_us)
VALUES ($1, $2, $3, $4, NULL, NULL, FALSE, NULL)",
)
.bind(bytes(1))
.bind(account.clone())
.bind(bytes(1))
.bind(vec![0u8; 32])
.execute(&mut *first)
.await
.unwrap();
let seen_by_second: i64 = sqlx::query_scalar(count_sql)
.bind(account.clone())
.bind(0i64)
.bind(0i16)
.fetch_one(&mut *second)
.await
.unwrap();
first.commit().await.unwrap();
second.rollback().await.unwrap();
assert_eq!(seen_by_first, 0);
assert_eq!(
seen_by_second, 0,
"both issuers read a count below a bound of one, so both would insert — \
which is why the account row is taken FOR UPDATE before the count"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn mixed_issuers_report_duplicates() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(GrantPolicy::default(), 1_000).await else {
return;
};
account_keys::mixed_issuers_report_duplicates(store).await;
}
#[tokio::test]
async fn mixed_issuers_waiting_on_an_account_report_duplicates() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(GrantPolicy::default(), 1_000).await else {
return;
};
let pool = tollgate_store_postgres::test_support::pool(&store);
for same_key in [true, false] {
let id = if same_key { 10 } else { 20 };
let first = key(id, id, 0xa1, None);
let second = key(
if same_key { id } else { id + 1 },
if same_key { id + 1 } else { id },
0xb2,
None,
);
let mut blocker = pool.begin().await.unwrap();
sqlx::query(tollgate_store_postgres::test_support::issuance_sql::ACCOUNT_LOCK)
.bind(ACCOUNT.0.to_be_bytes().to_vec())
.fetch_one(&mut *blocker)
.await
.unwrap();
let blocker_pid: i32 = sqlx::query_scalar("SELECT pg_backend_pid()")
.fetch_one(&mut *blocker)
.await
.unwrap();
let mut attempts = tokio::task::JoinSet::new();
let bounded_store = Arc::clone(&store);
attempts.spawn(async move {
bounded_store
.insert_key_within(first, NonZeroUsize::new(32).unwrap(), t(100))
.await
});
wait_for_blocked_operations(&pool, blocker_pid, 1).await;
let unbounded_store = Arc::clone(&store);
attempts.spawn(async move { unbounded_store.insert_key(second).await });
wait_for_blocked_operations(&pool, blocker_pid, 2).await;
blocker.commit().await.unwrap();
let outcomes = tokio::time::timeout(std::time::Duration::from_secs(10), async {
[
attempts.join_next().await.unwrap().unwrap(),
attempts.join_next().await.unwrap().unwrap(),
]
})
.await
.expect("both issuers finish after the account lock is released");
assert!(
matches!(
outcomes,
[Ok(()), Err(KeyError::AlreadyExists)] | [Err(KeyError::AlreadyExists), Ok(())]
),
"one success and one duplicate, including principal collisions: {outcomes:?}"
);
}
assert_eq!(store.active_keys(t(100)).await.unwrap().len(), 2);
}
async fn wait_for_blocked_operations(pool: &sqlx::PgPool, blocker_pid: i32, expected: i64) {
tokio::time::timeout(std::time::Duration::from_secs(10), async {
loop {
let waiting: i64 = sqlx::query_scalar(
"WITH RECURSIVE blocked(pid) AS (
SELECT $1::integer
UNION
SELECT a.pid FROM pg_stat_activity a JOIN blocked b
ON b.pid = ANY(pg_blocking_pids(a.pid))
) SELECT count(*) - 1 FROM blocked",
)
.bind(blocker_pid)
.fetch_one(pool)
.await
.unwrap();
if waiting == expected {
break;
}
tokio::task::yield_now().await;
}
})
.await
.expect("issuers reach the account lock queue");
}
#[path = "../../tollgate-store/tests/support/account_view.rs"]
mod account_view;
#[tokio::test]
async fn an_unknown_account_view_is_absent_not_empty() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(GrantPolicy::default(), 1_000).await else {
return;
};
account_view::an_unknown_account_is_absent_not_empty(&*store).await;
}
#[tokio::test]
async fn the_account_view_reports_what_was_administered() {
let _guard = DB_LOCK.lock().await;
let Some(store) = empty_store(GrantPolicy::default()).await else {
return;
};
account_view::the_view_reports_what_was_administered(&*store).await;
}
#[tokio::test]
async fn funding_out_on_lease_is_not_reported_as_usage() {
let _guard = DB_LOCK.lock().await;
let Some(store) = empty_store(full_grant_policy()).await else {
return;
};
account_view::funding_out_on_lease_is_not_reported_as_usage(&*store).await;
}
#[tokio::test]
async fn setting_a_budget_schedule_reports_what_it_replaced() {
let _guard = DB_LOCK.lock().await;
let Some(store) = empty_store(GrantPolicy::default()).await else {
return;
};
account_view::setting_a_schedule_reports_what_it_replaced(&*store).await;
}
#[path = "../../tollgate-store/tests/support/admin_receipts.rs"]
mod admin_receipts;
#[tokio::test]
async fn concurrent_budget_receipts_form_one_serial_history() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(GrantPolicy::default(), 1_000).await else {
return;
};
admin_receipts::budget_receipts_form_a_serial_history(store).await;
}
#[tokio::test]
async fn credential_receipts_identify_issuance_and_concurrent_revocation() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(GrantPolicy::default(), 1_000).await else {
return;
};
admin_receipts::credential_receipts_capture_lifecycle(store).await;
}
#[tokio::test]
async fn queued_budget_updates_report_the_locked_predecessor() {
use tollgate_store::AdminState;
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(GrantPolicy::default(), 1_000).await else {
return;
};
store
.set_budget_schedule(ACCOUNT, Some(admin_receipts::schedule(100)))
.await
.unwrap();
let pool = tollgate_store_postgres::test_support::pool(&store);
let mut blocker = pool.begin().await.unwrap();
sqlx::query(tollgate_store_postgres::test_support::issuance_sql::ACCOUNT_LOCK)
.bind(ACCOUNT.0.to_be_bytes().to_vec())
.fetch_one(&mut *blocker)
.await
.unwrap();
let pid: i32 = sqlx::query_scalar("SELECT pg_backend_pid()")
.fetch_one(&mut *blocker)
.await
.unwrap();
let mut tasks = tokio::task::JoinSet::new();
for (waiting, allowance) in [(1, 200), (2, 300)] {
let store = Arc::clone(&store);
tasks.spawn(async move {
(
allowance,
store
.set_budget_schedule(ACCOUNT, Some(admin_receipts::schedule(allowance)))
.await,
)
});
wait_for_blocked_operations(&pool, pid, waiting).await;
}
blocker.commit().await.unwrap();
tokio::time::timeout(std::time::Duration::from_secs(10), async {
while let Some(result) = tasks.join_next().await {
let (allowance, receipt) = result.unwrap();
let receipt = receipt.unwrap();
assert_eq!(
receipt.before,
AdminState::Budget {
schedule: Some(admin_receipts::schedule(allowance - 100))
}
);
assert_eq!(
receipt.after,
AdminState::Budget {
schedule: Some(admin_receipts::schedule(allowance))
}
);
}
})
.await
.expect("queued budget mutations complete");
assert_eq!(
store.account_view(ACCOUNT).await.unwrap().unwrap().schedule,
Some(admin_receipts::schedule(300))
);
}
#[tokio::test]
async fn credential_auditing_refuses_a_corrupt_owner_without_retiring_it() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(GrantPolicy::default(), 1_000).await else {
return;
};
let pool = tollgate_store_postgres::test_support::pool(&store);
sqlx::query(
"INSERT INTO tollgate_accounts
(account_id, balance, deposited, status, capacity_class, next_fence,
usage_recorded, settlement_loss, overage_recorded)
VALUES ($1, 0, 0, 'Active', 'Assured', 1, 0, 0, 0)",
)
.bind(vec![1u8])
.execute(&pool)
.await
.unwrap();
let key = KeyId(900);
sqlx::query(
"INSERT INTO tollgate_credential_keys
(key_id, account_id, principal, digest, not_after_is_lower_bound)
VALUES ($1, $2, $1, $3, FALSE)",
)
.bind(key.0.to_be_bytes().to_vec())
.bind(vec![1u8])
.bind(vec![0u8; 32])
.execute(&pool)
.await
.unwrap();
assert!(matches!(
store.revoke_key_audited(key, t(100)).await,
Err(KeyError::Storage(_))
));
let revoked: bool = sqlx::query_scalar(
"SELECT revoked_at_us IS NOT NULL
FROM tollgate_credential_keys WHERE key_id = $1",
)
.bind(key.0.to_be_bytes().to_vec())
.fetch_one(&pool)
.await
.unwrap();
assert!(
!revoked,
"invalid audit evidence must not commit a retirement"
);
}
fn held_elsewhere(remaining: u64) -> AllocateError {
AllocateError::BalanceInsufficient(tollgate_core::BalanceShortfall {
remaining: CostUnits(remaining),
period_end: None,
})
}
#[tokio::test]
async fn grants_report_ledger_remaining_including_outstanding_leases() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 100).await else {
return;
};
let first = store
.acquire(ACCOUNT, CostUnits(40), TTL, t(0))
.await
.unwrap();
assert_eq!(
first.funding,
Some(tollgate_core::BalanceShortfall {
remaining: CostUnits(100),
period_end: None,
})
);
store
.ingest(&[usage(&first.grant, 1, 30, 0)], t(0))
.await
.unwrap();
let second = store
.acquire(ACCOUNT, CostUnits(60), TTL, t(1))
.await
.unwrap();
assert_eq!(second.grant.units, CostUnits(60));
assert_eq!(
second.funding.map(|funding| funding.remaining),
Some(CostUnits(70)),
"recorded usage is consumed; both leases' units are still the account's"
);
assert_conserved(&store).await;
}
#[tokio::test]
async fn a_partial_balance_refusal_reports_remaining_funding() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 100).await else {
return;
};
let held = store
.acquire(ACCOUNT, CostUnits(100), TTL, t(0))
.await
.unwrap()
.grant;
store.ingest(&[usage(&held, 1, 99, 0)], t(0)).await.unwrap();
let tail = store
.consolidate(
held.lease_id,
held.fencing_token,
CostUnits(1),
CostUnits(252),
CostUnits::ZERO,
TTL,
t(1),
)
.await
.unwrap();
assert_eq!(tail.grant.units, CostUnits(1));
assert_eq!(
tail.funding.map(|funding| funding.remaining),
Some(CostUnits(1))
);
assert_eq!(
store
.acquire(ACCOUNT, CostUnits(252), TTL, t(2))
.await
.unwrap_err(),
held_elsewhere(1)
);
assert_conserved(&store).await;
}
#[tokio::test]
async fn shortfall_evidence_names_the_stored_period() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 0).await else {
return;
};
AdminStore::set_budget_schedule(&*store, ACCOUNT, Some(monthly(100)))
.await
.unwrap();
roll(&store, FEB).await.unwrap();
let evidence = tollgate_core::BalanceShortfall {
remaining: CostUnits(100),
period_end: Some(t(MAR)),
};
let held = store
.acquire(ACCOUNT, CostUnits(100), TTL, t(FEB))
.await
.unwrap();
assert_eq!(held.funding, Some(evidence));
assert_eq!(
store
.acquire(ACCOUNT, CostUnits(1), TTL, t(FEB + 1))
.await
.unwrap_err(),
AllocateError::BalanceInsufficient(evidence)
);
assert_conserved(&store).await;
}
#[tokio::test]
async fn an_insufficient_balance_recovers_when_another_instance_returns_its_lease() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 100).await else {
return;
};
let held = store
.acquire(ACCOUNT, CostUnits(100), TTL, t(0))
.await
.unwrap()
.grant;
assert_eq!(held.units, CostUnits(100));
assert_eq!(
store
.acquire(ACCOUNT, CostUnits(100), TTL, t(1))
.await
.unwrap_err(),
held_elsewhere(100),
);
store
.release(held.lease_id, held.fencing_token, held.units, t(2))
.await
.unwrap();
let recovered = store
.acquire(ACCOUNT, CostUnits(100), TTL, t(3))
.await
.unwrap()
.grant;
assert_eq!(recovered.units, held.units);
assert_conserved(&store).await;
}
#[tokio::test]
async fn a_reclaimed_lease_forfeits_its_funding_instead_of_restoring_it() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 100).await else {
return;
};
let held = store
.acquire(ACCOUNT, CostUnits(100), TTL, t(0))
.await
.unwrap()
.grant;
assert_eq!(
store
.acquire(ACCOUNT, CostUnits(100), TTL, t(1))
.await
.unwrap_err(),
held_elsewhere(100),
);
let reclaimed = store
.reclaim_expired_batch(t(61), NonZeroUsize::new(1).unwrap())
.await
.unwrap();
assert_eq!(reclaimed.reclaimed().len(), 1);
assert_eq!(reclaimed.reclaimed()[0].forfeited, held.units);
assert_eq!(
store
.acquire(ACCOUNT, CostUnits(100), TTL, t(62))
.await
.unwrap_err(),
AllocateError::BalanceExhausted(tollgate_core::BalanceExhaustion { period_end: None })
);
assert_conserved(&store).await;
}
#[tokio::test]
async fn exhaustion_requires_recorded_consumption_and_survives_lease_expiry() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 100).await else {
return;
};
let held = store
.acquire(ACCOUNT, CostUnits(100), TTL, t(0))
.await
.unwrap()
.grant;
assert_eq!(
store
.acquire(ACCOUNT, CostUnits(1), TTL, t(1))
.await
.unwrap_err(),
held_elsewhere(100)
);
store
.ingest(&[usage(&held, 1, 100, 1)], t(1))
.await
.unwrap();
let exhausted =
AllocateError::BalanceExhausted(tollgate_core::BalanceExhaustion { period_end: None });
assert_eq!(
store
.acquire(ACCOUNT, CostUnits(1), TTL, t(2))
.await
.unwrap_err(),
exhausted
);
assert_eq!(
store
.consolidate(
held.lease_id,
held.fencing_token,
CostUnits::ZERO,
CostUnits(1),
CostUnits::ZERO,
TTL,
t(2)
)
.await
.unwrap_err(),
exhausted
);
store
.reclaim_expired_batch(t(61), NonZeroUsize::new(1).unwrap())
.await
.unwrap();
assert_eq!(
store
.acquire(ACCOUNT, CostUnits(1), TTL, t(62))
.await
.unwrap_err(),
exhausted
);
AdminStore::deposit(&*store, ACCOUNT, CostUnits(10))
.await
.unwrap();
assert_eq!(
store
.acquire(ACCOUNT, CostUnits(10), TTL, t(63))
.await
.unwrap()
.grant
.units,
CostUnits(10)
);
assert_conserved(&store).await;
}
#[tokio::test]
async fn zero_requested_units_never_produce_exhaustion_evidence() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 0).await else {
return;
};
assert_eq!(
store
.acquire(ACCOUNT, CostUnits::ZERO, TTL, t(0))
.await
.unwrap_err(),
AllocateError::InsufficientBalance
);
assert_eq!(
store
.acquire(ACCOUNT, CostUnits(1), TTL, t(0))
.await
.unwrap_err(),
AllocateError::BalanceExhausted(tollgate_core::BalanceExhaustion { period_end: None })
);
assert_conserved(&store).await;
}
#[tokio::test]
async fn exhaustion_evidence_names_the_stored_period_and_rollover_restores_funding() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 0).await else {
return;
};
AdminStore::set_budget_schedule(&*store, ACCOUNT, Some(monthly(100)))
.await
.unwrap();
roll(&store, FEB).await.unwrap();
let held = store
.acquire(ACCOUNT, CostUnits(100), TTL, t(FEB))
.await
.unwrap()
.grant;
store
.ingest(&[usage(&held, 1, 100, FEB)], t(FEB))
.await
.unwrap();
let exhausted = AllocateError::BalanceExhausted(tollgate_core::BalanceExhaustion {
period_end: Some(t(MAR)),
});
assert_eq!(
store
.acquire(ACCOUNT, CostUnits(1), TTL, t(FEB + 1))
.await
.unwrap_err(),
exhausted
);
assert_eq!(
store
.acquire(ACCOUNT, CostUnits(1), TTL, t(MAR))
.await
.unwrap_err(),
exhausted
);
roll(&store, MAR).await.unwrap();
assert_eq!(
store
.acquire(ACCOUNT, CostUnits(100), TTL, t(MAR))
.await
.unwrap()
.grant
.units,
CostUnits(100)
);
assert_conserved(&store).await;
}
#[tokio::test]
async fn a_provisioned_account_is_born_unfunded_suspended_and_best_effort() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(GrantPolicy::default(), 1_000).await else {
return;
};
let provisioned = AccountId(39);
let receipt = AdminStore::create_provisioned_account(&*store, provisioned)
.await
.unwrap();
assert_eq!(receipt.before, AdminState::Absent);
assert_eq!(
receipt.after,
AdminState::AccountCreated {
initial_balance: CostUnits::ZERO,
status: AccountStatus::Suspended,
capacity_class: CapacityClass::BestEffort,
origin: AdminAuthority::Provisioner,
}
);
let view = AdminStore::account_view(&*store, provisioned)
.await
.unwrap()
.unwrap();
assert_eq!(view.status, AccountStatus::Suspended);
assert_eq!(view.capacity_class, CapacityClass::BestEffort);
assert_eq!(view.origin, AdminAuthority::Provisioner);
assert_eq!(view.status_set_by, AdminAuthority::Provisioner);
assert_eq!(view.conservation.deposited, CostUnits::ZERO);
assert_eq!(view.conservation.balance, CostUnits::ZERO);
assert_eq!(
AdminStore::create_provisioned_account(&*store, provisioned)
.await
.unwrap_err(),
CreateAccountError::AlreadyExists
);
assert_eq!(
AdminStore::create_provisioned_account(&*store, ACCOUNT)
.await
.unwrap_err(),
CreateAccountError::AlreadyExists,
"an operator's account is never re-created as a provisioner's"
);
let operator = AdminStore::account_view(&*store, ACCOUNT)
.await
.unwrap()
.unwrap();
assert_eq!(operator.origin, AdminAuthority::Operator);
assert_eq!(operator.status_set_by, AdminAuthority::Operator);
}
#[tokio::test]
async fn provisioned_activation_honours_provenance_and_operator_holds() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(GrantPolicy::default(), 1_000).await else {
return;
};
assert_eq!(
AdminStore::activate_provisioned(&*store, AccountId(999))
.await
.unwrap_err(),
SetStatusError::UnknownAccount
);
AdminStore::set_account_status(&*store, ACCOUNT, AccountStatus::Suspended)
.await
.unwrap();
assert_eq!(
AdminStore::activate_provisioned(&*store, ACCOUNT)
.await
.unwrap_err(),
SetStatusError::NotProvisioned,
"an operator's account is out of a provisioner's reach"
);
assert_eq!(
AdminStore::account_view(&*store, ACCOUNT)
.await
.unwrap()
.unwrap()
.status,
AccountStatus::Suspended
);
let provisioned = AccountId(39);
AdminStore::create_provisioned_account(&*store, provisioned)
.await
.unwrap();
let activated = AdminStore::activate_provisioned(&*store, provisioned)
.await
.unwrap();
assert_eq!(
activated.before,
AdminState::Status {
status: AccountStatus::Suspended,
set_by: AdminAuthority::Provisioner,
}
);
assert_eq!(
activated.after,
AdminState::Status {
status: AccountStatus::Active,
set_by: AdminAuthority::Provisioner,
}
);
let repeated = AdminStore::activate_provisioned(&*store, provisioned)
.await
.unwrap();
assert_eq!(repeated.before, repeated.after);
let suspended = AdminStore::set_account_status(&*store, provisioned, AccountStatus::Suspended)
.await
.unwrap();
assert_eq!(
suspended.after,
AdminState::Status {
status: AccountStatus::Suspended,
set_by: AdminAuthority::Operator,
}
);
assert_eq!(
AdminStore::activate_provisioned(&*store, provisioned)
.await
.unwrap_err(),
SetStatusError::OperatorHold
);
let held = AdminStore::account_view(&*store, provisioned)
.await
.unwrap()
.unwrap();
assert_eq!(held.status, AccountStatus::Suspended);
assert_eq!(held.status_set_by, AdminAuthority::Operator);
AdminStore::set_account_status(&*store, provisioned, AccountStatus::Active)
.await
.unwrap();
let noop = AdminStore::activate_provisioned(&*store, provisioned)
.await
.unwrap();
assert_eq!(noop.before, noop.after);
assert_eq!(
noop.after,
AdminState::Status {
status: AccountStatus::Active,
set_by: AdminAuthority::Operator,
}
);
AdminStore::set_account_status(&*store, provisioned, AccountStatus::Closed)
.await
.unwrap();
assert_eq!(
AdminStore::activate_provisioned(&*store, provisioned)
.await
.unwrap_err(),
SetStatusError::AccountClosed
);
}
#[tokio::test]
async fn accounts_predating_provenance_belong_to_operators() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(GrantPolicy::default(), 1_000).await else {
return;
};
let legacy = AccountId(40);
sqlx::query(
"INSERT INTO tollgate_accounts
(account_id, balance, deposited, status, next_fence, usage_recorded, settlement_loss)
VALUES ($1, 0, 0, 'Suspended', 1, 0, 0)",
)
.bind(legacy.0.to_be_bytes().to_vec())
.execute(&corruption_pool().await)
.await
.unwrap();
let view = AdminStore::account_view(&*store, legacy)
.await
.unwrap()
.unwrap();
assert_eq!(view.origin, AdminAuthority::Operator);
assert_eq!(view.status_set_by, AdminAuthority::Operator);
assert_eq!(
AdminStore::activate_provisioned(&*store, legacy)
.await
.unwrap_err(),
SetStatusError::NotProvisioned
);
}
#[tokio::test]
async fn account_key_listing_distinguishes_unknown_from_empty() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(GrantPolicy::default(), 1_000).await else {
return;
};
account_keys::listing_distinguishes_unknown_from_empty(&*store).await;
}
#[tokio::test]
async fn deposits_accept_the_unit_ceiling_and_refuse_overflow() {
let _guard = DB_LOCK.lock().await;
let ceiling = u64::try_from(i64::MAX).unwrap();
let Some(store) = store_with_balance(full_grant_policy(), ceiling - 1).await else {
return;
};
AdminStore::deposit(&*store, ACCOUNT, CostUnits(1))
.await
.unwrap();
let before = store.conservation(ACCOUNT).await.unwrap().unwrap();
assert_eq!(before.balance, CostUnits(ceiling));
assert_eq!(before.deposited, CostUnits(ceiling));
assert!(before.holds());
for units in [1, u64::MAX] {
for _ in 0..2 {
assert_eq!(
AdminStore::deposit(&*store, ACCOUNT, CostUnits(units))
.await
.unwrap_err(),
AllocateError::BalanceOverflow,
);
assert_eq!(store.conservation(ACCOUNT).await.unwrap().unwrap(), before);
}
}
assert_eq!(
AdminStore::deposit(&*store, AccountId(999), CostUnits(1))
.await
.unwrap_err(),
AllocateError::UnknownAccount,
);
}
#[tokio::test]
async fn a_deposit_outside_bigint_is_a_permanent_refusal() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(full_grant_policy(), 0).await else {
return;
};
let before = store.conservation(ACCOUNT).await.unwrap().unwrap();
for units in [i64::MAX as u64 + 1, u64::MAX] {
assert_eq!(
AdminStore::deposit(&*store, ACCOUNT, CostUnits(units))
.await
.unwrap_err(),
AllocateError::BalanceOverflow,
);
assert_eq!(store.conservation(ACCOUNT).await.unwrap().unwrap(), before);
}
}
#[tokio::test]
async fn repeating_operator_suspension_establishes_a_hold() {
let _guard = DB_LOCK.lock().await;
let Some(store) = store_with_balance(GrantPolicy::default(), 1_000).await else {
return;
};
let provisioned = AccountId(390);
AdminStore::create_provisioned_account(&*store, provisioned)
.await
.unwrap();
let receipt = AdminStore::set_account_status(&*store, provisioned, AccountStatus::Suspended)
.await
.unwrap();
assert_eq!(
receipt.before,
AdminState::Status {
status: AccountStatus::Suspended,
set_by: AdminAuthority::Provisioner,
}
);
assert_eq!(
receipt.after,
AdminState::Status {
status: AccountStatus::Suspended,
set_by: AdminAuthority::Operator,
}
);
assert_eq!(
AdminStore::activate_provisioned(&*store, provisioned)
.await
.unwrap_err(),
SetStatusError::OperatorHold
);
let view = AdminStore::account_view(&*store, provisioned)
.await
.unwrap()
.unwrap();
assert_eq!(view.status, AccountStatus::Suspended);
assert_eq!(view.status_set_by, AdminAuthority::Operator);
}