1pub 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
67pub struct ServerState<S> {
70 pub store: Arc<S>,
71 pub clock: Arc<dyn Clock>,
72 pub security: Arc<ServerSecurity>,
73 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
94pub trait Backend:
96 LeaseAllocator
97 + SnapshotSource
98 + UsageSink
99 + AdminStore
100 + KeySource
101 + 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
126pub 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
202pub 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 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 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 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
281async 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 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 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 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
445async 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 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 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
540async 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
594async 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
640async 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 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
781async 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
820async 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 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
863async 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 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 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 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
950async 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 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
985async 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
1024async 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
1053async 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
1073fn 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 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}