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