1use std::collections::{BTreeSet, HashMap, HashSet};
3use std::sync::{Arc, Mutex};
4
5use jiff::Timestamp;
6use tokio::sync::watch;
7use tollgate_admission::{LeaseSlot, MapEntry};
8use tollgate_core::{
9 AccountId, AccountSnapshot, AccountStatus, EnforcementMode, LocalSharding, Principal,
10};
11
12const _: () = assert!(usize::BITS <= 64);
16
17#[derive(Debug, Clone, Default)]
21pub struct RuntimeFundingReport {
22 pub total_lease_remaining: Option<u128>,
25 pub total_overage_spent: u128,
29 pub total_overage_cap: Option<u128>,
34 pub earliest_lease_usable_until: Option<Timestamp>,
38}
39
40pub struct SlotRegistry {
43 inner: Mutex<Registry>,
44 sharding: LocalSharding,
45 changed: Option<watch::Sender<()>>,
46}
47
48#[derive(Default)]
49struct Registry {
50 slots: HashMap<AccountId, Arc<LeaseSlot>>,
51 principals: HashMap<Principal, Arc<AccountSnapshot>>,
52 members: HashMap<AccountId, HashMap<Principal, Arc<AccountSnapshot>>>,
53 deadlines: BTreeSet<(Timestamp, Principal)>,
54 dirty: BTreeSet<AccountId>,
55 tracked: HashSet<Principal>,
56 resolutions: HashMap<Principal, Timestamp>,
57}
58
59#[derive(Clone)]
60pub(crate) struct AccountBinding {
61 pub account: AccountId,
62 pub slot: Arc<LeaseSlot>,
63 pub snapshots: Vec<Arc<AccountSnapshot>>,
64}
65
66impl AccountBinding {
67 pub fn eligible(&self, now: Timestamp) -> bool {
68 self.snapshots
69 .iter()
70 .any(|snapshot| snapshot.status == AccountStatus::Active && now < snapshot.valid_until)
71 }
72
73 pub fn fundable(&self, now: Timestamp) -> bool {
74 self.snapshots.iter().any(|snapshot| {
75 snapshot.status == AccountStatus::Active
76 && now < snapshot.valid_until
77 && quota_usable(&self.slot, snapshot.enforcement_mode, now)
78 })
79 }
80}
81
82impl Default for SlotRegistry {
83 fn default() -> Self {
84 Self {
85 inner: Mutex::new(Registry::default()),
86 sharding: LocalSharding::SINGLE,
87 changed: None,
88 }
89 }
90}
91
92impl SlotRegistry {
93 pub(crate) fn track(&self, principal: Principal) {
94 if self.observes() {
95 self.inner
96 .lock()
97 .expect("slot registry poisoned")
98 .tracked
99 .insert(principal);
100 }
101 }
102
103 pub(crate) fn retain(&self, principals: &HashSet<Principal>) {
104 if self.observes() {
105 let mut inner = self.inner.lock().expect("slot registry poisoned");
106 inner.tracked.clone_from(principals);
107 #[allow(
110 clippy::disallowed_methods,
111 reason = "pure, total predicate: the visit order cannot change which entries survive"
112 )]
113 inner
114 .resolutions
115 .retain(|principal, _| principals.contains(principal));
116 }
117 }
118
119 pub(crate) fn resolution_counts(&self, now: Timestamp) -> (usize, usize) {
120 let inner = self.inner.lock().expect("slot registry poisoned");
121 #[allow(
122 clippy::disallowed_methods,
123 reason = "counts unresolved principals; a count does not depend on the order they are counted in"
124 )]
125 let unresolved = inner
126 .tracked
127 .iter()
128 .filter(|principal| {
129 inner
130 .resolutions
131 .get(principal)
132 .is_none_or(|until| now >= *until)
133 })
134 .count();
135 (inner.tracked.len(), unresolved)
136 }
137
138 pub(crate) fn funding(&self, now: Timestamp) -> RuntimeFundingReport {
139 let inner = self.inner.lock().expect("slot registry poisoned");
140 let mut report = RuntimeFundingReport::default();
141 #[allow(
142 clippy::disallowed_methods,
143 reason = "sums overage across slots with checked addition; the total does not depend on the order the addends arrive in"
144 )]
145 for slot in inner.slots.values() {
146 report.total_overage_spent = report
147 .total_overage_spent
148 .checked_add(u128::from(slot.overage().spent().get()))
149 .expect("at most usize::MAX u64 contributions fit u128");
150 if let Some(lease) = slot.load_observed() {
151 let sum = report
152 .total_lease_remaining
153 .unwrap_or(0)
154 .checked_add(u128::from(lease.remaining().get()))
155 .expect("at most usize::MAX u64 contributions fit u128");
156 report.total_lease_remaining = Some(sum);
157 }
158 }
159 #[allow(
160 clippy::disallowed_methods,
161 reason = "reduces each account's members to `any(..)` and `max(..)`; neither depends on the order they are visited in"
162 )]
163 for (account, members) in &inner.members {
164 if members
165 .values()
166 .any(|s| s.status == AccountStatus::Active && now < s.valid_until)
167 && let Some(lease) = inner.slots[account].load_observed()
168 {
169 let until = lease.usable_until();
170 report.earliest_lease_usable_until = Some(
171 report
172 .earliest_lease_usable_until
173 .map_or(until, |old| old.min(until)),
174 );
175 }
176 let cap = members
177 .values()
178 .filter(|s| s.status == AccountStatus::Active && now < s.valid_until)
179 .filter_map(|s| s.enforcement_mode.overage_cap())
180 .max();
181 if let Some(cap) = cap {
182 let sum = report
183 .total_overage_cap
184 .unwrap_or(0)
185 .checked_add(u128::from(cap.get()))
186 .expect("at most usize::MAX u64 contributions fit u128");
187 report.total_overage_cap = Some(sum);
188 }
189 }
190 report
191 }
192 #[must_use]
194 pub fn new() -> Arc<Self> {
195 Arc::new(Self::default())
196 }
197
198 #[must_use]
202 pub fn with_sharding(sharding: LocalSharding) -> Arc<Self> {
203 Arc::new(Self {
204 sharding,
205 ..Self::default()
206 })
207 }
208
209 pub(crate) fn observed(sharding: LocalSharding) -> (Arc<Self>, watch::Receiver<()>) {
210 let (changed, receiver) = watch::channel(());
211 (
212 Arc::new(Self {
213 sharding,
214 changed: Some(changed),
215 ..Self::default()
216 }),
217 receiver,
218 )
219 }
220
221 #[must_use]
224 pub fn slot(&self, account: AccountId) -> Arc<LeaseSlot> {
225 Arc::clone(
226 self.inner
227 .lock()
228 .expect("slot registry poisoned")
229 .slots
230 .entry(account)
231 .or_insert_with(|| LeaseSlot::with_sharding(account, self.sharding)),
232 )
233 }
234
235 #[must_use]
237 pub fn sharding(&self) -> LocalSharding {
238 self.sharding
239 }
240
241 pub(crate) fn observes(&self) -> bool {
242 self.changed.is_some()
243 }
244
245 pub(crate) fn observe_many(
249 &self,
250 entries: impl IntoIterator<Item = (Principal, Option<MapEntry>)>,
251 ) {
252 let mut inner = self.inner.lock().expect("slot registry poisoned");
253 for (principal, entry) in entries {
254 let until = match &entry {
255 Some(MapEntry::Present(state)) => Some(state.snapshot.valid_until),
256 Some(MapEntry::NegativeUntil { until }) => Some(*until),
257 None => None,
258 };
259 if let Some(until) = until {
260 inner.resolutions.insert(principal, until);
261 } else {
262 inner.resolutions.remove(&principal);
263 }
264 if let Some(previous) = inner.principals.remove(&principal) {
265 inner.deadlines.remove(&(previous.valid_until, principal));
266 let account = previous.account_id;
267 if let Some(members) = inner.members.get_mut(&account) {
268 members.remove(&principal);
269 if members.is_empty() {
270 inner.members.remove(&account);
271 }
272 }
273 inner.dirty.insert(account);
274 }
275 if let Some(MapEntry::Present(state)) = entry {
276 let snapshot = Arc::clone(&state.snapshot);
277 let account = snapshot.account_id;
278 inner.deadlines.insert((snapshot.valid_until, principal));
279 inner.principals.insert(principal, Arc::clone(&snapshot));
280 inner
281 .members
282 .entry(account)
283 .or_default()
284 .insert(principal, snapshot);
285 inner.dirty.insert(account);
286 }
287 }
288 drop(inner);
289 self.wake();
290 }
291
292 fn wake(&self) {
293 if let Some(changed) = &self.changed {
294 changed.send_replace(());
295 }
296 }
297
298 pub(crate) fn drain_changes(&self, now: Timestamp) -> Vec<AccountBinding> {
299 let mut inner = self.inner.lock().expect("slot registry poisoned");
300 while let Some(&(at, principal)) = inner.deadlines.first() {
301 if at > now {
302 break;
303 }
304 inner.deadlines.pop_first();
305 if let Some(snapshot) = inner.principals.get(&principal) {
306 let account = snapshot.account_id;
307 inner.dirty.insert(account);
308 }
309 }
310 let dirty = std::mem::take(&mut inner.dirty);
311 dirty
312 .into_iter()
313 .map(|account| inner.binding(account))
314 .collect()
315 }
316
317 pub(crate) fn next_expiry(&self) -> Option<Timestamp> {
318 self.inner
319 .lock()
320 .expect("slot registry poisoned")
321 .deadlines
322 .first()
323 .map(|(at, _)| *at)
324 }
325
326 #[allow(
329 clippy::disallowed_methods,
330 reason = "callers count or re-key into a BTreeMap; the key order never reaches an output"
331 )]
332 pub(crate) fn bindings(&self) -> Vec<AccountBinding> {
333 let inner = self.inner.lock().expect("slot registry poisoned");
334 inner
335 .members
336 .keys()
337 .map(|&account| inner.binding(account))
338 .collect()
339 }
340
341 #[allow(
347 clippy::disallowed_methods,
348 reason = "the contention report sorts by (count, account) before any output"
349 )]
350 pub(crate) fn slots(&self) -> Vec<(AccountId, Arc<LeaseSlot>)> {
351 self.inner
352 .lock()
353 .expect("slot registry poisoned")
354 .slots
355 .iter()
356 .map(|(&account, slot)| (account, Arc::clone(slot)))
357 .collect()
358 }
359
360 pub(crate) fn retained_slots(&self) -> usize {
361 self.inner
362 .lock()
363 .expect("slot registry poisoned")
364 .slots
365 .len()
366 }
367}
368
369impl Registry {
370 #[allow(
373 clippy::disallowed_methods,
374 reason = "the snapshots are only reduced with any(..); their order never reaches an output"
375 )]
376 fn binding(&self, account: AccountId) -> AccountBinding {
377 AccountBinding {
378 account,
379 slot: Arc::clone(
380 self.slots
381 .get(&account)
382 .expect("published membership has a slot"),
383 ),
384 snapshots: self
385 .members
386 .get(&account)
387 .map(|members| members.values().cloned().collect())
388 .unwrap_or_default(),
389 }
390 }
391}
392
393fn quota_usable(slot: &LeaseSlot, mode: EnforcementMode, now: Timestamp) -> bool {
394 let lease_usable = slot
395 .load_observed()
396 .is_some_and(|lease| now < lease.usable_until() && !lease.remaining().is_zero());
397 lease_usable
398 || mode
399 .overage_cap()
400 .is_some_and(|cap| !slot.overage().headroom(cap).is_zero())
401}
402
403#[cfg(test)]
404mod tests {
405 #![allow(
406 clippy::disallowed_methods,
407 reason = "unit tests that build an arbitrary `now` the assertions are relative to; \
408 no assertion here depends on what the clock actually said"
409 )]
410 use super::*;
411 use jiff::SignedDuration;
412 use tollgate_core::{
413 CommitFunding, CostUnits, EnforcementMode, FencingToken, LeaseGrant, LeaseId, LocalLease,
414 Reservation,
415 };
416 fn lease_expiring_at(expires_at: Timestamp, units: u64) -> Arc<LocalLease> {
417 Arc::new(LocalLease::new(
418 LeaseGrant {
419 lease_id: LeaseId(1),
420 account_id: AccountId(1),
421 fencing_token: FencingToken(1),
422 units: CostUnits(units),
423 expires_at,
424 },
425 CostUnits(0),
426 ))
427 }
428
429 #[test]
430 fn configured_sharding_is_preserved_by_stable_account_slots() {
431 let sharding = LocalSharding::new(std::num::NonZeroUsize::new(8).unwrap());
432 let registry = SlotRegistry::with_sharding(sharding);
433 let slot = registry.slot(AccountId(1));
434 assert!(!registry.observes());
435 assert_eq!(registry.retained_slots(), 1);
436 assert_eq!(registry.sharding(), sharding);
437 assert_eq!(slot.sharding(), sharding);
438 assert!(Arc::ptr_eq(&slot, ®istry.slot(AccountId(1))));
439 assert!(!Arc::ptr_eq(&slot, ®istry.slot(AccountId(2))));
440 assert_eq!(registry.retained_slots(), 2);
441 }
442
443 #[tokio::test]
444 async fn membership_wakes_once_per_publication_and_expires_at_its_exact_deadline() {
445 use tollgate_admission::{ArcSwapSnapshotMap, SnapshotMap};
446 use tollgate_core::{CostTable, Generation, PermissionBits, ResolvedLimits};
447 let (registry, mut changed) = SlotRegistry::observed(LocalSharding::SINGLE);
448 let map = ArcSwapSnapshotMap::default();
449 let now = Timestamp::from_second(100).unwrap();
450 let until = Timestamp::from_second(101).unwrap();
451 let principal = Principal(11);
452 let snapshot = Arc::new(
453 AccountSnapshot::builder(
454 AccountId(1),
455 Generation(1),
456 AccountStatus::Active,
457 until,
458 PermissionBits::bit(0),
459 ResolvedLimits::new(100),
460 Arc::new(CostTable::builder(CostUnits(0), CostUnits(0)).build()),
461 )
462 .build(),
463 );
464 map.install(principal, snapshot, registry.slot(AccountId(1)))
465 .unwrap();
466 registry.observe_many([(principal, map.get(&principal))]);
467 assert!(changed.has_changed().unwrap());
468 changed.borrow_and_update();
469 assert_eq!(registry.drain_changes(now).len(), 1);
470 assert!(registry.drain_changes(now).is_empty());
471 assert_eq!(registry.next_expiry(), Some(until));
472 let expired = registry.drain_changes(until);
473 assert_eq!(expired.len(), 1);
474 assert!(!expired[0].eligible(until));
475 assert_eq!(registry.next_expiry(), None);
476 assert!(registry.drain_changes(until).is_empty());
477 registry.observe_many([(principal, None)]);
478 assert!(changed.has_changed().unwrap());
479 assert!(registry.bindings().is_empty());
480 assert_eq!(registry.retained_slots(), 1);
481 }
482
483 #[test]
484 fn readiness_closes_the_lease_window_exactly_when_debits_do() {
485 let now = Timestamp::now();
486 let slot = LeaseSlot::for_account(AccountId(1));
487 assert!(
488 !quota_usable(&slot, EnforcementMode::Strict, now),
489 "an empty slot funds nothing"
490 );
491
492 let lease = lease_expiring_at(now, 100);
493 drop(slot.replace(Arc::clone(&lease)));
494 assert!(
495 lease.try_debit(CostUnits(1), now).is_err(),
496 "the request path denies at the boundary",
497 );
498 assert!(
499 !quota_usable(&slot, EnforcementMode::Strict, now),
500 "so readiness must not still be advertising at it",
501 );
502
503 drop(slot.replace(lease_expiring_at(
504 now.checked_add(SignedDuration::from_secs(60)).unwrap(),
505 100,
506 )));
507 assert!(quota_usable(&slot, EnforcementMode::Strict, now));
508
509 drop(slot.replace(lease_expiring_at(
510 now.checked_add(SignedDuration::from_secs(60)).unwrap(),
511 0,
512 )));
513 assert!(
514 !quota_usable(&slot, EnforcementMode::Strict, now),
515 "a live lease with nothing left funds nothing either",
516 );
517 }
518
519 #[test]
524 fn readiness_counts_overage_headroom_for_an_elastic_account() {
525 let now = Timestamp::now();
526 let elastic = EnforcementMode::Elastic {
527 overage_cap: CostUnits(100),
528 };
529 let slot = LeaseSlot::for_account(AccountId(1));
530
531 assert!(!quota_usable(&slot, EnforcementMode::Strict, now));
534 assert!(quota_usable(&slot, elastic, now));
535
536 drop(slot.replace(lease_expiring_at(
537 now.checked_add(SignedDuration::from_secs(60)).unwrap(),
538 0,
539 )));
540 assert!(!quota_usable(&slot, EnforcementMode::Strict, now));
541 assert!(quota_usable(&slot, elastic, now));
542
543 let overage =
546 Reservation::reserve_overage(slot.overage(), CostUnits(100), CostUnits(100)).unwrap();
547 overage
548 .commit_at_execution_start(now, CommitFunding::LeaseOnly)
549 .unwrap();
550 assert!(
551 !quota_usable(&slot, elastic, now),
552 "a spent cap is not admissible, and readiness must say so"
553 );
554
555 assert!(quota_usable(
557 &slot,
558 EnforcementMode::Elastic {
559 overage_cap: CostUnits(200)
560 },
561 now
562 ));
563 }
564}