tollgate_store/memory.rs
1//! The in-memory reference backend.
2//!
3//! A single mutex over plain maps: writes happen at control-plane frequency,
4//! so contention is irrelevant, and the simplicity makes the settlement rules
5//! auditable. A real backend (Postgres) must reproduce exactly these rules —
6//! the shared correctness suite in `tests/` runs against both.
7//!
8//! [`MemoryStore::conservation`] returns the per-account ledger checked by
9//! [`Conservation::holds`], including overage funding and expired allowances.
10//! Usage recorded against a still-active lease lives *inside* that lease's
11//! grant (the grant was debited whole at acquire), so it only stands alone in
12//! the equation once the lease settles. A settlement loss is billing a
13//! released/expired lease could not account for (usage that never arrived
14//! before settlement).
15//!
16//! # This backend never forgets
17//!
18//! Nothing here is ever deleted, and two of the maps therefore grow with what
19//! the process has *done* rather than with what it currently holds:
20//!
21//! - `usage` is the idempotency index and keeps one event per request served,
22//! for the life of the process. It is the unbounded term, and bounding it
23//! needs a dedup-window retention decision rather than a deletion — see the
24//! deferred list in `docs/DESIGN.md`.
25//! - `leases` keeps settled records, so it grows with lease rotations. They
26//! are retained because a straggling usage event must still be matched to
27//! its lease capability and checked against settlement capacity (see
28//! [`UsageSink::ingest`]).
29//! - `snapshots` keeps revoked principals as tombstones, deliberately: that is
30//! INVARIANTS.md GL-15's anti-resurrection watermark, and the population is
31//! bounded by the number of principals rather than by traffic.
32//!
33//! The *sweep* cost does not follow that growth. Active leases are indexed
34//! in the private `leases` module, so reclaim and [`MemoryStore::conservation`] walk
35//! the live population, not the historical one (GL-23). Memory still does grow,
36//! so this backend suits development and demos but not soak or load testing;
37//! [`MemoryStore::stored_records`] reports the numbers, and the server logs
38//! them once per sweep.
39
40use crate::{AccountView, AdminAuthority, AdminReceipt, AdminState};
41
42use std::collections::{HashMap, HashSet};
43use std::num::NonZeroUsize;
44use std::sync::Arc;
45use std::sync::Mutex;
46
47use async_trait::async_trait;
48use jiff::{SignedDuration, Timestamp};
49use tokio::sync::broadcast;
50
51use tollgate_core::{
52 AccountId, AccountStatus, BudgetSchedule, BudgetView, CapacityClass, CostUnits, FencingToken,
53 Generation, KeyId, LeaseGrant, LeaseId, Principal, PublishableSnapshot, UsageEvent,
54 UsageSource,
55};
56
57use crate::leases::{LeaseRecord, Leases, Settled};
58pub use crate::traits::{AccountConfig, Conservation, StatusChange};
59use crate::traits::{
60 AdminStore, AllocateError, Allocation, BudgetError, CreateAccountError, GrantPolicy,
61 GrantPolicyError, IngestError, IngestReport, KeyDirectory, KeyError, KeyRecord,
62 KeySnapshotError, KeySummary, LeaseAllocator, PUSH_CHANNEL_CAPACITY, PublishSnapshotError,
63 ReclaimBatch, ReclaimedLease, Revocation, RolledAccount, RolloverBatch, SetStatusError,
64 SnapshotPush, SnapshotResolution, SnapshotSource, StoreError, StoreHealth, UsageSink,
65 pushes_exceed_capacity,
66};
67
68/// An account's balance, split by what expires and what does not (GL-97).
69///
70/// One number could not carry this. At a period boundary the allowance's
71/// remainder is expired and manual credits are kept, and a single balance can
72/// only either expire the credits with it or resurrect allowance units that
73/// were already spent — there is no arithmetic on one counter that separates
74/// "unspent allowance" from "unspent top-up" after the fact.
75///
76/// Spend order is allowance first. A credit bought or granted out of band
77/// should outlive the monthly allowance sitting beside it, so the units with
78/// an expiry date are the ones consumed first.
79#[derive(Debug, Clone, Copy, Default)]
80struct Balance {
81 /// Unspent units from the current period's allowance. Expired whole at
82 /// the next boundary under `Rollover::None`.
83 allowance: CostUnits,
84 /// Unspent units from manual deposits. Never expired by a rollover.
85 topup: CostUnits,
86}
87
88impl Balance {
89 /// What the account can spend, and what every reader outside this module
90 /// means by "balance".
91 fn total(self) -> CostUnits {
92 self.allowance
93 .checked_add(self.topup)
94 .expect("a balance that was funded in halves fits the sum it came from")
95 }
96
97 /// Take `units`, allowance first, reporting the split so the lease can
98 /// give each half back to the bucket it came from.
99 ///
100 /// Returning the split rather than crediting allowance-first on release is
101 /// what stops a top-up being silently expired: a lease funded entirely
102 /// from credits, released after a boundary, would otherwise return its
103 /// units to the allowance bucket and have them expired at the next one.
104 fn take(&mut self, units: CostUnits) -> Option<Drawn> {
105 let from_allowance = self.allowance.min(units);
106 let from_topup = units.checked_sub(from_allowance)?;
107 self.allowance = self.allowance.checked_sub(from_allowance)?;
108 self.topup = self.topup.checked_sub(from_topup)?;
109 Some(Drawn {
110 from_allowance,
111 from_topup,
112 })
113 }
114}
115
116/// Which buckets a grant drew from, carried by the lease so settlement can
117/// return each half to where it came from.
118#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
119pub(crate) struct Drawn {
120 pub(crate) from_allowance: CostUnits,
121 pub(crate) from_topup: CostUnits,
122}
123
124#[derive(Debug)]
125struct AccountRecord {
126 balance: Balance,
127 deposited: CostUnits,
128 /// The account's periodic allowance, if it has one. `None` keeps the
129 /// manual-deposit behaviour the ledger has always had, and a rollover pass
130 /// skips the account entirely — which is why this is an `Option` rather
131 /// than a schedule with a zero allowance, a very different thing.
132 schedule: Option<BudgetSchedule>,
133 /// The first instant of the period this account is currently in, and the
134 /// marker rollover is idempotent against: a boundary is crossed exactly
135 /// once because the update that crosses it is conditional on this value
136 /// still being the old one.
137 period_start: Timestamp,
138 /// Units funded but never spendable again, because the period that funded
139 /// them closed. Monotonic.
140 expired: CostUnits,
141 /// Mirrors `tollgate_accounts.status`. An [`AccountStatus`] rather than a
142 /// bool so `Closed` is representable and terminality can be checked here
143 /// instead of inferred from snapshots (GL-51).
144 status: AccountStatus,
145 /// Mirrors `tollgate_accounts.capacity_class`. The account-owned fact a
146 /// snapshot's `capacity_class` is a copy of, and the reason a publish
147 /// carrying a different one is refused: two writers for one fact is the
148 /// divergence GL-51 abolished for status (GL-99).
149 capacity_class: CapacityClass,
150 /// Mirrors `tollgate_accounts.origin`: which authority created the
151 /// account. Written once, never changed (#39).
152 origin: AdminAuthority,
153 /// Mirrors `tollgate_accounts.status_set_by`: which authority wrote
154 /// `status` last, the fact an operator hold is read from (#39).
155 status_set_by: AdminAuthority,
156 next_fence: u64,
157 /// Usage accepted into the billing ledger.
158 usage_recorded: CostUnits,
159 /// Unfunded units billed under elastic enforcement: the second funding
160 /// term of the conservation equation. Monotonic, like `deposited` and
161 /// `usage_recorded`; only a deposit settles it, and settling it does not
162 /// reduce it.
163 overage_recorded: CostUnits,
164 /// Billing lost at settlement: usage that had not arrived when a lease
165 /// was released/reclaimed. Bounded by construction; reconciliation
166 /// watches it.
167 settlement_loss: CostUnits,
168}
169
170impl AccountRecord {
171 /// What this account could still spend, for a snapshot's budget view
172 /// (GL-97).
173 ///
174 /// Balance *plus* the unspent remainder of every active lease, because
175 /// units out on lease are still the account's — an instance holding a
176 /// 500-unit lease has not lost those units, and a figure that excluded
177 /// them would tell a customer their quota had halved the moment a lease
178 /// was taken.
179 ///
180 /// Derived from the account row alone, with no walk over the leases. The
181 /// conservation equation is what makes that possible: `balance + active
182 /// grants` is `funded - consumed`, so what an account can still spend is
183 /// everything it was funded with minus everything it has consumed or lost.
184 /// A lease scan would be O(the account's leases) at every publication and
185 /// would answer the same number.
186 ///
187 /// The `expect`s are the split `PostgresStore::conservation` documents:
188 /// here the counters are maintained by one process under one lock, so an
189 /// underflow is unrepresentable rather than corruption a caller could
190 /// have caused. The PostgreSQL backend reads a database it does not
191 /// exclusively own and reports the same condition as an error.
192 fn budget_view(&self) -> BudgetView {
193 let funded = self
194 .deposited
195 .checked_add(self.overage_recorded)
196 .expect("an account cannot be funded past what it was funded with");
197 let consumed = self
198 .usage_recorded
199 .checked_add(self.settlement_loss)
200 .and_then(|spent| spent.checked_add(self.expired))
201 .expect("consumption cannot exceed the funding it came from");
202 BudgetView {
203 balance_at_publish: funded
204 .checked_sub(consumed)
205 .expect("conservation keeps consumption within funding"),
206 // The period the account is *in*, which is what it can spend
207 // against. An account whose schedule was set but whose first
208 // rollover has not run yet is still in its previous period, and
209 // says so, until the next sweep tick moves it.
210 period_end: self
211 .schedule
212 .map(|schedule| schedule.period.end_after(self.period_start)),
213 }
214 }
215}
216
217/// Whether the period that funded a lease had already closed by the time the
218/// lease settles.
219///
220/// Both settlement and consolidation use this decision so the credited bucket
221/// and the replacement grant's floor agree about which units are spendable.
222fn allowance_lapsed(record: &AccountRecord, period_start: Timestamp) -> bool {
223 period_start < record.period_start
224}
225
226/// Return a settled lease's unspent units to the account — expiring the
227/// allowance half, if the period that funded it has closed (GL-97).
228///
229/// This is the whole of the "drain then expire" decision. An active lease at a
230/// period boundary keeps serving to its own TTL, so there is no admission gap
231/// and no clock read on the request path; the boundary shows up here instead,
232/// when the lease finally settles. Unspent units funded by an allowance that
233/// no longer exists cannot go back to a balance — that would resurrect an
234/// expired allowance, and the account would carry units its schedule says it
235/// should not have.
236///
237/// The split is charged in the order the account spends, allowance first, so
238/// the top-up half is what survives a partly-spent lease. Only that half is
239/// unconditional: a top-up never expires, boundary or not, which is what makes
240/// "manual credits persist across rollover" true even for a lease that
241/// straddles one.
242///
243/// Usage is untouched either way, which is what makes a straggler across the
244/// boundary bill against the period the lease was granted in.
245fn credit_settlement(
246 record: &mut AccountRecord,
247 funding: Drawn,
248 period_start: Timestamp,
249 unspent: CostUnits,
250) {
251 let to_topup = funding.from_topup.min(unspent);
252 let to_allowance = unspent
253 .checked_sub(to_topup)
254 .expect("the top-up half never exceeds the grant it was drawn from");
255 record.balance.topup = record
256 .balance
257 .topup
258 .checked_add(to_topup)
259 .expect("settlement credit overflow");
260 // The only thing the boundary decides.
261 let bucket = if allowance_lapsed(record, period_start) {
262 &mut record.expired
263 } else {
264 &mut record.balance.allowance
265 };
266 *bucket = bucket
267 .checked_add(to_allowance)
268 .expect("settlement credit overflow");
269}
270
271/// A principal's stored snapshot state.
272///
273/// An enum rather than `{ generation, snapshot: Option<_> }` so the generation
274/// is stored exactly once (GL-54). The struct held it twice whenever a snapshot
275/// was present — once in the field and once inside the snapshot — with nothing
276/// but caller discipline keeping them equal, the same shape PostgreSQL had
277/// between its column and its JSONB.
278///
279/// The field could not simply be deleted: revoking sets the snapshot aside, and
280/// its generation is then the only surviving watermark, without which
281/// INVARIANTS.md GL-15's anti-resurrection rule would be unimplementable here. So
282/// each state carries the generation in exactly one place instead.
283///
284/// Variants named for the [`SnapshotResolution`] they map onto, since
285/// `SnapshotSource::snapshot` is the only place a reader meets them.
286///
287/// Note this makes memory's inability to un-revoke structural, where a
288/// PostgreSQL tombstone keeps its JSON and could in principle be resurrected.
289/// Identical behaviour today; if un-revoking is ever added, this is the backend
290/// that changes shape.
291#[derive(Debug)]
292enum SnapshotRecord {
293 /// Live: the generation is the snapshot's own.
294 Present(PublishableSnapshot),
295 /// Revoked: the snapshot is gone and the watermark is all that remains.
296 Revoked(Generation),
297}
298
299impl SnapshotRecord {
300 /// The record's generation, wherever this state keeps it.
301 fn generation(&self) -> Generation {
302 match self {
303 SnapshotRecord::Present(snapshot) => snapshot.generation,
304 SnapshotRecord::Revoked(generation) => *generation,
305 }
306 }
307}
308
309#[derive(Default)]
310struct Inner {
311 accounts: HashMap<AccountId, AccountRecord>,
312 leases: Leases,
313 snapshots: HashMap<Principal, SnapshotRecord>,
314 usage: HashMap<tollgate_core::RequestId, UsageEvent>,
315 credential_activity: HashMap<KeyId, Timestamp>,
316 /// Credential records, live and retired alike. A revocation sets
317 /// `revoked_at` rather than removing the row: the tombstone is what makes
318 /// "already retired" distinguishable from "never existed", and a removed
319 /// row would let a replayed issuance resurrect the credential.
320 keys: HashMap<KeyId, StoredKey>,
321 /// Ordered unrevoked ids avoid scanning retired history on every page.
322 unrevoked_keys: std::collections::BTreeSet<KeyId>,
323 /// Retained through retirement, like the durable principal UNIQUE index.
324 key_principals: HashMap<Principal, KeyId>,
325 credential_revision: u64,
326 next_lease_id: u128,
327}
328
329/// A credential as this backend holds it: the durable record plus its
330/// retirement, which is the one mutable field.
331#[derive(Debug, Clone)]
332struct StoredKey {
333 record: KeyRecord,
334 revoked_at: Option<Timestamp>,
335}
336
337/// What the in-memory ledger is currently holding.
338///
339/// This backend never forgets, so the first two numbers climb with traffic for
340/// the life of the process while the third returns to the live population.
341/// Reported so "grows without bound" is a number an operator can watch rather
342/// than a claim in a doc comment (GL-23).
343#[derive(Debug, Clone, Copy, PartialEq, Eq)]
344pub struct StoredRecords {
345 /// Credentials with recorded attributable commitments, including retired keys.
346 pub credential_activity: usize,
347 /// Usage events retained for idempotency: one per request served, kept
348 /// forever. This is the term that grows without bound.
349 pub usage_events: usize,
350 /// Lease records, active and settled together — grows with rotations.
351 pub leases: usize,
352 /// Of those, the ones still active: the only records the reclaim sweep and
353 /// [`MemoryStore::conservation`] examine. Steady-state, unlike the two
354 /// above, which is what makes the sweep's cost independent of them.
355 pub active_leases: usize,
356}
357
358/// The in-memory reference backend: every store trait behind one mutex.
359///
360/// It implements [`LeaseAllocator`], [`SnapshotSource`], [`UsageSink`],
361/// [`AdminStore`], [`KeyDirectory`], [`KeySource`](crate::KeySource) and
362/// [`StoreHealth`] with the settlement rules every backend must reproduce. It
363/// never deletes a record, so it suits development, tests and demos but not
364/// long-running load; see the [module documentation](crate::memory).
365pub struct MemoryStore {
366 inner: Mutex<Inner>,
367 policy: GrantPolicy,
368 push: broadcast::Sender<SnapshotPush>,
369}
370
371impl MemoryStore {
372 fn publish_key_snapshot_with_generation(
373 &self,
374 account: AccountId,
375 key: KeyId,
376 snapshot: PublishableSnapshot,
377 allocate_generation: bool,
378 ) -> Result<crate::AdminReceipt<()>, KeySnapshotError> {
379 // One guard across resolution, the retirement check and publication,
380 // so a revocation cannot land between them.
381 let (principal, published, before, after) = {
382 let mut inner = self.lock();
383 let stored = account_key(&inner, account, key)?;
384 if stored.revoked_at.is_some() {
385 return Err(KeySnapshotError::Retired { key_id: key });
386 }
387 let principal = stored.record.principal;
388 if snapshot.key_id != Some(key) {
389 return Err(PublishSnapshotError::CredentialMismatch { key_id: key }.into());
390 }
391 let before = snapshot_audit(inner.snapshots.get(&principal));
392 let snapshot = if allocate_generation {
393 let previous = inner
394 .snapshots
395 .get(&principal)
396 .map_or(0, |s| s.generation().0);
397 let next = previous
398 .checked_add(1)
399 .ok_or_else(|| StoreError("snapshot generation overflow".into()))?;
400 snapshot.restamped(snapshot.status, Generation(next))
401 } else {
402 snapshot
403 };
404 let published = publish_locked(&mut inner, principal, snapshot)?;
405 let after = snapshot_audit(inner.snapshots.get(&principal));
406 (principal, published, before, after)
407 };
408 if let Some(snapshot) = published {
409 self.push_to_subscribers(SnapshotPush {
410 principal,
411 resolution: SnapshotResolution::Present(snapshot),
412 });
413 }
414 Ok(AdminReceipt::new((), before, after))
415 }
416
417 /// Create an empty store that sizes grants with `policy`.
418 ///
419 /// # Errors
420 ///
421 /// [`GrantPolicyError`] when `policy` fails [`GrantPolicy::validate`].
422 pub fn new(policy: GrantPolicy) -> Result<Arc<Self>, GrantPolicyError> {
423 policy.validate()?;
424 let (push, _) = broadcast::channel(PUSH_CHANNEL_CAPACITY);
425 Ok(Arc::new(MemoryStore {
426 inner: Mutex::new(Inner::default()),
427 policy,
428 push,
429 }))
430 }
431
432 /// The expiry a grant issued now would carry, clamping the request to the
433 /// policy's `max_ttl`.
434 ///
435 /// Every fallible value is computed before the ledger moves: an overflow
436 /// must not debit a balance without creating a lease, nor settle a lease
437 /// without replacing it.
438 fn grant_expiry(
439 &self,
440 ttl: SignedDuration,
441 now: Timestamp,
442 ) -> Result<Timestamp, AllocateError> {
443 if ttl <= SignedDuration::ZERO {
444 return Err(AllocateError::InvalidTtl);
445 }
446 let ttl = ttl.min(self.policy.max_ttl);
447 now.checked_add(ttl)
448 .map_err(|e| AllocateError::Storage(StoreError(format!("ttl overflow: {e}"))))
449 }
450
451 fn lock(&self) -> std::sync::MutexGuard<'_, Inner> {
452 // Lock poisoning would mean a panic while holding the ledger; the
453 // reference backend treats that as unrecoverable.
454 self.inner.lock().expect("memory store lock poisoned")
455 }
456
457 // ---- admin / control-plane surface -------------------------------
458
459 /// Test/bootstrap convenience: create a fresh account, panicking on a
460 /// duplicate. Production paths use [`AdminStore::create_account`].
461 pub fn create_account(&self, config: AccountConfig) {
462 self.try_create_account(config)
463 .expect("account already exists");
464 }
465
466 /// Create an account. Never destructive: an existing account (with its
467 /// balance, ledger totals, fencing sequence, and leases) is left
468 /// untouched and the caller told (review finding GL-7).
469 pub fn try_create_account(&self, config: AccountConfig) -> Result<(), CreateAccountError> {
470 self.create_account_audited(config, AdminAuthority::Operator)
471 .map(|receipt| receipt.outcome)
472 }
473
474 fn create_account_audited(
475 &self,
476 config: AccountConfig,
477 origin: AdminAuthority,
478 ) -> Result<AdminReceipt<()>, CreateAccountError> {
479 let mut inner = self.lock();
480 if inner.accounts.contains_key(&config.account_id) {
481 return Err(CreateAccountError::AlreadyExists);
482 }
483 inner.accounts.insert(
484 config.account_id,
485 AccountRecord {
486 // The opening balance is a top-up, not an allowance: it was
487 // deposited by whoever created the account, and nothing has
488 // scheduled it to expire. A schedule set later starts its
489 // first period at the next rollover.
490 balance: Balance {
491 allowance: CostUnits::ZERO,
492 topup: config.initial_balance,
493 },
494 deposited: config.initial_balance,
495 schedule: None,
496 period_start: Timestamp::UNIX_EPOCH,
497 expired: CostUnits::ZERO,
498 status: config.status,
499 capacity_class: config.capacity_class,
500 origin,
501 status_set_by: origin,
502 next_fence: 1,
503 usage_recorded: CostUnits::ZERO,
504 overage_recorded: CostUnits::ZERO,
505 settlement_loss: CostUnits::ZERO,
506 },
507 );
508 Ok(AdminReceipt::new(
509 (),
510 AdminState::Absent,
511 AdminState::AccountCreated {
512 initial_balance: config.initial_balance,
513 status: config.status,
514 capacity_class: config.capacity_class,
515 origin,
516 },
517 ))
518 }
519
520 /// The one status transition, for either authority (#39). Operator
521 /// writes and provisioner activations share the lock, the refusal phase
522 /// and the republication, so the hold is checked at the same point the
523 /// write is made.
524 fn set_status_as(
525 &self,
526 account: AccountId,
527 status: AccountStatus,
528 authority: AdminAuthority,
529 ) -> Result<AdminReceipt<StatusChange>, SetStatusError> {
530 // One lock for both records. The inherent `publish_snapshot` takes the
531 // lock itself, so it cannot be reused here: the whole point is that no
532 // observer sees the ledger moved and the snapshots not.
533 let (republished, before, set_by) = {
534 let mut inner = self.lock();
535 // Phase 1, read-only: every refusal happens before anything moves.
536 let record = inner
537 .accounts
538 .get(&account)
539 .ok_or(SetStatusError::UnknownAccount)?;
540 if authority == AdminAuthority::Provisioner
541 && record.origin != AdminAuthority::Provisioner
542 {
543 return Err(SetStatusError::NotProvisioned);
544 }
545 if record.status == AccountStatus::Closed && status != AccountStatus::Closed {
546 return Err(SetStatusError::AccountClosed);
547 }
548 if authority == AdminAuthority::Provisioner
549 && record.status != status
550 && record.status_set_by == AdminAuthority::Operator
551 {
552 return Err(SetStatusError::OperatorHold);
553 }
554 // A provisioner repeating an activation changes nothing, not even
555 // the author: an operator who reactivated keeps that authorship.
556 let set_by = if authority == AdminAuthority::Provisioner && record.status == status {
557 record.status_set_by
558 } else {
559 authority
560 };
561 // Phase 2, still read-only: plan every snapshot write, so a
562 // generation overflow surfaces here rather than after the ledger
563 // has already moved. Writing the ledger first and failing here
564 // would leave the account suspended with Active snapshots -- the
565 // divergence INVARIANTS.md GL-22 forbids, in the very backend that
566 // serves as its reference.
567 let before = AdminState::Status {
568 status: record.status,
569 set_by: record.status_set_by,
570 };
571 let planned = plan_republish(&inner, account, Restamp::Status(status))?;
572 // Phase 3: apply. Nothing below this line can fail.
573 let record = inner
574 .accounts
575 .get_mut(&account)
576 .expect("the account was found under this same guard");
577 record.status = status;
578 record.status_set_by = set_by;
579 (apply_republish(&mut inner, planned), before, set_by)
580 };
581 // Outside the guard: `push_to_subscribers` must not run under it, and
582 // a subscriber must never observe a push for a transition that is
583 // still mid-flight.
584 if pushes_exceed_capacity(republished.len()) {
585 tracing::warn!(
586 %account,
587 principals = republished.len(),
588 capacity = PUSH_CHANNEL_CAPACITY,
589 "status change emitted more pushes than the channel holds; subscribers will resync"
590 );
591 }
592 let planned_pushes = republished;
593 let republished = planned_pushes.len();
594 for (principal, snapshot) in planned_pushes {
595 self.push_to_subscribers(SnapshotPush {
596 principal,
597 resolution: SnapshotResolution::Present(snapshot),
598 });
599 }
600 Ok(AdminReceipt::new(
601 StatusChange {
602 republished,
603 // This backend holds validated snapshots rather than encoded ones,
604 // so there is nothing here that can fail to decode.
605 unreadable: 0,
606 },
607 before,
608 AdminState::Status { status, set_by },
609 ))
610 }
611
612 /// Add balance to an existing account (top-up).
613 pub fn deposit(&self, account: AccountId, units: CostUnits) -> Result<(), AllocateError> {
614 self.deposit_audited(account, units)
615 .map(|receipt| receipt.outcome)
616 }
617
618 fn deposit_audited(
619 &self,
620 account: AccountId,
621 units: CostUnits,
622 ) -> Result<AdminReceipt<()>, AllocateError> {
623 let mut inner = self.lock();
624 let record = inner
625 .accounts
626 .get_mut(&account)
627 .ok_or(AllocateError::UnknownAccount)?;
628 // Both sums are computed before either lands, the rule `acquire`
629 // states and `set_account_status` splits into plan/apply. Assigning
630 // as they were computed left the ledger permanently short when the
631 // second overflowed: conservation keeps `balance <= deposited`, so
632 // `deposited` reaches the ceiling first, and a refused top-up on a
633 // fully-spent account still credited `balance` (GL-57). `PostgresStore`
634 // moves both columns in one statement, where an overflow aborts it
635 // and nothing moves.
636 // A manual deposit is a top-up: it survives a period boundary, which
637 // is the documented default (GL-97). An allowance only ever arrives
638 // through `roll_period`.
639 let topup = record
640 .balance
641 .topup
642 .checked_add(units)
643 .ok_or(AllocateError::BalanceOverflow)?;
644 let deposited = record
645 .deposited
646 .checked_add(units)
647 .ok_or(AllocateError::BalanceOverflow)?;
648 let before = AdminState::Funding {
649 topup: record.balance.topup,
650 deposited: record.deposited,
651 };
652 record.balance.topup = topup;
653 record.deposited = deposited;
654 Ok(AdminReceipt::new(
655 (),
656 before,
657 AdminState::Funding { topup, deposited },
658 ))
659 }
660
661 /// Bind (or replace) a principal's compiled snapshot and push it to
662 /// subscribers. Generation-monotonic: a replayed or reordered publish
663 /// carrying an older (or equal) generation is a no-op — matching the
664 /// Postgres backend, which enforces the same rule in its upsert (review
665 /// finding GL-5's backend-divergence note).
666 pub fn publish_snapshot(
667 &self,
668 principal: Principal,
669 snapshot: PublishableSnapshot,
670 ) -> Result<(), PublishSnapshotError> {
671 let published = {
672 let mut inner = self.lock();
673 publish_locked(&mut inner, principal, snapshot)?
674 };
675 if let Some(snapshot) = published {
676 self.push_to_subscribers(SnapshotPush {
677 principal,
678 resolution: SnapshotResolution::Present(snapshot),
679 });
680 }
681 Ok(())
682 }
683
684 /// Broadcast a control-plane change. No receivers is not a failure — a
685 /// pull via `snapshot()` still observes it — but how many instances the
686 /// push actually reached is the difference between "propagated in
687 /// milliseconds" and "propagated at the next refresh interval", so it is
688 /// reported rather than discarded.
689 fn push_to_subscribers(&self, push: SnapshotPush) {
690 let principal = push.principal;
691 let subscribers = self.push.send(push).unwrap_or(0);
692 tracing::debug!(
693 %principal,
694 subscribers,
695 "snapshot pushed to subscribers"
696 );
697 }
698
699 /// Tombstone a principal's live snapshot at its current generation and push
700 /// the revocation to subscribers. A principal with no snapshot, or one already
701 /// tombstoned, is left unchanged and nothing is pushed.
702 ///
703 /// Test/bootstrap convenience beside [`AdminStore::remove_snapshot`], which
704 /// also returns the audit receipt.
705 pub fn remove_snapshot(&self, principal: Principal) {
706 self.remove_snapshot_audited(principal);
707 }
708
709 fn remove_snapshot_audited(&self, principal: Principal) -> AdminReceipt<()> {
710 let receipt = remove_locked(&mut self.lock(), principal);
711 self.announce_removal(principal, &receipt);
712 receipt
713 }
714
715 /// Push a tombstone only when the removal changed something.
716 fn announce_removal(&self, principal: Principal, receipt: &AdminReceipt<()>) {
717 if receipt.before != receipt.after
718 && let AdminState::Snapshot { generation, .. } = receipt.after
719 {
720 self.push_to_subscribers(SnapshotPush {
721 principal,
722 resolution: SnapshotResolution::Revoked { generation },
723 });
724 }
725 }
726
727 // ---- reconciliation / test surface -------------------------------
728
729 /// The sums are recomputed from the lease records every time. The index
730 /// narrows *which* records are read (GL-23) and is never the source of the
731 /// numbers: this function exists to catch ledger bugs, and one that read a
732 /// running total maintained by the same writers that might be wrong could
733 /// not catch them.
734 #[must_use]
735 pub fn conservation(&self, account: AccountId) -> Option<Conservation> {
736 let inner = self.lock();
737 let record = inner.accounts.get(&account)?;
738 let mut active_grants = CostUnits::ZERO;
739 let mut active_used = CostUnits::ZERO;
740 for lease in inner.leases.active_of(account) {
741 active_grants = active_grants
742 .checked_add(lease.granted)
743 .expect("grant sum overflow");
744 active_used = active_used
745 .checked_add(lease.used)
746 .expect("used sum overflow");
747 }
748 Some(Conservation {
749 deposited: record.deposited,
750 overage_recorded: record.overage_recorded,
751 balance: record.balance.total(),
752 active_lease_grants: active_grants,
753 settled_usage: record
754 .usage_recorded
755 .checked_sub(active_used)
756 .expect("active usage never exceeds recorded usage"),
757 settlement_loss: record.settlement_loss,
758 expired: record.expired,
759 })
760 }
761
762 /// Every usage unit accepted for `account`: on active and settled leases, and
763 /// as overage. Zero for an unknown account.
764 #[must_use]
765 pub fn usage_recorded(&self, account: AccountId) -> CostUnits {
766 self.lock()
767 .accounts
768 .get(&account)
769 .map(|a| a.usage_recorded)
770 .unwrap_or(CostUnits::ZERO)
771 }
772
773 /// The event this store settled for `request_id`, if it settled one.
774 ///
775 /// The reference implementation keeps whole events — the idempotency map
776 /// is keyed by request and holds the value — so this reads back what was
777 /// actually billed rather than a projection of it. `PostgresStore` keeps
778 /// only the columns it needs and has no equivalent, which is why a
779 /// backend-parity assertion on a stored field reads that backend's column
780 /// directly instead.
781 ///
782 /// Exists for inspection and tests, beside [`usage_recorded`]. It is not
783 /// part of `UsageSink`: no request-path code reads settled events back.
784 ///
785 /// [`usage_recorded`]: MemoryStore::usage_recorded
786 #[must_use]
787 pub fn settled_event(&self, request_id: tollgate_core::RequestId) -> Option<UsageEvent> {
788 self.lock().usage.get(&request_id).copied()
789 }
790
791 /// `account`'s spendable balance, allowance plus top-up, excluding units out on
792 /// lease. Zero for an unknown account.
793 #[must_use]
794 pub fn balance(&self, account: AccountId) -> CostUnits {
795 self.lock()
796 .accounts
797 .get(&account)
798 .map(|a| a.balance.total())
799 .unwrap_or(CostUnits::ZERO)
800 }
801
802 /// Lease records examined by the sweep and by `conservation` since this
803 /// store was created. The bound GL-23 claims is about work, so only a count
804 /// of records actually looked at can witness it.
805 #[cfg(test)]
806 fn leases_examined(&self) -> usize {
807 self.lock().leases.examined()
808 }
809
810 /// What this backend is currently holding — see [`StoredRecords`], and the
811 /// module docs for why two of the three only ever climb.
812 #[must_use]
813 pub fn stored_records(&self) -> StoredRecords {
814 let inner = self.lock();
815 StoredRecords {
816 credential_activity: inner.credential_activity.len(),
817 usage_events: inner.usage.len(),
818 leases: inner.leases.len(),
819 active_leases: inner.leases.active_len(),
820 }
821 }
822}
823
824/// The part of a settled lease's `unspent` that returns to *spendable*
825/// balance.
826///
827/// Not always all of it: the allowance half of a lease funded by a period that
828/// has since closed expires instead of coming back (GL-97, and
829/// [`credit_settlement`], which applies the same rule). A caller sizing a
830/// grant against the credit it is about to make must ask this rather than
831/// assume `unspent`.
832fn spendable_credit(
833 record: &AccountRecord,
834 funding: Drawn,
835 period_start: Timestamp,
836 unspent: CostUnits,
837) -> CostUnits {
838 if allowance_lapsed(record, period_start) {
839 funding.from_topup.min(unspent)
840 } else {
841 unspent
842 }
843}
844
845/// A validated release, decided before anything moves.
846///
847/// The reference backend holds one mutex where the SQL backend holds a
848/// transaction, so it has nothing to roll back: every way an operation can
849/// refuse is established here, and applying is infallible. That is what makes
850/// a consolidation all-or-nothing in both backends for the same reason rather
851/// than by coincidence.
852struct ReleasePlan {
853 account_id: AccountId,
854 funding: Drawn,
855 period_start: Timestamp,
856 unspent: CostUnits,
857 loss: CostUnits,
858}
859
860/// A validated grant, decided before anything moves. See [`ReleasePlan`].
861struct GrantPlan {
862 granted: CostUnits,
863 fencing_token: FencingToken,
864 next_fence: u64,
865 next_lease_id: u128,
866}
867
868fn plan_release(
869 inner: &Inner,
870 policy: &GrantPolicy,
871 lease_id: LeaseId,
872 fencing_token: FencingToken,
873 unspent: CostUnits,
874 now: Timestamp,
875) -> Result<ReleasePlan, AllocateError> {
876 let lease = inner
877 .leases
878 .get(lease_id)
879 .ok_or(AllocateError::UnknownLease)?;
880 if lease.fencing_token != fencing_token {
881 return Err(AllocateError::Fenced);
882 }
883 // A lease that lapsed before the release arrived settles by expiry
884 // reclaim instead; the late releaser is told, not silently absorbed.
885 // Releases are accepted through the grace window: a holder shutting down
886 // slowly may reach here after `expires_at` but before the sweep settles
887 // the lease. Only a settled (or grace-exhausted) lease refuses.
888 if !lease.is_active()
889 || policy
890 .reclaim_cutoff(now)
891 .is_some_and(|cutoff| lease.expires_at <= cutoff)
892 {
893 return Err(AllocateError::LeaseNotActive);
894 }
895 // granted = used + unspent + loss; a claim that doesn't fit is a client
896 // accounting bug. The loss is *provisional*: usage events for this lease
897 // that were committed but not yet flushed at release time still fit in the
898 // gap and convert loss back into billed usage when they arrive (see
899 // `ingest`).
900 let spent_plus_unspent = lease
901 .used
902 .checked_add(unspent)
903 .ok_or(AllocateError::InvalidRelease)?;
904 let loss = lease
905 .granted
906 .checked_sub(spent_plus_unspent)
907 .ok_or(AllocateError::InvalidRelease)?;
908 Ok(ReleasePlan {
909 account_id: lease.account_id,
910 funding: lease.funding,
911 period_start: lease.period_start,
912 unspent,
913 loss,
914 })
915}
916
917fn apply_release(inner: &mut Inner, lease_id: LeaseId, plan: &ReleasePlan) {
918 assert!(
919 inner
920 .leases
921 .settle(lease_id, Settled::Released, plan.unspent),
922 "the plan validated this lease as active under this same lock"
923 );
924 let record = inner
925 .accounts
926 .get_mut(&plan.account_id)
927 .expect("lease account exists");
928 credit_settlement(record, plan.funding, plan.period_start, plan.unspent);
929 record.settlement_loss = record
930 .settlement_loss
931 .checked_add(plan.loss)
932 .expect("loss overflow");
933}
934
935/// Size one grant against `account`'s balance.
936///
937/// `incoming` is units this same operation is about to credit to this same
938/// account — a consolidation's returned tail. It counts twice, and
939/// deliberately: as part of the balance the policy sizes against, and as the
940/// floor that answer may not fall below, so the exchange can grow a holding or
941/// leave it alone but never shrink it (see [`LeaseAllocator::consolidate`]).
942/// A plain acquire credits nothing and passes zero, leaving the policy's
943/// answer exactly as it was.
944fn plan_grant(
945 inner: &Inner,
946 policy: &GrantPolicy,
947 account: AccountId,
948 requested: CostUnits,
949 incoming: CostUnits,
950 needed: CostUnits,
951 attest_refusal: bool,
952) -> Result<GrantPlan, AllocateError> {
953 let next_lease_id = inner
954 .next_lease_id
955 .checked_add(1)
956 .ok_or_else(|| AllocateError::Storage(StoreError("lease id overflow".into())))?;
957 let record = inner
958 .accounts
959 .get(&account)
960 .ok_or(AllocateError::UnknownAccount)?;
961 if record.status != AccountStatus::Active {
962 return Err(AllocateError::AccountInactive);
963 }
964 let balance = record
965 .balance
966 .total()
967 .checked_add(incoming)
968 .ok_or_else(|| AllocateError::Storage(StoreError("account balance overflow".into())))?;
969 // The floor and demand are applied to the policy's answer, never to the
970 // balance test: an account with nothing left still refuses, and both are
971 // capped by the balance `incoming` is part of, so they can only re-select
972 // capacity the account demonstrably has.
973 let granted = policy
974 .consolidation_grant(requested, balance, incoming, needed)
975 .ok_or_else(|| {
976 // A refused consolidation's settlement never applies. Attest
977 // only where that settlement would not itself have removed
978 // funding, the same rule the SQL backend needs because its
979 // transaction has already applied the settlement it rolls back.
980 if requested.is_zero() || !attest_refusal {
981 return AllocateError::InsufficientBalance;
982 }
983 funding_refusal(record.budget_view().shortfall())
984 })?;
985 let next_fence = record
986 .next_fence
987 .checked_add(1)
988 .ok_or_else(|| AllocateError::Storage(StoreError("fencing token overflow".into())))?;
989 Ok(GrantPlan {
990 granted,
991 fencing_token: FencingToken(record.next_fence),
992 next_fence,
993 next_lease_id,
994 })
995}
996
997/// The refusal a ledger with nothing allocatable attests: exhaustion when no
998/// funding remains, otherwise how much remains in other leases.
999fn funding_refusal(evidence: tollgate_core::BalanceShortfall) -> AllocateError {
1000 match evidence.exhaustion() {
1001 Some(exhausted) => AllocateError::BalanceExhausted(exhausted),
1002 None => AllocateError::BalanceInsufficient(evidence),
1003 }
1004}
1005
1006/// Read after every half of the exchange has applied: the evidence must
1007/// describe the ledger the grant committed into, settlement loss included.
1008fn allocation(inner: &Inner, grant: LeaseGrant) -> Allocation {
1009 let funding = inner
1010 .accounts
1011 .get(&grant.account_id)
1012 .expect("a granted account exists")
1013 .budget_view()
1014 .shortfall();
1015 Allocation {
1016 grant,
1017 funding: Some(funding),
1018 }
1019}
1020
1021fn apply_grant(
1022 inner: &mut Inner,
1023 account: AccountId,
1024 plan: GrantPlan,
1025 expires_at: Timestamp,
1026) -> LeaseGrant {
1027 let record = inner
1028 .accounts
1029 .get_mut(&account)
1030 .expect("the plan validated this account under this same lock");
1031 // Allowance first, and the split travels with the lease so settlement
1032 // returns each half where it came from (GL-97).
1033 let drawn = record
1034 .balance
1035 .take(plan.granted)
1036 .expect("grant never exceeds the balance the plan sized it against");
1037 let period_start = record.period_start;
1038 record.next_fence = plan.next_fence;
1039
1040 inner.next_lease_id = plan.next_lease_id;
1041 let lease_id = LeaseId(plan.next_lease_id);
1042 inner.leases.open(
1043 lease_id,
1044 LeaseRecord::opened(account, plan.fencing_token, plan.granted, expires_at)
1045 .funded_by(drawn, period_start),
1046 );
1047 LeaseGrant {
1048 lease_id,
1049 account_id: account,
1050 fencing_token: plan.fencing_token,
1051 units: plan.granted,
1052 expires_at,
1053 }
1054}
1055
1056#[async_trait]
1057impl LeaseAllocator for MemoryStore {
1058 async fn acquire(
1059 &self,
1060 account: AccountId,
1061 requested: CostUnits,
1062 ttl: SignedDuration,
1063 now: Timestamp,
1064 ) -> Result<Allocation, AllocateError> {
1065 let expires_at = self.grant_expiry(ttl, now)?;
1066 let mut inner = self.lock();
1067 let plan = plan_grant(
1068 &inner,
1069 &self.policy,
1070 account,
1071 requested,
1072 CostUnits::ZERO,
1073 CostUnits::ZERO,
1074 true,
1075 )?;
1076 let grant = apply_grant(&mut inner, account, plan, expires_at);
1077 Ok(allocation(&inner, grant))
1078 }
1079
1080 async fn release(
1081 &self,
1082 lease_id: LeaseId,
1083 fencing_token: FencingToken,
1084 unspent: CostUnits,
1085 now: Timestamp,
1086 ) -> Result<(), AllocateError> {
1087 let mut inner = self.lock();
1088 let plan = plan_release(&inner, &self.policy, lease_id, fencing_token, unspent, now)?;
1089 apply_release(&mut inner, lease_id, &plan);
1090 Ok(())
1091 }
1092
1093 async fn consolidate(
1094 &self,
1095 lease_id: LeaseId,
1096 fencing_token: FencingToken,
1097 unspent: CostUnits,
1098 requested: CostUnits,
1099 needed: CostUnits,
1100 ttl: SignedDuration,
1101 now: Timestamp,
1102 ) -> Result<Allocation, AllocateError> {
1103 let expires_at = self.grant_expiry(ttl, now)?;
1104 let mut inner = self.lock();
1105 // Both halves are planned before either is applied. The reference
1106 // backend has no rollback, so a grant that refuses after the release
1107 // had already moved units would settle a lease it cannot replace —
1108 // which is the one failure this operation exists to prevent, arriving
1109 // by a different door.
1110 let release = plan_release(&inner, &self.policy, lease_id, fencing_token, unspent, now)?;
1111 let account = release.account_id;
1112 let record = inner
1113 .accounts
1114 .get(&account)
1115 .ok_or(AllocateError::UnknownAccount)?;
1116 // Not `unspent`: the allowance half of a lease funded by a closed
1117 // period expires rather than returning, so the grant is sized against
1118 // what the credit will actually restore (GL-97).
1119 let restored = spendable_credit(record, release.funding, release.period_start, unspent);
1120 let preserves_funding = release.loss.is_zero() && restored == unspent;
1121 let grant = plan_grant(
1122 &inner,
1123 &self.policy,
1124 account,
1125 requested,
1126 restored,
1127 needed,
1128 preserves_funding,
1129 )?;
1130
1131 apply_release(&mut inner, lease_id, &release);
1132 let grant = apply_grant(&mut inner, account, grant, expires_at);
1133 Ok(allocation(&inner, grant))
1134 }
1135
1136 async fn reclaim_expired_batch(
1137 &self,
1138 now: Timestamp,
1139 limit: NonZeroUsize,
1140 ) -> Result<ReclaimBatch, StoreError> {
1141 let mut inner = self.lock();
1142 // Walks the active index, oldest expiry first, and stops at the first
1143 // lease that is not yet due — so the cost is what is being reclaimed,
1144 // not what the process has ever leased (GL-23).
1145 let expired = inner
1146 .leases
1147 .reclaimable(self.policy.reclaim_cutoff(now), limit.get());
1148 // Planned, validated, then applied. `ReclaimBatch::try_new` is the
1149 // last fallible step, and running it after a loop that had already
1150 // settled leases and credited balances would return `Err` over a
1151 // ledger that had moved. Unreachable while `reclaimable` respects the
1152 // limit, but the structure is the defect, and it is one refactor away
1153 // from being reachable (GL-57).
1154 let mut reclaimed = Vec::with_capacity(expired.len());
1155 // A holder that never released cannot prove any unit unspent: its
1156 // lease accepted commits until `usable_until`, and whatever it had
1157 // committed but not flushed died with it. So a sweep settles the lease
1158 // as a release claiming nothing would: no credit, and the remainder
1159 // recorded as provisional settlement loss (GL-136). Usage that arrives
1160 // later still fits in that gap and converts loss into billed usage
1161 // (see `ingest`), which is how a holder that outlived an outage is
1162 // billed rather than dropped.
1163 for &lease_id in &expired {
1164 let lease = inner.leases.get(lease_id).expect("just listed");
1165 let forfeited = lease
1166 .granted
1167 .checked_sub(lease.used)
1168 .expect("usage never exceeds grant");
1169 reclaimed.push(ReclaimedLease {
1170 lease_id,
1171 account_id: lease.account_id,
1172 forfeited,
1173 });
1174 }
1175 let batch = ReclaimBatch::try_new(reclaimed, limit)?;
1176
1177 for entry in batch.reclaimed() {
1178 assert!(
1179 inner
1180 .leases
1181 .settle(entry.lease_id, Settled::Expired, CostUnits::ZERO),
1182 "reclaimable only yields active leases"
1183 );
1184 let record = inner
1185 .accounts
1186 .get_mut(&entry.account_id)
1187 .expect("lease account exists");
1188 record.settlement_loss = record
1189 .settlement_loss
1190 .checked_add(entry.forfeited)
1191 .expect("loss overflow");
1192 }
1193 // One line per sweep rather than one per batch: a drain calls this
1194 // until a batch comes back unsaturated, and that last call carries the
1195 // post-sweep numbers.
1196 if batch.reclaimed().len() < limit.get() {
1197 let held = StoredRecords {
1198 credential_activity: inner.credential_activity.len(),
1199 usage_events: inner.usage.len(),
1200 leases: inner.leases.len(),
1201 active_leases: inner.leases.active_len(),
1202 };
1203 tracing::debug!(
1204 credential_activity = held.credential_activity,
1205 usage_events = held.usage_events,
1206 leases = held.leases,
1207 active_leases = held.active_leases,
1208 "memory store holdings; usage events and settled leases are never reclaimed"
1209 );
1210 }
1211 Ok(batch)
1212 }
1213}
1214
1215/// Insert a snapshot under a guard the caller already holds, returning what
1216/// to push once the guard is released.
1217///
1218/// Splitting the locked work from the push is what lets a caller hold one
1219/// guard across its own checks and this insert. `push_to_subscribers` must not
1220/// run under the lock, and a subscriber must never observe a push for a
1221/// publication that is still mid-flight — so the value to push comes back out
1222/// instead of being sent from in here.
1223///
1224/// Generation-monotonic: a replayed or reordered publish carrying an older or
1225/// equal generation is a no-op, and returns `None` because there is nothing to
1226/// announce.
1227fn publish_locked(
1228 inner: &mut Inner,
1229 principal: Principal,
1230 snapshot: PublishableSnapshot,
1231) -> Result<Option<PublishableSnapshot>, PublishSnapshotError> {
1232 if let Some(key_id) = snapshot.key_id
1233 && !inner.keys.get(&key_id).is_some_and(|stored| {
1234 stored.record.principal == principal && stored.record.account_id == snapshot.account_id
1235 })
1236 {
1237 return Err(PublishSnapshotError::CredentialMismatch { key_id });
1238 }
1239 if let Some(record) = inner.accounts.get(&snapshot.account_id) {
1240 if record.status != snapshot.status {
1241 return Err(PublishSnapshotError::StatusMismatch {
1242 ledger: record.status,
1243 submitted: snapshot.status,
1244 });
1245 }
1246 if record.capacity_class != snapshot.capacity_class {
1247 return Err(PublishSnapshotError::CapacityClassMismatch {
1248 ledger: record.capacity_class,
1249 submitted: snapshot.capacity_class,
1250 });
1251 }
1252 }
1253 if let Some(existing) = inner.snapshots.get(&principal)
1254 && existing.generation() >= snapshot.generation
1255 {
1256 return Ok(None);
1257 }
1258 // The store stamps the budget view; a publisher cannot supply one (GL-97).
1259 // Done here rather than at each caller so the two publication entry points
1260 // cannot drift, and under the same lock as the write so the number
1261 // published is the ledger as of that write. An account this store does not
1262 // hold publishes unstamped, which is the same "adds no account-existence
1263 // requirement" rule the status check follows.
1264 // Unconditional, including the `None` arm: a snapshot arrives here having
1265 // crossed a wire, where nothing stops a publisher putting a balance in the
1266 // JSON. Overwriting always is what makes the store the field's only
1267 // writer, rather than only usually.
1268 let view = inner
1269 .accounts
1270 .get(&snapshot.account_id)
1271 .map(AccountRecord::budget_view);
1272 let snapshot = snapshot.with_budget(view);
1273 inner
1274 .snapshots
1275 .insert(principal, SnapshotRecord::Present(snapshot.clone()));
1276 Ok(Some(snapshot))
1277}
1278
1279/// Tombstone a live snapshot, keeping its generation as the watermark, and
1280/// report the predecessor and result under the caller's guard.
1281fn remove_locked(inner: &mut Inner, principal: Principal) -> AdminReceipt<()> {
1282 let before = snapshot_audit(inner.snapshots.get(&principal));
1283 if let Some(SnapshotRecord::Present(snapshot)) = inner.snapshots.get(&principal) {
1284 let generation = snapshot.generation;
1285 inner
1286 .snapshots
1287 .insert(principal, SnapshotRecord::Revoked(generation));
1288 }
1289 AdminReceipt::new((), before, snapshot_audit(inner.snapshots.get(&principal)))
1290}
1291
1292/// Resolve `account`'s credential `key` to its stored record under the
1293/// caller's guard: the principal never leaves the store (GL-143).
1294fn account_key(
1295 inner: &Inner,
1296 account: AccountId,
1297 key: KeyId,
1298) -> Result<&StoredKey, KeySnapshotError> {
1299 inner
1300 .keys
1301 .get(&key)
1302 .filter(|stored| stored.record.account_id == account)
1303 .ok_or(KeySnapshotError::UnknownCredential)
1304}
1305
1306/// Which account-owned fact a republication is carrying.
1307///
1308/// The two operator actions — a status change and a capacity-class change —
1309/// differ only in the field they compare and set. Everything that makes the
1310/// republication *safe* is identical: the two-phase overflow check, the
1311/// tombstone skip, the already-at-target skip, and the deterministic ordering.
1312/// Writing that twice is how the two would drift, and the half that drifted
1313/// would be the half nobody was looking at.
1314#[derive(Debug, Clone, Copy)]
1315enum Restamp {
1316 Status(AccountStatus),
1317 CapacityClass(CapacityClass),
1318}
1319
1320impl Restamp {
1321 /// Whether this snapshot already carries the target, and so must not be
1322 /// rewritten — that is what makes a repeated call converge.
1323 fn already_applied(self, snapshot: &tollgate_core::AccountSnapshot) -> bool {
1324 match self {
1325 Restamp::Status(status) => snapshot.status == status,
1326 Restamp::CapacityClass(class) => snapshot.capacity_class == class,
1327 }
1328 }
1329
1330 fn apply(self, snapshot: &PublishableSnapshot, generation: Generation) -> PublishableSnapshot {
1331 match self {
1332 Restamp::Status(status) => snapshot.restamped(status, generation),
1333 Restamp::CapacityClass(class) => snapshot.reclassified(class, generation),
1334 }
1335 }
1336}
1337
1338/// Plan the re-stamping of every live snapshot of `account`, under a lock the
1339/// caller already holds. Mutates nothing: the caller applies
1340/// the plan only once every fallible step has succeeded.
1341///
1342/// **Two-phase on purpose.** Every new generation is computed, and every
1343/// overflow surfaced, *before* a single record is touched. A loop that
1344/// mutated as it went would leave an account half-republished behind a `u64`
1345/// overflow — one ledger status, two different snapshot statuses — which is
1346/// precisely the divergence GL-51 exists to abolish, reintroduced in the
1347/// backend that serves as the executable reference.
1348///
1349/// Tombstones are skipped: `snapshot: None` is a revoked principal, and
1350/// republishing it would resurrect it (INVARIANTS.md GL-15). Note this backend
1351/// skips them because the tombstone has *lost* its account attribution, while
1352/// PostgreSQL skips them by `deleted = FALSE` with the JSON still present —
1353/// different mechanisms, identical behaviour, which is what the mirrored
1354/// tests pin.
1355///
1356/// Rows already at the target are left alone, so a repeated call converges
1357/// and bumps no generation.
1358///
1359/// Returned sorted by principal so both backends emit pushes in the same
1360/// order and a mirrored test need not assert on incidental ordering.
1361fn plan_republish(
1362 inner: &Inner,
1363 account: AccountId,
1364 restamp: Restamp,
1365) -> Result<Vec<(Principal, PublishableSnapshot)>, SetStatusError> {
1366 let mut planned = Vec::new();
1367 for (principal, record) in &inner.snapshots {
1368 // Revoked principals are skipped: republishing one would resurrect it,
1369 // which INVARIANTS.md GL-15 forbids.
1370 let SnapshotRecord::Present(snapshot) = record else {
1371 continue;
1372 };
1373 if snapshot.account_id != account || restamp.already_applied(snapshot) {
1374 continue;
1375 }
1376 let generation = snapshot
1377 .generation
1378 .0
1379 .checked_add(1)
1380 .map(Generation)
1381 .ok_or_else(|| {
1382 SetStatusError::Storage(StoreError("snapshot generation overflow".into()))
1383 })?;
1384 // The restamped snapshot carries the new generation; the plan does not
1385 // carry a second copy of it (GL-54).
1386 planned.push((*principal, restamp.apply(snapshot, generation)));
1387 }
1388 planned.sort_unstable_by_key(|(principal, _)| *principal);
1389 Ok(planned)
1390}
1391
1392/// Apply a plan from [`plan_republish`]. Infallible by construction: every
1393/// value it writes was computed and checked before the first mutation, which
1394/// is what keeps the ledger and the snapshots from parting company when a
1395/// generation is about to overflow.
1396fn apply_republish(
1397 inner: &mut Inner,
1398 planned: Vec<(Principal, PublishableSnapshot)>,
1399) -> Vec<(Principal, PublishableSnapshot)> {
1400 let mut pushes = Vec::with_capacity(planned.len());
1401 for (principal, snapshot) in planned {
1402 inner
1403 .snapshots
1404 .insert(principal, SnapshotRecord::Present(snapshot.clone()));
1405 pushes.push((principal, snapshot));
1406 }
1407 pushes
1408}
1409
1410#[async_trait]
1411impl StoreHealth for MemoryStore {
1412 async fn ping(&self) -> Result<(), StoreError> {
1413 Ok(())
1414 }
1415}
1416
1417#[async_trait]
1418impl AdminStore for MemoryStore {
1419 async fn create_account(
1420 &self,
1421 config: AccountConfig,
1422 ) -> Result<crate::AdminReceipt<()>, CreateAccountError> {
1423 self.create_account_audited(config, AdminAuthority::Operator)
1424 }
1425
1426 async fn create_provisioned_account(
1427 &self,
1428 account: AccountId,
1429 ) -> Result<crate::AdminReceipt<()>, CreateAccountError> {
1430 self.create_account_audited(
1431 AccountConfig {
1432 account_id: account,
1433 initial_balance: CostUnits::ZERO,
1434 status: AccountStatus::Suspended,
1435 capacity_class: CapacityClass::BestEffort,
1436 },
1437 AdminAuthority::Provisioner,
1438 )
1439 }
1440
1441 async fn deposit(
1442 &self,
1443 account: AccountId,
1444 units: CostUnits,
1445 ) -> Result<crate::AdminReceipt<()>, AllocateError> {
1446 self.deposit_audited(account, units)
1447 }
1448
1449 async fn set_budget_schedule(
1450 &self,
1451 account: AccountId,
1452 schedule: Option<BudgetSchedule>,
1453 ) -> Result<crate::AdminReceipt<()>, BudgetError> {
1454 let mut inner = self.lock();
1455 let record = inner
1456 .accounts
1457 .get_mut(&account)
1458 .ok_or(BudgetError::UnknownAccount)?;
1459 // Read before write, under the one guard, so the receipt reports what
1460 // this call replaced rather than what a later reader happens to find.
1461 let before = AdminState::Budget {
1462 schedule: record.schedule,
1463 };
1464 record.schedule = schedule;
1465 Ok(crate::AdminReceipt::new(
1466 (),
1467 before,
1468 AdminState::Budget { schedule },
1469 ))
1470 }
1471
1472 async fn roll_due_periods(
1473 &self,
1474 now: Timestamp,
1475 limit: NonZeroUsize,
1476 ) -> Result<RolloverBatch, StoreError> {
1477 let mut inner = self.lock();
1478 let mut rolled = Vec::new();
1479 // Choose *which* accounts this bounded batch rolls before rolling any,
1480 // oldest period first, exactly as `DUE_PERIODS_SQL` does with
1481 // `ORDER BY period_start_us LIMIT $3`.
1482 //
1483 // Iterating `accounts` directly and breaking at `limit` selected by
1484 // hash order, so with more accounts due than one batch can take, two
1485 // runs over identical state rolled different accounts — and a different
1486 // set again from the backend this one is the reference for (GL-100). The
1487 // drain loop means every due account is rolled eventually, so this was
1488 // not a ledger defect; it was an unbounded-in-principle wait for any
1489 // particular account, and a divergence no test could see.
1490 //
1491 // The account id breaks ties so the order is total. PostgreSQL leaves
1492 // equal `period_start_us` rows in an arbitrary order; being stricter
1493 // than the contract is safe, and being unpredictable is what this
1494 // avoids.
1495 #[allow(
1496 clippy::disallowed_methods,
1497 reason = "sorted below before the limit is applied, so the batch's membership is a function of the stored state"
1498 )]
1499 let mut due: Vec<(Timestamp, AccountId)> = inner
1500 .accounts
1501 .iter()
1502 .filter_map(|(account_id, record)| {
1503 let schedule = record.schedule?;
1504 (schedule.period.start_of(now) > record.period_start)
1505 .then_some((record.period_start, *account_id))
1506 })
1507 .collect();
1508 due.sort_unstable();
1509 due.truncate(limit.get());
1510
1511 for (_, account_id) in due {
1512 let account_id = &account_id;
1513 let record = inner
1514 .accounts
1515 .get_mut(account_id)
1516 .expect("the due set was taken from this map under the same lock");
1517 let Some(schedule) = record.schedule else {
1518 continue;
1519 };
1520 // The boundary test and the crossing happen under one lock, which
1521 // is what makes two passes racing a boundary produce one roll. The
1522 // reference backend gets that from the mutex; PostgreSQL gets it
1523 // from a row lock over the same comparison.
1524 let boundary = schedule.period.start_of(now);
1525 if boundary <= record.period_start {
1526 continue;
1527 }
1528 // One allowance, not one per missed boundary: an account left
1529 // unrolled for two months is entitled to what it has now, not to a
1530 // backlog.
1531 let expired = record.balance.allowance;
1532 let deposited = record
1533 .deposited
1534 .checked_add(schedule.allowance)
1535 .ok_or_else(|| StoreError(format!("deposit overflow for account {account_id}")))?;
1536 let total_expired = record
1537 .expired
1538 .checked_add(expired)
1539 .ok_or_else(|| StoreError(format!("expiry overflow for account {account_id}")))?;
1540 // Written only after every fallible step has succeeded: a rollover
1541 // that overflowed halfway would leave the account funded but not
1542 // credited, or credited twice at the next pass.
1543 record.deposited = deposited;
1544 record.expired = total_expired;
1545 record.balance.allowance = schedule.allowance;
1546 record.period_start = boundary;
1547 rolled.push(RolledAccount {
1548 account_id: *account_id,
1549 deposited: schedule.allowance,
1550 expired,
1551 });
1552 }
1553 RolloverBatch::try_new(rolled, limit)
1554 }
1555
1556 async fn set_account_status(
1557 &self,
1558 account: AccountId,
1559 status: AccountStatus,
1560 ) -> Result<crate::AdminReceipt<StatusChange>, SetStatusError> {
1561 self.set_status_as(account, status, AdminAuthority::Operator)
1562 }
1563
1564 async fn activate_provisioned(
1565 &self,
1566 account: AccountId,
1567 ) -> Result<crate::AdminReceipt<StatusChange>, SetStatusError> {
1568 self.set_status_as(account, AccountStatus::Active, AdminAuthority::Provisioner)
1569 }
1570
1571 async fn set_capacity_class(
1572 &self,
1573 account: AccountId,
1574 class: CapacityClass,
1575 ) -> Result<crate::AdminReceipt<StatusChange>, SetStatusError> {
1576 // Structurally identical to `set_account_status`, deliberately: one
1577 // lock across both records, refusals before anything moves, the whole
1578 // snapshot plan computed before the ledger is touched. What differs is
1579 // one field, and that difference lives in `Restamp` rather than in a
1580 // second copy of this procedure.
1581 let (republished, before) = {
1582 let mut inner = self.lock();
1583 let record = inner
1584 .accounts
1585 .get(&account)
1586 .ok_or(SetStatusError::UnknownAccount)?;
1587 // A closed account is terminal. Reclassifying one is meaningless,
1588 // and refusing costs nothing that a caller wanted.
1589 if record.status == AccountStatus::Closed {
1590 return Err(SetStatusError::AccountClosed);
1591 }
1592 let before = AdminState::CapacityClass {
1593 capacity_class: record.capacity_class,
1594 };
1595 let planned = plan_republish(&inner, account, Restamp::CapacityClass(class))?;
1596 inner
1597 .accounts
1598 .get_mut(&account)
1599 .expect("the account was found under this same guard")
1600 .capacity_class = class;
1601 (apply_republish(&mut inner, planned), before)
1602 };
1603 if pushes_exceed_capacity(republished.len()) {
1604 tracing::warn!(
1605 %account,
1606 principals = republished.len(),
1607 capacity = PUSH_CHANNEL_CAPACITY,
1608 "capacity class change emitted more pushes than the channel holds; \
1609 subscribers will resync"
1610 );
1611 }
1612 let planned_pushes = republished;
1613 let republished = planned_pushes.len();
1614 for (principal, snapshot) in planned_pushes {
1615 self.push_to_subscribers(SnapshotPush {
1616 principal,
1617 resolution: SnapshotResolution::Present(snapshot),
1618 });
1619 }
1620 Ok(AdminReceipt::new(
1621 StatusChange {
1622 republished,
1623 unreadable: 0,
1624 },
1625 before,
1626 AdminState::CapacityClass {
1627 capacity_class: class,
1628 },
1629 ))
1630 }
1631
1632 async fn publish_snapshot(
1633 &self,
1634 principal: Principal,
1635 snapshot: PublishableSnapshot,
1636 ) -> Result<crate::AdminReceipt<()>, PublishSnapshotError> {
1637 // One guard for the check and the write. Releasing it between them
1638 // left a window `set_account_status` could run through entirely: a
1639 // publish that passed the status check, a suspension that restamped
1640 // every snapshot then existing, and finally the insert of a principal
1641 // `plan_republish` never saw — because it did not exist yet. The
1642 // ledger said suspended, that principal's live snapshot said active,
1643 // and it kept being admitted until someone repeated the transition.
1644 // That is GL-51's defect reintroduced in the backend that serves as
1645 // INVARIANTS.md GL-22's reference. `PostgresStore` holds `FOR SHARE` on
1646 // the account row across the same pair, which `set_account_status`'s
1647 // `FOR UPDATE` serialises against.
1648 let (published, before, after) = {
1649 let mut inner = self.lock();
1650 let before = snapshot_audit(inner.snapshots.get(&principal));
1651 let published = publish_locked(&mut inner, principal, snapshot)?;
1652 (
1653 published,
1654 before,
1655 snapshot_audit(inner.snapshots.get(&principal)),
1656 )
1657 };
1658 if let Some(snapshot) = published {
1659 self.push_to_subscribers(SnapshotPush {
1660 principal,
1661 resolution: SnapshotResolution::Present(snapshot),
1662 });
1663 }
1664 Ok(AdminReceipt::new((), before, after))
1665 }
1666
1667 async fn account_view(&self, account: AccountId) -> Result<Option<AccountView>, StoreError> {
1668 // One guard over status, schedule and the funding equation. Reading
1669 // them separately could straddle a rollover or a status change and
1670 // describe a state this account was never in.
1671 let inner = self.lock();
1672 let Some(record) = inner.accounts.get(&account) else {
1673 return Ok(None);
1674 };
1675 let mut active_grants = CostUnits::ZERO;
1676 let mut active_used = CostUnits::ZERO;
1677 for lease in inner.leases.active_of(account) {
1678 active_grants = active_grants
1679 .checked_add(lease.granted)
1680 .ok_or_else(|| StoreError(format!("grant sum overflow for account {account}")))?;
1681 active_used = active_used
1682 .checked_add(lease.used)
1683 .ok_or_else(|| StoreError(format!("used sum overflow for account {account}")))?;
1684 }
1685 Ok(Some(AccountView {
1686 account_id: account,
1687 status: record.status,
1688 capacity_class: record.capacity_class,
1689 origin: record.origin,
1690 status_set_by: record.status_set_by,
1691 schedule: record.schedule,
1692 period_start: record.period_start,
1693 conservation: Conservation {
1694 deposited: record.deposited,
1695 overage_recorded: record.overage_recorded,
1696 balance: record.balance.total(),
1697 active_lease_grants: active_grants,
1698 settled_usage: record
1699 .usage_recorded
1700 .checked_sub(active_used)
1701 .ok_or_else(|| {
1702 StoreError(format!("active usage exceeds recorded for {account}"))
1703 })?,
1704 settlement_loss: record.settlement_loss,
1705 expired: record.expired,
1706 },
1707 }))
1708 }
1709
1710 async fn remove_snapshot(
1711 &self,
1712 principal: Principal,
1713 ) -> Result<crate::AdminReceipt<()>, StoreError> {
1714 Ok(self.remove_snapshot_audited(principal))
1715 }
1716}
1717
1718#[async_trait]
1719impl SnapshotSource for MemoryStore {
1720 async fn snapshot(&self, principal: Principal) -> Result<SnapshotResolution, StoreError> {
1721 Ok(match self.lock().snapshots.get(&principal) {
1722 Some(SnapshotRecord::Present(snapshot)) => {
1723 SnapshotResolution::Present(snapshot.clone())
1724 }
1725 Some(SnapshotRecord::Revoked(generation)) => SnapshotResolution::Revoked {
1726 generation: *generation,
1727 },
1728 None => SnapshotResolution::Unknown,
1729 })
1730 }
1731
1732 fn subscribe(&self) -> broadcast::Receiver<SnapshotPush> {
1733 self.push.subscribe()
1734 }
1735
1736 /// Tombstones included: a revoked principal is one an instance must keep
1737 /// tracking so it keeps *knowing* about the revocation. Dropping it from
1738 /// the catalogue would make it indistinguishable from a principal that
1739 /// never existed, which is the resurrection INVARIANTS.md GL-15 forbids.
1740 async fn principals(&self) -> Result<Option<Vec<Principal>>, StoreError> {
1741 // Sorted, because the catalogue is an output and `snapshots` is a
1742 // `HashMap`: returning its iteration order would make this call's
1743 // result depend on the hash seed rather than on the stored state, and
1744 // differ run to run. `PostgresStore` already answers
1745 // `ORDER BY principal`, so the reference backend was the one diverging
1746 // (GL-100). The suite's assertion sorted before comparing, which is how
1747 // it stayed invisible.
1748 #[allow(
1749 clippy::disallowed_methods,
1750 reason = "sorted on the next line before it leaves the function, so the hash order never reaches the caller"
1751 )]
1752 let mut principals: Vec<Principal> = self.lock().snapshots.keys().copied().collect();
1753 principals.sort_unstable_by_key(|principal| principal.0);
1754 Ok(Some(principals))
1755 }
1756}
1757
1758#[async_trait]
1759impl UsageSink for MemoryStore {
1760 async fn ingest(
1761 &self,
1762 events: &[UsageEvent],
1763 _now: Timestamp,
1764 ) -> Result<IngestReport, IngestError> {
1765 let mut inner = self.lock();
1766
1767 // Planned first, applied second. The overage branch can fail the whole
1768 // batch on an accounting overflow, and returning from the middle of an
1769 // applying loop left every earlier event of the batch committed while
1770 // the caller was told the batch failed (GL-57): `UsageWriter` reads a
1771 // `StoreError` as an outage and retries, the replay counts the applied
1772 // events as duplicates, and retry accounting then describes a batch
1773 // that partially succeeded as a total failure. `PostgresStore` runs
1774 // the batch in one transaction and rolls it back, so this is the shape
1775 // that makes the two backends agree (INVARIANTS.md GL-7).
1776 //
1777 // The plan carries overlays rather than reading `inner` twice, because
1778 // events in one batch see each other: two charges against the same
1779 // lease consume its capacity in order, and a request id repeated
1780 // inside a batch is a duplicate of the earlier one.
1781 let mut report = IngestReport::default();
1782 let mut accepted: Vec<UsageEvent> = Vec::new();
1783 let mut lease_used: HashMap<LeaseId, CostUnits> = HashMap::new();
1784 let mut usage_recorded: HashMap<AccountId, CostUnits> = HashMap::new();
1785 let mut overage_recorded: HashMap<AccountId, CostUnits> = HashMap::new();
1786 let mut settlement_loss: HashMap<AccountId, CostUnits> = HashMap::new();
1787 let mut planned: HashSet<tollgate_core::RequestId> = HashSet::new();
1788
1789 for event in events {
1790 if inner.usage.contains_key(&event.request_id) || planned.contains(&event.request_id) {
1791 report.duplicate += 1;
1792 continue;
1793 }
1794 // Overage has no lease, so it has neither a capability to verify
1795 // nor lease capacity to fit inside. It is recorded against the
1796 // account directly, moving two columns at once: `usage_recorded`,
1797 // so it is billed, and `overage_recorded`, so it is funded. Moving
1798 // only the first would break conservation by exactly these units.
1799 //
1800 // Deliberately not conditioned on the account's current
1801 // enforcement mode. The ledger does not carry the mode the
1802 // request was admitted under, and an account switched back to
1803 // `Strict` between admission and flush would otherwise have this
1804 // charge discarded — fail-open on accounting, to protect a
1805 // fail-closed decision that was already made correctly. The work
1806 // happened; it gets billed.
1807 let UsageSource::Leased {
1808 lease_id,
1809 fencing_token,
1810 } = event.source
1811 else {
1812 let Some(record) = inner.accounts.get(&event.account_id) else {
1813 report.rejected += 1;
1814 continue;
1815 };
1816 let recorded = *usage_recorded
1817 .get(&event.account_id)
1818 .unwrap_or(&record.usage_recorded);
1819 let overage = *overage_recorded
1820 .get(&event.account_id)
1821 .unwrap_or(&record.overage_recorded);
1822 // Both terms must move or neither does, or the equation is
1823 // left open (INVARIANTS.md GL-11). A total that cannot be
1824 // represented is corruption of a monotonic column rather than
1825 // a problem with this event, so it is surfaced as a store
1826 // error and the whole batch fails — and because nothing has
1827 // been applied yet, the batch fails whole.
1828 let (Some(next_recorded), Some(next_overage)) = (
1829 recorded.checked_add(event.units),
1830 overage.checked_add(event.units),
1831 ) else {
1832 // Refused, not unavailable: a monotonic column that cannot
1833 // absorb these units will not absorb them on the next
1834 // attempt either, so retrying this batch forever would
1835 // wedge the writer behind an arithmetic fact (GL-61).
1836 return Err(IngestError::Refused(StoreError(format!(
1837 "overage accounting overflow for account {:#034x}: recorded usage {} \
1838 and overage {} cannot absorb {}",
1839 event.account_id.0, recorded, overage, event.units
1840 ))));
1841 };
1842 usage_recorded.insert(event.account_id, next_recorded);
1843 overage_recorded.insert(event.account_id, next_overage);
1844 planned.insert(event.request_id);
1845 accepted.push(*event);
1846 report.accepted += 1;
1847 continue;
1848 };
1849 // Capability check: the (lease, token, account) triple must name
1850 // a known lease. Then the conservation check: the event must fit in
1851 // `granted - used - credited`. For an active lease `credited` is
1852 // zero (plain capacity check). For a released lease the gap is
1853 // exactly the provisional settlement loss, so a straggler that
1854 // was committed before release converts loss into billed usage.
1855 // For an expired lease the reclaim credited the full remainder —
1856 // nothing fits, so stragglers stay rejected (they'd double-count).
1857 let Some(lease) = inner.leases.get(lease_id) else {
1858 report.rejected += 1;
1859 continue;
1860 };
1861 if lease.fencing_token != fencing_token || lease.account_id != event.account_id {
1862 report.rejected += 1;
1863 continue;
1864 }
1865 let was_settled = !lease.is_active();
1866 let used = *lease_used.get(&lease_id).unwrap_or(&lease.used);
1867 let capacity = used
1868 .checked_add(lease.credited())
1869 .and_then(|committed| lease.granted.checked_sub(committed));
1870 let fits = matches!(capacity, Some(cap) if event.units <= cap);
1871 if !fits {
1872 report.rejected += 1;
1873 continue;
1874 }
1875 let account_id = lease.account_id;
1876 let record = inner
1877 .accounts
1878 .get(&account_id)
1879 .expect("lease account exists");
1880 let recorded = *usage_recorded
1881 .get(&account_id)
1882 .unwrap_or(&record.usage_recorded);
1883 lease_used.insert(
1884 lease_id,
1885 used.checked_add(event.units).expect("fits within grant"),
1886 );
1887 usage_recorded.insert(
1888 account_id,
1889 recorded.checked_add(event.units).ok_or_else(|| {
1890 IngestError::Refused(StoreError(format!(
1891 "usage accounting overflow for account {:#034x}: recorded usage {} \
1892 cannot absorb {}",
1893 account_id.0, recorded, event.units
1894 )))
1895 })?,
1896 );
1897 if was_settled {
1898 // The units move from provisional loss to billed usage.
1899 let loss = *settlement_loss
1900 .get(&account_id)
1901 .unwrap_or(&record.settlement_loss);
1902 settlement_loss.insert(
1903 account_id,
1904 loss.checked_sub(event.units)
1905 .expect("straggler fits within recorded loss"),
1906 );
1907 }
1908 planned.insert(event.request_id);
1909 accepted.push(*event);
1910 report.accepted += 1;
1911 }
1912
1913 // Attribute only the canonical accepted set. Replays were classified
1914 // above before their payload was trusted, and may never change history.
1915 let mut activity = HashMap::<KeyId, Timestamp>::new();
1916 let mut unattributed = 0;
1917 for event in &accepted {
1918 if let Some(key_id) = event.key_id
1919 && inner
1920 .keys
1921 .get(&key_id)
1922 .is_some_and(|key| key.record.account_id == event.account_id)
1923 {
1924 let at = crate::clock::timestamp_from_micros(event.occurred_at.as_microsecond())
1925 .expect("truncating a valid Timestamp to microseconds stays representable");
1926 activity
1927 .entry(key_id)
1928 .and_modify(|old| *old = (*old).max(at))
1929 .or_insert(at);
1930 } else {
1931 unattributed += 1;
1932 }
1933 }
1934 report.unattributed = Some(unattributed);
1935
1936 // Apply. Every value here was computed above, so nothing in this block
1937 // can fail and leave the ledger half-moved.
1938 for (lease_id, used) in lease_used {
1939 inner
1940 .leases
1941 .get_mut(lease_id)
1942 .expect("planned against a lease that exists")
1943 .used = used;
1944 }
1945 for (account_id, value) in usage_recorded {
1946 inner
1947 .accounts
1948 .get_mut(&account_id)
1949 .expect("planned against an account that exists")
1950 .usage_recorded = value;
1951 }
1952 for (account_id, value) in overage_recorded {
1953 inner
1954 .accounts
1955 .get_mut(&account_id)
1956 .expect("planned against an account that exists")
1957 .overage_recorded = value;
1958 }
1959 for (account_id, value) in settlement_loss {
1960 inner
1961 .accounts
1962 .get_mut(&account_id)
1963 .expect("planned against an account that exists")
1964 .settlement_loss = value;
1965 }
1966 for event in accepted {
1967 inner.usage.insert(event.request_id, event);
1968 }
1969 for (key_id, at) in activity {
1970 inner
1971 .credential_activity
1972 .entry(key_id)
1973 .and_modify(|old| *old = (*old).max(at))
1974 .or_insert(at);
1975 }
1976 Ok(report)
1977 }
1978}
1979
1980#[async_trait]
1981impl KeyDirectory for MemoryStore {
1982 async fn credential_activity(
1983 &self,
1984 keys: &[KeyId],
1985 ) -> Result<Vec<crate::CredentialActivity>, StoreError> {
1986 let inner = self.lock();
1987 Ok(keys
1988 .iter()
1989 .map(|&key_id| crate::CredentialActivity {
1990 key_id,
1991 state: if !inner.keys.contains_key(&key_id) {
1992 crate::CredentialActivityState::Unknown
1993 } else if let Some(&last_committed_at) = inner.credential_activity.get(&key_id) {
1994 crate::CredentialActivityState::Committed { last_committed_at }
1995 } else {
1996 crate::CredentialActivityState::Unobserved
1997 },
1998 })
1999 .collect())
2000 }
2001
2002 async fn insert_key(&self, record: KeyRecord) -> Result<(), KeyError> {
2003 let mut inner = self.lock();
2004 insert_key_locked(&mut inner, record).map(|receipt| receipt.outcome)
2005 }
2006
2007 async fn revoke_key(&self, key_id: KeyId, now: Timestamp) -> Result<Revocation, KeyError> {
2008 self.revoke_key_audited(key_id, now)
2009 .await
2010 .map(|receipt| receipt.outcome)
2011 }
2012
2013 async fn publish_key_snapshot(
2014 &self,
2015 account: AccountId,
2016 key: KeyId,
2017 snapshot: PublishableSnapshot,
2018 ) -> Result<AdminReceipt<()>, KeySnapshotError> {
2019 self.publish_key_snapshot_with_generation(account, key, snapshot, false)
2020 }
2021
2022 async fn publish_key_snapshot_next(
2023 &self,
2024 account: AccountId,
2025 key: KeyId,
2026 snapshot: PublishableSnapshot,
2027 ) -> Result<AdminReceipt<()>, KeySnapshotError> {
2028 self.publish_key_snapshot_with_generation(account, key, snapshot, true)
2029 }
2030
2031 async fn remove_key_snapshot(
2032 &self,
2033 account: AccountId,
2034 key: KeyId,
2035 ) -> Result<crate::AdminReceipt<()>, KeySnapshotError> {
2036 let (principal, receipt) = {
2037 let mut inner = self.lock();
2038 let principal = account_key(&inner, account, key)?.record.principal;
2039 (principal, remove_locked(&mut inner, principal))
2040 };
2041 self.announce_removal(principal, &receipt);
2042 Ok(receipt)
2043 }
2044
2045 async fn revoke_key_audited(
2046 &self,
2047 key_id: KeyId,
2048 now: Timestamp,
2049 ) -> Result<crate::AdminReceipt<Revocation>, KeyError> {
2050 let mut inner = self.lock();
2051 let Some(stored) = inner.keys.get(&key_id) else {
2052 return Err(KeyError::UnknownKey);
2053 };
2054 let account_id = stored.record.account_id;
2055 let before = AdminState::Credential {
2056 account_id,
2057 key_id,
2058 revoked: stored.revoked_at.is_some(),
2059 };
2060 if stored.revoked_at.is_some() {
2061 return Ok(crate::AdminReceipt::new(
2062 Revocation::AlreadyRetired,
2063 before,
2064 before,
2065 ));
2066 }
2067 let revision = next_credential_revision(inner.credential_revision)?;
2068 inner
2069 .keys
2070 .get_mut(&key_id)
2071 .expect("key exists under the same guard")
2072 .revoked_at = Some(now);
2073 inner.unrevoked_keys.remove(&key_id);
2074 inner.credential_revision = revision;
2075 Ok(crate::AdminReceipt::new(
2076 Revocation::Retired,
2077 before,
2078 AdminState::Credential {
2079 account_id,
2080 key_id,
2081 revoked: true,
2082 },
2083 ))
2084 }
2085
2086 async fn active_keys(&self, now: Timestamp) -> Result<Vec<KeyRecord>, StoreError> {
2087 let inner = self.lock();
2088 let active: Vec<KeyRecord> = inner
2089 .unrevoked_keys
2090 .iter()
2091 .map(|id| &inner.keys[id])
2092 // Expiry is decided here rather than by each reader, so every
2093 // backend answers "active" the same way and a projection cannot
2094 // disagree with the ledger about which credentials are live.
2095 // Compare the full Timestamp: truncating either side can change
2096 // an exclusive expiry boundary or the proof sent to a verifier.
2097 .filter(|stored| {
2098 stored
2099 .record
2100 .not_after
2101 .is_none_or(|not_after| now < not_after)
2102 })
2103 .map(|stored| stored.record.clone())
2104 .collect();
2105 // Deterministic order: the projection is rebuilt from this, and a
2106 // backend whose output order wanders makes two instances' tables
2107 // differ in a way no test would reproduce.
2108 Ok(active)
2109 }
2110
2111 async fn account_keys(
2112 &self,
2113 account: AccountId,
2114 after: Option<KeyId>,
2115 limit: NonZeroUsize,
2116 ) -> Result<Vec<KeySummary>, KeyError> {
2117 crate::validate_key_page_limit(limit)?;
2118 let inner = self.lock();
2119 if !inner.accounts.contains_key(&account) {
2120 return Err(KeyError::UnknownAccount);
2121 }
2122 // Revoked credentials are included, so `unrevoked_keys` is not the
2123 // index to read here; `keys` is a `HashMap`, so the order is imposed
2124 // rather than inherited. Sorting before truncating is what makes the
2125 // page a function of the stored state instead of the hash seed — the
2126 // same rule `principals` follows.
2127 #[allow(
2128 clippy::disallowed_methods,
2129 reason = "sorted below before the page is truncated, so the hash order never reaches the caller"
2130 )]
2131 let mut summaries: Vec<KeySummary> = inner
2132 .keys
2133 .values()
2134 .filter(|stored| stored.record.account_id == account)
2135 .filter(|stored| after.is_none_or(|cursor| stored.record.key_id > cursor))
2136 .map(|stored| KeySummary {
2137 key_id: stored.record.key_id,
2138 not_after: stored.record.not_after,
2139 revoked_at: stored.revoked_at,
2140 })
2141 .collect();
2142 summaries.sort_unstable_by_key(|summary| summary.key_id);
2143 summaries.truncate(limit.get());
2144 Ok(summaries)
2145 }
2146
2147 async fn insert_key_within(
2148 &self,
2149 record: KeyRecord,
2150 max_active: NonZeroUsize,
2151 now: Timestamp,
2152 ) -> Result<(), KeyError> {
2153 self.insert_key_within_audited(record, max_active, now)
2154 .await
2155 .map(|receipt| receipt.outcome)
2156 }
2157
2158 async fn insert_key_within_audited(
2159 &self,
2160 record: KeyRecord,
2161 max_active: NonZeroUsize,
2162 now: Timestamp,
2163 ) -> Result<crate::AdminReceipt<()>, KeyError> {
2164 // One guard across the count and the write. That is the whole point of
2165 // the method: releasing it between them would let two issuers each see
2166 // room for one more and both take it.
2167 let mut inner = self.lock();
2168 // Identity is decided *before* the bound, and the order is load-bearing.
2169 // A caller that lost the response resends the same `key_id`; by then
2170 // its own successful write may have filled the bound, and answering
2171 // `ActiveKeyLimit` would tell it to retire a credential when what
2172 // actually happened is that its first call worked. `AlreadyExists` is
2173 // both the true answer and the one that makes the retry safe (GL-121).
2174 if !inner.accounts.contains_key(&record.account_id) {
2175 return Err(KeyError::UnknownAccount);
2176 }
2177 if inner.keys.contains_key(&record.key_id)
2178 || inner.key_principals.contains_key(&record.principal)
2179 {
2180 return Err(KeyError::AlreadyExists);
2181 }
2182 #[allow(
2183 clippy::disallowed_methods,
2184 reason = "counts live credentials; a count does not depend on the order they are counted in"
2185 )]
2186 let live = inner
2187 .keys
2188 .values()
2189 .filter(|stored| stored.record.account_id == record.account_id)
2190 .filter(|stored| {
2191 KeySummary {
2192 key_id: stored.record.key_id,
2193 not_after: stored.record.not_after,
2194 revoked_at: stored.revoked_at,
2195 }
2196 .is_live(now)
2197 })
2198 .count();
2199 if live >= max_active.get() {
2200 return Err(KeyError::ActiveKeyLimit { limit: max_active });
2201 }
2202 insert_key_locked(&mut inner, record)
2203 }
2204}
2205
2206/// The credential write both [`KeyDirectory::insert_key`] and
2207/// [`KeyDirectory::insert_key_within`] perform, with the guard already held.
2208///
2209/// Shared rather than copied so the bounded path cannot drift from the
2210/// unbounded one on what it validates — and taken by `&mut Inner` so it
2211/// *cannot* acquire a lock, which is what makes the caller's guard span the
2212/// count and the write.
2213fn insert_key_locked(
2214 inner: &mut Inner,
2215 record: KeyRecord,
2216) -> Result<crate::AdminReceipt<()>, KeyError> {
2217 if !inner.accounts.contains_key(&record.account_id) {
2218 return Err(KeyError::UnknownAccount);
2219 }
2220 // Never destructive, for the reason `create_account` is not: an
2221 // overwrite would retire a live credential without saying so, and the
2222 // digest it replaced is unrecoverable.
2223 if inner.keys.contains_key(&record.key_id) {
2224 return Err(KeyError::AlreadyExists);
2225 }
2226 // Two credentials cannot share a principal. That value is the identity
2227 // admission decides with, so a collision would make one account's
2228 // revocation withdraw another's credential. At 128 bits of HMAC output
2229 // it means secret reuse or corruption rather than chance, and the
2230 // stored backend enforces it with a UNIQUE index — this is the
2231 // reference implementation of the same rule.
2232 if inner.key_principals.contains_key(&record.principal) {
2233 return Err(KeyError::AlreadyExists);
2234 }
2235 let revision = next_credential_revision(inner.credential_revision)?;
2236 inner.key_principals.insert(record.principal, record.key_id);
2237 inner.unrevoked_keys.insert(record.key_id);
2238 let after = AdminState::Credential {
2239 account_id: record.account_id,
2240 key_id: record.key_id,
2241 revoked: false,
2242 };
2243 inner.keys.insert(
2244 record.key_id,
2245 StoredKey {
2246 record,
2247 revoked_at: None,
2248 },
2249 );
2250 inner.credential_revision = revision;
2251 Ok(crate::AdminReceipt::new((), AdminState::Absent, after))
2252}
2253
2254fn next_credential_revision(revision: u64) -> Result<u64, KeyError> {
2255 revision
2256 .checked_add(1)
2257 .filter(|next| *next <= crate::MAX_KEY_REVISION)
2258 .ok_or_else(|| KeyError::Storage(StoreError("credential revision exhausted".into())))
2259}
2260
2261#[async_trait]
2262impl crate::KeySource for MemoryStore {
2263 async fn active_keys_page(
2264 &self,
2265 now: Timestamp,
2266 after: Option<KeyId>,
2267 limit: NonZeroUsize,
2268 ) -> Result<crate::KeyPage, StoreError> {
2269 use std::ops::Bound::{Excluded, Unbounded};
2270 crate::validate_key_page_limit(limit)?;
2271 let inner = self.lock();
2272 let lower = after.map_or(Unbounded, Excluded);
2273 let mut records: Vec<crate::CredentialRecord> = inner
2274 .unrevoked_keys
2275 .range((lower, Unbounded))
2276 .map(|id| &inner.keys[id].record)
2277 .filter(|record| record.not_after.is_none_or(|end| now < end))
2278 .take(limit.get() + 1)
2279 .cloned()
2280 .map(Into::into)
2281 .collect();
2282 let next_after = if records.len() > limit.get() {
2283 records.pop();
2284 records.last().map(|key| key.key_id)
2285 } else {
2286 None
2287 };
2288 crate::KeyPage::try_new(
2289 inner.credential_revision,
2290 now,
2291 after,
2292 limit,
2293 records,
2294 next_after,
2295 )
2296 }
2297}
2298
2299fn snapshot_audit(record: Option<&SnapshotRecord>) -> AdminState {
2300 match record {
2301 None => AdminState::Absent,
2302 Some(SnapshotRecord::Present(snapshot)) => AdminState::Snapshot {
2303 generation: snapshot.generation,
2304 revoked: false,
2305 },
2306 Some(SnapshotRecord::Revoked(generation)) => AdminState::Snapshot {
2307 generation: *generation,
2308 revoked: true,
2309 },
2310 }
2311}
2312
2313#[cfg(test)]
2314mod tests {
2315 use super::*;
2316 use std::num::NonZeroUsize;
2317
2318 #[tokio::test]
2319 async fn credential_revision_overflow_preserves_records_and_both_indexes() {
2320 use crate::KeySource;
2321 let store = store_with(100);
2322 let record = |id: u128| {
2323 let mut digest = [0; 32];
2324 digest[..16].copy_from_slice(&id.to_be_bytes());
2325 KeyRecord {
2326 key_id: KeyId(id),
2327 account_id: ACCOUNT,
2328 principal: Principal(id),
2329 digest,
2330 not_after: None,
2331 }
2332 };
2333 store.insert_key(record(1)).await.unwrap();
2334 store.inner.lock().unwrap().credential_revision = crate::MAX_KEY_REVISION - 1;
2335 store.insert_key(record(2)).await.unwrap();
2336 assert!(matches!(
2337 store.insert_key(record(3)).await,
2338 Err(KeyError::Storage(_))
2339 ));
2340 assert!(matches!(
2341 store.revoke_key(KeyId(1), t(100)).await,
2342 Err(KeyError::Storage(_))
2343 ));
2344 let page = store
2345 .active_keys_page(t(100), None, crate::DEFAULT_KEY_PAGE_LIMIT)
2346 .await
2347 .unwrap();
2348 assert_eq!(page.revision(), crate::MAX_KEY_REVISION);
2349 assert_eq!(
2350 page.records().iter().map(|k| k.key_id).collect::<Vec<_>>(),
2351 vec![KeyId(1), KeyId(2)]
2352 );
2353 store.inner.lock().unwrap().credential_revision = 2;
2354 store.insert_key(record(3)).await.unwrap(); // failure left no duplicate-principal index entry
2355 }
2356
2357 const ACCOUNT: AccountId = AccountId(1);
2358 const TTL: SignedDuration = SignedDuration::from_secs(60);
2359
2360 fn t(secs: i64) -> Timestamp {
2361 Timestamp::from_second(secs).unwrap()
2362 }
2363
2364 /// Grants exactly what is asked for, so a test can open a precise number
2365 /// of leases without the default policy's halving getting in the way.
2366 fn exact_grants() -> GrantPolicy {
2367 GrantPolicy {
2368 shrink_divisor: 1,
2369 min_grant: CostUnits(1),
2370 max_ttl: SignedDuration::from_secs(300),
2371 reclaim_grace: SignedDuration::ZERO,
2372 }
2373 }
2374
2375 fn store_with(balance: u64) -> Arc<MemoryStore> {
2376 let store = MemoryStore::new(exact_grants()).expect("policy is valid");
2377 store.create_account(AccountConfig {
2378 account_id: ACCOUNT,
2379 initial_balance: CostUnits(balance),
2380 status: AccountStatus::Active,
2381 capacity_class: CapacityClass::Assured,
2382 });
2383 store
2384 }
2385
2386 /// Build a store holding `settled` long-settled leases plus one that is
2387 /// due for reclaim, sweep it, and report how many lease records the sweep
2388 /// had to look at.
2389 async fn examined_by_a_sweep_past(settled: usize) -> usize {
2390 let store = store_with(1_000_000);
2391 for _ in 0..settled {
2392 let lease = store
2393 .acquire(ACCOUNT, CostUnits(1), TTL, t(0))
2394 .await
2395 .expect("funded")
2396 .grant;
2397 store
2398 .release(lease.lease_id, lease.fencing_token, lease.units, t(1))
2399 .await
2400 .expect("active");
2401 }
2402 let due = store
2403 .acquire(ACCOUNT, CostUnits(1), TTL, t(0))
2404 .await
2405 .expect("funded")
2406 .grant;
2407
2408 let before = store.leases_examined();
2409 let batch = store
2410 .reclaim_expired_batch(t(120), NonZeroUsize::new(64).unwrap())
2411 .await
2412 .expect("sweep");
2413 assert_eq!(batch.len(), 1, "exactly the one expired lease is settled");
2414 assert_eq!(batch.reclaimed()[0].lease_id, due.lease_id);
2415 store.leases_examined() - before
2416 }
2417
2418 /// GL-23's whole claim: the sweep's cost follows the live population, not
2419 /// what the process has ever leased. Two settled populations two orders of
2420 /// magnitude apart must cost the sweep the same, which is a stronger
2421 /// statement than any absolute constant — and the one that fails if the
2422 /// index is removed and the filter goes back over the whole table.
2423 #[tokio::test]
2424 async fn sweeping_examines_only_live_leases() {
2425 let few = examined_by_a_sweep_past(10).await;
2426 let many = examined_by_a_sweep_past(1_000).await;
2427 assert_eq!(
2428 few, many,
2429 "1,000 settled leases cost the sweep {many} record reads against \
2430 {few} for 10; the sweep is walking history again"
2431 );
2432 assert!(
2433 many <= 4,
2434 "one live lease should not cost {many} record reads"
2435 );
2436 }
2437
2438 /// The same for the other scan. `conservation` is the ledger checker, so
2439 /// it stays a recomputation from the lease records — the index only
2440 /// narrows which records it reads.
2441 #[tokio::test]
2442 async fn conservation_examines_only_live_leases() {
2443 async fn examined_by_conservation_past(settled: usize) -> usize {
2444 let store = store_with(1_000_000);
2445 for _ in 0..settled {
2446 let lease = store
2447 .acquire(ACCOUNT, CostUnits(1), TTL, t(0))
2448 .await
2449 .expect("funded")
2450 .grant;
2451 store
2452 .release(lease.lease_id, lease.fencing_token, lease.units, t(1))
2453 .await
2454 .expect("active");
2455 }
2456 let _live = store.acquire(ACCOUNT, CostUnits(5), TTL, t(0)).await;
2457
2458 let before = store.leases_examined();
2459 let conservation = store.conservation(ACCOUNT).expect("account exists");
2460 assert!(
2461 conservation.holds(),
2462 "conservation violated: {conservation:?}"
2463 );
2464 assert_eq!(conservation.active_lease_grants, CostUnits(5));
2465 store.leases_examined() - before
2466 }
2467
2468 let few = examined_by_conservation_past(10).await;
2469 let many = examined_by_conservation_past(1_000).await;
2470 assert_eq!(
2471 few, many,
2472 "conservation read {many} records against {few}; it is summing \
2473 over settled leases again"
2474 );
2475 }
2476}