use std::sync::Arc;
use std::time::SystemTime;
use serde::Serialize;
use tracing::{debug, warn};
use super::auth::{AdminAction, AdminAuthError, AdminGrant};
use super::diff::SemanticDiff;
use super::error::AdminError;
use super::protocol::{MutationRequest, WriteMode};
use super::reads::{
AuditPage, ConvergenceResult, HistoryRequest, RevisionPage, RevisionRecord, StateView,
};
use crate::backends::control_plane::{ControlPlaneError, ControlPlaneStore};
use crate::config::Mode;
use crate::convergence::RevisionReport;
use crate::desired_state::{
AuditEvent, AuditEventId, DesiredState, ExpectedRevision, LoadedRevision, Mutation, MutationId,
ResourceScope, RevisionCandidate, RevisionId, Uuid7Generator, ValidationError,
};
fn scope_covers(granted: &ResourceScope, resource: &ResourceScope) -> bool {
match granted {
ResourceScope::Deployment => true,
ResourceScope::Tenant(tenant) => resource.tenant() == Some(*tenant),
ResourceScope::Project { tenant, project } => matches!(
resource,
ResourceScope::Project {
tenant: other_tenant,
project: other_project,
} if other_tenant == tenant && other_project == project
),
}
}
pub trait DesiredStateEdit: Send + Sync {
fn edit(&self, state: &mut DesiredState) -> Result<(), ValidationError>;
}
impl<F> DesiredStateEdit for F
where
F: Fn(&mut DesiredState) -> Result<(), ValidationError> + Send + Sync,
{
fn edit(&self, state: &mut DesiredState) -> Result<(), ValidationError> {
self(state)
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
#[serde(tag = "result", rename_all = "snake_case")]
pub enum MutationResult {
Published { revision: String },
Replayed { revision: String },
DryRun,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
pub struct MutationOutcome {
#[serde(flatten)]
pub result: MutationResult,
#[serde(skip_serializing_if = "Option::is_none")]
pub base: Option<String>,
pub checksum: String,
pub mode: &'static str,
pub diff: SemanticDiff,
}
impl MutationOutcome {
pub fn revision(&self) -> Option<&str> {
match &self.result {
MutationResult::Published { revision } | MutationResult::Replayed { revision } => {
Some(revision)
}
MutationResult::DryRun => None,
}
}
}
pub struct AdminService {
store: Option<Arc<dyn ControlPlaneStore>>,
ids: Uuid7Generator,
}
impl AdminService {
pub fn stateless() -> Self {
Self {
store: None,
ids: Uuid7Generator::new(),
}
}
pub fn stateful(store: Arc<dyn ControlPlaneStore>) -> Self {
Self {
store: Some(store),
ids: Uuid7Generator::new(),
}
}
pub fn for_mode(mode: Mode, store: Option<Arc<dyn ControlPlaneStore>>) -> Self {
match (mode, store) {
(Mode::Stateful, Some(store)) => Self::stateful(store),
_ => Self::stateless(),
}
}
pub const fn mode(&self) -> Mode {
if self.store.is_some() {
Mode::Stateful
} else {
Mode::Stateless
}
}
fn store(&self) -> Result<&Arc<dyn ControlPlaneStore>, AdminError> {
self.store.as_ref().ok_or(AdminError::StatefulModeRequired)
}
fn within_scope(
granted: &ResourceScope,
before: &DesiredState,
after: &DesiredState,
) -> Result<(), AdminError> {
let touched = before
.resources()
.filter(|resource| after.get(&resource.reference) != Some(resource))
.chain(
after
.resources()
.filter(|resource| before.get(&resource.reference) != Some(resource)),
);
for resource in touched {
if !scope_covers(granted, &resource.scope) {
return Err(AdminError::Forbidden(AdminAuthError::ScopeNotPermitted));
}
}
Ok(())
}
fn permits(grant: &AdminGrant, action: AdminAction) -> Result<(), AdminError> {
if grant.action() != action {
return Err(AdminError::Forbidden(AdminAuthError::ActionNotPermitted {
action,
}));
}
Ok(())
}
fn permits_deployment_read(grant: &AdminGrant, action: AdminAction) -> Result<(), AdminError> {
Self::permits(grant, action)?;
if grant.scope() != &ResourceScope::Deployment {
return Err(AdminError::Forbidden(AdminAuthError::ScopeNotPermitted));
}
Ok(())
}
pub async fn desired_state(&self, grant: &AdminGrant) -> Result<StateView, AdminError> {
Self::permits_deployment_read(grant, AdminAction::ReadState)?;
let store = self.store()?;
let revision = store.load_desired_revision().await.map_err(log_store)?;
StateView::of(revision.as_ref())
}
pub async fn history(
&self,
grant: &AdminGrant,
request: HistoryRequest,
) -> Result<RevisionPage, AdminError> {
Self::permits_deployment_read(grant, AdminAction::ReadHistory)?;
let store = self.store()?;
let mut next = match request.start {
Some(start) => Some(start),
None => store.desired_revision().await.map_err(log_store)?,
};
let mut revisions = Vec::new();
while let Some(id) = next {
if revisions.len() >= request.limit.get() as usize {
break;
}
let manifest = match store.load_manifest(id).await {
Ok(manifest) => manifest,
Err(ControlPlaneError::RevisionNotFound(_)) if !revisions.is_empty() => break,
Err(error) => return Err(log_store(error)),
};
next = manifest.parent;
revisions.push(RevisionRecord::of(&manifest));
}
let next_start = next
.filter(|_| revisions.len() >= request.limit.get() as usize)
.map(|id| id.to_string());
Ok(RevisionPage {
revisions,
limit: request.limit.get(),
next_start,
})
}
pub async fn audit(
&self,
grant: &AdminGrant,
revision: RevisionId,
) -> Result<AuditPage, AdminError> {
Self::permits_deployment_read(grant, AdminAction::ReadAudit)?;
let store = self.store()?;
let events = store.audit_trail(revision).await.map_err(log_store)?;
Ok(AuditPage::of(revision, &events))
}
pub fn convergence(
&self,
grant: &AdminGrant,
report: &RevisionReport,
) -> Result<ConvergenceResult, AdminError> {
Self::permits_deployment_read(grant, AdminAction::ReadConvergence)?;
self.store()?;
Ok(ConvergenceResult::of(report))
}
pub async fn apply(
&self,
grant: &AdminGrant,
request: &MutationRequest,
edit: &dyn DesiredStateEdit,
) -> Result<MutationOutcome, AdminError> {
let store = self.store()?;
Self::permits(grant, AdminAction::for_mutation(request.kind))?;
if grant.scope() != &request.scope {
return Err(AdminError::Forbidden(AdminAuthError::ScopeNotPermitted));
}
let head = store.desired_revision().await.map_err(log_store)?;
let expected = request.preconditions.expected;
if !expected.matches(head) && request.mode().is_dry_run() {
return Err(AdminError::RevisionConflict {
expected,
actual: head,
});
}
let base = match expected {
ExpectedRevision::Empty => None,
ExpectedRevision::Exactly(id) => Some(id),
};
let current = match base {
Some(id) => match store.load_revision(id).await {
Ok(revision) => Some(revision),
Err(ControlPlaneError::RevisionNotFound(_)) if !expected.matches(head) => {
return Err(AdminError::RevisionConflict {
expected,
actual: head,
});
}
Err(error) => return Err(log_store(error)),
},
None => None,
};
let empty = DesiredState::new();
let current_state = current
.as_ref()
.map_or(&empty, |revision: &LoadedRevision| revision.state());
let mut candidate_state = current_state.clone();
edit.edit(&mut candidate_state)?;
Self::within_scope(grant.scope(), current_state, &candidate_state)?;
let mutation = MutationId::new(self.ids.next());
let submitted_at = SystemTime::now();
let identity = grant.identity();
let candidate = RevisionCandidate {
expected,
state: candidate_state,
mutation: Mutation {
id: mutation,
actor: identity.actor(),
kind: request.kind,
scope: request.scope.clone(),
idempotency_key: request.preconditions.idempotency_key.clone(),
submitted_at,
},
audit: AuditEvent {
id: AuditEventId::new(self.ids.next()),
mutation,
actor: identity.actor(),
kind: request.kind,
target: None,
summary: identity.audit_summary(request.summary.as_str()),
recorded_at: submitted_at,
},
};
let checksum = candidate.validated_checksum()?;
let diff = SemanticDiff::between(Some(current_state), &candidate.state)?;
let base = base.map(|id| id.to_string());
if request.mode().is_dry_run() {
debug!(
action = grant.action().as_str(),
breakglass = identity.is_breakglass(),
added = diff.summary.added,
removed = diff.summary.removed,
updated = diff.summary.updated,
"administrative dry run validated a candidate"
);
return Ok(MutationOutcome {
result: MutationResult::DryRun,
base,
checksum: checksum.to_string(),
mode: WriteMode::DryRun.as_str(),
diff,
});
}
let manifest = store.publish_revision(candidate).await.map_err(log_store)?;
let result = if manifest.mutation == mutation {
MutationResult::Published {
revision: manifest.id.to_string(),
}
} else {
MutationResult::Replayed {
revision: manifest.id.to_string(),
}
};
debug!(
action = grant.action().as_str(),
breakglass = identity.is_breakglass(),
revision = %manifest.id,
replayed = matches!(result, MutationResult::Replayed { .. }),
"administrative mutation published"
);
Ok(MutationOutcome {
result,
base,
checksum: checksum.to_string(),
mode: WriteMode::Apply.as_str(),
diff,
})
}
}
fn log_store(error: ControlPlaneError) -> AdminError {
let error = AdminError::from_control_plane(error);
if let Some(detail) = error.operator_detail() {
warn!(
code = error.code(),
detail, "control-plane operation failed"
);
}
error
}