Skip to main content

tollgate_server/
lib.rs

1//! The control-plane HTTP service.
2//!
3//! Latency-tolerant by design — its contract is correctness (no double-spend,
4//! fencing, idempotent ingest), enforced by whatever [`tollgate_store`] backend
5//! it is built over. The service is generic over that backend: anything
6//! satisfying [`Backend`] serves identically, which is the
7//! pluggability seam the Postgres backend drops into.
8//!
9//! Surface (`/v1`):
10//! - `POST /v1/leases/acquire`, `POST /v1/leases/release`,
11//!   `POST /v1/leases/consolidate`, `POST /v1/leases/reclaim` — fenced lease
12//!   lifecycle.
13//! - `GET /v1/keys` — revisioned active credential pages (instance authority only).
14//! - `GET /v1/snapshots` — the principal catalogue (GL-48).
15//! - `GET /v1/snapshots/{principal}` — compiled snapshot fetch.
16//! - `POST /v1/usage/ingest` — idempotent usage batches.
17//! - `POST /v1/admin/accounts`, `POST /v1/admin/accounts/{id}/deposit`,
18//!   `POST /v1/admin/accounts/{id}/status`,
19//!   `POST /v1/admin/accounts/{id}/capacity-class`,
20//!   `POST/GET /v1/admin/accounts/{id}/keys`,
21//!   `DELETE /v1/admin/accounts/{id}/keys/{key}`,
22//!   `PUT/DELETE /v1/admin/accounts/{id}/keys/{key}/snapshot` (GL-143),
23//!   `PUT/DELETE /v1/admin/snapshots/{principal}` — administration.
24//! - `GET /livez`, `GET /readyz` — probes.
25//!
26//! All protected handlers require verified instance or operator evidence.
27//! [`serve`] owns TLS and refuses exposed plaintext; [`security::ServerSecurity`]
28//! atomically rotates roles, credentials and TLS. Administrative audit events
29//! carry receipts captured by the backend at mutation. See
30//! `docs/CONTROL_PLANE_SECURITY.md` for deployment and audit collection.
31
32#![deny(missing_docs)]
33
34pub mod config;
35pub mod error;
36pub mod google;
37mod maintenance;
38pub mod security;
39pub mod transport;
40
41use std::sync::Arc;
42
43use axum::extract::{DefaultBodyLimit, State};
44use axum::http::StatusCode;
45use axum::routing::{get, post, put};
46use axum::{Json, Router};
47use tracing::Instrument as _;
48
49use tollgate_auth::CredentialIssuer;
50use tollgate_core::{AccountId, CapacityClass, CostUnits, KeyId, Principal, PublishableSnapshot};
51use tollgate_store::wire::{
52    API_PREFIX, AccountKeyResponse, AccountKeysResponse, AccountResponse, AcquireRequest,
53    AcquireResponse, ConsolidateRequest, ConsolidateResponse, CreateAccountRequest, DepositRequest,
54    IngestRequest, IssueKeyRequest, IssuedKeyResponse, MAX_INGEST_BODY_BYTES,
55    MAX_SNAPSHOT_BODY_BYTES, PrincipalsResponse, PublishSnapshotRequest, ReleaseRequest,
56    RevokeKeyResponse, SetBudgetRequest, SetBudgetResponse, SetCapacityClassRequest,
57    SetStatusRequest, SetStatusResponse,
58};
59use tollgate_store::{
60    AccountConfig, AdminState, AdminStore, Clock, DEFAULT_KEY_PAGE_LIMIT,
61    DEFAULT_RECLAIM_BATCH_LIMIT, DEFAULT_ROLLOVER_BATCH_LIMIT, IngestReport, KeyDirectory,
62    KeyRecord, KeySource, LeaseAllocator, ReclaimBatch, ReclaimedLease, SnapshotResolution,
63    SnapshotSource, StoreError, StoreHealth, UsageSink,
64};
65
66use crate::error::{ApiError, ApiJson, ApiPath, ApiQuery};
67use crate::security::{Authorization, InstanceIdentity, OperatorIdentity, Role, ServerSecurity};
68
69/// Everything the handlers need. `S` is the storage backend; the clock is the
70/// single place wall time enters the server.
71pub struct ServerState<S> {
72    /// The backend every handler and the [`serve`] maintenance sweep act on.
73    pub store: Arc<S>,
74    /// The time source for lease and ingest calls, credential page reads,
75    /// credential and client-certificate validity checks, audit event times,
76    /// and the maintenance sweep.
77    pub clock: Arc<dyn Clock>,
78    /// The live verifier, role map and TLS generation, shared by the
79    /// authorization middleware and the listener [`serve`] builds. Its
80    /// transport mode decides whether `serve` accepts a non-loopback
81    /// listener.
82    pub security: Arc<ServerSecurity>,
83    /// Who may mint credentials, if this deployment issues them at all (GL-121).
84    ///
85    /// `None` is a deliberate answer, not an omission: an instance that only
86    /// verifies has no business holding the capability to create credentials,
87    /// and the issuance routes answer `501` rather than pretending. That is the
88    /// shape `list_principals` already uses for a backend that cannot
89    /// enumerate — an unsupported capability is reported, never faked.
90    pub issuer: Option<Arc<dyn CredentialIssuer + Send + Sync>>,
91}
92
93impl<S> Clone for ServerState<S> {
94    fn clone(&self) -> Self {
95        ServerState {
96            store: Arc::clone(&self.store),
97            clock: Arc::clone(&self.clock),
98            security: Arc::clone(&self.security),
99            issuer: self.issuer.clone(),
100        }
101    }
102}
103
104/// The full store bound the server needs from a backend.
105pub trait Backend:
106    LeaseAllocator
107    + SnapshotSource
108    + UsageSink
109    + AdminStore
110    + KeySource
111    // A server administers credentials as well as projecting them (GL-121), so
112    // its backend must be the credential *directory* and not only a source.
113    // `HttpStore` implements `KeySource` but not this — correctly: it is a
114    // client of a server, never the authority behind one.
115    + KeyDirectory
116    + StoreHealth
117    + Send
118    + Sync
119    + 'static
120{
121}
122impl<T> Backend for T where
123    T: LeaseAllocator
124        + SnapshotSource
125        + UsageSink
126        + AdminStore
127        + KeySource
128        + KeyDirectory
129        + StoreHealth
130        + Send
131        + Sync
132        + 'static
133{
134}
135
136/// In-process handler router. Without an owned maintenance task `/readyz`
137/// returns 503 even when the store answers; use [`serve`] for a ready service.
138pub fn router<S: Backend>(state: ServerState<S>) -> Router {
139    router_with_maintenance(state, maintenance::Monitor::unmanaged())
140}
141
142fn router_with_maintenance<S: Backend>(
143    state: ServerState<S>,
144    maintenance: maintenance::Monitor,
145) -> Router {
146    let authorization = |role| Authorization {
147        security: Arc::clone(&state.security),
148        clock: Arc::clone(&state.clock),
149        role,
150    };
151    let instance = Router::new()
152        .route("/leases/acquire", post(acquire::<S>))
153        .route("/leases/release", post(release::<S>))
154        .route("/leases/consolidate", post(consolidate::<S>))
155        .route("/leases/reclaim", post(reclaim::<S>))
156        .route("/snapshots", get(list_principals::<S>))
157        .route("/keys", get(active_keys::<S>))
158        .route("/snapshots/{principal}", get(fetch_snapshot::<S>))
159        .route(
160            "/usage/ingest",
161            post(ingest::<S>).layer(DefaultBodyLimit::max(MAX_INGEST_BODY_BYTES)),
162        )
163        .route_layer(axum::middleware::from_fn_with_state(
164            authorization(Role::Instance),
165            security::authorize,
166        ));
167    let operator = Router::new()
168        .route("/accounts", post(create_account::<S>))
169        .route("/accounts/{account}", get(account::<S>))
170        .route("/accounts/{account}/budget", put(set_budget::<S>))
171        .route(
172            "/accounts/{account}/keys",
173            post(issue_key::<S>).get(list_account_keys::<S>),
174        )
175        .route(
176            "/accounts/{account}/keys/{key}",
177            axum::routing::delete(revoke_account_key::<S>),
178        )
179        .route(
180            "/accounts/{account}/keys/{key}/snapshot",
181            put(publish_key_snapshot::<S>)
182                .delete(remove_key_snapshot::<S>)
183                .layer(DefaultBodyLimit::max(MAX_SNAPSHOT_BODY_BYTES)),
184        )
185        .route("/accounts/{account}/deposit", post(deposit::<S>))
186        .route("/accounts/{account}/status", post(set_status::<S>))
187        .route(
188            "/accounts/{account}/capacity-class",
189            post(set_capacity_class::<S>),
190        )
191        .route(
192            "/snapshots/{principal}",
193            put(publish_snapshot::<S>)
194                .delete(remove_snapshot::<S>)
195                .layer(DefaultBodyLimit::max(MAX_SNAPSHOT_BODY_BYTES)),
196        )
197        .route_layer(axum::middleware::from_fn_with_state(
198            authorization(Role::Operator),
199            security::authorize,
200        ));
201    Router::new()
202        .route("/livez", get(async || StatusCode::OK))
203        .route(
204            "/readyz",
205            get(readyz::<S>).layer(axum::Extension(maintenance)),
206        )
207        .nest(API_PREFIX, instance.nest("/admin", operator))
208        .with_state(state)
209        .layer(axum::middleware::from_fn(error::report_http_failure))
210}
211
212/// Serve until `shutdown` resolves, running the maintenance sweep every
213/// `reclaim_interval` in the background (INVARIANTS.md GL-9's server half).
214pub async fn serve<S: Backend>(
215    listener: tokio::net::TcpListener,
216    state: ServerState<S>,
217    reclaim_interval: std::time::Duration,
218    shutdown: impl Future<Output = ()> + Send + 'static,
219) -> std::io::Result<()> {
220    if reclaim_interval.is_zero()
221        || tokio::time::Instant::now()
222            .checked_add(reclaim_interval)
223            .is_none()
224    {
225        return Err(std::io::Error::new(
226            std::io::ErrorKind::InvalidInput,
227            "reclaim_interval must be positive and fit the monotonic clock",
228        ));
229    }
230    let listener = transport::SecureListener::new(listener, Arc::clone(&state.security))?;
231    let sweep_store = Arc::clone(&state.store);
232    let sweep_clock = Arc::clone(&state.clock);
233    let mut sweeper = maintenance::Task::spawn(move |health| {
234        maintenance_sweep(sweep_store, sweep_clock, reclaim_interval, health).instrument(
235            tracing::info_span!(
236                "maintenance_sweep",
237                interval_ms = reclaim_interval.as_millis()
238            ),
239        )
240    });
241    // Axum spawns its signal waiter. Hand it only a receiver; retaining the
242    // sender and the caller's future here makes every serve exit close the
243    // signal, including cancellation before the caller requests shutdown.
244    let (signal, signalled) = tokio::sync::oneshot::channel::<()>();
245    let mut signal = Some(signal);
246    let server = axum::serve(
247        listener,
248        router_with_maintenance(state, sweeper.monitor.clone())
249            .into_make_service_with_connect_info::<transport::PeerIdentity>(),
250    )
251    .with_graceful_shutdown(async move {
252        // Sender drop is also a shutdown signal: it means serve exited.
253        let Ok(()) = signalled.await else { return };
254    });
255    let server = std::future::IntoFuture::into_future(server);
256    tokio::pin!(server, shutdown);
257    loop {
258        tokio::select! {
259            biased;
260            result = &mut sweeper.join => {
261                if result.as_ref().is_err_and(|error| error.is_cancelled())
262                    && sweeper.stop.requested()
263                {
264                    // Expected cancellation after the graceful-shutdown signal.
265                    // Readiness was withdrawn before the abort; HTTP may drain.
266                    return server.await;
267                }
268                sweeper.stop.stop();
269                drop(signal.take());
270                let reason = match result {
271                    Err(error) if error.is_panic() => "panic",
272                    Err(_) => "cancelled",
273                    Ok(()) => "unexpected-return",
274                };
275                tracing::error!(
276                    operation = "maintenance",
277                    reason,
278                    "server maintenance task stopped; stopping the listener"
279                );
280                return Err(std::io::Error::other("server maintenance task stopped unexpectedly"));
281            }
282            result = &mut server => return result,
283            () = &mut shutdown, if signal.is_some() => {
284                sweeper.stop.stop();
285                drop(signal.take());
286            }
287        }
288    }
289}
290
291/// Reclaim expired leases and roll due budget periods forever, on the
292/// configured interval.
293///
294/// This is INVARIANTS.md GL-9's server half, and it is the only thing that
295/// returns units stranded by a crashed holder. A store can serve `ping` and
296/// fail maintenance, so each outcome is published to readiness and reported
297/// independently. Silence about success would
298/// leave "the sweep is running but finding nothing" and "the sweep stopped"
299/// indistinguishable.
300///
301/// The rollover pass (GL-97) shares this tick rather than owning a timer,
302/// because it needs the same three things and gets them here already: a frozen
303/// cutoff, a bounded drain, and a failure that is reported rather than
304/// swallowed. It is also the only trigger — without it a schedule is a stored
305/// intention nothing ever acts on — and a boundary the pass misses is not lost
306/// but late, since the store crosses it on the next tick that sees it. The two
307/// passes are independent: one failing must not stop the other, because a
308/// stuck rollover stranding quota would be a strictly worse outcome than a
309/// late allowance.
310async fn maintenance_sweep<S: Backend>(
311    store: Arc<S>,
312    clock: Arc<dyn Clock>,
313    interval: std::time::Duration,
314    mut health: maintenance::Publisher,
315) {
316    #[derive(Default)]
317    struct Progress {
318        leases: u128,
319        units: u128,
320        batches: u64,
321    }
322
323    impl Progress {
324        fn record(&mut self, batch: &ReclaimBatch) -> Result<(), StoreError> {
325            let batch_leases =
326                u128::from(u64::try_from(batch.len()).map_err(|_| {
327                    StoreError("reclaim batch lease count exceeds u64 range".into())
328                })?);
329            let batch_units = batch.reclaimed().iter().try_fold(0u128, |total, lease| {
330                total
331                    .checked_add(u128::from(lease.forfeited.get()))
332                    .ok_or_else(|| StoreError("reclaim batch unit total overflow".into()))
333            })?;
334            let leases = self
335                .leases
336                .checked_add(batch_leases)
337                .ok_or_else(|| StoreError("reclaim cycle lease count overflow".into()))?;
338            let units = self
339                .units
340                .checked_add(batch_units)
341                .ok_or_else(|| StoreError("reclaim cycle unit total overflow".into()))?;
342            let batches = self
343                .batches
344                .checked_add(1)
345                .ok_or_else(|| StoreError("reclaim cycle batch count overflow".into()))?;
346            *self = Progress {
347                leases,
348                units,
349                batches,
350            };
351            Ok(())
352        }
353    }
354
355    let mut tick = tokio::time::interval(interval);
356    tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
357    loop {
358        tick.tick().await;
359        // Freeze the cutoff for this cycle. A continuously advancing cutoff
360        // could make a busy allocator feed the drain forever; a fixed one is
361        // a finite backlog and still returns all quota expired at this tick.
362        let now = clock.now();
363        let mut progress = Progress::default();
364        let outcome = loop {
365            match store
366                .reclaim_expired_batch(now, DEFAULT_RECLAIM_BATCH_LIMIT)
367                .await
368            {
369                Ok(batch) => {
370                    if progress.record(&batch).is_err() {
371                        tracing::error!(
372                            operation = "reclaim",
373                            "reclaim progress counter overflowed; stopping maintenance"
374                        );
375                        return;
376                    }
377                    if !batch.is_saturated() {
378                        break Ok(());
379                    }
380                    // MemoryStore can finish a batch without an .await that
381                    // yields. Let request handlers run between backlog chunks;
382                    // PostgreSQL benefits from the same explicit fairness.
383                    tokio::task::yield_now().await;
384                }
385                Err(error) => break Err(error),
386            }
387        };
388
389        let Ok(report) = health.record(maintenance::Operation::Reclaim, outcome.is_ok()) else {
390            tracing::error!(
391                operation = "reclaim",
392                "maintenance failure counter exhausted"
393            );
394            return;
395        };
396        match outcome {
397            Ok(()) => {
398                if report.recovered_after > 0 {
399                    tracing::info!(
400                        operation = "reclaim",
401                        after_failures = report.recovered_after,
402                        "reclaim sweep recovered"
403                    );
404                }
405                // A swept lease is one its holder never released, so its
406                // remainder was forfeited rather than returned (GL-136). Units
407                // forfeited are an operator signal: a crash, or a shutdown
408                // whose release deadline lapsed.
409                if progress.units > 0 {
410                    tracing::warn!(
411                        leases = %progress.leases,
412                        forfeited_units = %progress.units,
413                        batches = progress.batches,
414                        "swept unreleased leases; their remainders are forfeited as provisional loss"
415                    );
416                } else if progress.leases > 0 {
417                    tracing::info!(
418                        leases = %progress.leases,
419                        forfeited_units = %progress.units,
420                        batches = progress.batches,
421                        "swept unreleased leases with nothing left to forfeit"
422                    );
423                }
424            }
425            Err(_) => {
426                macro_rules! report_failure {
427                    ($level:ident) => {
428                        tracing::$level!(
429                            operation = "reclaim",
430                            code = "storage",
431                            consecutive_failures = report.failures,
432                            reclaimed_leases = %progress.leases,
433                            forfeited_units = %progress.units,
434                            completed_batches = progress.batches,
435                            "reclaim sweep failed; completed batches stay committed and remaining expired leases stay stranded until it recovers"
436                        );
437                    };
438                }
439                if report.failures >= maintenance::PERSISTENT_FAILURES {
440                    report_failure!(error);
441                } else {
442                    report_failure!(warn);
443                }
444            }
445        }
446        if roll_due_periods(store.as_ref(), now, &mut health)
447            .await
448            .is_break()
449        {
450            return;
451        }
452    }
453}
454
455/// Drain the rollover pass for one tick.
456///
457/// Bounded per transaction and drained to a partial batch, exactly as the
458/// reclaim loop above is, and for a sharper reason: every scheduled account is
459/// due at the same instant, so the first tick after midnight on the 1st has
460/// the whole scheduled population to cross. A single unbounded statement there
461/// would hold locks across the entire account table.
462///
463/// Failures are reported and the tick ends. Retrying inside the tick would
464/// spin against a store that is down; the next tick is the retry, and until it
465/// succeeds the affected accounts keep spending last period's allowance, which
466/// is late rather than wrong. Only counter exhaustion breaks maintenance;
467/// a reported backend failure continues to the next scheduled retry.
468async fn roll_due_periods<S: Backend>(
469    store: &S,
470    now: jiff::Timestamp,
471    health: &mut maintenance::Publisher,
472) -> std::ops::ControlFlow<()> {
473    let mut accounts: u64 = 0;
474    let mut batches: u64 = 0;
475    loop {
476        match store
477            .roll_due_periods(now, DEFAULT_ROLLOVER_BATCH_LIMIT)
478            .await
479        {
480            Ok(batch) => {
481                let Some(total) = u64::try_from(batch.len())
482                    .ok()
483                    .and_then(|rolled| accounts.checked_add(rolled))
484                    .zip(batches.checked_add(1))
485                else {
486                    // Unreachable short of a backend ignoring the batch limit
487                    // forever, and still not a place to wrap: a counter that
488                    // silently restarts would under-report a boundary that
489                    // rolled more accounts than it claimed.
490                    tracing::error!(
491                        accounts,
492                        batches,
493                        "budget rollover progress counter overflowed; stopping maintenance"
494                    );
495                    return std::ops::ControlFlow::Break(());
496                };
497                (accounts, batches) = (total.0, total.1);
498                if !batch.is_saturated() {
499                    break;
500                }
501                // MemoryStore can finish a batch without an .await that
502                // yields; let request handlers run between chunks.
503                tokio::task::yield_now().await;
504            }
505            Err(_) => {
506                let Ok(report) = health.record(maintenance::Operation::Rollover, false) else {
507                    tracing::error!(
508                        operation = "budget-rollover",
509                        "maintenance failure counter exhausted"
510                    );
511                    return std::ops::ControlFlow::Break(());
512                };
513                macro_rules! report_failure {
514                    ($level:ident) => {
515                        tracing::$level!(
516                            operation = "budget-rollover",
517                            code = "storage",
518                            consecutive_failures = report.failures,
519                            rolled_accounts = accounts,
520                            completed_batches = batches,
521                            "budget rollover failed; accounts past their boundary keep last period's allowance until it recovers"
522                        );
523                    };
524                }
525                if report.failures >= maintenance::PERSISTENT_FAILURES {
526                    report_failure!(error);
527                } else {
528                    report_failure!(warn);
529                }
530                return std::ops::ControlFlow::Continue(());
531            }
532        }
533    }
534    let report = health
535        .record(maintenance::Operation::Rollover, true)
536        .expect("a successful pass clears its failure counter without arithmetic");
537    if report.recovered_after > 0 {
538        tracing::info!(
539            operation = "budget-rollover",
540            after_failures = report.recovered_after,
541            "budget rollover recovered"
542        );
543    }
544    if accounts > 0 {
545        tracing::info!(accounts, batches, "rolled budget periods");
546    }
547    std::ops::ControlFlow::Continue(())
548}
549
550/// Ready only when the backing store answers (review finding GL-11): a server
551/// whose source of truth or maintenance is unavailable must not attract traffic.
552async fn readyz<S: Backend>(
553    State(state): State<ServerState<S>>,
554    axum::Extension(maintenance): axum::Extension<maintenance::Monitor>,
555) -> StatusCode {
556    match state.store.ping().await {
557        Ok(()) if maintenance.healthy() => StatusCode::OK,
558        Ok(()) => StatusCode::SERVICE_UNAVAILABLE,
559        Err(_) => {
560            tracing::warn!(
561                operation = "readiness",
562                code = "storage",
563                "readiness probe failed: the store did not answer"
564            );
565            StatusCode::SERVICE_UNAVAILABLE
566        }
567    }
568}
569
570async fn acquire<S: Backend>(
571    _identity: InstanceIdentity,
572    State(state): State<ServerState<S>>,
573    ApiJson(request): ApiJson<AcquireRequest>,
574) -> Result<Json<AcquireResponse>, ApiError> {
575    let grant = state
576        .store
577        .acquire(
578            request.account_id,
579            request.requested,
580            request.ttl.duration()?,
581            state.clock.now(),
582        )
583        .await?;
584    Ok(Json(grant))
585}
586
587async fn release<S: Backend>(
588    _identity: InstanceIdentity,
589    State(state): State<ServerState<S>>,
590    ApiJson(request): ApiJson<ReleaseRequest>,
591) -> Result<StatusCode, ApiError> {
592    state
593        .store
594        .release(
595            request.lease_id,
596            request.fencing_token,
597            request.unspent,
598            state.clock.now(),
599        )
600        .await?;
601    Ok(StatusCode::NO_CONTENT)
602}
603
604/// Return a lease's unspent units and re-grant against the restored balance,
605/// in one backend transaction.
606///
607/// It is a route rather than two calls the instance could make itself for the
608/// reason the trait method exists: split across the wire, another instance can
609/// take the returned units in the gap, and the grant policy can hand back less
610/// than was returned (`LeaseAllocator::consolidate`).
611async fn consolidate<S: Backend>(
612    _identity: InstanceIdentity,
613    State(state): State<ServerState<S>>,
614    ApiJson(request): ApiJson<ConsolidateRequest>,
615) -> Result<Json<ConsolidateResponse>, ApiError> {
616    let grant = state
617        .store
618        .consolidate(
619            request.lease_id,
620            request.fencing_token,
621            request.unspent,
622            request.requested,
623            request.needed,
624            request.ttl.duration()?,
625            state.clock.now(),
626        )
627        .await?;
628    Ok(Json(grant))
629}
630
631async fn reclaim<S: Backend>(
632    _identity: InstanceIdentity,
633    State(state): State<ServerState<S>>,
634) -> Result<Json<Vec<ReclaimedLease>>, ApiError> {
635    Ok(Json(state.store.reclaim_expired(state.clock.now()).await?))
636}
637
638async fn fetch_snapshot<S: Backend>(
639    _identity: InstanceIdentity,
640    State(state): State<ServerState<S>>,
641    ApiPath(principal): ApiPath<Principal>,
642) -> Result<Json<Arc<tollgate_core::AccountSnapshot>>, ApiError> {
643    match state.store.snapshot(principal).await? {
644        SnapshotResolution::Present(snapshot) => Ok(Json(snapshot.into_inner())),
645        SnapshotResolution::Revoked { generation } => Err(ApiError::revoked(generation)),
646        SnapshotResolution::Unknown => Err(ApiError::not_found("unknown-principal", "no snapshot")),
647    }
648}
649
650/// The principal catalogue, for instances that track every customer rather
651/// than a configured slice (GL-48).
652///
653/// A backend that cannot enumerate answers 501 rather than an empty list: an
654/// empty catalogue and an unsupported one lead an instance to opposite
655/// conclusions — forget everything, or keep what you were configured with —
656/// and conflating them would silently strand it on a stale set.
657async fn list_principals<S: Backend>(
658    _identity: InstanceIdentity,
659    State(state): State<ServerState<S>>,
660) -> Result<Json<PrincipalsResponse>, ApiError> {
661    match state.store.principals().await? {
662        Some(principals) => Ok(Json(PrincipalsResponse { principals })),
663        None => Err(ApiError::not_implemented(
664            "enumeration-unsupported",
665            "this backend cannot list principals",
666        )),
667    }
668}
669
670async fn ingest<S: Backend>(
671    _identity: InstanceIdentity,
672    State(state): State<ServerState<S>>,
673    body: Result<ApiJson<IngestRequest>, ApiError>,
674) -> Result<Json<IngestReport>, ApiError> {
675    let ApiJson(request) = body.map_err(|mut error| {
676        // Only this route accepts usage batches. The shared JSON rejection
677        // describes the byte limit without advice about another route's data.
678        if error.status == StatusCode::PAYLOAD_TOO_LARGE {
679            error.title.push_str(&format!(
680                "; usage batches are capped at {} events",
681                tollgate_store::MAX_INGEST_BATCH
682            ));
683        }
684        error
685    })?;
686    Ok(Json(
687        state
688            .store
689            .ingest(&request.events, state.clock.now())
690            .await?,
691    ))
692}
693
694#[derive(serde::Deserialize)]
695#[serde(deny_unknown_fields)]
696struct KeyQuery {
697    after: Option<tollgate_core::KeyId>,
698    limit: Option<usize>,
699}
700
701async fn active_keys<S: Backend>(
702    _instance: InstanceIdentity,
703    State(state): State<ServerState<S>>,
704    ApiQuery(query): ApiQuery<KeyQuery>,
705) -> Result<impl axum::response::IntoResponse, ApiError> {
706    let limit = credential_page_limit(query.limit)?;
707    let page = state
708        .store
709        .active_keys_page(state.clock.now(), query.after, limit)
710        .await
711        .and_then(|page| {
712            page.validate_request(query.after, limit)?;
713            Ok(page)
714        })
715        .map_err(|_| ApiError {
716            status: StatusCode::SERVICE_UNAVAILABLE,
717            code: "credential-source-unavailable",
718            title: "credential source is unavailable".into(),
719            generation: None,
720            balance_exhaustion: None,
721            balance_shortfall: None,
722        })?;
723    let body = tollgate_store::wire::KeysResponse {
724        revision: page.revision(),
725        as_of: page.as_of(),
726        next_after: page.next_after(),
727        keys: page.into_records(),
728    };
729    Ok((
730        [(axum::http::header::CACHE_CONTROL, "no-store")],
731        Json(body),
732    ))
733}
734
735async fn create_account<S: Backend>(
736    operator: OperatorIdentity,
737    State(state): State<ServerState<S>>,
738    ApiJson(request): ApiJson<CreateAccountRequest>,
739) -> Result<StatusCode, ApiError> {
740    operator
741        .run(
742            "create_account",
743            request.account_id,
744            state.clock.as_ref(),
745            state.store.create_account(AccountConfig {
746                account_id: request.account_id,
747                initial_balance: request.initial_balance,
748                status: request.status,
749                capacity_class: CapacityClass::Assured,
750            }),
751        )
752        .await?;
753    Ok(StatusCode::CREATED)
754}
755
756async fn deposit<S: Backend>(
757    operator: OperatorIdentity,
758    State(state): State<ServerState<S>>,
759    ApiPath(account): ApiPath<AccountId>,
760    ApiJson(request): ApiJson<DepositRequest>,
761) -> Result<StatusCode, ApiError> {
762    if request.units == CostUnits::ZERO {
763        return Err(ApiError::bad_request(
764            "zero-deposit",
765            "deposit must be positive",
766        ));
767    }
768    operator
769        .run(
770            "deposit",
771            account,
772            state.clock.as_ref(),
773            state.store.deposit(account, request.units),
774        )
775        .await?;
776    Ok(StatusCode::NO_CONTENT)
777}
778
779async fn set_status<S: Backend>(
780    operator: OperatorIdentity,
781    State(state): State<ServerState<S>>,
782    ApiPath(account): ApiPath<AccountId>,
783    ApiJson(request): ApiJson<SetStatusRequest>,
784) -> Result<Json<SetStatusResponse>, ApiError> {
785    // 200 with the blast radius, not 204: the ledger half is one row, but the
786    // snapshot half is however many credentials the account has, and an
787    // operator has no other way to learn which (GL-51).
788    let change = operator
789        .run(
790            "set_account_status",
791            account,
792            state.clock.as_ref(),
793            state.store.set_account_status(account, request.status),
794        )
795        .await?;
796    Ok(Json(SetStatusResponse {
797        republished: change.republished,
798        unreadable: change.unreadable,
799    }))
800}
801
802/// One account's administrative state, for an operator or an application
803/// backend acting for its customer (GL-121).
804///
805/// A read, so it takes no audit receipt: the receipt convention describes
806/// committed *outcomes*, and this commits nothing. It is still operator-only —
807/// an instance credential cannot reach it, because it exposes an account's
808/// funding position.
809///
810/// 404 for an unknown account rather than a zeroed body, so a caller cannot
811/// read "does not exist" as "exists with no funding".
812async fn account<S: Backend>(
813    _operator: OperatorIdentity,
814    State(state): State<ServerState<S>>,
815    ApiPath(account): ApiPath<AccountId>,
816) -> Result<Json<AccountResponse>, ApiError> {
817    let as_of = state.clock.now();
818    let view = state
819        .store
820        .account_view(account)
821        .await?
822        .ok_or_else(|| ApiError::not_found("unknown-account", "no such account"))?;
823    let funding = view.conservation;
824    Ok(Json(AccountResponse {
825        account_id: view.account_id,
826        as_of,
827        status: view.status,
828        capacity_class: view.capacity_class,
829        budget: view.schedule,
830        period_start: view.period_start,
831        balance: funding.balance,
832        outstanding_lease_grants: funding.active_lease_grants,
833        settled_usage: funding.settled_usage,
834        expired_allowance: funding.expired,
835        settlement_loss: funding.settlement_loss,
836        deposited: funding.deposited,
837        overage_recorded: funding.overage_recorded,
838    }))
839}
840
841/// Set or clear an account's periodic allowance (GL-121).
842///
843/// Audited through the receipt convention, so the response reports what the
844/// call actually replaced rather than echoing the request back. A repeat is a
845/// success with equal `previous` and `current`: the caller's intent is
846/// satisfied, and saying so is more useful than a conflict.
847async fn set_budget<S: Backend>(
848    operator: OperatorIdentity,
849    State(state): State<ServerState<S>>,
850    ApiPath(account): ApiPath<AccountId>,
851    ApiJson(request): ApiJson<SetBudgetRequest>,
852) -> Result<Json<SetBudgetResponse>, ApiError> {
853    // The receipt carries what this call actually replaced, so the outcome is
854    // mapped to it rather than re-read. Reporting `request.budget` as
855    // `previous` would echo the caller's own input back as history, and a
856    // separate read afterwards would be a different snapshot — either way the
857    // response would stop describing what committed here.
858    let previous = operator
859        .run(
860            "set_budget_schedule",
861            account,
862            state.clock.as_ref(),
863            async {
864                state
865                    .store
866                    .set_budget_schedule(account, request.budget)
867                    .await
868                    .map(|receipt| {
869                        let replaced = match receipt.before {
870                            AdminState::Budget { schedule } => schedule,
871                            _ => None,
872                        };
873                        tollgate_store::AdminReceipt::new(replaced, receipt.before, receipt.after)
874                    })
875            },
876        )
877        .await?;
878    Ok(Json(SetBudgetResponse {
879        previous,
880        current: request.budget,
881    }))
882}
883
884/// Issue one credential for an account, disclosing its secret exactly once.
885///
886/// **Order matters and is the contract.** The credential is minted, then
887/// stored, and only then returned. A crash between minting and storing loses a
888/// secret nobody has — harmless. The reverse order would hand out a credential
889/// the server has never heard of, and no later reconciliation could repair it,
890/// because what persists is a digest and the secret cannot be derived from it.
891///
892/// A lost *response* is recoverable without reissuing: the caller chose the
893/// `key_id`, so resending the same request answers `409` — the credential
894/// exists, and its secret is gone. The remedy is to revoke and issue a new
895/// one, never to ask for the same secret again.
896async fn issue_key<S: Backend>(
897    operator: OperatorIdentity,
898    State(state): State<ServerState<S>>,
899    ApiPath(account): ApiPath<AccountId>,
900    ApiJson(request): ApiJson<IssueKeyRequest>,
901) -> Result<(StatusCode, Json<IssuedKeyResponse>), ApiError> {
902    let issuer = state.issuer.as_ref().ok_or_else(|| {
903        ApiError::not_implemented(
904            "issuance-unsupported",
905            "this deployment does not issue credentials",
906        )
907    })?;
908    let minted = issuer.mint(request.key_id).map_err(|_| {
909        // 503, not 500: entropy exhaustion is an operational condition the
910        // caller should retry, not a defect in the request.
911        ApiError {
912            status: axum::http::StatusCode::SERVICE_UNAVAILABLE,
913            code: "entropy-unavailable",
914            title: "credential entropy unavailable".into(),
915            generation: None,
916            balance_exhaustion: None,
917            balance_shortfall: None,
918        }
919    })?;
920    // Disclosed exactly as minted: the text is what the digest covers, so
921    // the credential the owner presents is the credential that verifies.
922    // An embedder's issuer that breaks that contract is refused before the
923    // record is stored, so no unpresentable credential ever exists.
924    let secret = std::str::from_utf8(&minted.secret)
925        .ok()
926        .filter(|text| !text.is_empty() && text.bytes().all(|b| b.is_ascii_graphic()))
927        .ok_or_else(|| ApiError {
928            status: axum::http::StatusCode::INTERNAL_SERVER_ERROR,
929            code: "issuer-misconfigured",
930            title: "the configured issuer minted a credential that cannot be presented".into(),
931            generation: None,
932            balance_exhaustion: None,
933            balance_shortfall: None,
934        })?
935        .to_owned();
936
937    let record = KeyRecord {
938        key_id: minted.key_id,
939        account_id: account,
940        principal: minted.principal,
941        digest: minted.digest,
942        not_after: request.not_after,
943    };
944    operator
945        .run(
946            "issue_key",
947            format!("{account}/keys/{}", minted.key_id),
948            state.clock.as_ref(),
949            async {
950                state
951                    .store
952                    .insert_key_within_audited(record, request.max_active_keys, state.clock.now())
953                    .await
954            },
955        )
956        .await?;
957
958    // 201, as `create_account` answers: this made a credential that did not
959    // exist. It is also what distinguishes the successful call from the 409 a
960    // resent request gets, which is the signal a retrying caller reads.
961    Ok((
962        StatusCode::CREATED,
963        Json(IssuedKeyResponse {
964            key_id: minted.key_id,
965            secret,
966            not_after: request.not_after,
967        }),
968    ))
969}
970
971/// One page of an account's credentials — metadata only.
972///
973/// Never the secret, never the digest, and never the principal: the principal
974/// is the digest's leading 128 bits, so a listing carrying it would leak half
975/// of what the verifier compares against.
976async fn list_account_keys<S: Backend>(
977    _operator: OperatorIdentity,
978    State(state): State<ServerState<S>>,
979    ApiPath(account): ApiPath<AccountId>,
980    ApiQuery(page): ApiQuery<KeyQuery>,
981) -> Result<Json<AccountKeysResponse>, ApiError> {
982    let limit = credential_page_limit(page.limit)?;
983    let as_of = state.clock.now();
984    let summaries = state.store.account_keys(account, page.after, limit).await?;
985    // A full page may or may not be the last; a short one certainly is. Saying
986    // "there may be more" costs the caller one extra request and never skips a
987    // credential, which is the safe direction for a cursor.
988    let next_after = (summaries.len() == limit.get())
989        .then(|| summaries.last().map(|summary| summary.key_id))
990        .flatten();
991    Ok(Json(AccountKeysResponse {
992        as_of,
993        keys: summaries
994            .into_iter()
995            .map(|summary| AccountKeyResponse {
996                key_id: summary.key_id,
997                not_after: summary.not_after,
998                revoked_at: summary.revoked_at,
999                live: summary.is_live(as_of),
1000            })
1001            .collect(),
1002        next_after,
1003    }))
1004}
1005
1006/// Revoke one of an account's credentials.
1007///
1008/// **Bound to the account in the path.** A `key_id` belonging to someone else
1009/// is answered 404, not revoked: an application backend administering its own
1010/// customer must not be able to retire another customer's credential by
1011/// guessing or mistyping an id. The check is a read of the account's own
1012/// listing, so it cannot be satisfied by a credential the account does not own.
1013async fn revoke_account_key<S: Backend>(
1014    operator: OperatorIdentity,
1015    State(state): State<ServerState<S>>,
1016    ApiPath((account, key)): ApiPath<(AccountId, KeyId)>,
1017) -> Result<Json<RevokeKeyResponse>, ApiError> {
1018    let owned = match state
1019        .store
1020        .account_keys(account, key.0.checked_sub(1).map(KeyId), one())
1021        .await
1022    {
1023        Ok(keys) => keys.into_iter().any(|summary| summary.key_id == key),
1024        // Keep revocation credential-scoped, including a missing owner.
1025        Err(tollgate_store::KeyError::UnknownAccount) => false,
1026        Err(error) => return Err(error.into()),
1027    };
1028    if !owned {
1029        return Err(ApiError::not_found(
1030            "unknown-credential",
1031            "no such credential for this account",
1032        ));
1033    }
1034    let now = state.clock.now();
1035    let outcome = operator
1036        .run(
1037            "revoke_key",
1038            format!("{account}/keys/{key}"),
1039            state.clock.as_ref(),
1040            state.store.revoke_key_audited(key, now),
1041        )
1042        .await?;
1043    Ok(Json(RevokeKeyResponse {
1044        key_id: key,
1045        retired: matches!(outcome, tollgate_store::Revocation::Retired),
1046    }))
1047}
1048
1049/// Publish the policy for one of an account's credentials (GL-143).
1050///
1051/// The operator names the credential by the handles it already holds; the
1052/// store resolves its principal, which never appears in a request, response
1053/// or audit. The snapshot is bound to the key in the path: an unstated
1054/// `key_id` is filled in. One naming a different credential is left as
1055/// stated, and the store refuses it rather than this handler rewriting it.
1056async fn publish_key_snapshot<S: Backend>(
1057    operator: OperatorIdentity,
1058    State(state): State<ServerState<S>>,
1059    ApiPath((account, key)): ApiPath<(AccountId, KeyId)>,
1060    ApiJson(request): ApiJson<PublishSnapshotRequest>,
1061) -> Result<StatusCode, ApiError> {
1062    let mut snapshot = request.snapshot;
1063    if snapshot.key_id.is_none() {
1064        Arc::make_mut(&mut snapshot).key_id = Some(key);
1065    }
1066    let snapshot = PublishableSnapshot::try_new(snapshot)?;
1067    operator
1068        .run(
1069            "publish_key_snapshot",
1070            format!("{account}/keys/{key}/snapshot"),
1071            state.clock.as_ref(),
1072            state.store.publish_key_snapshot(account, key, snapshot),
1073        )
1074        .await?;
1075    Ok(StatusCode::NO_CONTENT)
1076}
1077
1078/// Withdraw the policy bound to one of an account's credentials (GL-143).
1079///
1080/// Revocation does not withdraw it, so this is part of revoking a key; it is
1081/// allowed for a revoked credential for exactly that reason.
1082async fn remove_key_snapshot<S: Backend>(
1083    operator: OperatorIdentity,
1084    State(state): State<ServerState<S>>,
1085    ApiPath((account, key)): ApiPath<(AccountId, KeyId)>,
1086) -> Result<StatusCode, ApiError> {
1087    operator
1088        .run(
1089            "remove_key_snapshot",
1090            format!("{account}/keys/{key}/snapshot"),
1091            state.clock.as_ref(),
1092            state.store.remove_key_snapshot(account, key),
1093        )
1094        .await?;
1095    Ok(StatusCode::NO_CONTENT)
1096}
1097
1098/// Both credential feeds validate input before any backend read.
1099fn credential_page_limit(limit: Option<usize>) -> Result<std::num::NonZeroUsize, ApiError> {
1100    std::num::NonZeroUsize::new(limit.unwrap_or(DEFAULT_KEY_PAGE_LIMIT.get()))
1101        .filter(|limit| limit.get() <= tollgate_store::MAX_KEY_PAGE_LIMIT)
1102        .ok_or_else(|| ApiError {
1103            status: StatusCode::UNPROCESSABLE_ENTITY,
1104            code: "invalid-limit",
1105            title: format!(
1106                "credential page limit must be between 1 and {}",
1107                tollgate_store::MAX_KEY_PAGE_LIMIT
1108            ),
1109            generation: None,
1110            balance_exhaustion: None,
1111            balance_shortfall: None,
1112        })
1113}
1114
1115fn one() -> std::num::NonZeroUsize {
1116    std::num::NonZeroUsize::new(1).expect("one is nonzero")
1117}
1118
1119async fn set_capacity_class<S: Backend>(
1120    operator: OperatorIdentity,
1121    State(state): State<ServerState<S>>,
1122    ApiPath(account): ApiPath<AccountId>,
1123    ApiJson(request): ApiJson<SetCapacityClassRequest>,
1124) -> Result<Json<SetStatusResponse>, ApiError> {
1125    // 200 with the blast radius, for the reason `set_status` returns one: the
1126    // ledger half is one row and the snapshot half is however many credentials
1127    // the account has (GL-99).
1128    let change = operator
1129        .run(
1130            "set_capacity_class",
1131            account,
1132            state.clock.as_ref(),
1133            state
1134                .store
1135                .set_capacity_class(account, request.capacity_class),
1136        )
1137        .await?;
1138    Ok(Json(SetStatusResponse {
1139        republished: change.republished,
1140        unreadable: change.unreadable,
1141    }))
1142}
1143
1144async fn publish_snapshot<S: Backend>(
1145    operator: OperatorIdentity,
1146    State(state): State<ServerState<S>>,
1147    ApiPath(principal): ApiPath<Principal>,
1148    ApiJson(request): ApiJson<PublishSnapshotRequest>,
1149) -> Result<StatusCode, ApiError> {
1150    let snapshot = PublishableSnapshot::try_new(request.snapshot)?;
1151    operator
1152        .run(
1153            "publish_snapshot",
1154            principal,
1155            state.clock.as_ref(),
1156            state.store.publish_snapshot(principal, snapshot),
1157        )
1158        .await?;
1159    Ok(StatusCode::NO_CONTENT)
1160}
1161
1162async fn remove_snapshot<S: Backend>(
1163    operator: OperatorIdentity,
1164    State(state): State<ServerState<S>>,
1165    ApiPath(principal): ApiPath<Principal>,
1166) -> Result<StatusCode, ApiError> {
1167    operator
1168        .run(
1169            "remove_snapshot",
1170            principal,
1171            state.clock.as_ref(),
1172            state.store.remove_snapshot(principal),
1173        )
1174        .await?;
1175    Ok(StatusCode::NO_CONTENT)
1176}