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