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::{
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
73pub struct ServerState<S> {
76 pub store: Arc<S>,
78 pub clock: Arc<dyn Clock>,
82 pub security: Arc<ServerSecurity>,
87 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
108pub trait Backend:
110 LeaseAllocator
111 + SnapshotSource
112 + UsageSink
113 + AdminStore
114 + KeySource
115 + 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
140pub 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 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
226pub 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 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 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 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
305async 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 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 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 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
469async 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 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 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
564async 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
618async 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
664async 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 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
749async 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 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 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
881async 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 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
935async 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 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 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
1004async 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 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 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 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 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
1102async 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 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
1146async 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 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
1198async 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 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
1263async 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
1293fn 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 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 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}