Skip to main content

tollgate_store/
traits.rs

1//! The storage traits and their shared vocabulary.
2
3use std::num::NonZeroUsize;
4
5use async_trait::async_trait;
6use jiff::{SignedDuration, Timestamp};
7use tokio::sync::broadcast;
8
9use tollgate_core::{
10    AccountId, AccountStatus, BudgetSchedule, CapacityClass, CostUnits, FencingToken, Generation,
11    KeyId, LeaseGrant, LeaseId, Principal, PublishableSnapshot, UsageEvent,
12};
13
14/// Backend failure unrelated to domain rules (connection lost, transaction
15/// aborted). Callers treat it as retryable-with-backoff; it must never be
16/// conflated with a domain refusal.
17#[derive(Debug, Clone, PartialEq, Eq)]
18pub struct StoreError(pub String);
19
20impl std::fmt::Display for StoreError {
21    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
22        write!(f, "store error: {}", self.0)
23    }
24}
25
26impl std::error::Error for StoreError {}
27
28/// Domain refusals from the allocator.
29#[derive(Debug, Clone, PartialEq, Eq)]
30pub enum AllocateError {
31    UnknownAccount,
32    /// The account exists but is not in a state that may spend.
33    AccountInactive,
34    /// Nothing left to lease. Distinct from `AccountInactive`: the client
35    /// should keep polling, because usage settlement or a top-up can restore
36    /// balance.
37    InsufficientBalance,
38    /// The ledger confirms no funding remains, including outstanding leases.
39    BalanceExhausted(tollgate_core::BalanceExhaustion),
40    /// No grant is possible, and the ledger confirms how much funding remains
41    /// outside this instance's reach — all of it held in other leases.
42    /// `remaining` is never zero; zero is [`Self::BalanceExhausted`].
43    /// `InsufficientBalance` stays the refusal that carries no attestation.
44    BalanceInsufficient(tollgate_core::BalanceShortfall),
45    /// A lease must specify one unambiguous, strictly positive lifetime.
46    InvalidTtl,
47    UnknownLease,
48    /// The fencing token does not match the lease record named by `lease_id`.
49    /// Token ordering across different active leases is irrelevant
50    /// (INVARIANTS.md GL-4).
51    Fenced,
52    /// The lease exists but is no longer active (already released, expired,
53    /// or reclaimed).
54    LeaseNotActive,
55    /// A release claimed more unspent units than the lease can still hold
56    /// (`unspent + recorded usage > granted`) — a client accounting bug,
57    /// surfaced rather than absorbed.
58    InvalidRelease,
59    Storage(StoreError),
60}
61
62impl AllocateError {
63    /// Stable metric labels, in [`index`](AllocateError::index) order.
64    ///
65    /// Variant names, not [`Display`](std::fmt::Display) output: `Storage`
66    /// wraps a backend message that varies per failure, so labelling by
67    /// rendered text would mint a fresh time series per connection error.
68    pub const NAMES: [&'static str; Self::COUNT] = [
69        "unknown_account",
70        "account_inactive",
71        "insufficient_balance",
72        "invalid_ttl",
73        "unknown_lease",
74        "fenced",
75        "lease_not_active",
76        "invalid_release",
77        "storage",
78        "balance_exhausted",
79        "balance_insufficient",
80    ];
81
82    /// How many distinct refusals exist — the width of a per-reason tally.
83    pub const COUNT: usize = 11;
84
85    /// This refusal's dense slot, for direct-indexed per-reason counters.
86    ///
87    /// `Storage` carries data, so there is no discriminant to cast; the
88    /// mapping is written out and the match is exhaustive, so a new variant
89    /// fails to compile until it is given a slot rather than silently landing
90    /// in another's bucket. Mirrors `DenyReason::index`.
91    #[must_use]
92    pub const fn index(&self) -> usize {
93        match self {
94            AllocateError::UnknownAccount => 0,
95            AllocateError::AccountInactive => 1,
96            AllocateError::InsufficientBalance => 2,
97            AllocateError::InvalidTtl => 3,
98            AllocateError::UnknownLease => 4,
99            AllocateError::Fenced => 5,
100            AllocateError::LeaseNotActive => 6,
101            AllocateError::InvalidRelease => 7,
102            AllocateError::Storage(_) => 8,
103            AllocateError::BalanceExhausted(_) => 9,
104            AllocateError::BalanceInsufficient(_) => 10,
105        }
106    }
107
108    /// This refusal's stable metric label.
109    #[must_use]
110    pub const fn name(&self) -> &'static str {
111        Self::NAMES[self.index()]
112    }
113}
114
115impl std::fmt::Display for AllocateError {
116    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
117        match self {
118            AllocateError::UnknownAccount => f.write_str("unknown account"),
119            AllocateError::AccountInactive => f.write_str("account inactive"),
120            AllocateError::InsufficientBalance => f.write_str("insufficient balance"),
121            AllocateError::BalanceExhausted(_) => f.write_str("account balance exhausted"),
122            AllocateError::BalanceInsufficient(evidence) => write!(
123                f,
124                "insufficient balance ({} units remain, all held in leases)",
125                evidence.remaining
126            ),
127            AllocateError::InvalidTtl => {
128                f.write_str("lease TTL must specify one positive duration")
129            }
130            AllocateError::UnknownLease => f.write_str("unknown lease"),
131            AllocateError::Fenced => f.write_str("fencing token mismatch"),
132            AllocateError::LeaseNotActive => f.write_str("lease not active"),
133            AllocateError::InvalidRelease => f.write_str("invalid release"),
134            AllocateError::Storage(e) => write!(f, "{e}"),
135        }
136    }
137}
138
139impl std::error::Error for AllocateError {}
140
141impl From<StoreError> for AllocateError {
142    fn from(error: StoreError) -> Self {
143        AllocateError::Storage(error)
144    }
145}
146
147/// Admin-side inputs when creating an account.
148#[derive(Debug, Clone, Copy)]
149pub struct AccountConfig {
150    pub account_id: AccountId,
151    pub initial_balance: CostUnits,
152    /// The account's administrative status at birth. Anything but
153    /// [`AccountStatus::Active`] refuses leases while keeping the ledger, and
154    /// `Closed` is terminal from creation onwards (INVARIANTS.md GL-22).
155    ///
156    /// An [`AccountStatus`] rather than a bool so creation and
157    /// [`AdminStore::set_account_status`] speak one vocabulary about one
158    /// column; the bool could not express `Closed`, which is what let
159    /// terminality be a convention instead of a check (GL-51).
160    pub status: AccountStatus,
161    /// The account's execution-capacity class at birth (GL-99).
162    ///
163    /// Present at creation for the reason `status` is: creation and
164    /// [`AdminStore::set_capacity_class`] speak one vocabulary about one
165    /// column, so there is no window in which an account exists without a
166    /// class and no second place that decides the default.
167    pub capacity_class: CapacityClass,
168}
169
170/// Per-account conservation view for reconciliation.
171///
172/// Read the equation as a funding statement: the left side is everything the
173/// account was ever funded with, the right side is where those units now sit.
174/// The two sources of funding are money in ([`deposited`](Self::deposited))
175/// and credit extended ([`overage_recorded`](Self::overage_recorded)); the
176/// three resting places are unspent balance, capacity currently out on lease,
177/// and units already consumed or written off.
178#[derive(Debug, Clone, Copy, PartialEq, Eq)]
179pub struct Conservation {
180    pub deposited: CostUnits,
181    /// Unfunded units billed under [`EnforcementMode::Elastic`]: spend no
182    /// deposit paid for and no lease debited.
183    ///
184    /// A *funding* term, on the left of the equation beside `deposited`, not
185    /// a bucket on the right. Overage usage also lands in `settled_usage`, so
186    /// without a matching term on the left the equation would fail by exactly
187    /// the overage — which is the whole reason this field exists rather than
188    /// the ledger simply recording the usage and saying nothing else.
189    ///
190    /// [`EnforcementMode::Elastic`]: tollgate_core::EnforcementMode::Elastic
191    pub overage_recorded: CostUnits,
192    pub balance: CostUnits,
193    pub active_lease_grants: CostUnits,
194    /// Usage billed against leases that have settled (released or expired),
195    /// plus all overage usage — which belongs to no lease and is therefore
196    /// settled the moment it is recorded. Usage on active leases is inside
197    /// `active_lease_grants`.
198    pub settled_usage: CostUnits,
199    pub settlement_loss: CostUnits,
200    /// Units that were funded but will never be spent, because the period
201    /// that funded them ended (GL-97).
202    ///
203    /// A resting place on the right of the equation, beside `settlement_loss`
204    /// and for the same reason: both are units the account was funded with
205    /// that no longer sit in a balance, on a lease, or in billed usage.
206    /// Without it, an allowance that resets each month would make the equation
207    /// fail by exactly the unspent remainder — the ledger reporting corruption
208    /// every time a budget did the one thing it exists to do.
209    ///
210    /// Monotonic, like `deposited` and `overage_recorded`: expiry is a fact
211    /// about a period that has closed, and closing a period is not reversible.
212    pub expired: CostUnits,
213}
214
215impl Conservation {
216    /// `deposited + overage == balance + active grants + settled usage + loss
217    /// + expired`, exactly.
218    ///
219    /// Both sides accumulate with checked arithmetic and an overflow answers
220    /// `false`, never a wrap or a panic: this function exists to *detect*
221    /// corrupt ledger state, so arithmetic that could not represent the state
222    /// must report a violation rather than quietly produce a total that
223    /// happens to match (INVARIANTS.md GL-11).
224    #[must_use]
225    pub fn holds(&self) -> bool {
226        let Some(funded) = self.deposited.checked_add(self.overage_recorded) else {
227            return false;
228        };
229        let mut sum = self.balance;
230        for part in [
231            self.active_lease_grants,
232            self.settled_usage,
233            self.settlement_loss,
234            self.expired,
235        ] {
236            match sum.checked_add(part) {
237                Some(next) => sum = next,
238                None => return false,
239            }
240        }
241        sum == funded
242    }
243}
244
245/// How an allocator sizes grants as an account's balance shrinks.
246///
247/// Deep balances grant the full request; near exhaustion the grant is capped
248/// at `balance / shrink_divisor` (floored at `min_grant`, and never above the
249/// remaining balance). This is the design-review answer to the quota-edge
250/// problem: N instances can no longer strand a small balance behind one
251/// holder's oversized lease. A consolidation may exceed the cap only to fund
252/// a quote the holder already refused ([`Self::consolidation_grant`]), which
253/// is demand rather than hoarding.
254#[derive(Debug, Clone, Copy)]
255pub struct GrantPolicy {
256    pub shrink_divisor: u64,
257    pub min_grant: CostUnits,
258    /// Hard cap on any single lease's TTL; requests beyond it are clamped.
259    pub max_ttl: SignedDuration,
260    /// How long past a lease's `expires_at` the allocator waits before
261    /// reclaiming its unspent units. Holders stop spending at
262    /// `expires_at - safety margin` (their side of the protocol), so work
263    /// committed inside the usability window has `margin + grace` to be
264    /// flushed and billed before settlement could reject it. Releases are
265    /// also accepted through the grace window (review finding GL-1).
266    pub reclaim_grace: SignedDuration,
267}
268
269/// Invalid allocator policy. Duration signs are part of the lease safety
270/// protocol, so invalid values are rejected at backend construction rather
271/// than normalized into a potentially unsafe policy.
272#[derive(Debug, Clone, Copy, PartialEq, Eq)]
273pub struct GrantPolicyError(pub &'static str);
274
275impl std::fmt::Display for GrantPolicyError {
276    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
277        f.write_str(self.0)
278    }
279}
280
281impl std::error::Error for GrantPolicyError {}
282
283impl Default for GrantPolicy {
284    fn default() -> Self {
285        GrantPolicy {
286            shrink_divisor: 2,
287            min_grant: CostUnits(1),
288            max_ttl: SignedDuration::from_secs(300),
289            reclaim_grace: SignedDuration::from_secs(30),
290        }
291    }
292}
293
294impl GrantPolicy {
295    /// Greatest expiry whose full grace window has elapsed at `now`.
296    /// A validated policy has nonnegative grace; subtraction underflow means
297    /// no representable expiry is due. In particular, a deadline beyond
298    /// Timestamp::MAX is never shortened to that last representable instant.
299    pub fn reclaim_cutoff(&self, now: Timestamp) -> Option<Timestamp> {
300        now.checked_sub(self.reclaim_grace).ok()
301    }
302
303    pub fn validate(&self) -> Result<(), GrantPolicyError> {
304        if self.shrink_divisor == 0 {
305            return Err(GrantPolicyError("shrink_divisor must be positive"));
306        }
307        if self.min_grant.is_zero() {
308            return Err(GrantPolicyError("min_grant must be positive"));
309        }
310        if self.max_ttl <= SignedDuration::ZERO {
311            return Err(GrantPolicyError("max_ttl must be positive"));
312        }
313        if self.reclaim_grace < SignedDuration::ZERO {
314            return Err(GrantPolicyError("reclaim_grace must not be negative"));
315        }
316        Ok(())
317    }
318
319    /// The units a request for `requested` receives from `balance`, or `None`
320    /// when the balance cannot fund any grant.
321    #[must_use]
322    pub fn grant(&self, requested: CostUnits, balance: CostUnits) -> Option<CostUnits> {
323        // Constructors reject these policy/request states, but `grant` is a
324        // public pure helper too. Keep direct use fail-closed instead of
325        // panicking on a zero divisor or manufacturing a unit for a zero
326        // request.
327        if requested.is_zero()
328            || balance.is_zero()
329            || self.shrink_divisor == 0
330            || self.min_grant.is_zero()
331        {
332            return None;
333        }
334        let cap = (balance.get() / self.shrink_divisor).max(self.min_grant.get());
335        Some(CostUnits(
336            requested.get().min(cap).min(balance.get()).max(1),
337        ))
338    }
339
340    /// The units a consolidation re-grants, or `None` exactly when
341    /// [`Self::grant`] refuses.
342    ///
343    /// `balance` already includes the credit the exchange restores, and
344    /// `floor` is that credit: the ordinary answer may not shrink the holding
345    /// (GL-109). `needed` is the largest quote the returned lease refused, and
346    /// the answer grows to it when `balance` can fund it (GL-131). The shrink
347    /// cap stops one holder hoarding a small balance ahead of demand; a quote
348    /// the holder already failed to fund is demand, and the request that
349    /// proved it spends it. A quote `balance` cannot fund grows nothing,
350    /// because no grant would serve it. A plain acquire is `floor` and
351    /// `needed` both zero, which is [`Self::grant`] unchanged.
352    #[must_use]
353    pub fn consolidation_grant(
354        &self,
355        requested: CostUnits,
356        balance: CostUnits,
357        floor: CostUnits,
358        needed: CostUnits,
359    ) -> Option<CostUnits> {
360        let sized = self.grant(requested, balance)?.max(floor.min(balance));
361        Some(if needed <= balance {
362            sized.max(needed)
363        } else {
364            sized
365        })
366    }
367}
368
369/// One lease settled by an expiry sweep.
370#[derive(Debug, Clone, Copy, PartialEq, Eq)]
371#[cfg_attr(feature = "wire", derive(serde::Serialize, serde::Deserialize))]
372pub struct ReclaimedLease {
373    pub lease_id: LeaseId,
374    pub account_id: AccountId,
375    /// Units recorded as provisional settlement loss: `granted - recorded
376    /// usage` at the sweep. Nothing is credited back, because a holder that
377    /// never released cannot prove any unit unspent (GL-136); usage for the
378    /// lease that arrives later converts loss into billed usage.
379    pub forfeited: CostUnits,
380}
381
382/// The production-sized upper bound for one expiry-reclaim transaction.
383///
384/// The limit bounds locks, row materialization, and SQL parameters per
385/// transaction; it does not cap a legitimate backlog because the maintenance
386/// task drains saturated batches until it reaches a partial one.
387pub const DEFAULT_RECLAIM_BATCH_LIMIT: NonZeroUsize =
388    NonZeroUsize::new(256).expect("the reclaim batch limit is nonzero");
389
390/// Verified evidence returned by one bounded expiry-reclaim transaction.
391///
392/// The fields are private so `saturated` cannot disagree with the requested
393/// limit. Callers may therefore use it to decide whether another batch is
394/// required without re-deriving the backend's result.
395#[derive(Debug, Clone, PartialEq, Eq)]
396pub struct ReclaimBatch {
397    reclaimed: Vec<ReclaimedLease>,
398    saturated: bool,
399}
400
401impl ReclaimBatch {
402    /// Build a batch and derive its saturation evidence from `limit`.
403    pub fn try_new(
404        reclaimed: Vec<ReclaimedLease>,
405        limit: NonZeroUsize,
406    ) -> Result<Self, StoreError> {
407        if reclaimed.len() > limit.get() {
408            return Err(StoreError(format!(
409                "reclaim backend returned {} leases for a batch limit of {}",
410                reclaimed.len(),
411                limit
412            )));
413        }
414        Ok(ReclaimBatch {
415            saturated: reclaimed.len() == limit.get(),
416            reclaimed,
417        })
418    }
419
420    #[must_use]
421    pub fn reclaimed(&self) -> &[ReclaimedLease] {
422        &self.reclaimed
423    }
424
425    #[must_use]
426    pub fn len(&self) -> usize {
427        self.reclaimed.len()
428    }
429
430    #[must_use]
431    pub fn is_empty(&self) -> bool {
432        self.reclaimed.is_empty()
433    }
434
435    #[must_use]
436    pub fn is_saturated(&self) -> bool {
437        self.saturated
438    }
439
440    #[must_use]
441    pub fn into_reclaimed(self) -> Vec<ReclaimedLease> {
442        self.reclaimed
443    }
444}
445
446/// A grant, and the ledger's remaining funding as of the transaction that
447/// made it.
448///
449/// `funding` is read from the committed ledger after the grant and any
450/// consolidation settlement, so it counts the new lease's units. It is an
451/// upper bound on what the account can still spend (see
452/// [`tollgate_core::BalanceShortfall`]), carried apart from the grant because
453/// the grant is a capability and this is evidence about the account. `None`
454/// means the answering allocator attested nothing, as an older server does.
455#[derive(Debug, Clone, Copy, PartialEq, Eq)]
456#[cfg_attr(feature = "wire", derive(serde::Serialize, serde::Deserialize))]
457pub struct Allocation {
458    #[cfg_attr(feature = "wire", serde(flatten))]
459    pub grant: LeaseGrant,
460    #[cfg_attr(
461        feature = "wire",
462        serde(default, skip_serializing_if = "Option::is_none")
463    )]
464    pub funding: Option<tollgate_core::BalanceShortfall>,
465}
466
467/// Atomic lease allocation against the account balance — the amortization
468/// point: one `acquire` funds thousands of local reservations.
469#[async_trait]
470pub trait LeaseAllocator: Send + Sync {
471    /// Atomically debit a grant from the account. The granted size follows
472    /// the backend's [`GrantPolicy`] and may be smaller than `requested`;
473    /// the fencing token comes from a strictly increasing per-account
474    /// sequence. It remains a capability for this lease only; allocating a
475    /// newer token does not invalidate another active lease.
476    async fn acquire(
477        &self,
478        account: AccountId,
479        requested: CostUnits,
480        ttl: SignedDuration,
481        now: Timestamp,
482    ) -> Result<Allocation, AllocateError>;
483
484    /// Graceful return: require the stored `(lease_id, fencing_token)` pair,
485    /// credit `unspent` back, and close the lease. Callers should flush usage
486    /// first when possible; events arriving after release are accepted only
487    /// when they fit its provisional settlement loss.
488    async fn release(
489        &self,
490        lease_id: LeaseId,
491        fencing_token: FencingToken,
492        unspent: CostUnits,
493        now: Timestamp,
494    ) -> Result<(), AllocateError>;
495
496    /// Atomically return an active lease's `unspent` units and re-grant
497    /// against the restored balance: [`release`](Self::release) followed by
498    /// [`acquire`](Self::acquire), in one transaction, for the same account
499    /// the lease names.
500    ///
501    /// This exists because the holder cannot compose it from the two calls.
502    /// A holder whose grant is too small for the work it is being offered is
503    /// holding exactly the units the next grant needs, and separating the
504    /// return from the request loses them twice over: the
505    /// [`GrantPolicy`] re-sizes against a balance the returned units have
506    /// already rejoined, so a `shrink_divisor` above one can hand back
507    /// *less* than was returned (49 units returned into a balance of 58
508    /// re-grants 29 under the default policy), and in the gap between the two
509    /// calls another instance can take them. Neither is recoverable by the
510    /// holder, which is why the exchange belongs to the component that owns
511    /// both the policy and the transaction (INVARIANTS.md GL-1, GL-6).
512    ///
513    /// **The grant is never smaller than the credit actually restored.**
514    /// Allowance funded by a closed period expires at settlement; only its
515    /// surviving top-up credit supplies a floor. Otherwise all `unspent`
516    /// units return. The result is the larger of this floor and the ordinary
517    /// policy grant, so it can exceed `requested` when preserving a larger
518    /// holding. A zero request or a balance unable to fund any grant refuses
519    /// the whole exchange.
520    ///
521    /// **The grant grows to `needed` when the restored balance can fund it.**
522    /// `needed` is the largest quote the returned lease refused, zero when
523    /// none. The shrink cap exists so one holder cannot hoard a small balance
524    /// ahead of demand; a refused quote is demand already proven, so under a
525    /// `shrink_divisor` above one it is what lets a single holder reach a
526    /// quote above `balance / shrink_divisor` at all. A `needed` the balance
527    /// cannot fund changes nothing. Sizing is
528    /// [`GrantPolicy::consolidation_grant`] in every backend.
529    ///
530    /// The transaction applies both halves or neither. A domain refusal
531    /// (`InsufficientBalance`, `BalanceExhausted`, `BalanceInsufficient`,
532    /// `UnknownAccount`, `AccountInactive`, or `InvalidTtl`) leaves the
533    /// original lease unchanged.
534    /// `InvalidRelease` also leaves it unchanged but reports an
535    /// accounting-integrity fault.
536    /// `UnknownLease`, `Fenced`, and `LeaseNotActive` provide no authority to
537    /// resume spending from the old lease.
538    ///
539    /// **`Storage`, timeouts, and cancellation have an ambiguous outcome.**
540    /// The transaction may have committed before its reply was lost, including
541    /// during an HTTP response or transaction-commit failure. A holder must
542    /// keep the old lease out of service and attempt its release; reinstating
543    /// it could spend credited units twice. The unanswered replacement grant
544    /// cannot be recovered through the old capability, must be reported as
545    /// uncertain, and remains bounded by TTL reclaim.
546    #[allow(
547        clippy::too_many_arguments,
548        reason = "one transactional exchange: the release half's capability and credit, the \
549                  grant half's size and demand, and the shared lifetime and clock"
550    )]
551    async fn consolidate(
552        &self,
553        lease_id: LeaseId,
554        fencing_token: FencingToken,
555        unspent: CostUnits,
556        requested: CostUnits,
557        needed: CostUnits,
558        ttl: SignedDuration,
559        now: Timestamp,
560    ) -> Result<Allocation, AllocateError>;
561
562    /// Settle at most `limit` active leases whose TTL (plus the policy's
563    /// reclaim grace) has lapsed and whose holder never released them.
564    ///
565    /// Nothing is credited back. Each lease's `granted - recorded usage` is
566    /// recorded as provisional settlement loss, exactly as a release claiming
567    /// nothing unspent would record it, because a holder that never released
568    /// cannot prove any unit unspent: it may have committed work it never
569    /// flushed (INVARIANTS.md GL-9, GL-136). Usage for the lease that arrives
570    /// later fits in that loss and converts it into billed usage.
571    ///
572    /// One call is one bounded atomic
573    /// transaction; [`ReclaimBatch::is_saturated`] is verified evidence that
574    /// the caller should immediately run another batch. It must be safe to
575    /// run concurrently with everything else.
576    async fn reclaim_expired_batch(
577        &self,
578        now: Timestamp,
579        limit: NonZeroUsize,
580    ) -> Result<ReclaimBatch, StoreError>;
581
582    /// Settle every currently expired lease through bounded transactions.
583    ///
584    /// This preserves the original full-drain caller API. If a later batch
585    /// fails, earlier batches are already committed, so the returned error
586    /// explicitly reports that partial progress rather than presenting the
587    /// operation as all-or-nothing.
588    async fn reclaim_expired(&self, now: Timestamp) -> Result<Vec<ReclaimedLease>, StoreError> {
589        drain_reclaim_expired(self, now).await
590    }
591}
592
593/// The drain loop [`LeaseAllocator::reclaim_expired`] performs, as a free
594/// function so that an override can reuse it instead of re-deriving it.
595///
596/// Rust has no `super` for a trait default, so a wrapper that overrides
597/// `reclaim_expired` cannot call the body it overrides. Without this it must
598/// choose between forwarding to an inner allocator — which silently discards
599/// the wrapper's own `reclaim_expired_batch` override, and so discards any
600/// failure that override injects — and copying this loop, which is how two
601/// copies drift apart. Calling this keeps one body and re-dispatches every
602/// batch through `allocator`, whatever `allocator` is (GL-83).
603///
604/// `#[doc(hidden)]` marks it cross-crate-visible for that purpose rather than
605/// part of the documented surface, as [`Reservation::reserve_at_locality`] is
606/// in `tollgate-core`.
607///
608/// [`Reservation::reserve_at_locality`]: https://docs.rs/tollgate-core
609#[doc(hidden)]
610pub async fn drain_reclaim_expired<A>(
611    allocator: &A,
612    now: Timestamp,
613) -> Result<Vec<ReclaimedLease>, StoreError>
614where
615    A: LeaseAllocator + ?Sized,
616{
617    let mut reclaimed: Vec<ReclaimedLease> = Vec::new();
618    loop {
619        let batch = match allocator
620            .reclaim_expired_batch(now, DEFAULT_RECLAIM_BATCH_LIMIT)
621            .await
622        {
623            Ok(batch) => batch,
624            Err(error) if reclaimed.is_empty() => return Err(error),
625            Err(error) => {
626                let units: u128 = reclaimed
627                    .iter()
628                    .map(|lease| u128::from(lease.forfeited.get()))
629                    .sum();
630                return Err(StoreError(format!(
631                    "reclaim drain failed after {} leases forfeiting {units} units were committed: {error}",
632                    reclaimed.len()
633                )));
634            }
635        };
636        let saturated = batch.is_saturated();
637        reclaimed.extend(batch.into_reclaimed());
638        if !saturated {
639            return Ok(reclaimed);
640        }
641        tokio::task::yield_now().await;
642    }
643}
644
645/// Authoritative state returned by a snapshot pull or push.
646#[derive(Debug, Clone)]
647pub enum SnapshotResolution {
648    /// A compiled snapshot is currently authoritative.
649    Present(PublishableSnapshot),
650    /// The principal existed but was revoked at this generation. Sources must
651    /// retain this watermark so a delayed older positive cannot resurrect it.
652    Revoked { generation: Generation },
653    /// The source has never observed this principal.
654    Unknown,
655}
656
657/// Slots in a backend's snapshot push channel.
658///
659/// Shared so the two backends cannot drift, and named rather than inlined
660/// because two things must agree on it: the channel, and the warning that
661/// fires when one operation would out-run it. Each slot retains an
662/// `Arc<AccountSnapshot>`, so this is also a bound on how much snapshot memory
663/// one slow subscriber can pin.
664pub const PUSH_CHANNEL_CAPACITY: usize = 256;
665
666/// Whether pushing `principals` updates at once will out-run the push channel.
667///
668/// A status change republishes every live snapshot of an account (GL-51), so a
669/// wide account can exceed the channel in one operation. Past this point every
670/// subscriber lags and resyncs its whole tracked set — correct, and bounded by
671/// the client's `max_concurrent_fetches`, but expensive enough that an
672/// operator should not have to infer it from a latency graph.
673///
674/// Strictly greater: a batch that exactly fills the channel is delivered, so
675/// warning at equality would cry wolf on the largest successful case. Pure, so
676/// the boundary is pinned by a test rather than by whichever backend is being
677/// read.
678#[must_use]
679pub fn pushes_exceed_capacity(principals: usize) -> bool {
680    principals > PUSH_CHANNEL_CAPACITY
681}
682
683/// One pushed snapshot update.
684#[derive(Debug, Clone)]
685pub struct SnapshotPush {
686    pub principal: Principal,
687    pub resolution: SnapshotResolution,
688}
689
690/// Where compiled snapshots come from.
691#[async_trait]
692pub trait SnapshotSource: Send + Sync {
693    /// Fetch the authoritative state for a principal. Revocation is distinct
694    /// from never-known so pull, lag recovery, and restart preserve the
695    /// generation watermark required for anti-resurrection semantics.
696    /// Reads used to reconstruct reclaimed local history must be linearizable
697    /// against durable publications/tombstones. Start a new source operation;
698    /// an earlier cached response or lagging replica cannot establish that
699    /// principal's forgotten generation floor. Return an error if this
700    /// authority is unavailable. MemoryStore and primary PostgresStore reads
701    /// supply this ordering; HTTP deployments must preserve it end to end.
702    async fn snapshot(&self, principal: Principal) -> Result<SnapshotResolution, StoreError>;
703
704    /// Subscribe to pushes. A lagging receiver may miss updates; the
705    /// contract is that a fresh `snapshot()` fetch after a lag error
706    /// observes at least the newest generation.
707    fn subscribe(&self) -> broadcast::Receiver<SnapshotPush>;
708
709    /// Every principal this source knows, including revoked ones — a
710    /// tombstone is still a principal an instance must track, so that it
711    /// knows the revocation (INVARIANTS.md GL-15).
712    ///
713    /// For instances that serve any customer rather than a configured slice
714    /// (GL-48). Pushes alone cannot answer this: they carry deltas from the
715    /// moment of subscribing, so a cold instance has no way to learn the set
716    /// that already exists.
717    ///
718    /// `Ok(None)` means this source cannot enumerate, and the manager stays
719    /// on its configured set — exactly today's behaviour. `Err` means
720    /// enumeration *failed* and is retried. The two are deliberately
721    /// distinct: collapsing them would let a broken source look like a
722    /// limited one, and an instance would quietly serve a stale set forever.
723    ///
724    /// Defaulted so a source that has no catalogue — a test double, an
725    /// embedder's own adapter — is unaffected.
726    async fn principals(&self) -> Result<Option<Vec<Principal>>, StoreError> {
727        Ok(None)
728    }
729}
730
731/// Liveness of the backing store, for readiness probes: a server must not
732/// report ready while its source of truth is unreachable (review finding
733/// GL-11).
734#[async_trait]
735pub trait StoreHealth: Send + Sync {
736    async fn ping(&self) -> Result<(), StoreError>;
737}
738
739/// Refusals from account creation (review finding GL-7): creation is never
740/// destructive and never silently idempotent — recreating an existing
741/// account is a surfaced error in every backend, because an overwrite would
742/// reset balances/fencing under live leases and a silent no-op would hide
743/// operator mistakes. Resetting an account is a deliberate, separate
744/// workflow, not a create.
745#[derive(Debug, Clone, PartialEq, Eq)]
746pub enum CreateAccountError {
747    AlreadyExists,
748    Storage(StoreError),
749}
750
751impl std::fmt::Display for CreateAccountError {
752    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
753        match self {
754            CreateAccountError::AlreadyExists => f.write_str("account already exists"),
755            CreateAccountError::Storage(e) => write!(f, "{e}"),
756        }
757    }
758}
759
760impl std::error::Error for CreateAccountError {}
761
762/// What a status transition actually did.
763///
764/// The blast radius of the operation, returned rather than logged, because an
765/// operator suspending an account has no other way to learn it: the ledger
766/// half is one row, but the snapshot half is however many credentials that
767/// account has, and 204 says nothing.
768///
769/// `republished == 0` is the interesting value. It means the account had no
770/// live snapshots to change — either it has no credentials yet, or the ones it
771/// has are all revoked, or a status change was repeated and everything was
772/// already at the target. All three are worth knowing at the moment of the
773/// call rather than from a later denial.
774#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
775pub struct StatusChange {
776    /// Live snapshots republished with the new status, at `generation + 1`.
777    /// Excludes tombstones, which are never republished, and snapshots already
778    /// carrying the target status, which are not rewritten.
779    pub republished: usize,
780    /// Rows that changed durably but could not be decoded well enough to push.
781    ///
782    /// Always zero in `MemoryStore`, which holds validated snapshots rather
783    /// than encoded ones. In a stored backend a row can be undecodable — it
784    /// already was before the transition touched it, and the request path
785    /// already refuses it — and the transition deliberately does not fail
786    /// whole over one corrupt credential. But it is not silently absorbed
787    /// either: those principals did not get a push, so they will not converge
788    /// until their next refresh.
789    pub unreadable: usize,
790}
791
792/// One account moved across a period boundary by a rollover pass.
793#[derive(Debug, Clone, Copy, PartialEq, Eq)]
794pub struct RolledAccount {
795    pub account_id: AccountId,
796    /// The new period's allowance, deposited by this pass.
797    pub deposited: CostUnits,
798    /// Unspent allowance from the period that just closed. Manual top-ups are
799    /// never included: they persist across a boundary (GL-97).
800    pub expired: CostUnits,
801}
802
803/// The production-sized upper bound for one rollover transaction.
804///
805/// Every scheduled account comes due at the same instant — that is what a
806/// calendar boundary means — so this is not a cap on a rare backlog but the
807/// normal shape of the first pass after midnight on the 1st. It bounds locks,
808/// row materialization, and transaction size per statement; the sweep drains
809/// saturated batches until it reaches a partial one, exactly as the expiry
810/// reclaim does.
811pub const DEFAULT_ROLLOVER_BATCH_LIMIT: NonZeroUsize =
812    NonZeroUsize::new(256).expect("the rollover batch limit is nonzero");
813
814/// Verified evidence returned by one bounded rollover transaction.
815///
816/// Private fields, so `saturated` cannot disagree with the requested limit —
817/// the same contract [`ReclaimBatch`] carries, and for the same reason: the
818/// caller decides whether to ask for another batch from this, without
819/// re-deriving the backend's result.
820#[derive(Debug, Clone, PartialEq, Eq)]
821pub struct RolloverBatch {
822    rolled: Vec<RolledAccount>,
823    saturated: bool,
824}
825
826impl RolloverBatch {
827    /// Build a batch and derive its saturation evidence from `limit`.
828    pub fn try_new(rolled: Vec<RolledAccount>, limit: NonZeroUsize) -> Result<Self, StoreError> {
829        if rolled.len() > limit.get() {
830            return Err(StoreError(format!(
831                "rollover backend returned {} accounts for a batch limit of {}",
832                rolled.len(),
833                limit
834            )));
835        }
836        Ok(RolloverBatch {
837            saturated: rolled.len() == limit.get(),
838            rolled,
839        })
840    }
841
842    #[must_use]
843    pub fn rolled(&self) -> &[RolledAccount] {
844        &self.rolled
845    }
846
847    #[must_use]
848    pub fn len(&self) -> usize {
849        self.rolled.len()
850    }
851
852    #[must_use]
853    pub fn is_empty(&self) -> bool {
854        self.rolled.is_empty()
855    }
856
857    #[must_use]
858    pub fn is_saturated(&self) -> bool {
859        self.saturated
860    }
861}
862
863/// Refusals from setting or rolling a budget schedule (GL-97).
864#[derive(Debug, Clone, PartialEq, Eq)]
865pub enum BudgetError {
866    /// No such account. Never a silent no-op: an operator setting a schedule
867    /// on a mistyped id has to learn it now rather than at the next boundary.
868    UnknownAccount,
869    Storage(StoreError),
870}
871
872impl std::fmt::Display for BudgetError {
873    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
874        match self {
875            BudgetError::UnknownAccount => f.write_str("no such account"),
876            BudgetError::Storage(e) => write!(f, "{e}"),
877        }
878    }
879}
880
881impl std::error::Error for BudgetError {}
882
883impl From<StoreError> for BudgetError {
884    fn from(error: StoreError) -> Self {
885        BudgetError::Storage(error)
886    }
887}
888
889/// Refusals from an account-status transition (GL-51).
890///
891/// Deliberately not an [`AllocateError`]: that enum's `NAMES`/`COUNT`/`index`
892/// are the width of `LeaseCounters`' per-reason tally, and a status refusal
893/// can never come out of `acquire`, so widening it would export a slot that
894/// is permanently zero in every deployment. [`CreateAccountError`] is the
895/// existing precedent for this shape — domain refusals plus `Storage`.
896#[derive(Debug, Clone, PartialEq, Eq)]
897pub enum SetStatusError {
898    /// No such account. Never a silent no-op, and the same answer whichever
899    /// status was asked for.
900    UnknownAccount,
901    /// [`AccountStatus::Closed`] is terminal: an account enters it from any
902    /// status and leaves it never. The refusal changes nothing — not the
903    /// ledger, not one snapshot, not one generation.
904    AccountClosed,
905    Storage(StoreError),
906}
907
908impl std::fmt::Display for SetStatusError {
909    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
910        match self {
911            SetStatusError::UnknownAccount => f.write_str("unknown account"),
912            SetStatusError::AccountClosed => f.write_str("account is closed"),
913            SetStatusError::Storage(e) => write!(f, "{e}"),
914        }
915    }
916}
917
918impl std::error::Error for SetStatusError {}
919
920impl From<StoreError> for SetStatusError {
921    fn from(error: StoreError) -> Self {
922        SetStatusError::Storage(error)
923    }
924}
925
926/// Refusals from publishing a snapshot (GL-51).
927///
928/// `publish_snapshot` used to return a bare [`StoreError`], which left it free
929/// to write a status contradicting the ledger and recreate the divergence
930/// [`AdminStore::set_account_status`] exists to abolish.
931#[derive(Debug, Clone, PartialEq, Eq)]
932pub enum PublishSnapshotError {
933    /// The stated credential is missing or belongs to another principal/account.
934    CredentialMismatch {
935        key_id: KeyId,
936    },
937    /// The snapshot's status disagrees with the account ledger. An account's
938    /// status is changed through [`AdminStore::set_account_status`], which
939    /// republishes; a publish may carry the current status but may not change
940    /// it.
941    StatusMismatch {
942        ledger: AccountStatus,
943        submitted: AccountStatus,
944    },
945    /// The snapshot's execution-capacity class disagrees with the account
946    /// ledger (GL-99). The class is an account-owned fact changed through
947    /// [`AdminStore::set_capacity_class`], which republishes; a publish may
948    /// carry the current class but may not change it. Two writers for one
949    /// fact is the divergence the status guard above already exists to
950    /// abolish, and a second field must not reintroduce it.
951    CapacityClassMismatch {
952        ledger: CapacityClass,
953        submitted: CapacityClass,
954    },
955    Storage(StoreError),
956}
957
958impl std::fmt::Display for PublishSnapshotError {
959    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
960        match self {
961            PublishSnapshotError::CredentialMismatch { key_id } => {
962                write!(
963                    f,
964                    "credential {key_id} does not bind the published principal and account"
965                )
966            }
967            PublishSnapshotError::StatusMismatch { ledger, submitted } => write!(
968                f,
969                "snapshot status {} contradicts account status {}",
970                submitted.as_str(),
971                ledger.as_str()
972            ),
973            PublishSnapshotError::CapacityClassMismatch { ledger, submitted } => write!(
974                f,
975                "snapshot capacity class {} contradicts account capacity class {}",
976                submitted.as_str(),
977                ledger.as_str()
978            ),
979            PublishSnapshotError::Storage(e) => write!(f, "{e}"),
980        }
981    }
982}
983
984impl std::error::Error for PublishSnapshotError {}
985
986impl From<StoreError> for PublishSnapshotError {
987    fn from(error: StoreError) -> Self {
988        PublishSnapshotError::Storage(error)
989    }
990}
991
992/// One account's administrative state, for an operator read (GL-121).
993///
994/// Assembled from types that already exist rather than a parallel vocabulary,
995/// so the HTTP surface reports the same terms the ledger reasons in and a
996/// reader can hold a dashboard next to a conservation check.
997///
998/// **Funding is not billing, and the shape says so.** A falling `balance` does
999/// not mean units were billed: it also falls when they go out on a lease that
1000/// has not settled, and it falls when a budget period closes and takes its
1001/// unspent allowance with it. Those are three different facts, and
1002/// [`Conservation`] keeps them apart — `active_lease_grants` is capacity
1003/// currently out, `settled_usage` is what was actually consumed,
1004/// `settlement_loss` and `expired` are what will never be. A surface that
1005/// reported only a balance would let a customer read depletion as spend, which
1006/// is exactly what GL-121 asks not to do.
1007#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1008pub struct AccountView {
1009    pub account_id: AccountId,
1010    pub status: AccountStatus,
1011    pub capacity_class: CapacityClass,
1012    /// The periodic allowance, if this account has one. `None` is "no
1013    /// schedule, the balance does not expire" — not "unknown".
1014    pub schedule: Option<BudgetSchedule>,
1015    /// First instant of the period currently in force. Meaningful only
1016    /// alongside a `schedule`; it is the marker rollover is idempotent
1017    /// against.
1018    pub period_start: Timestamp,
1019    /// Every term of the funding equation, so a caller can distinguish
1020    /// remaining funding from outstanding grants from settled usage.
1021    pub conservation: Conservation,
1022}
1023
1024/// Administrative writes: the control plane's mutation surface. Kept apart
1025/// from the data-plane traits so a read-only replica can implement those
1026/// without this.
1027///
1028/// HTTP-facing mutations return [`crate::AdminReceipt`] captured at the same
1029/// serialization point as the write. Its `outcome` contains the operation's
1030/// result; before/after describe the fields owned by that operation. A separate
1031/// read before or after the transaction is not a valid receipt under concurrency.
1032/// No-op receipts have equal states; errors carry no confirmed transition.
1033/// Direct callers attach their own actor and audit delivery policy.
1034#[async_trait]
1035pub trait AdminStore: Send + Sync {
1036    async fn create_account(
1037        &self,
1038        config: AccountConfig,
1039    ) -> Result<crate::AdminReceipt<()>, CreateAccountError>;
1040    async fn deposit(
1041        &self,
1042        account: AccountId,
1043        units: CostUnits,
1044    ) -> Result<crate::AdminReceipt<()>, AllocateError>;
1045    /// Set an existing account's administrative status, in one transaction:
1046    /// the ledger's status, and a republication of every *live* snapshot of
1047    /// that account carrying the new status at `generation + 1`.
1048    ///
1049    /// This is the whole operator action. Before GL-51 the ledger flag and the
1050    /// published `AccountStatus` were two records with two propagation paths
1051    /// and nothing checking them against each other, so "deactivate" returned
1052    /// success while the request path kept admitting.
1053    ///
1054    /// Rules, all enforced here rather than by caller discipline:
1055    /// - A missing account is [`SetStatusError::UnknownAccount`], never a
1056    ///   silent no-op, whichever status was asked for.
1057    /// - [`AccountStatus::Closed`] is terminal
1058    ///   ([`SetStatusError::AccountClosed`]); `Closed` → `Closed` is a no-op.
1059    /// - Revoked principals are never republished: resurrecting a tombstone
1060    ///   is what INVARIANTS.md GL-15 forbids, and revocation stays a separate
1061    ///   per-credential mechanism.
1062    /// - Snapshots already at the target status are not rewritten, so a
1063    ///   repeat converges and bumps no generation.
1064    /// - Outstanding leases are **not** reclaimed. Lease acquisition refuses
1065    ///   at once, but admission stops only when the new snapshot installs —
1066    ///   one `SnapshotManager` refresh interval, and already-debited units
1067    ///   settle at release or TTL reclaim (GL-9).
1068    async fn set_account_status(
1069        &self,
1070        account: AccountId,
1071        status: AccountStatus,
1072    ) -> Result<crate::AdminReceipt<StatusChange>, SetStatusError>;
1073
1074    /// Set an existing account's execution-capacity class, in one
1075    /// transaction: the ledger's class, and a republication of every *live*
1076    /// snapshot of that account carrying the new class at `generation + 1`
1077    /// (GL-99).
1078    ///
1079    /// The same operator action, and the same ownership argument, as
1080    /// [`set_account_status`](Self::set_account_status): the class is one
1081    /// fact with one writer. A control plane that published it per credential
1082    /// instead would recreate exactly the divergence GL-51 abolished — some of
1083    /// an account's principals assured and some best-effort, with nothing
1084    /// checking them against each other, and a request's treatment depending
1085    /// on which credential it arrived with.
1086    ///
1087    /// Rules, all enforced here rather than by caller discipline:
1088    /// - A missing account is [`SetStatusError::UnknownAccount`].
1089    /// - A closed account is [`SetStatusError::AccountClosed`]. Reclassifying
1090    ///   a terminally closed account is meaningless and the refusal changes
1091    ///   nothing, exactly as it does for a status change.
1092    /// - Revoked principals are never republished (INVARIANTS.md GL-15).
1093    /// - Snapshots already at the target class are not rewritten, so a repeat
1094    ///   converges and bumps no generation.
1095    /// - Nothing about funding changes. The class decides whether an instance
1096    ///   starts an already-funded request, so quota, leases, and outstanding
1097    ///   usage are untouched — an account reclassified mid-flight keeps every
1098    ///   charge it has already committed.
1099    ///
1100    /// Reuses [`SetStatusError`] rather than declaring a near-identical twin:
1101    /// the two refusals are the same two conditions about the same ledger row,
1102    /// and a second enum would be two vocabularies for one answer.
1103    async fn set_capacity_class(
1104        &self,
1105        account: AccountId,
1106        class: CapacityClass,
1107    ) -> Result<crate::AdminReceipt<StatusChange>, SetStatusError>;
1108
1109    /// Give an account a periodic allowance, or take it away.
1110    ///
1111    /// Setting a schedule does not deposit anything: the first allowance
1112    /// arrives at the first [`roll_due_periods`](Self::roll_due_periods) pass after the
1113    /// schedule exists. Depositing here would make "set a schedule" and "give
1114    /// this account units now" the same operation, and an operator correcting
1115    /// a mistyped allowance would fund the account twice.
1116    ///
1117    /// `None` removes the schedule and leaves the balance alone — including
1118    /// any unspent allowance, which simply stops expiring. Removing a schedule
1119    /// is not a way to claw units back.
1120    /// Set or clear an account's periodic allowance.
1121    ///
1122    /// Returns a receipt rather than `()` so the change can be audited like
1123    /// every other administrative mutation: an operator surface has to be able
1124    /// to report what a call actually committed, and a bare `Ok` cannot say
1125    /// whether a schedule was introduced, replaced, or was already what the
1126    /// caller asked for (GL-121). A repeat returns equal before/after states,
1127    /// which is how the convention expresses an idempotent no-op.
1128    async fn set_budget_schedule(
1129        &self,
1130        account: AccountId,
1131        schedule: Option<BudgetSchedule>,
1132    ) -> Result<crate::AdminReceipt<()>, BudgetError>;
1133
1134    /// Cross the period boundary for up to `limit` accounts that are past it:
1135    /// expire each closed period's unspent allowance and deposit the next one,
1136    /// one transaction per batch.
1137    ///
1138    /// **Idempotency is this method's job, not its caller's.** The pass runs
1139    /// on every control-plane replica, so two of them will race a boundary;
1140    /// the backend crosses it under a row lock, and the second caller then
1141    /// reads the period the first one wrote and skips the account. A caller
1142    /// that read each period first and then rolled would produce two deposits
1143    /// under exactly the race this exists to survive.
1144    ///
1145    /// Safe to call at any cadence: before a boundary it selects nothing, and
1146    /// after one the first caller wins. It never rolls an account more than
1147    /// one period at a time — an account left unrolled for two months lands in
1148    /// the current period with one allowance, because an allowance is what the
1149    /// account is entitled to now, not a backlog to be paid out.
1150    ///
1151    /// Bounded, and saturating means there is more: every scheduled account
1152    /// comes due at the same instant, so a caller must drain saturated batches
1153    /// until one comes back partial. Unscheduled accounts are never selected.
1154    async fn roll_due_periods(
1155        &self,
1156        now: Timestamp,
1157        limit: NonZeroUsize,
1158    ) -> Result<RolloverBatch, StoreError>;
1159    /// Publish a principal's compiled snapshot.
1160    ///
1161    /// Refused with [`PublishSnapshotError::StatusMismatch`] when the
1162    /// snapshot's status contradicts the account ledger, so the two records
1163    /// [`set_account_status`](AdminStore::set_account_status) unifies cannot
1164    /// be pulled apart again one principal at a time. A snapshot whose
1165    /// account the ledger does not hold publishes unchanged, as before.
1166    async fn publish_snapshot(
1167        &self,
1168        principal: Principal,
1169        snapshot: PublishableSnapshot,
1170    ) -> Result<crate::AdminReceipt<()>, PublishSnapshotError>;
1171    async fn remove_snapshot(
1172        &self,
1173        principal: Principal,
1174    ) -> Result<crate::AdminReceipt<()>, StoreError>;
1175
1176    /// One account's administrative state, or `None` if no such account (GL-121).
1177    ///
1178    /// The read an operator surface needs and the traits did not have. Both
1179    /// backends already expose `conservation` as an inherent method, but with
1180    /// different signatures — one synchronous returning an `Option`, one
1181    /// asynchronous returning a `Result` — so nothing generic over a backend
1182    /// could read an account at all.
1183    ///
1184    /// A single call rather than several, because the terms have to agree with
1185    /// each other: status, schedule and the funding equation read separately
1186    /// can straddle a rollover or a suspension and describe a state the account
1187    /// was never in. A backend answers this from one consistent read.
1188    async fn account_view(&self, account: AccountId) -> Result<Option<AccountView>, StoreError>;
1189}
1190
1191/// Outcome of one ingest batch.
1192#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
1193#[cfg_attr(feature = "wire", derive(serde::Serialize, serde::Deserialize))]
1194pub struct IngestReport {
1195    /// Newly recorded events.
1196    pub accepted: u64,
1197    /// Events whose `request_id` was already recorded (idempotent replay —
1198    /// INVARIANTS.md GL-7).
1199    pub duplicate: u64,
1200    /// Events refused: unknown lease, lease-capability mismatch, or no
1201    /// remaining accounting capacity, or units outside the backend's storage
1202    /// domain. These are bounded billing loss,
1203    /// visible to reconciliation.
1204    pub rejected: u64,
1205    /// Newly accepted events with absent, unknown or different-account key
1206    /// attribution. Duplicates never supply fresh activity evidence. `None`
1207    /// means this sink does not report attribution (including older servers),
1208    /// not that every event was attributed.
1209    #[cfg_attr(feature = "wire", serde(default))]
1210    pub unattributed: Option<u64>,
1211}
1212
1213impl IngestReport {
1214    /// A complete acknowledgement partitions this batch exactly once. Validate
1215    /// before releasing queued evidence, including replies from custom sinks.
1216    pub fn validate(&self, submitted: usize) -> Result<(), StoreError> {
1217        let total = self
1218            .accepted
1219            .checked_add(self.duplicate)
1220            .and_then(|n| n.checked_add(self.rejected));
1221        if !total.is_some_and(|n| u64::try_from(submitted) == Ok(n))
1222            || self.unattributed.is_some_and(|n| n > self.accepted)
1223        {
1224            return Err(StoreError(
1225                "invalid usage acknowledgement cardinality".into(),
1226            ));
1227        }
1228        Ok(())
1229    }
1230}
1231
1232/// The billing ledger's write side.
1233#[async_trait]
1234pub trait UsageSink: Send + Sync {
1235    /// Record a batch. Idempotent on `request_id`; every event must match its
1236    /// stored `(lease_id, account_id, fencing_token)` capability before lease
1237    /// state and accounting capacity are checked. Partial acceptance is
1238    /// normal — the report says what happened.
1239    ///
1240    /// Duplicates are classified before inspecting their payload. A new event
1241    /// outside the backend's unit domain is rejected individually; it is not
1242    /// remembered as accepted and cannot poison otherwise valid neighbors.
1243    /// MemoryStore supports `u64` units; PostgreSQL supports nonnegative
1244    /// `BIGINT` units (`0..=i64::MAX`). A representable event that overflows an
1245    /// accumulated accounting total refuses the entire batch atomically.
1246    async fn ingest(
1247        &self,
1248        events: &[UsageEvent],
1249        now: Timestamp,
1250    ) -> Result<IngestReport, IngestError>;
1251}
1252
1253/// One credential's durable record.
1254///
1255/// The `digest` is opaque here on purpose. The HMAC secret that produced it
1256/// lives with the verifier (`tollgate-auth`) and never reaches a store, so a
1257/// backend holds material that verifies nothing on its own — the property
1258/// `HmacRegistry` is built around, preserved across the persistence boundary.
1259/// A store that could compute a digest would be a store whose compromise is
1260/// sufficient to mint credentials.
1261///
1262/// `principal` is the digest's own truncation, so it is derived rather than
1263/// assigned: the request path is keyed by it, and per-credential revocation
1264/// is `install_revoked` for exactly this value.
1265#[derive(Clone, PartialEq, Eq)]
1266pub struct KeyRecord {
1267    pub key_id: KeyId,
1268    pub account_id: AccountId,
1269    /// The principal this credential authenticates as: the leading 128 bits
1270    /// of `digest`.
1271    pub principal: Principal,
1272    /// HMAC-SHA256 of the secret under the verifier's server secret.
1273    pub digest: [u8; 32],
1274    /// When the credential stops being valid of its own accord, independent
1275    /// of revocation. Surfaced to the request path through
1276    /// `Verified::reusable_until`, so a session cache cannot outlive it.
1277    pub not_after: Option<Timestamp>,
1278}
1279
1280impl std::fmt::Debug for KeyRecord {
1281    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1282        f.debug_struct("KeyRecord")
1283            .field("key_id", &self.key_id)
1284            .field("account_id", &self.account_id)
1285            .field("principal", &self.principal)
1286            .field("not_after", &self.not_after)
1287            .finish_non_exhaustive()
1288    }
1289}
1290
1291/// One requested key's activity. No observation is not proof of non-use.
1292#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1293pub struct CredentialActivity {
1294    pub key_id: KeyId,
1295    pub state: CredentialActivityState,
1296}
1297
1298#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1299pub enum CredentialActivityState {
1300    Unknown,
1301    Unobserved,
1302    /// Maximum accepted, attributable execution-start time, at microsecond
1303    /// precision. Never an authorization or independent server-clock fact.
1304    Committed {
1305        last_committed_at: Timestamp,
1306    },
1307}
1308
1309/// What a revocation actually did.
1310///
1311/// Returned rather than inferred, for the reason [`StatusChange`] is: an
1312/// operator retiring a suspicious credential needs to know whether they
1313/// retired anything. "Already revoked" and "no such key" are different
1314/// answers to the same request and only one of them is a mistake.
1315#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1316pub enum Revocation {
1317    /// The credential was live and is now retired.
1318    Retired,
1319    /// The credential was already retired; nothing changed.
1320    AlreadyRetired,
1321}
1322
1323/// The largest usage batch any sink accepts in one call, in events.
1324///
1325/// A batch cap has to admit the largest batch the system can legitimately
1326/// produce, not a comfortable one: the shipped example writes 256 per flush,
1327/// and this leaves a factor of sixteen for an embedder that batches harder to
1328/// cut round-trips. Beyond it the answer is more flushes, not a bigger body —
1329/// an ingest is idempotent, so splitting costs a round trip and risks nothing.
1330///
1331/// `UsageWriterConfig::validate` refuses a `max_batch` above this, so the
1332/// misconfiguration is a startup error rather than a permanently-rejected
1333/// batch discovered in production (GL-61).
1334pub const MAX_INGEST_BATCH: usize = 4_096;
1335
1336/// Why an ingest attempt failed, and whether replaying it unchanged could
1337/// ever succeed.
1338///
1339/// The distinction exists because the writer's correct response to the two is
1340/// opposite. An unreachable sink is a *duration*: retrying the same batch is
1341/// the designed behaviour, and giving up would lose billable events over a
1342/// blip. A refused batch is a *fact about the batch*: retrying it unchanged
1343/// gets the same answer forever, and every event queued behind it waits for a
1344/// recovery that cannot come — a permanent, deterministic error laundered
1345/// into an unbounded billing and availability outage (GL-61).
1346///
1347/// [`From<StoreError>`] yields [`Unavailable`](Self::Unavailable), so a
1348/// backend that does not classify keeps the retry-forever behaviour it had.
1349/// Terminality is asserted, never assumed.
1350#[derive(Debug, Clone, PartialEq, Eq)]
1351pub enum IngestError {
1352    /// The sink could not be reached, or could not answer in time. The batch
1353    /// is unchanged and will be retried.
1354    Unavailable(StoreError),
1355    /// The sink refused this batch and will refuse it again unchanged: a body
1356    /// over the endpoint's limit, an event it cannot decode, a contract it
1357    /// does not implement. Retrying cannot help.
1358    Refused(StoreError),
1359}
1360
1361impl IngestError {
1362    /// Whether replaying this batch unchanged could ever succeed.
1363    #[must_use]
1364    pub const fn is_retryable(&self) -> bool {
1365        matches!(self, IngestError::Unavailable(_))
1366    }
1367}
1368
1369impl From<StoreError> for IngestError {
1370    fn from(error: StoreError) -> Self {
1371        IngestError::Unavailable(error)
1372    }
1373}
1374
1375impl std::fmt::Display for IngestError {
1376    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1377        match self {
1378            IngestError::Unavailable(e) => write!(f, "{e}"),
1379            IngestError::Refused(e) => write!(f, "refused: {e}"),
1380        }
1381    }
1382}
1383
1384impl std::error::Error for IngestError {}
1385
1386/// Refusals from credential lifecycle operations.
1387#[derive(Debug, Clone, PartialEq, Eq)]
1388pub enum KeyError {
1389    /// No such credential. Never a silent no-op — an operator revoking a key
1390    /// that does not exist has either the wrong id or a false belief about
1391    /// what is live, and both are worth surfacing.
1392    UnknownKey,
1393    /// The credential's account does not exist, so nothing could authenticate
1394    /// as it. Refused at issuance rather than producing a key that verifies
1395    /// and is then denied by every admission.
1396    UnknownAccount,
1397    /// This `key_id` is already recorded. Issuance is never destructive, for
1398    /// the reason account creation is not ([`CreateAccountError`]): an
1399    /// overwrite would silently retire a live credential.
1400    ///
1401    /// This is also the retry answer. A caller that supplies the `key_id` and
1402    /// loses the response resends the same one and is told the credential
1403    /// exists — which is the truth, and which discloses no secret. That is why
1404    /// issuance must never become an upsert (GL-121).
1405    AlreadyExists,
1406    /// The account already holds `limit` live credentials, so issuing another
1407    /// would exceed the bound the caller supplied.
1408    ///
1409    /// "Live" excludes revoked keys and keys whose `not_after` has passed: a
1410    /// bound that counted expired credentials would strand an account behind
1411    /// keys nobody can authenticate with.
1412    ActiveKeyLimit {
1413        limit: NonZeroUsize,
1414    },
1415    Storage(StoreError),
1416}
1417
1418impl std::fmt::Display for KeyError {
1419    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1420        match self {
1421            KeyError::UnknownKey => f.write_str("no such credential"),
1422            KeyError::UnknownAccount => f.write_str("no such account"),
1423            KeyError::AlreadyExists => f.write_str("credential already exists"),
1424            KeyError::ActiveKeyLimit { limit } => {
1425                write!(f, "account already holds {limit} live credentials")
1426            }
1427            KeyError::Storage(e) => write!(f, "{e}"),
1428        }
1429    }
1430}
1431
1432impl std::error::Error for KeyError {}
1433
1434/// Refusals from binding or withdrawing a snapshot by the credential it was
1435/// issued as (GL-143).
1436#[derive(Debug, Clone, PartialEq, Eq)]
1437pub enum KeySnapshotError {
1438    /// No such credential, or it belongs to another account. One answer for
1439    /// both, as revocation gives: a foreign `key_id` discloses nothing.
1440    UnknownCredential,
1441    /// The credential was revoked. Revocation is terminal (INVARIANTS.md
1442    /// GL-27), so it is never granted positive authorization again; withdrawal
1443    /// remains allowed.
1444    Retired {
1445        key_id: KeyId,
1446    },
1447    /// Publication refused as [`AdminStore::publish_snapshot`] would refuse it.
1448    Publish(PublishSnapshotError),
1449    Storage(StoreError),
1450}
1451
1452impl std::fmt::Display for KeySnapshotError {
1453    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1454        match self {
1455            KeySnapshotError::UnknownCredential => f.write_str("no such credential"),
1456            KeySnapshotError::Retired { key_id } => {
1457                write!(f, "credential {key_id} is revoked")
1458            }
1459            KeySnapshotError::Publish(e) => write!(f, "{e}"),
1460            KeySnapshotError::Storage(e) => write!(f, "{e}"),
1461        }
1462    }
1463}
1464
1465impl std::error::Error for KeySnapshotError {}
1466
1467impl From<PublishSnapshotError> for KeySnapshotError {
1468    fn from(error: PublishSnapshotError) -> Self {
1469        Self::Publish(error)
1470    }
1471}
1472
1473impl From<StoreError> for KeySnapshotError {
1474    fn from(error: StoreError) -> Self {
1475        Self::Storage(error)
1476    }
1477}
1478
1479impl From<StoreError> for KeyError {
1480    fn from(error: StoreError) -> Self {
1481        Self::Storage(error)
1482    }
1483}
1484
1485/// One credential as an *administrator* sees it (GL-121).
1486///
1487/// Deliberately not a [`KeyRecord`]. A record carries `digest` — the HMAC the
1488/// verifier compares against — and `principal`, documented there as "the
1489/// leading 128 bits of `digest`". Both are digest material, and an
1490/// account-scoped listing is reachable by an application backend and from
1491/// there a browser, so neither may appear in it. `key_id` is the non-secret
1492/// handle: what the caller chose, what revocation names, what an audit shows.
1493///
1494/// Expiry and revocation are surfaced separately. A credential that lapsed on
1495/// its own is a different operational fact from one an operator withdrew, and
1496/// collapsing both into "inactive" loses the distinction exactly where someone
1497/// is deciding whether to issue a replacement.
1498#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1499pub struct KeySummary {
1500    pub key_id: KeyId,
1501    /// When the credential stops being valid of its own accord, if ever.
1502    pub not_after: Option<Timestamp>,
1503    /// When an operator withdrew it, if they did. Terminal.
1504    pub revoked_at: Option<Timestamp>,
1505}
1506
1507impl KeySummary {
1508    /// Whether this credential can still authenticate at `now`.
1509    ///
1510    /// The predicate the issuance bound counts with, written once so a listing
1511    /// and a limit cannot disagree about what "live" means.
1512    #[must_use]
1513    pub fn is_live(&self, now: Timestamp) -> bool {
1514        self.revoked_at.is_none() && self.not_after.is_none_or(|until| now < until)
1515    }
1516}
1517
1518/// Durable credential lifecycle: the half of key management that outlives a
1519/// process and is shared by a fleet.
1520///
1521/// **Why this is a store trait rather than registry state.** Every other
1522/// control-plane fact here — accounts, balances, leases, statuses, snapshots —
1523/// is transactional in a backend with a local read projection, and credentials
1524/// are the same class of fact. An in-memory-only registry would lose issued
1525/// keys on restart, keep verifying a credential another instance revoked, and
1526/// make a fleet-wide active-key limit unenforceable, because no instance sees
1527/// the fleet.
1528///
1529/// **Durability before disclosure.** [`insert_key`](Self::insert_key) must
1530/// commit before its caller returns the secret to anyone. A crash between the
1531/// two hands out a credential the server has never heard of, which no later
1532/// reconciliation can repair: the digest is unrecoverable from the record,
1533/// which is the point of storing digests.
1534///
1535/// The verifier's projection is rebuilt from [`active_keys`](Self::active_keys),
1536/// so a backend decides what "active" means once, here, rather than in each
1537/// reader.
1538#[async_trait]
1539pub trait KeyDirectory: crate::KeySource {
1540    /// Inspect requested keys in input order, including repeated IDs and
1541    /// retired keys. Every input has one explicit result or the read fails.
1542    /// This operator read uses O(keys.len()) output memory; backends bound
1543    /// individual queries internally. Multiple chunks need not share an instant.
1544    async fn credential_activity(
1545        &self,
1546        keys: &[KeyId],
1547    ) -> Result<Vec<CredentialActivity>, StoreError>;
1548
1549    /// Record a minted credential. The caller has already generated the
1550    /// secret and computed its digest; this stores what remains.
1551    async fn insert_key(&self, record: KeyRecord) -> Result<(), KeyError>;
1552
1553    /// Retire one credential, reporting whether it was live.
1554    ///
1555    /// Revocation is durable and terminal: a retired credential is never
1556    /// resurrected, for the same reason a snapshot tombstone is not
1557    /// (INVARIANTS.md GL-15).
1558    async fn revoke_key(&self, key_id: KeyId, now: Timestamp) -> Result<Revocation, KeyError>;
1559
1560    /// Unbounded operator read of every credential valid at `now`. Serving
1561    /// instances use `KeySource` pages and an owned, bounded drain instead.
1562    /// Retained for existing direct-store lifecycle tooling; no hidden page cap.
1563    async fn active_keys(&self, now: Timestamp) -> Result<Vec<KeyRecord>, StoreError>;
1564
1565    /// One account's credentials, ordered by `key_id`, for an operator
1566    /// listing (GL-121).
1567    ///
1568    /// Distinct from [`active_keys`](Self::active_keys), which is the
1569    /// fleet-wide, digest-bearing projection an *instance* pulls: this is
1570    /// account-scoped, bounded, and carries no digest material, because it
1571    /// answers a different question for a different caller.
1572    ///
1573    /// Paginate with `after` — the greatest `key_id` already seen, exclusive.
1574    /// Revoked and expired credentials are included, because an administrator
1575    /// deciding whether to issue a replacement needs to see what became of the
1576    /// last one; [`KeySummary::is_live`] separates them.
1577    async fn account_keys(
1578        &self,
1579        account: AccountId,
1580        after: Option<KeyId>,
1581        limit: NonZeroUsize,
1582    ) -> Result<Vec<KeySummary>, StoreError>;
1583
1584    /// Record a credential only if the account holds fewer than `max_active`
1585    /// live ones, counting and inserting indivisibly (GL-121).
1586    ///
1587    /// The bound is supplied per call rather than stored: what counts as a
1588    /// reasonable number of credentials belongs to the application's plan, not
1589    /// to Tollgate, and a value in the request is one the caller can change
1590    /// without a migration.
1591    ///
1592    /// **Why this is not [`active_keys`](Self::active_keys) then
1593    /// [`insert_key`](Self::insert_key).** Two issuers racing that pair both
1594    /// read `max_active - 1`, both insert, and the account ends up over the
1595    /// bound with no error raised anywhere — the more replicas, the likelier.
1596    /// A backend must make the count and the insert one indivisible step:
1597    /// `MemoryStore` holds a single lock across both, and `PostgresStore`
1598    /// takes the account row `FOR UPDATE` first — the row
1599    /// `set_account_status` already serialises against.
1600    ///
1601    /// Live excludes revoked credentials, and those whose `not_after` has
1602    /// passed at `now`, so an account cannot be stranded behind keys that can
1603    /// no longer authenticate.
1604    async fn insert_key_within(
1605        &self,
1606        record: KeyRecord,
1607        max_active: NonZeroUsize,
1608        now: Timestamp,
1609    ) -> Result<(), KeyError>;
1610
1611    /// Bounded issuance with lifecycle evidence captured under the mutation lock.
1612    /// HTTP administrators must use this receipt rather than synthesize history.
1613    async fn insert_key_within_audited(
1614        &self,
1615        record: KeyRecord,
1616        max_active: NonZeroUsize,
1617        now: Timestamp,
1618    ) -> Result<crate::AdminReceipt<()>, KeyError>;
1619
1620    /// Retire a credential and capture its actual owner, key and predecessor
1621    /// under the mutation lock. A repeated revocation returns equal states.
1622    async fn revoke_key_audited(
1623        &self,
1624        key_id: KeyId,
1625        now: Timestamp,
1626    ) -> Result<crate::AdminReceipt<Revocation>, KeyError>;
1627
1628    /// Publish `snapshot` for the principal of `account`'s credential `key`,
1629    /// resolved inside the store (GL-143).
1630    ///
1631    /// An operator holds `(account, key)`; the principal is digest material
1632    /// and never leaves the server. Resolution, the retirement check and the
1633    /// publication are one indivisible step, so a concurrent revocation
1634    /// either precedes it — and the publish is refused as
1635    /// [`KeySnapshotError::Retired`] — or follows it.
1636    ///
1637    /// `snapshot` must state `key_id == Some(key)`; anything else is
1638    /// [`PublishSnapshotError::CredentialMismatch`]. Every other rule is
1639    /// [`AdminStore::publish_snapshot`]'s, including the generation no-op.
1640    async fn publish_key_snapshot(
1641        &self,
1642        account: AccountId,
1643        key: KeyId,
1644        snapshot: PublishableSnapshot,
1645    ) -> Result<crate::AdminReceipt<()>, KeySnapshotError>;
1646
1647    /// Withdraw the snapshot of `account`'s credential `key`, tombstoning it
1648    /// as [`AdminStore::remove_snapshot`] does. Allowed for a revoked
1649    /// credential: revocation does not withdraw its snapshot, and withdrawal
1650    /// is the safe direction.
1651    async fn remove_key_snapshot(
1652        &self,
1653        account: AccountId,
1654        key: KeyId,
1655    ) -> Result<crate::AdminReceipt<()>, KeySnapshotError>;
1656}
1657
1658#[cfg(test)]
1659mod tests {
1660    use std::sync::atomic::{AtomicUsize, Ordering};
1661
1662    use super::KeySummary;
1663    use jiff::Timestamp;
1664    use tollgate_core::KeyId;
1665
1666    fn at(seconds: i64) -> Timestamp {
1667        Timestamp::from_second(seconds).expect("a test instant")
1668    }
1669
1670    /// `not_after` is exclusive, and the instant itself is the whole question.
1671    ///
1672    /// Both backends filter with `now < not_after` — `memory.rs` in three
1673    /// places, and four SQL predicates written `not_after > $now`. `is_live`
1674    /// is the summary of exactly those queries, so `<=` here would not merely
1675    /// be off by an instant: a listing would report a credential live for the
1676    /// one instant at which every query that selects credentials has already
1677    /// dropped it, and the issuance bound counts with this predicate.
1678    #[test]
1679    fn a_credential_is_dead_at_its_expiry_instant_not_after_it() {
1680        let expiring = |not_after| KeySummary {
1681            key_id: KeyId(1),
1682            not_after: Some(not_after),
1683            revoked_at: None,
1684        };
1685        assert!(
1686            expiring(at(100)).is_live(at(99)),
1687            "live up to the instant before"
1688        );
1689        assert!(
1690            !expiring(at(100)).is_live(at(100)),
1691            "dead *at* the boundary: expiry is exclusive, as both backends filter it"
1692        );
1693        assert!(!expiring(at(100)).is_live(at(101)), "and dead after it");
1694    }
1695
1696    use super::*;
1697
1698    #[derive(Clone, Copy)]
1699    enum ReclaimScript {
1700        FailFirst,
1701        FullBatchThenFail,
1702    }
1703
1704    struct ScriptedReclaimer {
1705        script: ReclaimScript,
1706        calls: AtomicUsize,
1707    }
1708
1709    #[async_trait]
1710    impl LeaseAllocator for ScriptedReclaimer {
1711        async fn acquire(
1712            &self,
1713            _account: AccountId,
1714            _requested: CostUnits,
1715            _ttl: SignedDuration,
1716            _now: Timestamp,
1717        ) -> Result<Allocation, AllocateError> {
1718            unreachable!("the full-drain tests only reclaim")
1719        }
1720
1721        async fn release(
1722            &self,
1723            _lease_id: LeaseId,
1724            _fencing_token: FencingToken,
1725            _unspent: CostUnits,
1726            _now: Timestamp,
1727        ) -> Result<(), AllocateError> {
1728            unreachable!("the full-drain tests only reclaim")
1729        }
1730
1731        async fn consolidate(
1732            &self,
1733            _lease_id: LeaseId,
1734            _fencing_token: FencingToken,
1735            _unspent: CostUnits,
1736            _requested: CostUnits,
1737            _needed: CostUnits,
1738            _ttl: SignedDuration,
1739            _now: Timestamp,
1740        ) -> Result<Allocation, AllocateError> {
1741            unreachable!("the reclaim drain never consolidates")
1742        }
1743
1744        async fn reclaim_expired_batch(
1745            &self,
1746            _now: Timestamp,
1747            limit: NonZeroUsize,
1748        ) -> Result<ReclaimBatch, StoreError> {
1749            let call = self.calls.fetch_add(1, Ordering::AcqRel);
1750            if matches!(self.script, ReclaimScript::FailFirst) || call > 0 {
1751                return Err(StoreError("scripted reclaim failure".into()));
1752            }
1753            let reclaimed = (0..limit.get())
1754                .map(|id| ReclaimedLease {
1755                    lease_id: LeaseId(u128::try_from(id).unwrap()),
1756                    account_id: AccountId(1),
1757                    forfeited: CostUnits(1),
1758                })
1759                .collect();
1760            ReclaimBatch::try_new(reclaimed, limit)
1761        }
1762    }
1763
1764    /// GL-131: under the default divisor of 2, a 60-unit balance re-granted as
1765    /// 30 forever, however often a 51-unit quote was refused.
1766    #[test]
1767    fn consolidation_grows_only_to_a_fundable_needed_quote() {
1768        let policy = GrantPolicy::default();
1769        let size = |requested, balance, floor, needed| {
1770            policy.consolidation_grant(
1771                CostUnits(requested),
1772                CostUnits(balance),
1773                CostUnits(floor),
1774                CostUnits(needed),
1775            )
1776        };
1777        assert_eq!(
1778            size(1_000, 60, 30, 0),
1779            Some(CostUnits(30)),
1780            "the GL-109 floor"
1781        );
1782        assert_eq!(
1783            size(1_000, 60, 30, 51),
1784            Some(CostUnits(51)),
1785            "proven demand"
1786        );
1787        assert_eq!(
1788            size(1_000, 60, 30, 60),
1789            Some(CostUnits(60)),
1790            "all of it, inclusive"
1791        );
1792        assert_eq!(
1793            size(1_000, 60, 30, 61),
1794            Some(CostUnits(30)),
1795            "an unfundable quote grows nothing"
1796        );
1797        assert_eq!(
1798            size(1_000, 60, 40, 35),
1799            Some(CostUnits(40)),
1800            "demand never shrinks the floor"
1801        );
1802        assert_eq!(
1803            size(10, 60, 0, 51),
1804            Some(CostUnits(51)),
1805            "past a small target"
1806        );
1807        assert_eq!(
1808            size(1_000, 0, 0, 51),
1809            None,
1810            "an empty balance still refuses"
1811        );
1812        assert_eq!(size(0, 60, 0, 51), None, "a zero request still refuses");
1813        for (requested, balance) in [(1_000, 60), (7, 60), (1_000, 1)] {
1814            assert_eq!(
1815                size(requested, balance, 0, 0),
1816                policy.grant(CostUnits(requested), CostUnits(balance)),
1817                "a plain acquire is the ordinary policy"
1818            );
1819        }
1820    }
1821
1822    /// Every variant, once. Sized by `COUNT`, so adding a refusal without
1823    /// widening this array fails to compile.
1824    fn all() -> [AllocateError; AllocateError::COUNT] {
1825        [
1826            AllocateError::UnknownAccount,
1827            AllocateError::AccountInactive,
1828            AllocateError::InsufficientBalance,
1829            AllocateError::InvalidTtl,
1830            AllocateError::UnknownLease,
1831            AllocateError::Fenced,
1832            AllocateError::LeaseNotActive,
1833            AllocateError::InvalidRelease,
1834            AllocateError::Storage(StoreError("connection reset".into())),
1835            AllocateError::BalanceExhausted(tollgate_core::BalanceExhaustion { period_end: None }),
1836            AllocateError::BalanceInsufficient(tollgate_core::BalanceShortfall {
1837                remaining: CostUnits(1),
1838                period_end: None,
1839            }),
1840        ]
1841    }
1842
1843    /// The indices must be a permutation of `0..COUNT`: two refusals sharing
1844    /// a slot would silently merge their tallies, and a slot no refusal maps
1845    /// to would export a counter that can never move.
1846    #[test]
1847    fn indices_cover_every_slot_exactly_once() {
1848        let mut seen = [false; AllocateError::COUNT];
1849        for error in all() {
1850            let index = error.index();
1851            assert!(index < AllocateError::COUNT, "{error} indexes out of range");
1852            assert!(!seen[index], "{error} shares slot {index}");
1853            seen[index] = true;
1854        }
1855        assert!(seen.iter().all(|hit| *hit), "every slot must be claimed");
1856    }
1857
1858    /// Labels are read by whatever scrapes the counters, so they are a
1859    /// contract: distinct, and free of the backend text `Display` carries.
1860    /// `Storage` wraps an arbitrary message, so labelling by rendered text
1861    /// would mint a fresh time series per connection failure.
1862    #[test]
1863    fn labels_are_distinct_and_free_of_backend_text() {
1864        for (position, error) in all().iter().enumerate() {
1865            assert_eq!(error.name(), AllocateError::NAMES[position]);
1866        }
1867        let storage = AllocateError::Storage(StoreError("connection reset".into()));
1868        assert_eq!(storage.name(), "storage");
1869        assert!(
1870            !storage.name().contains("connection"),
1871            "the label must not carry the backend's message"
1872        );
1873        let mut names = AllocateError::NAMES;
1874        names.sort_unstable();
1875        names.iter().reduce(|previous, next| {
1876            assert_ne!(previous, next, "duplicate label {next}");
1877            next
1878        });
1879    }
1880
1881    /// The payload must not affect the slot, or one refusal would scatter
1882    /// across slots as its message varied.
1883    #[test]
1884    fn payload_does_not_affect_the_slot() {
1885        assert_eq!(
1886            AllocateError::Storage(StoreError("a".into())).index(),
1887            AllocateError::Storage(StoreError("b".into())).index()
1888        );
1889    }
1890
1891    #[test]
1892    fn conservation_requires_an_exact_equation_without_overflow() {
1893        let balanced = Conservation {
1894            deposited: CostUnits(10),
1895            overage_recorded: CostUnits::ZERO,
1896            balance: CostUnits(1),
1897            active_lease_grants: CostUnits(2),
1898            settled_usage: CostUnits(3),
1899            settlement_loss: CostUnits(4),
1900            expired: CostUnits::ZERO,
1901        };
1902        assert!(balanced.holds());
1903
1904        let drifted = Conservation {
1905            deposited: CostUnits(11),
1906            ..balanced
1907        };
1908        assert!(!drifted.holds());
1909
1910        let overflowing = Conservation {
1911            deposited: CostUnits(u64::MAX),
1912            overage_recorded: CostUnits::ZERO,
1913            balance: CostUnits(u64::MAX),
1914            active_lease_grants: CostUnits(1),
1915            settled_usage: CostUnits::ZERO,
1916            settlement_loss: CostUnits::ZERO,
1917            expired: CostUnits::ZERO,
1918        };
1919        assert!(!overflowing.holds());
1920    }
1921
1922    /// Overage funds the left side, so usage it paid for closes the equation
1923    /// rather than breaking it — and the same numbers without the funding term
1924    /// must *not* balance, or the field would be decorative.
1925    #[test]
1926    fn overage_funds_the_usage_it_bills() {
1927        let elastic = Conservation {
1928            deposited: CostUnits(10),
1929            overage_recorded: CostUnits(5),
1930            balance: CostUnits(1),
1931            active_lease_grants: CostUnits(2),
1932            settled_usage: CostUnits(8),
1933            settlement_loss: CostUnits(4),
1934            expired: CostUnits::ZERO,
1935        };
1936        assert!(elastic.holds());
1937        assert!(
1938            !Conservation {
1939                overage_recorded: CostUnits::ZERO,
1940                ..elastic
1941            }
1942            .holds(),
1943            "the same ledger without the funding term must fail by exactly the overage"
1944        );
1945    }
1946
1947    /// The left side is checked too. Overflowing the funding sum answers
1948    /// `false` rather than wrapping to a total that might coincidentally match
1949    /// the right side (INVARIANTS.md GL-11).
1950    #[test]
1951    fn overflowing_the_funding_sum_is_a_violation_not_a_wrap() {
1952        let overflowing = Conservation {
1953            deposited: CostUnits(u64::MAX),
1954            overage_recorded: CostUnits(1),
1955            balance: CostUnits::ZERO,
1956            active_lease_grants: CostUnits::ZERO,
1957            settled_usage: CostUnits::ZERO,
1958            settlement_loss: CostUnits::ZERO,
1959            expired: CostUnits::ZERO,
1960        };
1961        assert!(!overflowing.holds());
1962    }
1963
1964    /// An allowance that expired at a period boundary left the balance without
1965    /// being spent, so the equation only closes if `expired` is on the right
1966    /// side — and the same ledger without the term must fail by exactly the
1967    /// units that expired, or the field would be decorative (GL-97).
1968    #[test]
1969    fn expiry_accounts_for_an_allowance_that_was_never_spent() {
1970        let rolled = Conservation {
1971            deposited: CostUnits(10),
1972            overage_recorded: CostUnits::ZERO,
1973            balance: CostUnits(1),
1974            active_lease_grants: CostUnits(2),
1975            settled_usage: CostUnits(3),
1976            settlement_loss: CostUnits::ZERO,
1977            expired: CostUnits(4),
1978        };
1979        assert!(rolled.holds());
1980        assert!(
1981            !Conservation {
1982                expired: CostUnits::ZERO,
1983                ..rolled
1984            }
1985            .holds(),
1986            "the same ledger without the expiry term must fail by exactly the expired units"
1987        );
1988    }
1989
1990    /// A rollover pass reports what it did, and its caller decides whether to
1991    /// ask for another batch from `is_saturated` alone. Every accessor is
1992    /// asserted directly: a `len` that always answered 1, or an `is_empty`
1993    /// stuck either way, would send the server's drain loop into an endless
1994    /// round of empty batches or stop it one batch short of the boundary it
1995    /// was crossing.
1996    #[test]
1997    fn rollover_batch_reports_what_it_rolled() {
1998        let limit = NonZeroUsize::new(2).unwrap();
1999        let rolled = |account| RolledAccount {
2000            account_id: AccountId(account),
2001            deposited: CostUnits(100),
2002            expired: CostUnits::ZERO,
2003        };
2004
2005        let empty = RolloverBatch::try_new(Vec::new(), limit).unwrap();
2006        assert!(empty.is_empty());
2007        assert_eq!(empty.len(), 0);
2008        assert!(!empty.is_saturated());
2009        assert!(empty.rolled().is_empty());
2010
2011        let partial = RolloverBatch::try_new(vec![rolled(1)], limit).unwrap();
2012        assert!(!partial.is_empty());
2013        assert_eq!(partial.len(), 1);
2014        assert!(
2015            !partial.is_saturated(),
2016            "a partial batch is what ends the drain"
2017        );
2018        assert_eq!(partial.rolled(), &[rolled(1)]);
2019
2020        let full = RolloverBatch::try_new(vec![rolled(1), rolled(2)], limit).unwrap();
2021        assert_eq!(full.len(), 2);
2022        assert!(
2023            full.is_saturated(),
2024            "a batch at the limit means there may be more"
2025        );
2026    }
2027
2028    /// Saturation is derived from the limit the caller asked for, so a backend
2029    /// returning more than it was allowed is corruption to refuse rather than
2030    /// a batch to trust — the same contract `ReclaimBatch::try_new` carries.
2031    #[test]
2032    fn a_rollover_batch_beyond_its_limit_is_refused() {
2033        let rolled = |account| RolledAccount {
2034            account_id: AccountId(account),
2035            deposited: CostUnits(100),
2036            expired: CostUnits::ZERO,
2037        };
2038        assert!(
2039            RolloverBatch::try_new(vec![rolled(1), rolled(2)], NonZeroUsize::new(1).unwrap())
2040                .is_err()
2041        );
2042    }
2043
2044    #[test]
2045    fn reclaim_batch_reports_an_empty_result() {
2046        let batch = ReclaimBatch::try_new(Vec::new(), NonZeroUsize::new(2).unwrap()).unwrap();
2047        assert!(batch.is_empty());
2048        assert_eq!(batch.len(), 0);
2049        assert!(!batch.is_saturated());
2050        assert!(batch.reclaimed().is_empty());
2051    }
2052
2053    #[tokio::test]
2054    async fn full_drain_preserves_an_initial_batch_error() {
2055        let allocator = ScriptedReclaimer {
2056            script: ReclaimScript::FailFirst,
2057            calls: AtomicUsize::new(0),
2058        };
2059        let expected = StoreError("scripted reclaim failure".into());
2060        assert_eq!(
2061            allocator.reclaim_expired(Timestamp::MIN).await,
2062            Err(expected)
2063        );
2064        assert_eq!(allocator.calls.load(Ordering::Acquire), 1);
2065    }
2066
2067    #[tokio::test]
2068    async fn full_drain_reports_progress_before_a_later_batch_error() {
2069        let allocator = ScriptedReclaimer {
2070            script: ReclaimScript::FullBatchThenFail,
2071            calls: AtomicUsize::new(0),
2072        };
2073        let original = StoreError("scripted reclaim failure".into());
2074        let error = allocator.reclaim_expired(Timestamp::MIN).await.unwrap_err();
2075        assert_ne!(error, original);
2076        assert!(error.0.contains("256 leases"));
2077        assert!(error.0.contains("256 units"));
2078        assert_eq!(allocator.calls.load(Ordering::Acquire), 2);
2079    }
2080
2081    /// The push-capacity boundary, pinned where both backends read it.
2082    ///
2083    /// Strictly greater is the whole content of the rule: a batch that exactly
2084    /// fills the channel is delivered, so warning at equality would fire on the
2085    /// largest successful case and train an operator to ignore it. Mutation
2086    /// testing found this untested — the comparison could be flipped to `<`,
2087    /// `<=` or `>=` and every scenario stayed green, because nothing observed
2088    /// the warning at all (GL-51).
2089    #[test]
2090    fn the_push_capacity_warning_fires_only_above_the_channel() {
2091        assert!(
2092            !pushes_exceed_capacity(0),
2093            "an empty batch is not a capacity problem"
2094        );
2095        assert!(!pushes_exceed_capacity(PUSH_CHANNEL_CAPACITY - 1));
2096        assert!(
2097            !pushes_exceed_capacity(PUSH_CHANNEL_CAPACITY),
2098            "a batch that exactly fills the channel is still delivered"
2099        );
2100        assert!(
2101            pushes_exceed_capacity(PUSH_CHANNEL_CAPACITY + 1),
2102            "one more than the channel holds is what makes a subscriber lag"
2103        );
2104    }
2105}