use std::sync::Arc;
use std::time::Duration;
use jiff::SignedDuration;
use tollgate_admission::NoGate;
use tollgate_client::{
AccountLeaseConfig, Clock, InstanceRuntime, InstanceRuntimeConfig, SnapshotManagerConfig,
SystemClock, TrackedPrincipals, UsageWriterConfig,
};
use tollgate_core::{
AccountId, AccountSnapshot, AccountStatus, CapacityClass, CostTable, CostUnits, DenyReason,
Generation, LocalSharding, OpIndex, PermissionBits, Principal, PublishableSnapshot, RequestId,
ResolvedLimits,
};
use tollgate_store::{AccountConfig, GrantPolicy, MemoryStore};
#[derive(Clone, Copy)]
enum Op {
Quote,
Report,
}
impl OpIndex for Op {
fn index(&self) -> usize {
*self as usize
}
}
const CALL: PermissionBits = PermissionBits::bit(0);
const EXPORT: PermissionBits = PermissionBits::bit(1);
#[tokio::test]
async fn getting_started() {
let store = MemoryStore::new(GrantPolicy::default()).expect("default grant policy is valid");
let clock: Arc<dyn Clock> = Arc::new(SystemClock);
let account = AccountId(1);
store.create_account(AccountConfig {
account_id: account,
initial_balance: CostUnits(10_000),
status: AccountStatus::Active,
capacity_class: CapacityClass::Assured,
});
let costs = CostTable::builder(CostUnits(5), CostUnits(5))
.weight(&Op::Quote, CostUnits(1))
.weight(&Op::Report, CostUnits(20))
.build();
let valid_until = clock
.now()
.checked_add(SignedDuration::from_hours(1))
.expect("an hour from now is a valid timestamp");
let snapshot = AccountSnapshot::builder(
account,
Generation(1),
AccountStatus::Active,
valid_until,
CALL,
ResolvedLimits::new(64),
Arc::new(costs),
)
.build();
let caller = Principal(42);
store
.publish_snapshot(
caller,
PublishableSnapshot::try_new(Arc::new(snapshot)).expect("a consistent snapshot"),
)
.expect("the account exists");
let config = InstanceRuntimeConfig {
snapshots: SnapshotManagerConfig {
principals: TrackedPrincipals::All { seed: vec![caller] },
refresh_interval: Duration::from_secs(30),
unknown_ttl: SignedDuration::from_secs(60),
revoked_ttl: SignedDuration::from_hours(1),
retry_backoff: Duration::from_millis(200),
max_concurrent_fetches: 16,
fetch_timeout: Duration::from_secs(5),
enumeration_timeout: Duration::from_secs(30),
},
leases: AccountLeaseConfig {
target_grant: CostUnits(1_000),
low_water: CostUnits(100),
lease_ttl: SignedDuration::from_secs(60),
expiry_safety_margin: SignedDuration::from_secs(2),
poll_interval: Duration::from_millis(20),
store_call_timeout: Duration::from_secs(5),
shutdown_release_deadline: Duration::from_secs(5),
},
usage: UsageWriterConfig {
queue_capacity: 4_096,
max_batch: 256,
flush_interval: Duration::from_millis(25),
retry_backoff: Duration::from_millis(50),
shutdown_drain_deadline: Duration::from_secs(5),
ingest_timeout: Duration::from_secs(5),
},
sharding: LocalSharding::SINGLE,
snapshot_history_capacity:
tollgate_admission::ArcSwapSnapshotMap::DEFAULT_GENERATION_CAPACITY,
idle_account_linger: Duration::from_secs(1),
manager_restart_backoff: Duration::from_millis(200),
shutdown_deadline: Duration::from_secs(15),
};
let (runtime, handle) = InstanceRuntime::spawn(
store.clone(),
store.clone(),
store.clone(),
Arc::clone(&clock),
config,
)
.expect("a valid runtime configuration");
while !handle.readiness(clock.now()).is_ready() {
tokio::time::sleep(Duration::from_millis(10)).await;
}
let now = clock.now();
let context = handle
.begin(caller, CALL, now)
.expect("a known, active caller with the permission");
let slot = handle
.recorder()
.try_reserve()
.expect("the usage queue has room");
let pending = context
.admit(&[(Op::Quote, 3), (Op::Report, 1)], slot, now)
.expect("within limits and funded");
let committed = pending
.acquire_capacity(&NoGate)
.expect("no capacity gate configured")
.commit(RequestId(1), clock.now())
.map_err(|(error, _released)| error)
.expect("the lease is still usable");
drop(committed);
let refused = handle.begin(caller, EXPORT, clock.now());
assert!(matches!(refused, Err(DenyReason::MissingPermission)));
let stranger = handle.begin(Principal(7), CALL, clock.now());
assert!(matches!(stranger, Err(DenyReason::UnknownPrincipal)));
let report = runtime.shutdown().await.expect("the runtime task finished");
assert!(!report.deadline_expired);
assert_eq!(report.usage.expect("usage drained").accepted, 1);
assert_eq!(store.usage_recorded(account), CostUnits(28));
assert_eq!(store.balance(account), CostUnits(9_972));
let ledger = store.conservation(account).expect("the account exists");
assert_eq!(ledger.deposited, CostUnits(10_000));
assert_eq!(ledger.settled_usage, CostUnits(28));
assert_eq!(ledger.active_lease_grants, CostUnits::ZERO);
}