Skip to main content

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}