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