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, Uri};
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, uuid_detail};
use super::router::{ADMIN_MAX_REQUEST_BYTES, AdminApi};
use super::secrets::{
self, RotateSecretRequest, SecretLifecycleRequest, SecretTransitionView, SecretVersionView,
SecretVersionsView, StageSecretRequest,
};
use super::service::{AvailabilityAuthority, MutationOutcome};
use crate::desired_state::{
InvalidId, 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)
}
pub(super) fn stage_secret_route() -> MethodRouter<Arc<AdminApi>> {
post(stage_secret)
}
pub(super) fn rotate_secret_route() -> MethodRouter<Arc<AdminApi>> {
post(rotate_secret)
}
pub(super) fn secret_lifecycle_route() -> MethodRouter<Arc<AdminApi>> {
post(secret_lifecycle)
}
pub(super) fn secret_versions_route() -> MethodRouter<Arc<AdminApi>> {
get(secret_versions)
}
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`: {}", id_detail(RevisionId::PREFIX, &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>,
}
pub(super) const CATALOGUE_MAX_QUERY_BYTES: usize = 2 * 1024;
pub(super) const CATALOGUE_MAX_QUERY_PARAMS: usize = 6;
fn validate_catalogue_query(uri: &Uri) -> Result<(), AdminError> {
let query = uri.query().unwrap_or_default();
if query.len() > CATALOGUE_MAX_QUERY_BYTES {
return Err(AdminError::RequestInvalid {
schema: "catalogue",
detail: format!("query exceeds the {CATALOGUE_MAX_QUERY_BYTES}-byte limit"),
});
}
let parameters = query
.split('&')
.filter(|parameter| !parameter.is_empty())
.count();
if parameters > CATALOGUE_MAX_QUERY_PARAMS {
return Err(AdminError::RequestInvalid {
schema: "catalogue",
detail: format!("query has more than {CATALOGUE_MAX_QUERY_PARAMS} parameters"),
});
}
Ok(())
}
async fn catalogue(
State(api): State<Arc<AdminApi>>,
identity: AdminIdentity,
headers: HeaderMap,
uri: Uri,
query: Result<Query<CatalogueQuery>, QueryRejection>,
) -> Result<Conditional<CatalogueView>, AdminError> {
const SCHEMA: &str = "catalogue";
validate_catalogue_query(&uri)?;
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`: {}", id_detail(RevisionId::PREFIX, &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`: {}", id_detail(RevisionId::PREFIX, &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 stage_secret(
State(api): State<Arc<AdminApi>>,
identity: AdminIdentity,
body: Result<Bytes, BytesRejection>,
) -> Result<Json<SecretVersionView>, AdminError> {
const SCHEMA: &str = "secret";
let body = document(SCHEMA, body)?;
let request: StageSecretRequest = secret_document(SCHEMA, &body)?;
let owner = secrets::owner_of(SCHEMA, &request.tenant, request.project.as_deref())?;
let material = secrets::material_of(SCHEMA, request.material)?;
let grant = api
.authorize(
&identity,
AdminAction::WriteSecrets,
Surface::Credential,
&owner.scope(),
)
.await?;
Ok(Json(
api.service.stage_secret(&grant, owner, material).await?,
))
}
async fn rotate_secret(
State(api): State<Arc<AdminApi>>,
identity: AdminIdentity,
body: Result<Bytes, BytesRejection>,
) -> Result<Json<SecretVersionView>, AdminError> {
const SCHEMA: &str = "secret_rotation";
let body = document(SCHEMA, body)?;
let request: RotateSecretRequest = secret_document(SCHEMA, &body)?;
let owner = secrets::owner_of(SCHEMA, &request.tenant, request.project.as_deref())?;
let reference = secrets::reference_of(SCHEMA, &request.reference)?;
let material = secrets::material_of(SCHEMA, request.material)?;
let grant = api
.authorize(
&identity,
AdminAction::WriteSecrets,
Surface::Credential,
&owner.scope(),
)
.await?;
Ok(Json(
api.service
.rotate_secret(&grant, owner, reference, material)
.await?,
))
}
async fn secret_lifecycle(
State(api): State<Arc<AdminApi>>,
identity: AdminIdentity,
body: Result<Bytes, BytesRejection>,
) -> Result<Json<SecretTransitionView>, AdminError> {
const SCHEMA: &str = "secret_lifecycle";
let body = document(SCHEMA, body)?;
let request: SecretLifecycleRequest = secret_document(SCHEMA, &body)?;
let owner = secrets::owner_of(SCHEMA, &request.tenant, request.project.as_deref())?;
let reference = secrets::reference_of(SCHEMA, &request.reference)?;
let next = secrets::lifecycle_of(SCHEMA, &request.lifecycle)?;
let grant = api
.authorize(
&identity,
AdminAction::WriteSecrets,
Surface::Credential,
&owner.scope(),
)
.await?;
Ok(Json(
api.service
.move_secret(&grant, owner, reference, next)
.await?,
))
}
#[derive(Debug, Deserialize)]
#[serde(deny_unknown_fields)]
struct SecretVersionsQuery {
tenant: 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))
}
async fn secret_versions(
State(api): State<Arc<AdminApi>>,
identity: AdminIdentity,
headers: HeaderMap,
Path(secret): Path<String>,
query: Result<Query<SecretVersionsQuery>, QueryRejection>,
) -> Result<Conditional<SecretVersionsView>, AdminError> {
const SCHEMA: &str = "secret_versions";
let Query(query) = query.map_err(|rejection| AdminError::RequestInvalid {
schema: SCHEMA,
detail: rejection.body_text(),
})?;
let owner = secrets::owner_of(SCHEMA, &query.tenant, query.project.as_deref())?;
let secret = secrets::secret_of(SCHEMA, &secret)?;
let grant = api
.authorize(
&identity,
AdminAction::ReadSecrets,
Surface::Credential,
&owner.scope(),
)
.await?;
Ok(Conditional::new(
&headers,
api.service.secret_versions(&grant, owner, secret).await?,
))
}
fn secret_document<T: serde::de::DeserializeOwned>(
schema: &'static str,
body: &Bytes,
) -> Result<T, AdminError> {
serde_json::from_slice(body).map_err(|error| AdminError::RequestInvalid {
schema,
detail: format!(
"the document is not a valid `{schema}` request (line {}, column {}); its text is not \
reported, because it carries material",
error.line(),
error.column()
),
})
}
pub(super) 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", id_detail(TenantId::PREFIX, &error)))?;
match project {
None => Ok(ResourceScope::Tenant(tenant)),
Some(project) => {
let project = ProjectId::parse(project).map_err(|error| {
invalid("project", id_detail(ProjectId::PREFIX, &error))
})?;
Ok(ResourceScope::Project { tenant, project })
}
}
}
}
}
fn id_detail(prefix: &'static str, error: &InvalidId) -> String {
match error {
InvalidId::Prefix { .. } => format!("is not a `{prefix}`-prefixed id"),
InvalidId::Uuid(uuid) => format!("has a uuid that {}", uuid_detail(uuid)),
}
}
#[cfg(test)]
mod tests {
use super::*;
const PASTED_MATERIAL: &str = "sk-axond-admin-sentinel-51H9xNEVERLOGME";
fn detail(error: AdminError) -> String {
error
.operator_detail()
.expect("a request refusal has operator detail")
.to_owned()
}
#[test]
fn a_malformed_scope_id_is_refused_without_echoing_what_arrived() {
let refusal = detail(
scope_of("rollback", Some(PASTED_MATERIAL), None).expect_err("a wrong-prefix tenant"),
);
assert_eq!(refusal, "`tenant`: is not a `ten_`-prefixed id");
let malformed = format!("{}not-a-uuid", TenantId::PREFIX);
let refusal =
detail(scope_of("rollback", Some(&malformed), None).expect_err("a malformed uuid"));
assert_eq!(
refusal,
"`tenant`: has a uuid that is not a hyphenated 8-4-4-4-12 uuid"
);
assert!(!refusal.contains(&malformed));
let tenant = format!("{}0189f8c1-2a3b-7c4d-8e5f-6a7b8c9d0e1f", TenantId::PREFIX);
let refusal = detail(
scope_of("rollback", Some(&tenant), Some(PASTED_MATERIAL))
.expect_err("a wrong-prefix project"),
);
assert_eq!(refusal, "`project`: is not a `prj_`-prefixed id");
assert!(!refusal.contains(PASTED_MATERIAL));
}
}
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)
}
}
}