1#![deny(missing_docs)]
33
34pub mod config;
35pub mod error;
36pub mod google;
37mod maintenance;
38pub mod security;
39pub mod transport;
40
41use std::sync::Arc;
42
43use axum::extract::{DefaultBodyLimit, State};
44use axum::http::StatusCode;
45use axum::routing::{get, post, put};
46use axum::{Json, Router};
47use tracing::Instrument as _;
48
49use tollgate_auth::CredentialIssuer;
50use tollgate_core::{AccountId, CapacityClass, CostUnits, KeyId, Principal, PublishableSnapshot};
51use tollgate_store::wire::{
52 API_PREFIX, AccountKeyResponse, AccountKeysResponse, AccountResponse, AcquireRequest,
53 AcquireResponse, ConsolidateRequest, ConsolidateResponse, CreateAccountRequest, DepositRequest,
54 IngestRequest, IssueKeyRequest, IssuedKeyResponse, MAX_INGEST_BODY_BYTES,
55 MAX_SNAPSHOT_BODY_BYTES, PrincipalsResponse, PublishSnapshotRequest, ReleaseRequest,
56 RevokeKeyResponse, SetBudgetRequest, SetBudgetResponse, SetCapacityClassRequest,
57 SetStatusRequest, SetStatusResponse,
58};
59use tollgate_store::{
60 AccountConfig, AdminState, AdminStore, Clock, DEFAULT_KEY_PAGE_LIMIT,
61 DEFAULT_RECLAIM_BATCH_LIMIT, DEFAULT_ROLLOVER_BATCH_LIMIT, IngestReport, KeyDirectory,
62 KeyRecord, KeySource, LeaseAllocator, ReclaimBatch, ReclaimedLease, SnapshotResolution,
63 SnapshotSource, StoreError, StoreHealth, UsageSink,
64};
65
66use crate::error::{ApiError, ApiJson, ApiPath, ApiQuery};
67use crate::security::{Authorization, InstanceIdentity, OperatorIdentity, Role, ServerSecurity};
68
69pub struct ServerState<S> {
72 pub store: Arc<S>,
74 pub clock: Arc<dyn Clock>,
78 pub security: Arc<ServerSecurity>,
83 pub issuer: Option<Arc<dyn CredentialIssuer + Send + Sync>>,
91}
92
93impl<S> Clone for ServerState<S> {
94 fn clone(&self) -> Self {
95 ServerState {
96 store: Arc::clone(&self.store),
97 clock: Arc::clone(&self.clock),
98 security: Arc::clone(&self.security),
99 issuer: self.issuer.clone(),
100 }
101 }
102}
103
104pub trait Backend:
106 LeaseAllocator
107 + SnapshotSource
108 + UsageSink
109 + AdminStore
110 + KeySource
111 + KeyDirectory
116 + StoreHealth
117 + Send
118 + Sync
119 + 'static
120{
121}
122impl<T> Backend for T where
123 T: LeaseAllocator
124 + SnapshotSource
125 + UsageSink
126 + AdminStore
127 + KeySource
128 + KeyDirectory
129 + StoreHealth
130 + Send
131 + Sync
132 + 'static
133{
134}
135
136pub fn router<S: Backend>(state: ServerState<S>) -> Router {
139 router_with_maintenance(state, maintenance::Monitor::unmanaged())
140}
141
142fn router_with_maintenance<S: Backend>(
143 state: ServerState<S>,
144 maintenance: maintenance::Monitor,
145) -> Router {
146 let authorization = |role| Authorization {
147 security: Arc::clone(&state.security),
148 clock: Arc::clone(&state.clock),
149 role,
150 };
151 let instance = Router::new()
152 .route("/leases/acquire", post(acquire::<S>))
153 .route("/leases/release", post(release::<S>))
154 .route("/leases/consolidate", post(consolidate::<S>))
155 .route("/leases/reclaim", post(reclaim::<S>))
156 .route("/snapshots", get(list_principals::<S>))
157 .route("/keys", get(active_keys::<S>))
158 .route("/snapshots/{principal}", get(fetch_snapshot::<S>))
159 .route(
160 "/usage/ingest",
161 post(ingest::<S>).layer(DefaultBodyLimit::max(MAX_INGEST_BODY_BYTES)),
162 )
163 .route_layer(axum::middleware::from_fn_with_state(
164 authorization(Role::Instance),
165 security::authorize,
166 ));
167 let operator = Router::new()
168 .route("/accounts", post(create_account::<S>))
169 .route("/accounts/{account}", get(account::<S>))
170 .route("/accounts/{account}/budget", put(set_budget::<S>))
171 .route(
172 "/accounts/{account}/keys",
173 post(issue_key::<S>).get(list_account_keys::<S>),
174 )
175 .route(
176 "/accounts/{account}/keys/{key}",
177 axum::routing::delete(revoke_account_key::<S>),
178 )
179 .route(
180 "/accounts/{account}/keys/{key}/snapshot",
181 put(publish_key_snapshot::<S>)
182 .delete(remove_key_snapshot::<S>)
183 .layer(DefaultBodyLimit::max(MAX_SNAPSHOT_BODY_BYTES)),
184 )
185 .route("/accounts/{account}/deposit", post(deposit::<S>))
186 .route("/accounts/{account}/status", post(set_status::<S>))
187 .route(
188 "/accounts/{account}/capacity-class",
189 post(set_capacity_class::<S>),
190 )
191 .route(
192 "/snapshots/{principal}",
193 put(publish_snapshot::<S>)
194 .delete(remove_snapshot::<S>)
195 .layer(DefaultBodyLimit::max(MAX_SNAPSHOT_BODY_BYTES)),
196 )
197 .route_layer(axum::middleware::from_fn_with_state(
198 authorization(Role::Operator),
199 security::authorize,
200 ));
201 Router::new()
202 .route("/livez", get(async || StatusCode::OK))
203 .route(
204 "/readyz",
205 get(readyz::<S>).layer(axum::Extension(maintenance)),
206 )
207 .nest(API_PREFIX, instance.nest("/admin", operator))
208 .with_state(state)
209 .layer(axum::middleware::from_fn(error::report_http_failure))
210}
211
212pub async fn serve<S: Backend>(
215 listener: tokio::net::TcpListener,
216 state: ServerState<S>,
217 reclaim_interval: std::time::Duration,
218 shutdown: impl Future<Output = ()> + Send + 'static,
219) -> std::io::Result<()> {
220 if reclaim_interval.is_zero()
221 || tokio::time::Instant::now()
222 .checked_add(reclaim_interval)
223 .is_none()
224 {
225 return Err(std::io::Error::new(
226 std::io::ErrorKind::InvalidInput,
227 "reclaim_interval must be positive and fit the monotonic clock",
228 ));
229 }
230 let listener = transport::SecureListener::new(listener, Arc::clone(&state.security))?;
231 let sweep_store = Arc::clone(&state.store);
232 let sweep_clock = Arc::clone(&state.clock);
233 let mut sweeper = maintenance::Task::spawn(move |health| {
234 maintenance_sweep(sweep_store, sweep_clock, reclaim_interval, health).instrument(
235 tracing::info_span!(
236 "maintenance_sweep",
237 interval_ms = reclaim_interval.as_millis()
238 ),
239 )
240 });
241 let (signal, signalled) = tokio::sync::oneshot::channel::<()>();
245 let mut signal = Some(signal);
246 let server = axum::serve(
247 listener,
248 router_with_maintenance(state, sweeper.monitor.clone())
249 .into_make_service_with_connect_info::<transport::PeerIdentity>(),
250 )
251 .with_graceful_shutdown(async move {
252 let Ok(()) = signalled.await else { return };
254 });
255 let server = std::future::IntoFuture::into_future(server);
256 tokio::pin!(server, shutdown);
257 loop {
258 tokio::select! {
259 biased;
260 result = &mut sweeper.join => {
261 if result.as_ref().is_err_and(|error| error.is_cancelled())
262 && sweeper.stop.requested()
263 {
264 return server.await;
267 }
268 sweeper.stop.stop();
269 drop(signal.take());
270 let reason = match result {
271 Err(error) if error.is_panic() => "panic",
272 Err(_) => "cancelled",
273 Ok(()) => "unexpected-return",
274 };
275 tracing::error!(
276 operation = "maintenance",
277 reason,
278 "server maintenance task stopped; stopping the listener"
279 );
280 return Err(std::io::Error::other("server maintenance task stopped unexpectedly"));
281 }
282 result = &mut server => return result,
283 () = &mut shutdown, if signal.is_some() => {
284 sweeper.stop.stop();
285 drop(signal.take());
286 }
287 }
288 }
289}
290
291async fn maintenance_sweep<S: Backend>(
311 store: Arc<S>,
312 clock: Arc<dyn Clock>,
313 interval: std::time::Duration,
314 mut health: maintenance::Publisher,
315) {
316 #[derive(Default)]
317 struct Progress {
318 leases: u128,
319 units: u128,
320 batches: u64,
321 }
322
323 impl Progress {
324 fn record(&mut self, batch: &ReclaimBatch) -> Result<(), StoreError> {
325 let batch_leases =
326 u128::from(u64::try_from(batch.len()).map_err(|_| {
327 StoreError("reclaim batch lease count exceeds u64 range".into())
328 })?);
329 let batch_units = batch.reclaimed().iter().try_fold(0u128, |total, lease| {
330 total
331 .checked_add(u128::from(lease.forfeited.get()))
332 .ok_or_else(|| StoreError("reclaim batch unit total overflow".into()))
333 })?;
334 let leases = self
335 .leases
336 .checked_add(batch_leases)
337 .ok_or_else(|| StoreError("reclaim cycle lease count overflow".into()))?;
338 let units = self
339 .units
340 .checked_add(batch_units)
341 .ok_or_else(|| StoreError("reclaim cycle unit total overflow".into()))?;
342 let batches = self
343 .batches
344 .checked_add(1)
345 .ok_or_else(|| StoreError("reclaim cycle batch count overflow".into()))?;
346 *self = Progress {
347 leases,
348 units,
349 batches,
350 };
351 Ok(())
352 }
353 }
354
355 let mut tick = tokio::time::interval(interval);
356 tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
357 loop {
358 tick.tick().await;
359 let now = clock.now();
363 let mut progress = Progress::default();
364 let outcome = loop {
365 match store
366 .reclaim_expired_batch(now, DEFAULT_RECLAIM_BATCH_LIMIT)
367 .await
368 {
369 Ok(batch) => {
370 if progress.record(&batch).is_err() {
371 tracing::error!(
372 operation = "reclaim",
373 "reclaim progress counter overflowed; stopping maintenance"
374 );
375 return;
376 }
377 if !batch.is_saturated() {
378 break Ok(());
379 }
380 tokio::task::yield_now().await;
384 }
385 Err(error) => break Err(error),
386 }
387 };
388
389 let Ok(report) = health.record(maintenance::Operation::Reclaim, outcome.is_ok()) else {
390 tracing::error!(
391 operation = "reclaim",
392 "maintenance failure counter exhausted"
393 );
394 return;
395 };
396 match outcome {
397 Ok(()) => {
398 if report.recovered_after > 0 {
399 tracing::info!(
400 operation = "reclaim",
401 after_failures = report.recovered_after,
402 "reclaim sweep recovered"
403 );
404 }
405 if progress.units > 0 {
410 tracing::warn!(
411 leases = %progress.leases,
412 forfeited_units = %progress.units,
413 batches = progress.batches,
414 "swept unreleased leases; their remainders are forfeited as provisional loss"
415 );
416 } else if progress.leases > 0 {
417 tracing::info!(
418 leases = %progress.leases,
419 forfeited_units = %progress.units,
420 batches = progress.batches,
421 "swept unreleased leases with nothing left to forfeit"
422 );
423 }
424 }
425 Err(_) => {
426 macro_rules! report_failure {
427 ($level:ident) => {
428 tracing::$level!(
429 operation = "reclaim",
430 code = "storage",
431 consecutive_failures = report.failures,
432 reclaimed_leases = %progress.leases,
433 forfeited_units = %progress.units,
434 completed_batches = progress.batches,
435 "reclaim sweep failed; completed batches stay committed and remaining expired leases stay stranded until it recovers"
436 );
437 };
438 }
439 if report.failures >= maintenance::PERSISTENT_FAILURES {
440 report_failure!(error);
441 } else {
442 report_failure!(warn);
443 }
444 }
445 }
446 if roll_due_periods(store.as_ref(), now, &mut health)
447 .await
448 .is_break()
449 {
450 return;
451 }
452 }
453}
454
455async fn roll_due_periods<S: Backend>(
469 store: &S,
470 now: jiff::Timestamp,
471 health: &mut maintenance::Publisher,
472) -> std::ops::ControlFlow<()> {
473 let mut accounts: u64 = 0;
474 let mut batches: u64 = 0;
475 loop {
476 match store
477 .roll_due_periods(now, DEFAULT_ROLLOVER_BATCH_LIMIT)
478 .await
479 {
480 Ok(batch) => {
481 let Some(total) = u64::try_from(batch.len())
482 .ok()
483 .and_then(|rolled| accounts.checked_add(rolled))
484 .zip(batches.checked_add(1))
485 else {
486 tracing::error!(
491 accounts,
492 batches,
493 "budget rollover progress counter overflowed; stopping maintenance"
494 );
495 return std::ops::ControlFlow::Break(());
496 };
497 (accounts, batches) = (total.0, total.1);
498 if !batch.is_saturated() {
499 break;
500 }
501 tokio::task::yield_now().await;
504 }
505 Err(_) => {
506 let Ok(report) = health.record(maintenance::Operation::Rollover, false) else {
507 tracing::error!(
508 operation = "budget-rollover",
509 "maintenance failure counter exhausted"
510 );
511 return std::ops::ControlFlow::Break(());
512 };
513 macro_rules! report_failure {
514 ($level:ident) => {
515 tracing::$level!(
516 operation = "budget-rollover",
517 code = "storage",
518 consecutive_failures = report.failures,
519 rolled_accounts = accounts,
520 completed_batches = batches,
521 "budget rollover failed; accounts past their boundary keep last period's allowance until it recovers"
522 );
523 };
524 }
525 if report.failures >= maintenance::PERSISTENT_FAILURES {
526 report_failure!(error);
527 } else {
528 report_failure!(warn);
529 }
530 return std::ops::ControlFlow::Continue(());
531 }
532 }
533 }
534 let report = health
535 .record(maintenance::Operation::Rollover, true)
536 .expect("a successful pass clears its failure counter without arithmetic");
537 if report.recovered_after > 0 {
538 tracing::info!(
539 operation = "budget-rollover",
540 after_failures = report.recovered_after,
541 "budget rollover recovered"
542 );
543 }
544 if accounts > 0 {
545 tracing::info!(accounts, batches, "rolled budget periods");
546 }
547 std::ops::ControlFlow::Continue(())
548}
549
550async fn readyz<S: Backend>(
553 State(state): State<ServerState<S>>,
554 axum::Extension(maintenance): axum::Extension<maintenance::Monitor>,
555) -> StatusCode {
556 match state.store.ping().await {
557 Ok(()) if maintenance.healthy() => StatusCode::OK,
558 Ok(()) => StatusCode::SERVICE_UNAVAILABLE,
559 Err(_) => {
560 tracing::warn!(
561 operation = "readiness",
562 code = "storage",
563 "readiness probe failed: the store did not answer"
564 );
565 StatusCode::SERVICE_UNAVAILABLE
566 }
567 }
568}
569
570async fn acquire<S: Backend>(
571 _identity: InstanceIdentity,
572 State(state): State<ServerState<S>>,
573 ApiJson(request): ApiJson<AcquireRequest>,
574) -> Result<Json<AcquireResponse>, ApiError> {
575 let grant = state
576 .store
577 .acquire(
578 request.account_id,
579 request.requested,
580 request.ttl.duration()?,
581 state.clock.now(),
582 )
583 .await?;
584 Ok(Json(grant))
585}
586
587async fn release<S: Backend>(
588 _identity: InstanceIdentity,
589 State(state): State<ServerState<S>>,
590 ApiJson(request): ApiJson<ReleaseRequest>,
591) -> Result<StatusCode, ApiError> {
592 state
593 .store
594 .release(
595 request.lease_id,
596 request.fencing_token,
597 request.unspent,
598 state.clock.now(),
599 )
600 .await?;
601 Ok(StatusCode::NO_CONTENT)
602}
603
604async fn consolidate<S: Backend>(
612 _identity: InstanceIdentity,
613 State(state): State<ServerState<S>>,
614 ApiJson(request): ApiJson<ConsolidateRequest>,
615) -> Result<Json<ConsolidateResponse>, ApiError> {
616 let grant = state
617 .store
618 .consolidate(
619 request.lease_id,
620 request.fencing_token,
621 request.unspent,
622 request.requested,
623 request.needed,
624 request.ttl.duration()?,
625 state.clock.now(),
626 )
627 .await?;
628 Ok(Json(grant))
629}
630
631async fn reclaim<S: Backend>(
632 _identity: InstanceIdentity,
633 State(state): State<ServerState<S>>,
634) -> Result<Json<Vec<ReclaimedLease>>, ApiError> {
635 Ok(Json(state.store.reclaim_expired(state.clock.now()).await?))
636}
637
638async fn fetch_snapshot<S: Backend>(
639 _identity: InstanceIdentity,
640 State(state): State<ServerState<S>>,
641 ApiPath(principal): ApiPath<Principal>,
642) -> Result<Json<Arc<tollgate_core::AccountSnapshot>>, ApiError> {
643 match state.store.snapshot(principal).await? {
644 SnapshotResolution::Present(snapshot) => Ok(Json(snapshot.into_inner())),
645 SnapshotResolution::Revoked { generation } => Err(ApiError::revoked(generation)),
646 SnapshotResolution::Unknown => Err(ApiError::not_found("unknown-principal", "no snapshot")),
647 }
648}
649
650async fn list_principals<S: Backend>(
658 _identity: InstanceIdentity,
659 State(state): State<ServerState<S>>,
660) -> Result<Json<PrincipalsResponse>, ApiError> {
661 match state.store.principals().await? {
662 Some(principals) => Ok(Json(PrincipalsResponse { principals })),
663 None => Err(ApiError::not_implemented(
664 "enumeration-unsupported",
665 "this backend cannot list principals",
666 )),
667 }
668}
669
670async fn ingest<S: Backend>(
671 _identity: InstanceIdentity,
672 State(state): State<ServerState<S>>,
673 ApiJson(request): ApiJson<IngestRequest>,
674) -> Result<Json<IngestReport>, ApiError> {
675 Ok(Json(
676 state
677 .store
678 .ingest(&request.events, state.clock.now())
679 .await?,
680 ))
681}
682
683#[derive(serde::Deserialize)]
684#[serde(deny_unknown_fields)]
685struct KeyQuery {
686 after: Option<tollgate_core::KeyId>,
687 limit: Option<usize>,
688}
689
690async fn active_keys<S: Backend>(
691 _instance: InstanceIdentity,
692 State(state): State<ServerState<S>>,
693 ApiQuery(query): ApiQuery<KeyQuery>,
694) -> Result<impl axum::response::IntoResponse, ApiError> {
695 let limit = credential_page_limit(query.limit)?;
696 let page = state
697 .store
698 .active_keys_page(state.clock.now(), query.after, limit)
699 .await
700 .and_then(|page| {
701 page.validate_request(query.after, limit)?;
702 Ok(page)
703 })
704 .map_err(|_| ApiError {
705 status: StatusCode::SERVICE_UNAVAILABLE,
706 code: "credential-source-unavailable",
707 title: "credential source is unavailable".into(),
708 generation: None,
709 balance_exhaustion: None,
710 balance_shortfall: None,
711 })?;
712 let body = tollgate_store::wire::KeysResponse {
713 revision: page.revision(),
714 as_of: page.as_of(),
715 next_after: page.next_after(),
716 keys: page.into_records(),
717 };
718 Ok((
719 [(axum::http::header::CACHE_CONTROL, "no-store")],
720 Json(body),
721 ))
722}
723
724async fn create_account<S: Backend>(
725 operator: OperatorIdentity,
726 State(state): State<ServerState<S>>,
727 ApiJson(request): ApiJson<CreateAccountRequest>,
728) -> Result<StatusCode, ApiError> {
729 operator
730 .run(
731 "create_account",
732 request.account_id,
733 state.clock.as_ref(),
734 state.store.create_account(AccountConfig {
735 account_id: request.account_id,
736 initial_balance: request.initial_balance,
737 status: request.status,
738 capacity_class: CapacityClass::Assured,
739 }),
740 )
741 .await?;
742 Ok(StatusCode::CREATED)
743}
744
745async fn deposit<S: Backend>(
746 operator: OperatorIdentity,
747 State(state): State<ServerState<S>>,
748 ApiPath(account): ApiPath<AccountId>,
749 ApiJson(request): ApiJson<DepositRequest>,
750) -> Result<StatusCode, ApiError> {
751 if request.units == CostUnits::ZERO {
752 return Err(ApiError::bad_request(
753 "zero-deposit",
754 "deposit must be positive",
755 ));
756 }
757 operator
758 .run(
759 "deposit",
760 account,
761 state.clock.as_ref(),
762 state.store.deposit(account, request.units),
763 )
764 .await?;
765 Ok(StatusCode::NO_CONTENT)
766}
767
768async fn set_status<S: Backend>(
769 operator: OperatorIdentity,
770 State(state): State<ServerState<S>>,
771 ApiPath(account): ApiPath<AccountId>,
772 ApiJson(request): ApiJson<SetStatusRequest>,
773) -> Result<Json<SetStatusResponse>, ApiError> {
774 let change = operator
778 .run(
779 "set_account_status",
780 account,
781 state.clock.as_ref(),
782 state.store.set_account_status(account, request.status),
783 )
784 .await?;
785 Ok(Json(SetStatusResponse {
786 republished: change.republished,
787 unreadable: change.unreadable,
788 }))
789}
790
791async fn account<S: Backend>(
802 _operator: OperatorIdentity,
803 State(state): State<ServerState<S>>,
804 ApiPath(account): ApiPath<AccountId>,
805) -> Result<Json<AccountResponse>, ApiError> {
806 let as_of = state.clock.now();
807 let view = state
808 .store
809 .account_view(account)
810 .await?
811 .ok_or_else(|| ApiError::not_found("unknown-account", "no such account"))?;
812 let funding = view.conservation;
813 Ok(Json(AccountResponse {
814 account_id: view.account_id,
815 as_of,
816 status: view.status,
817 capacity_class: view.capacity_class,
818 budget: view.schedule,
819 period_start: view.period_start,
820 balance: funding.balance,
821 outstanding_lease_grants: funding.active_lease_grants,
822 settled_usage: funding.settled_usage,
823 expired_allowance: funding.expired,
824 settlement_loss: funding.settlement_loss,
825 deposited: funding.deposited,
826 overage_recorded: funding.overage_recorded,
827 }))
828}
829
830async fn set_budget<S: Backend>(
837 operator: OperatorIdentity,
838 State(state): State<ServerState<S>>,
839 ApiPath(account): ApiPath<AccountId>,
840 ApiJson(request): ApiJson<SetBudgetRequest>,
841) -> Result<Json<SetBudgetResponse>, ApiError> {
842 let previous = operator
848 .run(
849 "set_budget_schedule",
850 account,
851 state.clock.as_ref(),
852 async {
853 state
854 .store
855 .set_budget_schedule(account, request.budget)
856 .await
857 .map(|receipt| {
858 let replaced = match receipt.before {
859 AdminState::Budget { schedule } => schedule,
860 _ => None,
861 };
862 tollgate_store::AdminReceipt::new(replaced, receipt.before, receipt.after)
863 })
864 },
865 )
866 .await?;
867 Ok(Json(SetBudgetResponse {
868 previous,
869 current: request.budget,
870 }))
871}
872
873async fn issue_key<S: Backend>(
886 operator: OperatorIdentity,
887 State(state): State<ServerState<S>>,
888 ApiPath(account): ApiPath<AccountId>,
889 ApiJson(request): ApiJson<IssueKeyRequest>,
890) -> Result<(StatusCode, Json<IssuedKeyResponse>), ApiError> {
891 let issuer = state.issuer.as_ref().ok_or_else(|| {
892 ApiError::not_implemented(
893 "issuance-unsupported",
894 "this deployment does not issue credentials",
895 )
896 })?;
897 let minted = issuer.mint(request.key_id).map_err(|_| {
898 ApiError {
901 status: axum::http::StatusCode::SERVICE_UNAVAILABLE,
902 code: "entropy-unavailable",
903 title: "credential entropy unavailable".into(),
904 generation: None,
905 balance_exhaustion: None,
906 balance_shortfall: None,
907 }
908 })?;
909 let secret = std::str::from_utf8(&minted.secret)
914 .ok()
915 .filter(|text| !text.is_empty() && text.bytes().all(|b| b.is_ascii_graphic()))
916 .ok_or_else(|| ApiError {
917 status: axum::http::StatusCode::INTERNAL_SERVER_ERROR,
918 code: "issuer-misconfigured",
919 title: "the configured issuer minted a credential that cannot be presented".into(),
920 generation: None,
921 balance_exhaustion: None,
922 balance_shortfall: None,
923 })?
924 .to_owned();
925
926 let record = KeyRecord {
927 key_id: minted.key_id,
928 account_id: account,
929 principal: minted.principal,
930 digest: minted.digest,
931 not_after: request.not_after,
932 };
933 operator
934 .run(
935 "issue_key",
936 format!("{account}/keys/{}", minted.key_id),
937 state.clock.as_ref(),
938 async {
939 state
940 .store
941 .insert_key_within_audited(record, request.max_active_keys, state.clock.now())
942 .await
943 },
944 )
945 .await?;
946
947 Ok((
951 StatusCode::CREATED,
952 Json(IssuedKeyResponse {
953 key_id: minted.key_id,
954 secret,
955 not_after: request.not_after,
956 }),
957 ))
958}
959
960async fn list_account_keys<S: Backend>(
966 _operator: OperatorIdentity,
967 State(state): State<ServerState<S>>,
968 ApiPath(account): ApiPath<AccountId>,
969 ApiQuery(page): ApiQuery<KeyQuery>,
970) -> Result<Json<AccountKeysResponse>, ApiError> {
971 let limit = credential_page_limit(page.limit)?;
972 let as_of = state.clock.now();
973 let summaries = state.store.account_keys(account, page.after, limit).await?;
974 let next_after = (summaries.len() == limit.get())
978 .then(|| summaries.last().map(|summary| summary.key_id))
979 .flatten();
980 Ok(Json(AccountKeysResponse {
981 as_of,
982 keys: summaries
983 .into_iter()
984 .map(|summary| AccountKeyResponse {
985 key_id: summary.key_id,
986 not_after: summary.not_after,
987 revoked_at: summary.revoked_at,
988 live: summary.is_live(as_of),
989 })
990 .collect(),
991 next_after,
992 }))
993}
994
995async fn revoke_account_key<S: Backend>(
1003 operator: OperatorIdentity,
1004 State(state): State<ServerState<S>>,
1005 ApiPath((account, key)): ApiPath<(AccountId, KeyId)>,
1006) -> Result<Json<RevokeKeyResponse>, ApiError> {
1007 let owned = state
1008 .store
1009 .account_keys(account, key.0.checked_sub(1).map(KeyId), one())
1010 .await?
1011 .into_iter()
1012 .any(|summary| summary.key_id == key);
1013 if !owned {
1014 return Err(ApiError::not_found(
1015 "unknown-credential",
1016 "no such credential for this account",
1017 ));
1018 }
1019 let now = state.clock.now();
1020 let outcome = operator
1021 .run(
1022 "revoke_key",
1023 format!("{account}/keys/{key}"),
1024 state.clock.as_ref(),
1025 state.store.revoke_key_audited(key, now),
1026 )
1027 .await?;
1028 Ok(Json(RevokeKeyResponse {
1029 key_id: key,
1030 retired: matches!(outcome, tollgate_store::Revocation::Retired),
1031 }))
1032}
1033
1034async fn publish_key_snapshot<S: Backend>(
1042 operator: OperatorIdentity,
1043 State(state): State<ServerState<S>>,
1044 ApiPath((account, key)): ApiPath<(AccountId, KeyId)>,
1045 ApiJson(request): ApiJson<PublishSnapshotRequest>,
1046) -> Result<StatusCode, ApiError> {
1047 let mut snapshot = request.snapshot;
1048 if snapshot.key_id.is_none() {
1049 Arc::make_mut(&mut snapshot).key_id = Some(key);
1050 }
1051 let snapshot = PublishableSnapshot::try_new(snapshot)?;
1052 operator
1053 .run(
1054 "publish_key_snapshot",
1055 format!("{account}/keys/{key}/snapshot"),
1056 state.clock.as_ref(),
1057 state.store.publish_key_snapshot(account, key, snapshot),
1058 )
1059 .await?;
1060 Ok(StatusCode::NO_CONTENT)
1061}
1062
1063async fn remove_key_snapshot<S: Backend>(
1068 operator: OperatorIdentity,
1069 State(state): State<ServerState<S>>,
1070 ApiPath((account, key)): ApiPath<(AccountId, KeyId)>,
1071) -> Result<StatusCode, ApiError> {
1072 operator
1073 .run(
1074 "remove_key_snapshot",
1075 format!("{account}/keys/{key}/snapshot"),
1076 state.clock.as_ref(),
1077 state.store.remove_key_snapshot(account, key),
1078 )
1079 .await?;
1080 Ok(StatusCode::NO_CONTENT)
1081}
1082
1083fn credential_page_limit(limit: Option<usize>) -> Result<std::num::NonZeroUsize, ApiError> {
1085 std::num::NonZeroUsize::new(limit.unwrap_or(DEFAULT_KEY_PAGE_LIMIT.get()))
1086 .filter(|limit| limit.get() <= tollgate_store::MAX_KEY_PAGE_LIMIT)
1087 .ok_or_else(|| ApiError {
1088 status: StatusCode::UNPROCESSABLE_ENTITY,
1089 code: "invalid-limit",
1090 title: format!(
1091 "credential page limit must be between 1 and {}",
1092 tollgate_store::MAX_KEY_PAGE_LIMIT
1093 ),
1094 generation: None,
1095 balance_exhaustion: None,
1096 balance_shortfall: None,
1097 })
1098}
1099
1100fn one() -> std::num::NonZeroUsize {
1101 std::num::NonZeroUsize::new(1).expect("one is nonzero")
1102}
1103
1104async fn set_capacity_class<S: Backend>(
1105 operator: OperatorIdentity,
1106 State(state): State<ServerState<S>>,
1107 ApiPath(account): ApiPath<AccountId>,
1108 ApiJson(request): ApiJson<SetCapacityClassRequest>,
1109) -> Result<Json<SetStatusResponse>, ApiError> {
1110 let change = operator
1114 .run(
1115 "set_capacity_class",
1116 account,
1117 state.clock.as_ref(),
1118 state
1119 .store
1120 .set_capacity_class(account, request.capacity_class),
1121 )
1122 .await?;
1123 Ok(Json(SetStatusResponse {
1124 republished: change.republished,
1125 unreadable: change.unreadable,
1126 }))
1127}
1128
1129async fn publish_snapshot<S: Backend>(
1130 operator: OperatorIdentity,
1131 State(state): State<ServerState<S>>,
1132 ApiPath(principal): ApiPath<Principal>,
1133 ApiJson(request): ApiJson<PublishSnapshotRequest>,
1134) -> Result<StatusCode, ApiError> {
1135 let snapshot = PublishableSnapshot::try_new(request.snapshot)?;
1136 operator
1137 .run(
1138 "publish_snapshot",
1139 principal,
1140 state.clock.as_ref(),
1141 state.store.publish_snapshot(principal, snapshot),
1142 )
1143 .await?;
1144 Ok(StatusCode::NO_CONTENT)
1145}
1146
1147async fn remove_snapshot<S: Backend>(
1148 operator: OperatorIdentity,
1149 State(state): State<ServerState<S>>,
1150 ApiPath(principal): ApiPath<Principal>,
1151) -> Result<StatusCode, ApiError> {
1152 operator
1153 .run(
1154 "remove_snapshot",
1155 principal,
1156 state.clock.as_ref(),
1157 state.store.remove_snapshot(principal),
1158 )
1159 .await?;
1160 Ok(StatusCode::NO_CONTENT)
1161}