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 backend failure unrelated to domain rules; see [`StoreError`].
972 Storage(StoreError),
973}
974
975impl std::fmt::Display for SetStatusError {
976 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
977 match self {
978 SetStatusError::UnknownAccount => f.write_str("unknown account"),
979 SetStatusError::AccountClosed => f.write_str("account is closed"),
980 SetStatusError::Storage(e) => write!(f, "{e}"),
981 }
982 }
983}
984
985impl std::error::Error for SetStatusError {}
986
987impl From<StoreError> for SetStatusError {
988 fn from(error: StoreError) -> Self {
989 SetStatusError::Storage(error)
990 }
991}
992
993/// Refusals from publishing a snapshot (GL-51).
994///
995/// `publish_snapshot` used to return a bare [`StoreError`], which left it free
996/// to write a status contradicting the ledger and recreate the divergence
997/// [`AdminStore::set_account_status`] exists to abolish.
998#[derive(Debug, Clone, PartialEq, Eq)]
999pub enum PublishSnapshotError {
1000 /// The stated credential is missing or belongs to another principal/account.
1001 CredentialMismatch {
1002 /// The credential the snapshot states.
1003 key_id: KeyId,
1004 },
1005 /// The snapshot's status disagrees with the account ledger. An account's
1006 /// status is changed through [`AdminStore::set_account_status`], which
1007 /// republishes; a publish may carry the current status but may not change
1008 /// it.
1009 StatusMismatch {
1010 /// The status the account ledger holds.
1011 ledger: AccountStatus,
1012 /// The status the snapshot carried.
1013 submitted: AccountStatus,
1014 },
1015 /// The snapshot's execution-capacity class disagrees with the account
1016 /// ledger (GL-99). The class is an account-owned fact changed through
1017 /// [`AdminStore::set_capacity_class`], which republishes; a publish may
1018 /// carry the current class but may not change it. Two writers for one
1019 /// fact is the divergence the status guard above already exists to
1020 /// abolish, and a second field must not reintroduce it.
1021 CapacityClassMismatch {
1022 /// The class the account ledger holds.
1023 ledger: CapacityClass,
1024 /// The class the snapshot carried.
1025 submitted: CapacityClass,
1026 },
1027 /// A backend failure unrelated to domain rules; see [`StoreError`].
1028 Storage(StoreError),
1029}
1030
1031impl std::fmt::Display for PublishSnapshotError {
1032 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1033 match self {
1034 PublishSnapshotError::CredentialMismatch { key_id } => {
1035 write!(
1036 f,
1037 "credential {key_id} does not bind the published principal and account"
1038 )
1039 }
1040 PublishSnapshotError::StatusMismatch { ledger, submitted } => write!(
1041 f,
1042 "snapshot status {} contradicts account status {}",
1043 submitted.as_str(),
1044 ledger.as_str()
1045 ),
1046 PublishSnapshotError::CapacityClassMismatch { ledger, submitted } => write!(
1047 f,
1048 "snapshot capacity class {} contradicts account capacity class {}",
1049 submitted.as_str(),
1050 ledger.as_str()
1051 ),
1052 PublishSnapshotError::Storage(e) => write!(f, "{e}"),
1053 }
1054 }
1055}
1056
1057impl std::error::Error for PublishSnapshotError {}
1058
1059impl From<StoreError> for PublishSnapshotError {
1060 fn from(error: StoreError) -> Self {
1061 PublishSnapshotError::Storage(error)
1062 }
1063}
1064
1065/// One account's administrative state, for an operator read (GL-121).
1066///
1067/// Assembled from types that already exist rather than a parallel vocabulary,
1068/// so the HTTP surface reports the same terms the ledger reasons in and a
1069/// reader can hold a dashboard next to a conservation check.
1070///
1071/// **Funding is not billing, and the shape says so.** A falling `balance` does
1072/// not mean units were billed: it also falls when they go out on a lease that
1073/// has not settled, and it falls when a budget period closes and takes its
1074/// unspent allowance with it. Those are three different facts, and
1075/// [`Conservation`] keeps them apart — `active_lease_grants` is capacity
1076/// currently out, `settled_usage` is what was actually consumed,
1077/// `settlement_loss` and `expired` are what will never be. A surface that
1078/// reported only a balance would let a customer read depletion as spend, which
1079/// is exactly what GL-121 asks not to do.
1080#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1081pub struct AccountView {
1082 /// The account read.
1083 pub account_id: AccountId,
1084 /// Administrative status, from the ledger.
1085 pub status: AccountStatus,
1086 /// Execution-capacity class, from the ledger.
1087 pub capacity_class: CapacityClass,
1088 /// The periodic allowance, if this account has one. `None` is "no
1089 /// schedule, the balance does not expire" — not "unknown".
1090 pub schedule: Option<BudgetSchedule>,
1091 /// First instant of the period currently in force. Meaningful only
1092 /// alongside a `schedule`; it is the marker rollover is idempotent
1093 /// against.
1094 pub period_start: Timestamp,
1095 /// Every term of the funding equation, so a caller can distinguish
1096 /// remaining funding from outstanding grants from settled usage.
1097 pub conservation: Conservation,
1098}
1099
1100/// Administrative writes: the control plane's mutation surface. Kept apart
1101/// from the data-plane traits so a read-only replica can implement those
1102/// without this.
1103///
1104/// HTTP-facing mutations return [`crate::AdminReceipt`] captured at the same
1105/// serialization point as the write. Its `outcome` contains the operation's
1106/// result; before/after describe the fields owned by that operation. A separate
1107/// read before or after the transaction is not a valid receipt under concurrency.
1108/// No-op receipts have equal states; errors carry no confirmed transition.
1109/// Direct callers attach their own actor and audit delivery policy.
1110#[async_trait]
1111pub trait AdminStore: Send + Sync {
1112 /// Create an account from `config`, with its opening balance deposited as a
1113 /// top-up and its fencing sequence starting at one.
1114 ///
1115 /// Never destructive and never silently idempotent (INVARIANTS.md 14): an
1116 /// existing account is refused with [`CreateAccountError::AlreadyExists`] and
1117 /// left untouched, including its balance, ledger totals, fencing sequence and
1118 /// leases. The receipt's `before` is [`AdminState::Absent`](crate::AdminState::Absent).
1119 async fn create_account(
1120 &self,
1121 config: AccountConfig,
1122 ) -> Result<crate::AdminReceipt<()>, CreateAccountError>;
1123 /// Add `units` to an existing account as a top-up, which survives period
1124 /// boundaries, raising its balance and its `deposited` total together.
1125 ///
1126 /// A missing account is [`AllocateError::UnknownAccount`]. An overflow of
1127 /// either counter, or units outside the backend's numeric domain, is
1128 /// [`AllocateError::BalanceOverflow`] and moves neither. Account
1129 /// status is not checked. The receipt carries
1130 /// [`AdminState::Funding`](crate::AdminState::Funding) before and after.
1131 async fn deposit(
1132 &self,
1133 account: AccountId,
1134 units: CostUnits,
1135 ) -> Result<crate::AdminReceipt<()>, AllocateError>;
1136 /// Set an existing account's administrative status, in one transaction:
1137 /// the ledger's status, and a republication of every *live* snapshot of
1138 /// that account carrying the new status at `generation + 1`.
1139 ///
1140 /// This is the whole operator action. Before GL-51 the ledger flag and the
1141 /// published `AccountStatus` were two records with two propagation paths
1142 /// and nothing checking them against each other, so "deactivate" returned
1143 /// success while the request path kept admitting.
1144 ///
1145 /// Rules, all enforced here rather than by caller discipline:
1146 /// - A missing account is [`SetStatusError::UnknownAccount`], never a
1147 /// silent no-op, whichever status was asked for.
1148 /// - [`AccountStatus::Closed`] is terminal
1149 /// ([`SetStatusError::AccountClosed`]); `Closed` → `Closed` is a no-op.
1150 /// - Revoked principals are never republished: resurrecting a tombstone
1151 /// is what INVARIANTS.md GL-15 forbids, and revocation stays a separate
1152 /// per-credential mechanism.
1153 /// - Snapshots already at the target status are not rewritten, so a
1154 /// repeat converges and bumps no generation.
1155 /// - Outstanding leases are **not** reclaimed. Lease acquisition refuses
1156 /// at once, but admission stops only when the new snapshot installs —
1157 /// one `SnapshotManager` refresh interval, and already-debited units
1158 /// settle at release or TTL reclaim (GL-9).
1159 async fn set_account_status(
1160 &self,
1161 account: AccountId,
1162 status: AccountStatus,
1163 ) -> Result<crate::AdminReceipt<StatusChange>, SetStatusError>;
1164
1165 /// Set an existing account's execution-capacity class, in one
1166 /// transaction: the ledger's class, and a republication of every *live*
1167 /// snapshot of that account carrying the new class at `generation + 1`
1168 /// (GL-99).
1169 ///
1170 /// The same operator action, and the same ownership argument, as
1171 /// [`set_account_status`](Self::set_account_status): the class is one
1172 /// fact with one writer. A control plane that published it per credential
1173 /// instead would recreate exactly the divergence GL-51 abolished — some of
1174 /// an account's principals assured and some best-effort, with nothing
1175 /// checking them against each other, and a request's treatment depending
1176 /// on which credential it arrived with.
1177 ///
1178 /// Rules, all enforced here rather than by caller discipline:
1179 /// - A missing account is [`SetStatusError::UnknownAccount`].
1180 /// - A closed account is [`SetStatusError::AccountClosed`]. Reclassifying
1181 /// a terminally closed account is meaningless and the refusal changes
1182 /// nothing, exactly as it does for a status change.
1183 /// - Revoked principals are never republished (INVARIANTS.md GL-15).
1184 /// - Snapshots already at the target class are not rewritten, so a repeat
1185 /// converges and bumps no generation.
1186 /// - Nothing about funding changes. The class decides whether an instance
1187 /// starts an already-funded request, so quota, leases, and outstanding
1188 /// usage are untouched — an account reclassified mid-flight keeps every
1189 /// charge it has already committed.
1190 ///
1191 /// Reuses [`SetStatusError`] rather than declaring a near-identical twin:
1192 /// the two refusals are the same two conditions about the same ledger row,
1193 /// and a second enum would be two vocabularies for one answer.
1194 async fn set_capacity_class(
1195 &self,
1196 account: AccountId,
1197 class: CapacityClass,
1198 ) -> Result<crate::AdminReceipt<StatusChange>, SetStatusError>;
1199
1200 /// Give an account a periodic allowance, or take it away.
1201 ///
1202 /// Setting a schedule does not deposit anything: the first allowance
1203 /// arrives at the first [`roll_due_periods`](Self::roll_due_periods) pass after the
1204 /// schedule exists. Depositing here would make "set a schedule" and "give
1205 /// this account units now" the same operation, and an operator correcting
1206 /// a mistyped allowance would fund the account twice.
1207 ///
1208 /// `None` removes the schedule and leaves the balance alone — including
1209 /// any unspent allowance, which simply stops expiring. Removing a schedule
1210 /// is not a way to claw units back.
1211 /// Set or clear an account's periodic allowance.
1212 ///
1213 /// Returns a receipt rather than `()` so the change can be audited like
1214 /// every other administrative mutation: an operator surface has to be able
1215 /// to report what a call actually committed, and a bare `Ok` cannot say
1216 /// whether a schedule was introduced, replaced, or was already what the
1217 /// caller asked for (GL-121). A repeat returns equal before/after states,
1218 /// which is how the convention expresses an idempotent no-op.
1219 async fn set_budget_schedule(
1220 &self,
1221 account: AccountId,
1222 schedule: Option<BudgetSchedule>,
1223 ) -> Result<crate::AdminReceipt<()>, BudgetError>;
1224
1225 /// Cross the period boundary for up to `limit` accounts that are past it:
1226 /// expire each closed period's unspent allowance and deposit the next one,
1227 /// one transaction per batch.
1228 ///
1229 /// **Idempotency is this method's job, not its caller's.** The pass runs
1230 /// on every control-plane replica, so two of them will race a boundary;
1231 /// the backend crosses it under a row lock, and the second caller then
1232 /// reads the period the first one wrote and skips the account. A caller
1233 /// that read each period first and then rolled would produce two deposits
1234 /// under exactly the race this exists to survive.
1235 ///
1236 /// Safe to call at any cadence: before a boundary it selects nothing, and
1237 /// after one the first caller wins. It never rolls an account more than
1238 /// one period at a time — an account left unrolled for two months lands in
1239 /// the current period with one allowance, because an allowance is what the
1240 /// account is entitled to now, not a backlog to be paid out.
1241 ///
1242 /// Bounded, and saturating means there is more: every scheduled account
1243 /// comes due at the same instant, so a caller must drain saturated batches
1244 /// until one comes back partial. Unscheduled accounts are never selected.
1245 async fn roll_due_periods(
1246 &self,
1247 now: Timestamp,
1248 limit: NonZeroUsize,
1249 ) -> Result<RolloverBatch, StoreError>;
1250 /// Publish a principal's compiled snapshot.
1251 ///
1252 /// Refused with [`PublishSnapshotError::StatusMismatch`] when the
1253 /// snapshot's status contradicts the account ledger, so the two records
1254 /// [`set_account_status`](AdminStore::set_account_status) unifies cannot
1255 /// be pulled apart again one principal at a time. A snapshot whose
1256 /// account the ledger does not hold publishes unchanged, as before.
1257 async fn publish_snapshot(
1258 &self,
1259 principal: Principal,
1260 snapshot: PublishableSnapshot,
1261 ) -> Result<crate::AdminReceipt<()>, PublishSnapshotError>;
1262 /// Withdraw a principal's snapshot by tombstoning it at its current
1263 /// generation, and push the revocation to subscribers. The tombstone is
1264 /// durable, so no positive snapshot at or below that generation can resurrect
1265 /// the principal (INVARIANTS.md 15).
1266 ///
1267 /// A principal with no snapshot, or one already tombstoned, is left unchanged
1268 /// and the receipt's states are equal.
1269 async fn remove_snapshot(
1270 &self,
1271 principal: Principal,
1272 ) -> Result<crate::AdminReceipt<()>, StoreError>;
1273
1274 /// One account's administrative state, or `None` if no such account (GL-121).
1275 ///
1276 /// The read an operator surface needs and the traits did not have. Both
1277 /// backends already expose `conservation` as an inherent method, but with
1278 /// different signatures — one synchronous returning an `Option`, one
1279 /// asynchronous returning a `Result` — so nothing generic over a backend
1280 /// could read an account at all.
1281 ///
1282 /// A single call rather than several, because the terms have to agree with
1283 /// each other: status, schedule and the funding equation read separately
1284 /// can straddle a rollover or a suspension and describe a state the account
1285 /// was never in. A backend answers this from one consistent read.
1286 async fn account_view(&self, account: AccountId) -> Result<Option<AccountView>, StoreError>;
1287}
1288
1289/// Outcome of one ingest batch.
1290#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
1291#[cfg_attr(feature = "wire", derive(serde::Serialize, serde::Deserialize))]
1292pub struct IngestReport {
1293 /// Newly recorded events.
1294 pub accepted: u64,
1295 /// Events whose `request_id` was already recorded (idempotent replay —
1296 /// INVARIANTS.md GL-7).
1297 pub duplicate: u64,
1298 /// Events refused: unknown lease, lease-capability mismatch, or no
1299 /// remaining accounting capacity, or units outside the backend's storage
1300 /// domain. These are bounded billing loss,
1301 /// visible to reconciliation.
1302 pub rejected: u64,
1303 /// Newly accepted events with absent, unknown or different-account key
1304 /// attribution. Duplicates never supply fresh activity evidence. `None`
1305 /// means this sink does not report attribution (including older servers),
1306 /// not that every event was attributed.
1307 #[cfg_attr(feature = "wire", serde(default))]
1308 pub unattributed: Option<u64>,
1309}
1310
1311impl IngestReport {
1312 /// A complete acknowledgement partitions this batch exactly once. Validate
1313 /// before releasing queued evidence, including replies from custom sinks.
1314 pub fn validate(&self, submitted: usize) -> Result<(), StoreError> {
1315 let total = self
1316 .accepted
1317 .checked_add(self.duplicate)
1318 .and_then(|n| n.checked_add(self.rejected));
1319 if !total.is_some_and(|n| u64::try_from(submitted) == Ok(n))
1320 || self.unattributed.is_some_and(|n| n > self.accepted)
1321 {
1322 return Err(StoreError(
1323 "invalid usage acknowledgement cardinality".into(),
1324 ));
1325 }
1326 Ok(())
1327 }
1328}
1329
1330/// The billing ledger's write side.
1331#[async_trait]
1332pub trait UsageSink: Send + Sync {
1333 /// Record a batch. Idempotent on `request_id`; every event must match its
1334 /// stored `(lease_id, account_id, fencing_token)` capability before lease
1335 /// state and accounting capacity are checked. Partial acceptance is
1336 /// normal — the report says what happened.
1337 ///
1338 /// Duplicates are classified before inspecting their payload. A new event
1339 /// outside the backend's unit domain is rejected individually; it is not
1340 /// remembered as accepted and cannot poison otherwise valid neighbors.
1341 /// MemoryStore supports `u64` units; PostgreSQL supports nonnegative
1342 /// `BIGINT` units (`0..=i64::MAX`). A representable event that overflows an
1343 /// accumulated accounting total refuses the entire batch atomically.
1344 async fn ingest(
1345 &self,
1346 events: &[UsageEvent],
1347 now: Timestamp,
1348 ) -> Result<IngestReport, IngestError>;
1349}
1350
1351/// One credential's durable record.
1352///
1353/// The `digest` is opaque here on purpose. The HMAC secret that produced it
1354/// lives with the verifier (`tollgate-auth`) and never reaches a store, so a
1355/// backend holds material that verifies nothing on its own — the property
1356/// `HmacRegistry` is built around, preserved across the persistence boundary.
1357/// A store that could compute a digest would be a store whose compromise is
1358/// sufficient to mint credentials.
1359///
1360/// `principal` is the digest's own truncation, so it is derived rather than
1361/// assigned: the request path is keyed by it, and per-credential revocation
1362/// is `install_revoked` for exactly this value.
1363#[derive(Clone, PartialEq, Eq)]
1364pub struct KeyRecord {
1365 /// The credential's non-secret identifier, chosen by its issuer.
1366 pub key_id: KeyId,
1367 /// The account the credential authenticates for.
1368 pub account_id: AccountId,
1369 /// The principal this credential authenticates as: the leading 128 bits
1370 /// of `digest`.
1371 pub principal: Principal,
1372 /// HMAC-SHA256 of the secret under the verifier's server secret.
1373 pub digest: [u8; 32],
1374 /// When the credential stops being valid of its own accord, independent
1375 /// of revocation. Surfaced to the request path through
1376 /// `Verified::reusable_until`, so a session cache cannot outlive it.
1377 pub not_after: Option<Timestamp>,
1378}
1379
1380impl std::fmt::Debug for KeyRecord {
1381 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1382 f.debug_struct("KeyRecord")
1383 .field("key_id", &self.key_id)
1384 .field("account_id", &self.account_id)
1385 .field("principal", &self.principal)
1386 .field("not_after", &self.not_after)
1387 .finish_non_exhaustive()
1388 }
1389}
1390
1391/// One requested key's activity. No observation is not proof of non-use.
1392#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1393pub struct CredentialActivity {
1394 /// The requested key.
1395 pub key_id: KeyId,
1396 /// What the store has recorded for it.
1397 pub state: CredentialActivityState,
1398}
1399
1400/// A credential's recorded commitment activity (INVARIANTS.md 35).
1401/// It is never authentication or authorization evidence.
1402#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1403pub enum CredentialActivityState {
1404 /// The directory holds no credential with this id.
1405 Unknown,
1406 /// The credential exists, but no accepted, attributable usage has been
1407 /// recorded for it. Not proof that it was never used.
1408 Unobserved,
1409 /// Maximum accepted, attributable execution-start time, at microsecond
1410 /// precision. Never an authorization or independent server-clock fact.
1411 Committed {
1412 /// The latest recorded execution-start time.
1413 last_committed_at: Timestamp,
1414 },
1415}
1416
1417/// What a revocation actually did.
1418///
1419/// Returned rather than inferred, for the reason [`StatusChange`] is: an
1420/// operator retiring a suspicious credential needs to know whether they
1421/// retired anything. "Already revoked" and "no such key" are different
1422/// answers to the same request and only one of them is a mistake.
1423#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1424pub enum Revocation {
1425 /// The credential was live and is now retired.
1426 Retired,
1427 /// The credential was already retired; nothing changed.
1428 AlreadyRetired,
1429}
1430
1431/// The largest usage batch any sink accepts in one call, in events.
1432///
1433/// A batch cap has to admit the largest batch the system can legitimately
1434/// produce, not a comfortable one: the shipped example writes 256 per flush,
1435/// and this leaves a factor of sixteen for an embedder that batches harder to
1436/// cut round-trips. Beyond it the answer is more flushes, not a bigger body —
1437/// an ingest is idempotent, so splitting costs a round trip and risks nothing.
1438///
1439/// `UsageWriterConfig::validate` refuses a `max_batch` above this, so the
1440/// misconfiguration is a startup error rather than a permanently-rejected
1441/// batch discovered in production (GL-61).
1442pub const MAX_INGEST_BATCH: usize = 4_096;
1443
1444/// Why an ingest attempt failed, and whether replaying it unchanged could
1445/// ever succeed.
1446///
1447/// The distinction exists because the writer's correct response to the two is
1448/// opposite. An unreachable sink is a *duration*: retrying the same batch is
1449/// the designed behaviour, and giving up would lose billable events over a
1450/// blip. A refused batch is a *fact about the batch*: retrying it unchanged
1451/// gets the same answer forever, and every event queued behind it waits for a
1452/// recovery that cannot come — a permanent, deterministic error laundered
1453/// into an unbounded billing and availability outage (GL-61).
1454///
1455/// [`From<StoreError>`] yields [`Unavailable`](Self::Unavailable), so a
1456/// backend that does not classify keeps the retry-forever behaviour it had.
1457/// Terminality is asserted, never assumed.
1458#[derive(Debug, Clone, PartialEq, Eq)]
1459pub enum IngestError {
1460 /// The sink could not be reached, or could not answer in time. The batch
1461 /// is unchanged and will be retried.
1462 Unavailable(StoreError),
1463 /// The sink refused this batch and will refuse it again unchanged: a body
1464 /// over the endpoint's limit, an event it cannot decode, a contract it
1465 /// does not implement. Retrying cannot help.
1466 Refused(StoreError),
1467}
1468
1469impl IngestError {
1470 /// Whether replaying this batch unchanged could ever succeed.
1471 #[must_use]
1472 pub const fn is_retryable(&self) -> bool {
1473 matches!(self, IngestError::Unavailable(_))
1474 }
1475}
1476
1477impl From<StoreError> for IngestError {
1478 fn from(error: StoreError) -> Self {
1479 IngestError::Unavailable(error)
1480 }
1481}
1482
1483impl std::fmt::Display for IngestError {
1484 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1485 match self {
1486 IngestError::Unavailable(e) => write!(f, "{e}"),
1487 IngestError::Refused(e) => write!(f, "refused: {e}"),
1488 }
1489 }
1490}
1491
1492impl std::error::Error for IngestError {}
1493
1494/// Refusals from credential lifecycle operations.
1495#[derive(Debug, Clone, PartialEq, Eq)]
1496pub enum KeyError {
1497 /// No such credential. Never a silent no-op — an operator revoking a key
1498 /// that does not exist has either the wrong id or a false belief about
1499 /// what is live, and both are worth surfacing.
1500 UnknownKey,
1501 /// The credential's account does not exist, so nothing could authenticate
1502 /// as it. Refused at issuance rather than producing a key that verifies
1503 /// and is then denied by every admission.
1504 UnknownAccount,
1505 /// This `key_id` is already recorded. Issuance is never destructive, for
1506 /// the reason account creation is not ([`CreateAccountError`]): an
1507 /// overwrite would silently retire a live credential.
1508 ///
1509 /// This is also the retry answer. A caller that supplies the `key_id` and
1510 /// loses the response resends the same one and is told the credential
1511 /// exists — which is the truth, and which discloses no secret. That is why
1512 /// issuance must never become an upsert (GL-121).
1513 AlreadyExists,
1514 /// The account already holds `limit` live credentials, so issuing another
1515 /// would exceed the bound the caller supplied.
1516 ///
1517 /// "Live" excludes revoked keys and keys whose `not_after` has passed: a
1518 /// bound that counted expired credentials would strand an account behind
1519 /// keys nobody can authenticate with.
1520 ActiveKeyLimit {
1521 /// The bound the caller supplied.
1522 limit: NonZeroUsize,
1523 },
1524 /// A backend failure unrelated to domain rules; see [`StoreError`].
1525 Storage(StoreError),
1526}
1527
1528impl std::fmt::Display for KeyError {
1529 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1530 match self {
1531 KeyError::UnknownKey => f.write_str("no such credential"),
1532 KeyError::UnknownAccount => f.write_str("no such account"),
1533 KeyError::AlreadyExists => f.write_str("credential already exists"),
1534 KeyError::ActiveKeyLimit { limit } => {
1535 write!(f, "account already holds {limit} live credentials")
1536 }
1537 KeyError::Storage(e) => write!(f, "{e}"),
1538 }
1539 }
1540}
1541
1542impl std::error::Error for KeyError {}
1543
1544/// Refusals from binding or withdrawing a snapshot by the credential it was
1545/// issued as (GL-143).
1546#[derive(Debug, Clone, PartialEq, Eq)]
1547pub enum KeySnapshotError {
1548 /// No such credential, or it belongs to another account. One answer for
1549 /// both, as revocation gives: a foreign `key_id` discloses nothing.
1550 UnknownCredential,
1551 /// The credential was revoked. Revocation is terminal (INVARIANTS.md
1552 /// GL-27), so it is never granted positive authorization again; withdrawal
1553 /// remains allowed.
1554 Retired {
1555 /// The revoked credential.
1556 key_id: KeyId,
1557 },
1558 /// Publication refused as [`AdminStore::publish_snapshot`] would refuse it.
1559 Publish(PublishSnapshotError),
1560 /// A backend failure unrelated to domain rules; see [`StoreError`].
1561 Storage(StoreError),
1562}
1563
1564impl std::fmt::Display for KeySnapshotError {
1565 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1566 match self {
1567 KeySnapshotError::UnknownCredential => f.write_str("no such credential"),
1568 KeySnapshotError::Retired { key_id } => {
1569 write!(f, "credential {key_id} is revoked")
1570 }
1571 KeySnapshotError::Publish(e) => write!(f, "{e}"),
1572 KeySnapshotError::Storage(e) => write!(f, "{e}"),
1573 }
1574 }
1575}
1576
1577impl std::error::Error for KeySnapshotError {}
1578
1579impl From<PublishSnapshotError> for KeySnapshotError {
1580 fn from(error: PublishSnapshotError) -> Self {
1581 Self::Publish(error)
1582 }
1583}
1584
1585impl From<StoreError> for KeySnapshotError {
1586 fn from(error: StoreError) -> Self {
1587 Self::Storage(error)
1588 }
1589}
1590
1591impl From<StoreError> for KeyError {
1592 fn from(error: StoreError) -> Self {
1593 Self::Storage(error)
1594 }
1595}
1596
1597/// One credential as an *administrator* sees it (GL-121).
1598///
1599/// Deliberately not a [`KeyRecord`]. A record carries `digest` — the HMAC the
1600/// verifier compares against — and `principal`, documented there as "the
1601/// leading 128 bits of `digest`". Both are digest material, and an
1602/// account-scoped listing is reachable by an application backend and from
1603/// there a browser, so neither may appear in it. `key_id` is the non-secret
1604/// handle: what the caller chose, what revocation names, what an audit shows.
1605///
1606/// Expiry and revocation are surfaced separately. A credential that lapsed on
1607/// its own is a different operational fact from one an operator withdrew, and
1608/// collapsing both into "inactive" loses the distinction exactly where someone
1609/// is deciding whether to issue a replacement.
1610#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1611pub struct KeySummary {
1612 /// The credential's non-secret identifier.
1613 pub key_id: KeyId,
1614 /// When the credential stops being valid of its own accord, if ever.
1615 pub not_after: Option<Timestamp>,
1616 /// When an operator withdrew it, if they did. Terminal.
1617 pub revoked_at: Option<Timestamp>,
1618}
1619
1620impl KeySummary {
1621 /// Whether this credential can still authenticate at `now`.
1622 ///
1623 /// The predicate the issuance bound counts with, written once so a listing
1624 /// and a limit cannot disagree about what "live" means.
1625 #[must_use]
1626 pub fn is_live(&self, now: Timestamp) -> bool {
1627 self.revoked_at.is_none() && self.not_after.is_none_or(|until| now < until)
1628 }
1629}
1630
1631/// Durable credential lifecycle: the half of key management that outlives a
1632/// process and is shared by a fleet.
1633///
1634/// **Why this is a store trait rather than registry state.** Every other
1635/// control-plane fact here — accounts, balances, leases, statuses, snapshots —
1636/// is transactional in a backend with a local read projection, and credentials
1637/// are the same class of fact. An in-memory-only registry would lose issued
1638/// keys on restart, keep verifying a credential another instance revoked, and
1639/// make a fleet-wide active-key limit unenforceable, because no instance sees
1640/// the fleet.
1641///
1642/// **Durability before disclosure.** [`insert_key`](Self::insert_key) must
1643/// commit before its caller returns the secret to anyone. A crash between the
1644/// two hands out a credential the server has never heard of, which no later
1645/// reconciliation can repair: the digest is unrecoverable from the record,
1646/// which is the point of storing digests.
1647///
1648/// The verifier's projection is rebuilt from [`active_keys`](Self::active_keys),
1649/// so a backend decides what "active" means once, here, rather than in each
1650/// reader.
1651#[async_trait]
1652pub trait KeyDirectory: crate::KeySource {
1653 /// Inspect requested keys in input order, including repeated IDs and
1654 /// retired keys. Every input has one explicit result or the read fails.
1655 /// This operator read uses O(keys.len()) output memory; backends bound
1656 /// individual queries internally. Multiple chunks need not share an instant.
1657 async fn credential_activity(
1658 &self,
1659 keys: &[KeyId],
1660 ) -> Result<Vec<CredentialActivity>, StoreError>;
1661
1662 /// Record a minted credential. The caller has already generated the
1663 /// secret and computed its digest; this stores what remains.
1664 async fn insert_key(&self, record: KeyRecord) -> Result<(), KeyError>;
1665
1666 /// Retire one credential, reporting whether it was live.
1667 ///
1668 /// Revocation is durable and terminal: a retired credential is never
1669 /// resurrected, for the same reason a snapshot tombstone is not
1670 /// (INVARIANTS.md GL-15).
1671 async fn revoke_key(&self, key_id: KeyId, now: Timestamp) -> Result<Revocation, KeyError>;
1672
1673 /// Unbounded operator read of every credential valid at `now`. Serving
1674 /// instances use `KeySource` pages and an owned, bounded drain instead.
1675 /// Retained for existing direct-store lifecycle tooling; no hidden page cap.
1676 async fn active_keys(&self, now: Timestamp) -> Result<Vec<KeyRecord>, StoreError>;
1677
1678 /// One account's credentials, ordered by `key_id`, for an operator
1679 /// listing (GL-121).
1680 ///
1681 /// Distinct from [`active_keys`](Self::active_keys), which is the
1682 /// fleet-wide, digest-bearing projection an *instance* pulls: this is
1683 /// account-scoped, bounded, and carries no digest material, because it
1684 /// answers a different question for a different caller.
1685 ///
1686 /// Paginate with `after` — the greatest `key_id` already seen, exclusive.
1687 /// Revoked and expired credentials are included, because an administrator
1688 /// deciding whether to issue a replacement needs to see what became of the
1689 /// last one; [`KeySummary::is_live`] separates them.
1690 async fn account_keys(
1691 &self,
1692 account: AccountId,
1693 after: Option<KeyId>,
1694 limit: NonZeroUsize,
1695 ) -> Result<Vec<KeySummary>, StoreError>;
1696
1697 /// Record a credential only if the account holds fewer than `max_active`
1698 /// live ones, counting and inserting indivisibly (GL-121).
1699 ///
1700 /// The bound is supplied per call rather than stored: what counts as a
1701 /// reasonable number of credentials belongs to the application's plan, not
1702 /// to Tollgate, and a value in the request is one the caller can change
1703 /// without a migration.
1704 ///
1705 /// **Why this is not [`active_keys`](Self::active_keys) then
1706 /// [`insert_key`](Self::insert_key).** Two issuers racing that pair both
1707 /// read `max_active - 1`, both insert, and the account ends up over the
1708 /// bound with no error raised anywhere — the more replicas, the likelier.
1709 /// A backend must make the count and the insert one indivisible step:
1710 /// `MemoryStore` holds a single lock across both, and `PostgresStore`
1711 /// takes the account row `FOR UPDATE` first — the row
1712 /// `set_account_status` already serialises against.
1713 ///
1714 /// Live excludes revoked credentials, and those whose `not_after` has
1715 /// passed at `now`, so an account cannot be stranded behind keys that can
1716 /// no longer authenticate.
1717 async fn insert_key_within(
1718 &self,
1719 record: KeyRecord,
1720 max_active: NonZeroUsize,
1721 now: Timestamp,
1722 ) -> Result<(), KeyError>;
1723
1724 /// Bounded issuance with lifecycle evidence captured under the mutation lock.
1725 /// HTTP administrators must use this receipt rather than synthesize history.
1726 async fn insert_key_within_audited(
1727 &self,
1728 record: KeyRecord,
1729 max_active: NonZeroUsize,
1730 now: Timestamp,
1731 ) -> Result<crate::AdminReceipt<()>, KeyError>;
1732
1733 /// Retire a credential and capture its actual owner, key and predecessor
1734 /// under the mutation lock. A repeated revocation returns equal states.
1735 async fn revoke_key_audited(
1736 &self,
1737 key_id: KeyId,
1738 now: Timestamp,
1739 ) -> Result<crate::AdminReceipt<Revocation>, KeyError>;
1740
1741 /// Publish `snapshot` for the principal of `account`'s credential `key`,
1742 /// resolved inside the store (GL-143).
1743 ///
1744 /// An operator holds `(account, key)`; the principal is digest material
1745 /// and never leaves the server. Resolution, the retirement check and the
1746 /// publication are one indivisible step, so a concurrent revocation
1747 /// either precedes it — and the publish is refused as
1748 /// [`KeySnapshotError::Retired`] — or follows it.
1749 ///
1750 /// `snapshot` must state `key_id == Some(key)`; anything else is
1751 /// [`PublishSnapshotError::CredentialMismatch`]. Every other rule is
1752 /// [`AdminStore::publish_snapshot`]'s, including the generation no-op.
1753 async fn publish_key_snapshot(
1754 &self,
1755 account: AccountId,
1756 key: KeyId,
1757 snapshot: PublishableSnapshot,
1758 ) -> Result<crate::AdminReceipt<()>, KeySnapshotError>;
1759
1760 /// Withdraw the snapshot of `account`'s credential `key`, tombstoning it
1761 /// as [`AdminStore::remove_snapshot`] does. Allowed for a revoked
1762 /// credential: revocation does not withdraw its snapshot, and withdrawal
1763 /// is the safe direction.
1764 async fn remove_key_snapshot(
1765 &self,
1766 account: AccountId,
1767 key: KeyId,
1768 ) -> Result<crate::AdminReceipt<()>, KeySnapshotError>;
1769}
1770
1771#[cfg(test)]
1772mod tests {
1773 use std::sync::atomic::{AtomicUsize, Ordering};
1774
1775 use super::KeySummary;
1776 use jiff::Timestamp;
1777 use tollgate_core::KeyId;
1778
1779 fn at(seconds: i64) -> Timestamp {
1780 Timestamp::from_second(seconds).expect("a test instant")
1781 }
1782
1783 /// `not_after` is exclusive, and the instant itself is the whole question.
1784 ///
1785 /// Both backends filter with `now < not_after` — `memory.rs` in three
1786 /// places, and four SQL predicates written `not_after > $now`. `is_live`
1787 /// is the summary of exactly those queries, so `<=` here would not merely
1788 /// be off by an instant: a listing would report a credential live for the
1789 /// one instant at which every query that selects credentials has already
1790 /// dropped it, and the issuance bound counts with this predicate.
1791 #[test]
1792 fn a_credential_is_dead_at_its_expiry_instant_not_after_it() {
1793 let expiring = |not_after| KeySummary {
1794 key_id: KeyId(1),
1795 not_after: Some(not_after),
1796 revoked_at: None,
1797 };
1798 assert!(
1799 expiring(at(100)).is_live(at(99)),
1800 "live up to the instant before"
1801 );
1802 assert!(
1803 !expiring(at(100)).is_live(at(100)),
1804 "dead *at* the boundary: expiry is exclusive, as both backends filter it"
1805 );
1806 assert!(!expiring(at(100)).is_live(at(101)), "and dead after it");
1807 }
1808
1809 use super::*;
1810
1811 #[derive(Clone, Copy)]
1812 enum ReclaimScript {
1813 FailFirst,
1814 FullBatchThenFail,
1815 }
1816
1817 struct ScriptedReclaimer {
1818 script: ReclaimScript,
1819 calls: AtomicUsize,
1820 }
1821
1822 #[async_trait]
1823 impl LeaseAllocator for ScriptedReclaimer {
1824 async fn acquire(
1825 &self,
1826 _account: AccountId,
1827 _requested: CostUnits,
1828 _ttl: SignedDuration,
1829 _now: Timestamp,
1830 ) -> Result<Allocation, AllocateError> {
1831 unreachable!("the full-drain tests only reclaim")
1832 }
1833
1834 async fn release(
1835 &self,
1836 _lease_id: LeaseId,
1837 _fencing_token: FencingToken,
1838 _unspent: CostUnits,
1839 _now: Timestamp,
1840 ) -> Result<(), AllocateError> {
1841 unreachable!("the full-drain tests only reclaim")
1842 }
1843
1844 async fn consolidate(
1845 &self,
1846 _lease_id: LeaseId,
1847 _fencing_token: FencingToken,
1848 _unspent: CostUnits,
1849 _requested: CostUnits,
1850 _needed: CostUnits,
1851 _ttl: SignedDuration,
1852 _now: Timestamp,
1853 ) -> Result<Allocation, AllocateError> {
1854 unreachable!("the reclaim drain never consolidates")
1855 }
1856
1857 async fn reclaim_expired_batch(
1858 &self,
1859 _now: Timestamp,
1860 limit: NonZeroUsize,
1861 ) -> Result<ReclaimBatch, StoreError> {
1862 let call = self.calls.fetch_add(1, Ordering::AcqRel);
1863 if matches!(self.script, ReclaimScript::FailFirst) || call > 0 {
1864 return Err(StoreError("scripted reclaim failure".into()));
1865 }
1866 let reclaimed = (0..limit.get())
1867 .map(|id| ReclaimedLease {
1868 lease_id: LeaseId(u128::try_from(id).unwrap()),
1869 account_id: AccountId(1),
1870 forfeited: CostUnits(1),
1871 })
1872 .collect();
1873 ReclaimBatch::try_new(reclaimed, limit)
1874 }
1875 }
1876
1877 /// GL-131: under the default divisor of 2, a 60-unit balance re-granted as
1878 /// 30 forever, however often a 51-unit quote was refused.
1879 #[test]
1880 fn consolidation_grows_only_to_a_fundable_needed_quote() {
1881 let policy = GrantPolicy::default();
1882 let size = |requested, balance, floor, needed| {
1883 policy.consolidation_grant(
1884 CostUnits(requested),
1885 CostUnits(balance),
1886 CostUnits(floor),
1887 CostUnits(needed),
1888 )
1889 };
1890 assert_eq!(
1891 size(1_000, 60, 30, 0),
1892 Some(CostUnits(30)),
1893 "the GL-109 floor"
1894 );
1895 assert_eq!(
1896 size(1_000, 60, 30, 51),
1897 Some(CostUnits(51)),
1898 "proven demand"
1899 );
1900 assert_eq!(
1901 size(1_000, 60, 30, 60),
1902 Some(CostUnits(60)),
1903 "all of it, inclusive"
1904 );
1905 assert_eq!(
1906 size(1_000, 60, 30, 61),
1907 Some(CostUnits(30)),
1908 "an unfundable quote grows nothing"
1909 );
1910 assert_eq!(
1911 size(1_000, 60, 40, 35),
1912 Some(CostUnits(40)),
1913 "demand never shrinks the floor"
1914 );
1915 assert_eq!(
1916 size(10, 60, 0, 51),
1917 Some(CostUnits(51)),
1918 "past a small target"
1919 );
1920 assert_eq!(
1921 size(1_000, 0, 0, 51),
1922 None,
1923 "an empty balance still refuses"
1924 );
1925 assert_eq!(size(0, 60, 0, 51), None, "a zero request still refuses");
1926 for (requested, balance) in [(1_000, 60), (7, 60), (1_000, 1)] {
1927 assert_eq!(
1928 size(requested, balance, 0, 0),
1929 policy.grant(CostUnits(requested), CostUnits(balance)),
1930 "a plain acquire is the ordinary policy"
1931 );
1932 }
1933 }
1934
1935 /// Every variant, once. Sized by `COUNT`, so adding a refusal without
1936 /// widening this array fails to compile.
1937 fn all() -> [AllocateError; AllocateError::COUNT] {
1938 [
1939 AllocateError::UnknownAccount,
1940 AllocateError::AccountInactive,
1941 AllocateError::InsufficientBalance,
1942 AllocateError::InvalidTtl,
1943 AllocateError::UnknownLease,
1944 AllocateError::Fenced,
1945 AllocateError::LeaseNotActive,
1946 AllocateError::InvalidRelease,
1947 AllocateError::Storage(StoreError("connection reset".into())),
1948 AllocateError::BalanceExhausted(tollgate_core::BalanceExhaustion { period_end: None }),
1949 AllocateError::BalanceInsufficient(tollgate_core::BalanceShortfall {
1950 remaining: CostUnits(1),
1951 period_end: None,
1952 }),
1953 AllocateError::BalanceOverflow,
1954 ]
1955 }
1956
1957 /// The indices must be a permutation of `0..COUNT`: two refusals sharing
1958 /// a slot would silently merge their tallies, and a slot no refusal maps
1959 /// to would export a counter that can never move.
1960 #[test]
1961 fn indices_cover_every_slot_exactly_once() {
1962 let mut seen = [false; AllocateError::COUNT];
1963 for error in all() {
1964 let index = error.index();
1965 assert!(index < AllocateError::COUNT, "{error} indexes out of range");
1966 assert!(!seen[index], "{error} shares slot {index}");
1967 seen[index] = true;
1968 }
1969 assert!(seen.iter().all(|hit| *hit), "every slot must be claimed");
1970 }
1971
1972 /// Labels are read by whatever scrapes the counters, so they are a
1973 /// contract: distinct, and free of the backend text `Display` carries.
1974 /// `Storage` wraps an arbitrary message, so labelling by rendered text
1975 /// would mint a fresh time series per connection failure.
1976 #[test]
1977 fn labels_are_distinct_and_free_of_backend_text() {
1978 for (position, error) in all().iter().enumerate() {
1979 assert_eq!(error.name(), AllocateError::NAMES[position]);
1980 }
1981 let storage = AllocateError::Storage(StoreError("connection reset".into()));
1982 assert_eq!(storage.name(), "storage");
1983 assert!(
1984 !storage.name().contains("connection"),
1985 "the label must not carry the backend's message"
1986 );
1987 let mut names = AllocateError::NAMES;
1988 names.sort_unstable();
1989 names.iter().reduce(|previous, next| {
1990 assert_ne!(previous, next, "duplicate label {next}");
1991 next
1992 });
1993 }
1994
1995 /// The payload must not affect the slot, or one refusal would scatter
1996 /// across slots as its message varied.
1997 #[test]
1998 fn payload_does_not_affect_the_slot() {
1999 assert_eq!(
2000 AllocateError::Storage(StoreError("a".into())).index(),
2001 AllocateError::Storage(StoreError("b".into())).index()
2002 );
2003 }
2004
2005 #[test]
2006 fn conservation_requires_an_exact_equation_without_overflow() {
2007 let balanced = Conservation {
2008 deposited: CostUnits(10),
2009 overage_recorded: CostUnits::ZERO,
2010 balance: CostUnits(1),
2011 active_lease_grants: CostUnits(2),
2012 settled_usage: CostUnits(3),
2013 settlement_loss: CostUnits(4),
2014 expired: CostUnits::ZERO,
2015 };
2016 assert!(balanced.holds());
2017
2018 let drifted = Conservation {
2019 deposited: CostUnits(11),
2020 ..balanced
2021 };
2022 assert!(!drifted.holds());
2023
2024 let overflowing = Conservation {
2025 deposited: CostUnits(u64::MAX),
2026 overage_recorded: CostUnits::ZERO,
2027 balance: CostUnits(u64::MAX),
2028 active_lease_grants: CostUnits(1),
2029 settled_usage: CostUnits::ZERO,
2030 settlement_loss: CostUnits::ZERO,
2031 expired: CostUnits::ZERO,
2032 };
2033 assert!(!overflowing.holds());
2034 }
2035
2036 /// Overage funds the left side, so usage it paid for closes the equation
2037 /// rather than breaking it — and the same numbers without the funding term
2038 /// must *not* balance, or the field would be decorative.
2039 #[test]
2040 fn overage_funds_the_usage_it_bills() {
2041 let elastic = Conservation {
2042 deposited: CostUnits(10),
2043 overage_recorded: CostUnits(5),
2044 balance: CostUnits(1),
2045 active_lease_grants: CostUnits(2),
2046 settled_usage: CostUnits(8),
2047 settlement_loss: CostUnits(4),
2048 expired: CostUnits::ZERO,
2049 };
2050 assert!(elastic.holds());
2051 assert!(
2052 !Conservation {
2053 overage_recorded: CostUnits::ZERO,
2054 ..elastic
2055 }
2056 .holds(),
2057 "the same ledger without the funding term must fail by exactly the overage"
2058 );
2059 }
2060
2061 /// The left side is checked too. Overflowing the funding sum answers
2062 /// `false` rather than wrapping to a total that might coincidentally match
2063 /// the right side (INVARIANTS.md GL-11).
2064 #[test]
2065 fn overflowing_the_funding_sum_is_a_violation_not_a_wrap() {
2066 let overflowing = Conservation {
2067 deposited: CostUnits(u64::MAX),
2068 overage_recorded: CostUnits(1),
2069 balance: CostUnits::ZERO,
2070 active_lease_grants: CostUnits::ZERO,
2071 settled_usage: CostUnits::ZERO,
2072 settlement_loss: CostUnits::ZERO,
2073 expired: CostUnits::ZERO,
2074 };
2075 assert!(!overflowing.holds());
2076 }
2077
2078 /// An allowance that expired at a period boundary left the balance without
2079 /// being spent, so the equation only closes if `expired` is on the right
2080 /// side — and the same ledger without the term must fail by exactly the
2081 /// units that expired, or the field would be decorative (GL-97).
2082 #[test]
2083 fn expiry_accounts_for_an_allowance_that_was_never_spent() {
2084 let rolled = 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::ZERO,
2091 expired: CostUnits(4),
2092 };
2093 assert!(rolled.holds());
2094 assert!(
2095 !Conservation {
2096 expired: CostUnits::ZERO,
2097 ..rolled
2098 }
2099 .holds(),
2100 "the same ledger without the expiry term must fail by exactly the expired units"
2101 );
2102 }
2103
2104 /// A rollover pass reports what it did, and its caller decides whether to
2105 /// ask for another batch from `is_saturated` alone. Every accessor is
2106 /// asserted directly: a `len` that always answered 1, or an `is_empty`
2107 /// stuck either way, would send the server's drain loop into an endless
2108 /// round of empty batches or stop it one batch short of the boundary it
2109 /// was crossing.
2110 #[test]
2111 fn rollover_batch_reports_what_it_rolled() {
2112 let limit = NonZeroUsize::new(2).unwrap();
2113 let rolled = |account| RolledAccount {
2114 account_id: AccountId(account),
2115 deposited: CostUnits(100),
2116 expired: CostUnits::ZERO,
2117 };
2118
2119 let empty = RolloverBatch::try_new(Vec::new(), limit).unwrap();
2120 assert!(empty.is_empty());
2121 assert_eq!(empty.len(), 0);
2122 assert!(!empty.is_saturated());
2123 assert!(empty.rolled().is_empty());
2124
2125 let partial = RolloverBatch::try_new(vec![rolled(1)], limit).unwrap();
2126 assert!(!partial.is_empty());
2127 assert_eq!(partial.len(), 1);
2128 assert!(
2129 !partial.is_saturated(),
2130 "a partial batch is what ends the drain"
2131 );
2132 assert_eq!(partial.rolled(), &[rolled(1)]);
2133
2134 let full = RolloverBatch::try_new(vec![rolled(1), rolled(2)], limit).unwrap();
2135 assert_eq!(full.len(), 2);
2136 assert!(
2137 full.is_saturated(),
2138 "a batch at the limit means there may be more"
2139 );
2140 }
2141
2142 /// Saturation is derived from the limit the caller asked for, so a backend
2143 /// returning more than it was allowed is corruption to refuse rather than
2144 /// a batch to trust — the same contract `ReclaimBatch::try_new` carries.
2145 #[test]
2146 fn a_rollover_batch_beyond_its_limit_is_refused() {
2147 let rolled = |account| RolledAccount {
2148 account_id: AccountId(account),
2149 deposited: CostUnits(100),
2150 expired: CostUnits::ZERO,
2151 };
2152 assert!(
2153 RolloverBatch::try_new(vec![rolled(1), rolled(2)], NonZeroUsize::new(1).unwrap())
2154 .is_err()
2155 );
2156 }
2157
2158 #[test]
2159 fn reclaim_batch_reports_an_empty_result() {
2160 let batch = ReclaimBatch::try_new(Vec::new(), NonZeroUsize::new(2).unwrap()).unwrap();
2161 assert!(batch.is_empty());
2162 assert_eq!(batch.len(), 0);
2163 assert!(!batch.is_saturated());
2164 assert!(batch.reclaimed().is_empty());
2165 }
2166
2167 #[tokio::test]
2168 async fn full_drain_preserves_an_initial_batch_error() {
2169 let allocator = ScriptedReclaimer {
2170 script: ReclaimScript::FailFirst,
2171 calls: AtomicUsize::new(0),
2172 };
2173 let expected = StoreError("scripted reclaim failure".into());
2174 assert_eq!(
2175 allocator.reclaim_expired(Timestamp::MIN).await,
2176 Err(expected)
2177 );
2178 assert_eq!(allocator.calls.load(Ordering::Acquire), 1);
2179 }
2180
2181 #[tokio::test]
2182 async fn full_drain_reports_progress_before_a_later_batch_error() {
2183 let allocator = ScriptedReclaimer {
2184 script: ReclaimScript::FullBatchThenFail,
2185 calls: AtomicUsize::new(0),
2186 };
2187 let original = StoreError("scripted reclaim failure".into());
2188 let error = allocator.reclaim_expired(Timestamp::MIN).await.unwrap_err();
2189 assert_ne!(error, original);
2190 assert!(error.0.contains("256 leases"));
2191 assert!(error.0.contains("256 units"));
2192 assert_eq!(allocator.calls.load(Ordering::Acquire), 2);
2193 }
2194
2195 /// The push-capacity boundary, pinned where both backends read it.
2196 ///
2197 /// Strictly greater is the whole content of the rule: a batch that exactly
2198 /// fills the channel is delivered, so warning at equality would fire on the
2199 /// largest successful case and train an operator to ignore it. Mutation
2200 /// testing found this untested — the comparison could be flipped to `<`,
2201 /// `<=` or `>=` and every scenario stayed green, because nothing observed
2202 /// the warning at all (GL-51).
2203 #[test]
2204 fn the_push_capacity_warning_fires_only_above_the_channel() {
2205 assert!(
2206 !pushes_exceed_capacity(0),
2207 "an empty batch is not a capacity problem"
2208 );
2209 assert!(!pushes_exceed_capacity(PUSH_CHANNEL_CAPACITY - 1));
2210 assert!(
2211 !pushes_exceed_capacity(PUSH_CHANNEL_CAPACITY),
2212 "a batch that exactly fills the channel is still delivered"
2213 );
2214 assert!(
2215 pushes_exceed_capacity(PUSH_CHANNEL_CAPACITY + 1),
2216 "one more than the channel holds is what makes a subscriber lag"
2217 );
2218 }
2219}