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    ApiJson(request): ApiJson<IngestRequest>,
674) -> Result<Json<IngestReport>, ApiError> {
675    Ok(Json(
676        state
677            .store
678            .ingest(&request.events, state.clock.now())
679            .await?,
680    ))
681}
682
683#[derive(serde::Deserialize)]
684#[serde(deny_unknown_fields)]
685struct KeyQuery {
686    after: Option<tollgate_core::KeyId>,
687    limit: Option<usize>,
688}
689
690async fn active_keys<S: Backend>(
691    _instance: InstanceIdentity,
692    State(state): State<ServerState<S>>,
693    ApiQuery(query): ApiQuery<KeyQuery>,
694) -> Result<impl axum::response::IntoResponse, ApiError> {
695    let limit = credential_page_limit(query.limit)?;
696    let page = state
697        .store
698        .active_keys_page(state.clock.now(), query.after, limit)
699        .await
700        .and_then(|page| {
701            page.validate_request(query.after, limit)?;
702            Ok(page)
703        })
704        .map_err(|_| ApiError {
705            status: StatusCode::SERVICE_UNAVAILABLE,
706            code: "credential-source-unavailable",
707            title: "credential source is unavailable".into(),
708            generation: None,
709            balance_exhaustion: None,
710            balance_shortfall: None,
711        })?;
712    let body = tollgate_store::wire::KeysResponse {
713        revision: page.revision(),
714        as_of: page.as_of(),
715        next_after: page.next_after(),
716        keys: page.into_records(),
717    };
718    Ok((
719        [(axum::http::header::CACHE_CONTROL, "no-store")],
720        Json(body),
721    ))
722}
723
724async fn create_account<S: Backend>(
725    operator: OperatorIdentity,
726    State(state): State<ServerState<S>>,
727    ApiJson(request): ApiJson<CreateAccountRequest>,
728) -> Result<StatusCode, ApiError> {
729    operator
730        .run(
731            "create_account",
732            request.account_id,
733            state.clock.as_ref(),
734            state.store.create_account(AccountConfig {
735                account_id: request.account_id,
736                initial_balance: request.initial_balance,
737                status: request.status,
738                capacity_class: CapacityClass::Assured,
739            }),
740        )
741        .await?;
742    Ok(StatusCode::CREATED)
743}
744
745async fn deposit<S: Backend>(
746    operator: OperatorIdentity,
747    State(state): State<ServerState<S>>,
748    ApiPath(account): ApiPath<AccountId>,
749    ApiJson(request): ApiJson<DepositRequest>,
750) -> Result<StatusCode, ApiError> {
751    if request.units == CostUnits::ZERO {
752        return Err(ApiError::bad_request(
753            "zero-deposit",
754            "deposit must be positive",
755        ));
756    }
757    operator
758        .run(
759            "deposit",
760            account,
761            state.clock.as_ref(),
762            state.store.deposit(account, request.units),
763        )
764        .await?;
765    Ok(StatusCode::NO_CONTENT)
766}
767
768async fn set_status<S: Backend>(
769    operator: OperatorIdentity,
770    State(state): State<ServerState<S>>,
771    ApiPath(account): ApiPath<AccountId>,
772    ApiJson(request): ApiJson<SetStatusRequest>,
773) -> Result<Json<SetStatusResponse>, ApiError> {
774    // 200 with the blast radius, not 204: the ledger half is one row, but the
775    // snapshot half is however many credentials the account has, and an
776    // operator has no other way to learn which (GL-51).
777    let change = operator
778        .run(
779            "set_account_status",
780            account,
781            state.clock.as_ref(),
782            state.store.set_account_status(account, request.status),
783        )
784        .await?;
785    Ok(Json(SetStatusResponse {
786        republished: change.republished,
787        unreadable: change.unreadable,
788    }))
789}
790
791/// One account's administrative state, for an operator or an application
792/// backend acting for its customer (GL-121).
793///
794/// A read, so it takes no audit receipt: the receipt convention describes
795/// committed *outcomes*, and this commits nothing. It is still operator-only —
796/// an instance credential cannot reach it, because it exposes an account's
797/// funding position.
798///
799/// 404 for an unknown account rather than a zeroed body, so a caller cannot
800/// read "does not exist" as "exists with no funding".
801async fn account<S: Backend>(
802    _operator: OperatorIdentity,
803    State(state): State<ServerState<S>>,
804    ApiPath(account): ApiPath<AccountId>,
805) -> Result<Json<AccountResponse>, ApiError> {
806    let as_of = state.clock.now();
807    let view = state
808        .store
809        .account_view(account)
810        .await?
811        .ok_or_else(|| ApiError::not_found("unknown-account", "no such account"))?;
812    let funding = view.conservation;
813    Ok(Json(AccountResponse {
814        account_id: view.account_id,
815        as_of,
816        status: view.status,
817        capacity_class: view.capacity_class,
818        budget: view.schedule,
819        period_start: view.period_start,
820        balance: funding.balance,
821        outstanding_lease_grants: funding.active_lease_grants,
822        settled_usage: funding.settled_usage,
823        expired_allowance: funding.expired,
824        settlement_loss: funding.settlement_loss,
825        deposited: funding.deposited,
826        overage_recorded: funding.overage_recorded,
827    }))
828}
829
830/// Set or clear an account's periodic allowance (GL-121).
831///
832/// Audited through the receipt convention, so the response reports what the
833/// call actually replaced rather than echoing the request back. A repeat is a
834/// success with equal `previous` and `current`: the caller's intent is
835/// satisfied, and saying so is more useful than a conflict.
836async fn set_budget<S: Backend>(
837    operator: OperatorIdentity,
838    State(state): State<ServerState<S>>,
839    ApiPath(account): ApiPath<AccountId>,
840    ApiJson(request): ApiJson<SetBudgetRequest>,
841) -> Result<Json<SetBudgetResponse>, ApiError> {
842    // The receipt carries what this call actually replaced, so the outcome is
843    // mapped to it rather than re-read. Reporting `request.budget` as
844    // `previous` would echo the caller's own input back as history, and a
845    // separate read afterwards would be a different snapshot — either way the
846    // response would stop describing what committed here.
847    let previous = operator
848        .run(
849            "set_budget_schedule",
850            account,
851            state.clock.as_ref(),
852            async {
853                state
854                    .store
855                    .set_budget_schedule(account, request.budget)
856                    .await
857                    .map(|receipt| {
858                        let replaced = match receipt.before {
859                            AdminState::Budget { schedule } => schedule,
860                            _ => None,
861                        };
862                        tollgate_store::AdminReceipt::new(replaced, receipt.before, receipt.after)
863                    })
864            },
865        )
866        .await?;
867    Ok(Json(SetBudgetResponse {
868        previous,
869        current: request.budget,
870    }))
871}
872
873/// Issue one credential for an account, disclosing its secret exactly once.
874///
875/// **Order matters and is the contract.** The credential is minted, then
876/// stored, and only then returned. A crash between minting and storing loses a
877/// secret nobody has — harmless. The reverse order would hand out a credential
878/// the server has never heard of, and no later reconciliation could repair it,
879/// because what persists is a digest and the secret cannot be derived from it.
880///
881/// A lost *response* is recoverable without reissuing: the caller chose the
882/// `key_id`, so resending the same request answers `409` — the credential
883/// exists, and its secret is gone. The remedy is to revoke and issue a new
884/// one, never to ask for the same secret again.
885async fn issue_key<S: Backend>(
886    operator: OperatorIdentity,
887    State(state): State<ServerState<S>>,
888    ApiPath(account): ApiPath<AccountId>,
889    ApiJson(request): ApiJson<IssueKeyRequest>,
890) -> Result<(StatusCode, Json<IssuedKeyResponse>), ApiError> {
891    let issuer = state.issuer.as_ref().ok_or_else(|| {
892        ApiError::not_implemented(
893            "issuance-unsupported",
894            "this deployment does not issue credentials",
895        )
896    })?;
897    let minted = issuer.mint(request.key_id).map_err(|_| {
898        // 503, not 500: entropy exhaustion is an operational condition the
899        // caller should retry, not a defect in the request.
900        ApiError {
901            status: axum::http::StatusCode::SERVICE_UNAVAILABLE,
902            code: "entropy-unavailable",
903            title: "credential entropy unavailable".into(),
904            generation: None,
905            balance_exhaustion: None,
906            balance_shortfall: None,
907        }
908    })?;
909    // Disclosed exactly as minted: the text is what the digest covers, so
910    // the credential the owner presents is the credential that verifies.
911    // An embedder's issuer that breaks that contract is refused before the
912    // record is stored, so no unpresentable credential ever exists.
913    let secret = std::str::from_utf8(&minted.secret)
914        .ok()
915        .filter(|text| !text.is_empty() && text.bytes().all(|b| b.is_ascii_graphic()))
916        .ok_or_else(|| ApiError {
917            status: axum::http::StatusCode::INTERNAL_SERVER_ERROR,
918            code: "issuer-misconfigured",
919            title: "the configured issuer minted a credential that cannot be presented".into(),
920            generation: None,
921            balance_exhaustion: None,
922            balance_shortfall: None,
923        })?
924        .to_owned();
925
926    let record = KeyRecord {
927        key_id: minted.key_id,
928        account_id: account,
929        principal: minted.principal,
930        digest: minted.digest,
931        not_after: request.not_after,
932    };
933    operator
934        .run(
935            "issue_key",
936            format!("{account}/keys/{}", minted.key_id),
937            state.clock.as_ref(),
938            async {
939                state
940                    .store
941                    .insert_key_within_audited(record, request.max_active_keys, state.clock.now())
942                    .await
943            },
944        )
945        .await?;
946
947    // 201, as `create_account` answers: this made a credential that did not
948    // exist. It is also what distinguishes the successful call from the 409 a
949    // resent request gets, which is the signal a retrying caller reads.
950    Ok((
951        StatusCode::CREATED,
952        Json(IssuedKeyResponse {
953            key_id: minted.key_id,
954            secret,
955            not_after: request.not_after,
956        }),
957    ))
958}
959
960/// One page of an account's credentials — metadata only.
961///
962/// Never the secret, never the digest, and never the principal: the principal
963/// is the digest's leading 128 bits, so a listing carrying it would leak half
964/// of what the verifier compares against.
965async fn list_account_keys<S: Backend>(
966    _operator: OperatorIdentity,
967    State(state): State<ServerState<S>>,
968    ApiPath(account): ApiPath<AccountId>,
969    ApiQuery(page): ApiQuery<KeyQuery>,
970) -> Result<Json<AccountKeysResponse>, ApiError> {
971    let limit = credential_page_limit(page.limit)?;
972    let as_of = state.clock.now();
973    let summaries = state.store.account_keys(account, page.after, limit).await?;
974    // A full page may or may not be the last; a short one certainly is. Saying
975    // "there may be more" costs the caller one extra request and never skips a
976    // credential, which is the safe direction for a cursor.
977    let next_after = (summaries.len() == limit.get())
978        .then(|| summaries.last().map(|summary| summary.key_id))
979        .flatten();
980    Ok(Json(AccountKeysResponse {
981        as_of,
982        keys: summaries
983            .into_iter()
984            .map(|summary| AccountKeyResponse {
985                key_id: summary.key_id,
986                not_after: summary.not_after,
987                revoked_at: summary.revoked_at,
988                live: summary.is_live(as_of),
989            })
990            .collect(),
991        next_after,
992    }))
993}
994
995/// Revoke one of an account's credentials.
996///
997/// **Bound to the account in the path.** A `key_id` belonging to someone else
998/// is answered 404, not revoked: an application backend administering its own
999/// customer must not be able to retire another customer's credential by
1000/// guessing or mistyping an id. The check is a read of the account's own
1001/// listing, so it cannot be satisfied by a credential the account does not own.
1002async fn revoke_account_key<S: Backend>(
1003    operator: OperatorIdentity,
1004    State(state): State<ServerState<S>>,
1005    ApiPath((account, key)): ApiPath<(AccountId, KeyId)>,
1006) -> Result<Json<RevokeKeyResponse>, ApiError> {
1007    let owned = state
1008        .store
1009        .account_keys(account, key.0.checked_sub(1).map(KeyId), one())
1010        .await?
1011        .into_iter()
1012        .any(|summary| summary.key_id == key);
1013    if !owned {
1014        return Err(ApiError::not_found(
1015            "unknown-credential",
1016            "no such credential for this account",
1017        ));
1018    }
1019    let now = state.clock.now();
1020    let outcome = operator
1021        .run(
1022            "revoke_key",
1023            format!("{account}/keys/{key}"),
1024            state.clock.as_ref(),
1025            state.store.revoke_key_audited(key, now),
1026        )
1027        .await?;
1028    Ok(Json(RevokeKeyResponse {
1029        key_id: key,
1030        retired: matches!(outcome, tollgate_store::Revocation::Retired),
1031    }))
1032}
1033
1034/// Publish the policy for one of an account's credentials (GL-143).
1035///
1036/// The operator names the credential by the handles it already holds; the
1037/// store resolves its principal, which never appears in a request, response
1038/// or audit. The snapshot is bound to the key in the path: an unstated
1039/// `key_id` is filled in. One naming a different credential is left as
1040/// stated, and the store refuses it rather than this handler rewriting it.
1041async fn publish_key_snapshot<S: Backend>(
1042    operator: OperatorIdentity,
1043    State(state): State<ServerState<S>>,
1044    ApiPath((account, key)): ApiPath<(AccountId, KeyId)>,
1045    ApiJson(request): ApiJson<PublishSnapshotRequest>,
1046) -> Result<StatusCode, ApiError> {
1047    let mut snapshot = request.snapshot;
1048    if snapshot.key_id.is_none() {
1049        Arc::make_mut(&mut snapshot).key_id = Some(key);
1050    }
1051    let snapshot = PublishableSnapshot::try_new(snapshot)?;
1052    operator
1053        .run(
1054            "publish_key_snapshot",
1055            format!("{account}/keys/{key}/snapshot"),
1056            state.clock.as_ref(),
1057            state.store.publish_key_snapshot(account, key, snapshot),
1058        )
1059        .await?;
1060    Ok(StatusCode::NO_CONTENT)
1061}
1062
1063/// Withdraw the policy bound to one of an account's credentials (GL-143).
1064///
1065/// Revocation does not withdraw it, so this is part of revoking a key; it is
1066/// allowed for a revoked credential for exactly that reason.
1067async fn remove_key_snapshot<S: Backend>(
1068    operator: OperatorIdentity,
1069    State(state): State<ServerState<S>>,
1070    ApiPath((account, key)): ApiPath<(AccountId, KeyId)>,
1071) -> Result<StatusCode, ApiError> {
1072    operator
1073        .run(
1074            "remove_key_snapshot",
1075            format!("{account}/keys/{key}/snapshot"),
1076            state.clock.as_ref(),
1077            state.store.remove_key_snapshot(account, key),
1078        )
1079        .await?;
1080    Ok(StatusCode::NO_CONTENT)
1081}
1082
1083/// Both credential feeds validate input before any backend read.
1084fn credential_page_limit(limit: Option<usize>) -> Result<std::num::NonZeroUsize, ApiError> {
1085    std::num::NonZeroUsize::new(limit.unwrap_or(DEFAULT_KEY_PAGE_LIMIT.get()))
1086        .filter(|limit| limit.get() <= tollgate_store::MAX_KEY_PAGE_LIMIT)
1087        .ok_or_else(|| ApiError {
1088            status: StatusCode::UNPROCESSABLE_ENTITY,
1089            code: "invalid-limit",
1090            title: format!(
1091                "credential page limit must be between 1 and {}",
1092                tollgate_store::MAX_KEY_PAGE_LIMIT
1093            ),
1094            generation: None,
1095            balance_exhaustion: None,
1096            balance_shortfall: None,
1097        })
1098}
1099
1100fn one() -> std::num::NonZeroUsize {
1101    std::num::NonZeroUsize::new(1).expect("one is nonzero")
1102}
1103
1104async fn set_capacity_class<S: Backend>(
1105    operator: OperatorIdentity,
1106    State(state): State<ServerState<S>>,
1107    ApiPath(account): ApiPath<AccountId>,
1108    ApiJson(request): ApiJson<SetCapacityClassRequest>,
1109) -> Result<Json<SetStatusResponse>, ApiError> {
1110    // 200 with the blast radius, for the reason `set_status` returns one: the
1111    // ledger half is one row and the snapshot half is however many credentials
1112    // the account has (GL-99).
1113    let change = operator
1114        .run(
1115            "set_capacity_class",
1116            account,
1117            state.clock.as_ref(),
1118            state
1119                .store
1120                .set_capacity_class(account, request.capacity_class),
1121        )
1122        .await?;
1123    Ok(Json(SetStatusResponse {
1124        republished: change.republished,
1125        unreadable: change.unreadable,
1126    }))
1127}
1128
1129async fn publish_snapshot<S: Backend>(
1130    operator: OperatorIdentity,
1131    State(state): State<ServerState<S>>,
1132    ApiPath(principal): ApiPath<Principal>,
1133    ApiJson(request): ApiJson<PublishSnapshotRequest>,
1134) -> Result<StatusCode, ApiError> {
1135    let snapshot = PublishableSnapshot::try_new(request.snapshot)?;
1136    operator
1137        .run(
1138            "publish_snapshot",
1139            principal,
1140            state.clock.as_ref(),
1141            state.store.publish_snapshot(principal, snapshot),
1142        )
1143        .await?;
1144    Ok(StatusCode::NO_CONTENT)
1145}
1146
1147async fn remove_snapshot<S: Backend>(
1148    operator: OperatorIdentity,
1149    State(state): State<ServerState<S>>,
1150    ApiPath(principal): ApiPath<Principal>,
1151) -> Result<StatusCode, ApiError> {
1152    operator
1153        .run(
1154            "remove_snapshot",
1155            principal,
1156            state.clock.as_ref(),
1157            state.store.remove_snapshot(principal),
1158        )
1159        .await?;
1160    Ok(StatusCode::NO_CONTENT)
1161}