use std::sync::Arc;
use aion_core::{ClusterEvent, NamespacePlacementWire};
use aion_store::{MintOutcome, NamespaceOrigin, NamespacePlacement, NamespaceStore};
use super::route::{MintCredentials, MintRoute, NamespaceRouting};
use crate::cluster_publisher::ClusterEventPublisher;
use crate::config::AutoCreate;
use crate::error::ServerError;
#[derive(Clone)]
pub struct NamespaceMinter {
store: Arc<dyn NamespaceStore>,
policy: AutoCreate,
cluster_publisher: Option<ClusterEventPublisher>,
routing: Option<NamespaceRouting>,
}
impl std::fmt::Debug for NamespaceMinter {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter
.debug_struct("NamespaceMinter")
.field("policy", &self.policy)
.field("cluster_publisher", &self.cluster_publisher.is_some())
.field("routing", &self.routing.is_some())
.finish_non_exhaustive()
}
}
impl NamespaceMinter {
#[must_use]
pub fn new(store: Arc<dyn NamespaceStore>, policy: AutoCreate) -> Self {
Self {
store,
policy,
cluster_publisher: None,
routing: None,
}
}
#[must_use]
pub fn with_routing(mut self, routing: NamespaceRouting) -> Self {
self.routing = Some(routing);
self
}
#[must_use]
pub fn without_routing(mut self) -> Self {
self.routing = None;
self
}
#[must_use]
pub fn with_caller_credentials(mut self, credentials: MintCredentials) -> Self {
self.routing = self
.routing
.map(|routing| routing.with_credentials(credentials));
self
}
#[must_use]
pub fn with_cluster_publisher(mut self, publisher: ClusterEventPublisher) -> Self {
self.cluster_publisher = Some(publisher);
self
}
#[must_use]
pub fn policy(&self) -> AutoCreate {
self.policy
}
pub async fn mint_or_gate(
&self,
namespaces: &[String],
origin: NamespaceOrigin,
) -> Result<(), ServerError> {
for namespace in namespaces {
match self.policy {
AutoCreate::Open => self.mint(namespace, origin).await?,
AutoCreate::Closed => {
if self.store.get_namespace(namespace).await?.is_none() {
return Err(ServerError::namespace_denied(format!(
"namespace {namespace} does not exist and auto_create is closed"
)));
}
}
}
}
Ok(())
}
async fn mint(&self, namespace: &str, origin: NamespaceOrigin) -> Result<(), ServerError> {
if let Some(routing) = &self.routing
&& let MintRoute::Remote { shard, target } = routing.route_for(namespace)
{
return routing.forward(namespace, target, shard, origin).await;
}
if self.store.register_namespace(namespace, origin).await? == MintOutcome::Created {
self.announce_created(namespace, origin).await?;
}
Ok(())
}
pub async fn create_explicit(&self, name: &str) -> Result<MintOutcome, ServerError> {
let outcome = self
.store
.register_namespace(name, NamespaceOrigin::Explicit)
.await?;
if outcome == MintOutcome::Created {
self.announce_created(name, NamespaceOrigin::Explicit)
.await?;
}
Ok(outcome)
}
pub async fn set_placement(
&self,
name: &str,
placement: NamespacePlacement,
) -> Result<bool, ServerError> {
if self
.store
.set_namespace_placement(name, placement.clone())
.await?
.is_none()
{
return Ok(false);
}
self.announce_placement_changed(name, &placement);
Ok(true)
}
pub async fn placement_of(&self, name: &str) -> Result<NamespacePlacement, ServerError> {
Ok(self
.store
.get_namespace(name)
.await?
.map(|record| record.placement)
.unwrap_or_default())
}
fn announce_placement_changed(&self, name: &str, placement: &NamespacePlacement) {
let wire = placement_to_wire(placement);
tracing::info!(
namespace = %name,
placement_kind = %wire.kind,
"namespace placement changed"
);
let Some(publisher) = &self.cluster_publisher else {
return;
};
let name = name.to_owned();
drop(
publisher.emit(move |meta| ClusterEvent::NamespacePlacementChanged {
meta,
name,
placement: wire,
}),
);
}
async fn announce_created(
&self,
name: &str,
origin: NamespaceOrigin,
) -> Result<(), ServerError> {
tracing::info!(
namespace = %name,
origin = origin_label(origin),
"namespace created"
);
let Some(publisher) = &self.cluster_publisher else {
return Ok(());
};
let Some(record) = self.store.get_namespace(name).await? else {
return Ok(());
};
let name = record.name;
let created_at = record.created_at;
let label = origin_label(record.origin).to_owned();
drop(publisher.emit(move |meta| ClusterEvent::NamespaceCreated {
meta,
name,
created_at,
origin: label,
}));
Ok(())
}
}
fn placement_to_wire(placement: &NamespacePlacement) -> NamespacePlacementWire {
match placement {
NamespacePlacement::Unplaced => NamespacePlacementWire {
kind: "unplaced".to_owned(),
nodes: Vec::new(),
},
NamespacePlacement::Prefer { nodes } => NamespacePlacementWire {
kind: "prefer".to_owned(),
nodes: nodes.iter().cloned().collect(),
},
NamespacePlacement::Pinned { nodes } => NamespacePlacementWire {
kind: "pinned".to_owned(),
nodes: nodes.iter().cloned().collect(),
},
}
}
const fn origin_label(origin: NamespaceOrigin) -> &'static str {
match origin {
NamespaceOrigin::WorkerMint => "worker_mint",
NamespaceOrigin::StartMint => "start_mint",
NamespaceOrigin::Explicit => "explicit",
NamespaceOrigin::InferredFromState => "inferred_from_state",
}
}
#[cfg(test)]
#[path = "minter_tests.rs"]
mod tests;