use std::sync::Arc;
use std::time::SystemTime;
use axum::Json;
use axum::body::Bytes;
use axum::extract::rejection::{BytesRejection, QueryRejection};
use axum::extract::{Path, Query, State};
use axum::http::{HeaderMap, StatusCode};
use axum::routing::{MethodRouter, get, post};
use serde::Deserialize;
use super::auth::{AdminAction, AdminIdentity};
use super::catalogue::{CatalogueFilters, CatalogueRequest, CatalogueView};
use super::conditional::Conditional;
use super::error::AdminError;
use super::protocol::{AuditSummary, MutationPreconditions, MutationRequest};
use super::reads::{
AuditPage, AvailabilityResult, ConvergenceResult, HistoryLimit, HistoryRequest, RevisionPage,
StateView,
};
use super::resources::{AdminResourceRequest, MutationEnvelope, RollbackRequest};
use super::router::{ADMIN_MAX_REQUEST_BYTES, AdminApi};
use super::service::{AvailabilityAuthority, MutationOutcome};
use crate::desired_state::{
ModelLifecycle, MutationKind, OfferingId, ProjectId, ResourceScope, RevisionId, Surface,
TenantId, WireFamily,
};
pub(super) fn publish_route<R: AdminResourceRequest>() -> MethodRouter<Arc<AdminApi>> {
post(publish::<R>)
}
pub(super) fn rollback_route() -> MethodRouter<Arc<AdminApi>> {
post(rollback)
}
pub(super) fn state_route() -> MethodRouter<Arc<AdminApi>> {
get(state)
}
pub(super) fn catalogue_route() -> MethodRouter<Arc<AdminApi>> {
get(catalogue)
}
pub(super) fn history_route() -> MethodRouter<Arc<AdminApi>> {
get(history)
}
pub(super) fn audit_route() -> MethodRouter<Arc<AdminApi>> {
get(audit)
}
pub(super) fn convergence_route() -> MethodRouter<Arc<AdminApi>> {
get(convergence)
}
pub(super) fn availability_route() -> MethodRouter<Arc<AdminApi>> {
get(availability)
}
fn document(
schema: &'static str,
body: Result<Bytes, BytesRejection>,
) -> Result<Bytes, AdminError> {
body.map_err(|rejection| {
if rejection.status() == StatusCode::PAYLOAD_TOO_LARGE {
AdminError::RequestTooLarge {
limit: ADMIN_MAX_REQUEST_BYTES,
}
} else {
AdminError::RequestInvalid {
schema,
detail: rejection.body_text(),
}
}
})
}
async fn publish<R: AdminResourceRequest>(
State(api): State<Arc<AdminApi>>,
identity: AdminIdentity,
preconditions: MutationPreconditions,
body: Result<Bytes, BytesRejection>,
) -> Result<Json<MutationOutcome>, AdminError> {
let body = document(R::SCHEMA, body)?;
let envelope: MutationEnvelope<R> =
serde_json::from_slice(&body).map_err(|error| AdminError::RequestInvalid {
schema: R::SCHEMA,
detail: error.to_string(),
})?;
let summary = AuditSummary::parse(&envelope.summary)?;
let kind = envelope.mutation.kind();
let plan = envelope.resource.plan()?;
if kind == MutationKind::Delete && !plan.retires {
return Err(AdminError::RequestInvalid {
schema: R::SCHEMA,
detail: "`mutation: \"delete\"` requires a document that retires the resource: this \
surface removes nothing, so state the terminal lifecycle the resource \
supports (a tenant `deleted`, a credential `revoked`, an enablement or alias \
`disabled`) — or record the change as an update"
.to_owned(),
});
}
let grant = api
.authorize(&identity, AdminAction::Publish, R::SURFACE, &plan.scope)
.await?;
let request = MutationRequest {
preconditions,
kind,
surface: R::SURFACE,
scope: plan.scope.clone(),
summary,
};
let outcome = api
.service
.apply(&grant, &request, plan.edit.as_ref())
.await?;
Ok(Json(outcome))
}
async fn rollback(
State(api): State<Arc<AdminApi>>,
identity: AdminIdentity,
preconditions: MutationPreconditions,
body: Result<Bytes, BytesRejection>,
) -> Result<Json<MutationOutcome>, AdminError> {
let body = document("rollback", body)?;
let request: RollbackRequest =
serde_json::from_slice(&body).map_err(|error| AdminError::RequestInvalid {
schema: "rollback",
detail: error.to_string(),
})?;
let summary = AuditSummary::parse(&request.summary)?;
let target =
RevisionId::parse(&request.revision).map_err(|error| AdminError::RequestInvalid {
schema: "rollback",
detail: format!("`revision`: {error}"),
})?;
let scope = scope_of(
"rollback",
request.tenant.as_deref(),
request.project.as_deref(),
)?;
let grant = api
.authorize(
&identity,
AdminAction::Rollback,
Surface::AuditTrail,
&scope,
)
.await?;
let mutation = MutationRequest {
preconditions,
kind: MutationKind::Rollback,
surface: Surface::AuditTrail,
scope,
summary,
};
let outcome = api.service.rollback(&grant, &mutation, target).await?;
Ok(Json(outcome))
}
async fn state(
State(api): State<Arc<AdminApi>>,
identity: AdminIdentity,
headers: HeaderMap,
) -> Result<Conditional<StateView>, AdminError> {
let grant = api
.authorize(
&identity,
AdminAction::ReadState,
Surface::AuditTrail,
&ResourceScope::Deployment,
)
.await?;
Ok(Conditional::new(
&headers,
api.service.desired_state(&grant).await?,
))
}
#[derive(Debug, Default, Deserialize)]
#[serde(deny_unknown_fields)]
struct CatalogueQuery {
tenant: String,
#[serde(default)]
project: Option<String>,
#[serde(default)]
state: Option<String>,
#[serde(default)]
wire_family: Option<String>,
#[serde(default)]
offering: Option<String>,
#[serde(default)]
billable: Option<bool>,
}
async fn catalogue(
State(api): State<Arc<AdminApi>>,
identity: AdminIdentity,
headers: HeaderMap,
query: Result<Query<CatalogueQuery>, QueryRejection>,
) -> Result<Conditional<CatalogueView>, AdminError> {
const SCHEMA: &str = "catalogue";
let Query(query) = query.map_err(|rejection| AdminError::RequestInvalid {
schema: SCHEMA,
detail: rejection.body_text(),
})?;
let invalid = |field: &'static str, detail: String| AdminError::RequestInvalid {
schema: SCHEMA,
detail: format!("`{field}`: {detail}"),
};
let tenant =
TenantId::parse(&query.tenant).map_err(|error| invalid("tenant", error.to_string()))?;
let project = match query.project.as_deref() {
None => None,
Some(project) => {
Some(ProjectId::parse(project).map_err(|error| invalid("project", error.to_string()))?)
}
};
let state =
match query.state.as_deref() {
None => None,
Some(text) => Some(ModelLifecycle::parse(text).ok_or_else(|| {
invalid("state", format!("`{text}` is not a model lifecycle state"))
})?),
};
let wire_family = match query.wire_family.as_deref() {
None => None,
Some(text) => Some(
WireFamily::parse(text)
.ok_or_else(|| invalid("wire_family", format!("`{text}` is not a wire family")))?,
),
};
let offering = match query.offering.as_deref() {
None => None,
Some(text) => {
Some(OfferingId::parse(text).map_err(|error| invalid("offering", error.to_string()))?)
}
};
let request = CatalogueRequest {
tenant,
project,
filters: CatalogueFilters {
state,
wire_family,
offering,
billable: query.billable,
},
};
let grant = api
.authorize(
&identity,
AdminAction::ReadState,
Surface::Model,
&request.scope(),
)
.await?;
Ok(Conditional::new(
&headers,
api.service.model_catalogue(&grant, &request).await?,
))
}
#[derive(Debug, Default, Deserialize)]
#[serde(deny_unknown_fields)]
struct HistoryQuery {
#[serde(default)]
limit: Option<u32>,
#[serde(default)]
start: Option<String>,
}
async fn history(
State(api): State<Arc<AdminApi>>,
identity: AdminIdentity,
headers: HeaderMap,
query: Result<Query<HistoryQuery>, QueryRejection>,
) -> Result<Conditional<RevisionPage>, AdminError> {
let Query(query) = query.map_err(|rejection| AdminError::RequestInvalid {
schema: "history",
detail: rejection.body_text(),
})?;
let grant = api
.authorize(
&identity,
AdminAction::ReadHistory,
Surface::AuditTrail,
&ResourceScope::Deployment,
)
.await?;
let limit = match query.limit {
None => HistoryLimit::default(),
Some(limit) => HistoryLimit::parse(limit)?,
};
let start = match query.start.as_deref() {
None => None,
Some(text) => {
Some(
RevisionId::parse(text).map_err(|error| AdminError::RequestInvalid {
schema: "history",
detail: format!("`start`: {error}"),
})?,
)
}
};
let page = api
.service
.history(&grant, HistoryRequest { limit, start })
.await?;
Ok(Conditional::new(&headers, page))
}
async fn audit(
State(api): State<Arc<AdminApi>>,
identity: AdminIdentity,
headers: HeaderMap,
Path(revision): Path<String>,
) -> Result<Conditional<AuditPage>, AdminError> {
let grant = api
.authorize(
&identity,
AdminAction::ReadAudit,
Surface::AuditTrail,
&ResourceScope::Deployment,
)
.await?;
let revision = RevisionId::parse(&revision).map_err(|error| AdminError::RequestInvalid {
schema: "audit",
detail: format!("`revision`: {error}"),
})?;
Ok(Conditional::new(
&headers,
api.service.audit(&grant, revision).await?,
))
}
async fn convergence(
State(api): State<Arc<AdminApi>>,
identity: AdminIdentity,
headers: HeaderMap,
) -> Result<Conditional<ConvergenceResult>, AdminError> {
let grant = api
.authorize(
&identity,
AdminAction::ReadConvergence,
Surface::AuditTrail,
&ResourceScope::Deployment,
)
.await?;
let report = api.convergence_report();
let result = api.service.convergence(&grant, report.as_ref())?;
let identity = result.identity();
Ok(Conditional::identified_by(&headers, result, &identity))
}
#[derive(Debug, Default, Deserialize)]
#[serde(deny_unknown_fields)]
struct AvailabilityQuery {
#[serde(default)]
tenant: Option<String>,
#[serde(default)]
project: Option<String>,
}
async fn availability(
State(api): State<Arc<AdminApi>>,
identity: AdminIdentity,
headers: HeaderMap,
query: Result<Query<AvailabilityQuery>, QueryRejection>,
) -> Result<Conditional<AvailabilityResult>, AdminError> {
let Query(query) = query.map_err(|rejection| AdminError::RequestInvalid {
schema: "availability",
detail: rejection.body_text(),
})?;
if query.tenant.is_none() && query.project.is_none() {
return Err(AdminError::RequestInvalid {
schema: "availability",
detail: "`tenant`: an availability read must name the tenant it asks about".to_owned(),
});
}
let scope = scope_of(
"availability",
query.tenant.as_deref(),
query.project.as_deref(),
)?;
let grant = api
.authorize(
&identity,
AdminAction::ReadAvailability,
Surface::Model,
&scope,
)
.await?;
let authority = AvailabilityAuthority::of(
api.holds_deployment_authority(&identity, AdminAction::ReadAvailability),
);
let result = api.service.availability(
&grant,
&scope,
authority,
api.availability.as_deref(),
SystemTime::now(),
)?;
Ok(Conditional::new(&headers, result))
}
fn scope_of(
schema: &'static str,
tenant: Option<&str>,
project: Option<&str>,
) -> Result<ResourceScope, AdminError> {
let invalid = |field: &'static str, detail: String| AdminError::RequestInvalid {
schema,
detail: format!("`{field}`: {detail}"),
};
match (tenant, project) {
(None, None) => Ok(ResourceScope::Deployment),
(None, Some(_)) => Err(invalid(
"project",
"a project scope must name the tenant that owns it".to_owned(),
)),
(Some(tenant), project) => {
let tenant =
TenantId::parse(tenant).map_err(|error| invalid("tenant", error.to_string()))?;
match project {
None => Ok(ResourceScope::Tenant(tenant)),
Some(project) => {
let project = ProjectId::parse(project)
.map_err(|error| invalid("project", error.to_string()))?;
Ok(ResourceScope::Project { tenant, project })
}
}
}
}
}
mod extractors {
use super::{AdminError, AdminIdentity, MutationPreconditions};
use axum::extract::FromRequestParts;
use axum::http::request::Parts;
impl<S: Send + Sync> FromRequestParts<S> for AdminIdentity {
type Rejection = AdminError;
async fn from_request_parts(
parts: &mut Parts,
_state: &S,
) -> Result<Self, Self::Rejection> {
parts
.extensions
.get::<Self>()
.cloned()
.ok_or(AdminError::Unauthenticated(
crate::admin::auth::AdminAuthError::MissingCredential,
))
}
}
impl<S: Send + Sync> FromRequestParts<S> for MutationPreconditions {
type Rejection = AdminError;
async fn from_request_parts(
parts: &mut Parts,
_state: &S,
) -> Result<Self, Self::Rejection> {
if let Some(preconditions) = parts.extensions.get::<Self>() {
return Ok(preconditions.clone());
}
Self::from_headers(&parts.headers)
}
}
}