use std::sync::Arc;
use aion_core::{ClusterEvent, NamespacePlacementWire};
use aion_store::{MintOutcome, NamespaceOrigin, NamespacePlacement, NamespaceStore};
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>,
}
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())
.finish_non_exhaustive()
}
}
impl NamespaceMinter {
#[must_use]
pub fn new(store: Arc<dyn NamespaceStore>, policy: AutoCreate) -> Self {
Self {
store,
policy,
cluster_publisher: None,
}
}
#[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 => {
if self.store.register_namespace(namespace, origin).await?
== MintOutcome::Created
{
self.announce_created(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(())
}
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)]
mod tests {
#![allow(clippy::expect_used)]
use std::num::NonZeroUsize;
use std::sync::Arc;
use aion_core::ClusterEvent;
use aion_store::{InMemoryStore, NamespaceOrigin, NamespaceStore};
use futures::StreamExt;
use super::NamespaceMinter;
use crate::cluster_publisher::ClusterEventPublisher;
use crate::config::AutoCreate;
fn publisher() -> ClusterEventPublisher {
ClusterEventPublisher::new(NonZeroUsize::new(16).expect("non-zero capacity"))
}
fn open_minter(store: Arc<InMemoryStore>, publisher: ClusterEventPublisher) -> NamespaceMinter {
let store: Arc<dyn NamespaceStore> = store;
NamespaceMinter::new(store, AutoCreate::Open).with_cluster_publisher(publisher)
}
async fn next_created_name<S>(deltas: &mut S) -> Result<String, Box<dyn std::error::Error>>
where
S: futures::Stream<
Item = Result<ClusterEvent, crate::cluster_publisher::ClusterStreamLagged>,
> + Unpin,
{
let event = deltas
.next()
.await
.ok_or("expected a namespace-created delta")?
.map_err(|lag| format!("unexpected lag: {lag:?}"))?;
match event {
ClusterEvent::NamespaceCreated { name, origin, .. } => {
assert_eq!(origin, "explicit");
Ok(name)
}
other => Err(format!("expected NamespaceCreated, got {other:?}").into()),
}
}
#[tokio::test]
async fn namespace_created_delta_emits_once_on_created_and_not_on_already_existed()
-> Result<(), Box<dyn std::error::Error>> {
let store = Arc::new(InMemoryStore::default());
let publisher = publisher();
let mut deltas = publisher.subscribe(0);
let minter = open_minter(Arc::clone(&store), publisher);
minter
.mint_or_gate(&["orders".to_owned()], NamespaceOrigin::WorkerMint)
.await?;
let first = deltas
.next()
.await
.ok_or("expected one namespace-created delta")?
.map_err(|lag| format!("unexpected lag: {lag:?}"))?;
match first {
ClusterEvent::NamespaceCreated {
name,
origin,
created_at,
..
} => {
assert_eq!(name, "orders");
assert_eq!(origin, "worker_mint");
let record = store
.get_namespace("orders")
.await?
.ok_or("record must exist after a Created mint")?;
assert_eq!(created_at, record.created_at);
}
other => return Err(format!("expected NamespaceCreated, got {other:?}").into()),
}
minter
.mint_or_gate(&["orders".to_owned()], NamespaceOrigin::WorkerMint)
.await?;
minter
.mint_or_gate(&["billing".to_owned()], NamespaceOrigin::StartMint)
.await?;
let next = deltas
.next()
.await
.ok_or("expected the second namespace's delta")?
.map_err(|lag| format!("unexpected lag: {lag:?}"))?;
match next {
ClusterEvent::NamespaceCreated { name, origin, .. } => {
assert_eq!(name, "billing");
assert_eq!(origin, "start_mint");
}
other => return Err(format!("expected NamespaceCreated, got {other:?}").into()),
}
Ok(())
}
#[tokio::test]
async fn explicit_create_emits_created_delta_once_then_silent_on_recreate()
-> Result<(), Box<dyn std::error::Error>> {
let store = Arc::new(InMemoryStore::default());
let publisher = publisher();
let mut deltas = publisher.subscribe(0);
let minter = open_minter(Arc::clone(&store), publisher);
let created = minter.create_explicit("tenant-a").await?;
assert_eq!(created, aion_store::MintOutcome::Created);
let again = minter.create_explicit("tenant-a").await?;
assert_eq!(again, aion_store::MintOutcome::AlreadyExisted);
let _ = minter.create_explicit("tenant-b").await?;
let first = next_created_name(&mut deltas).await?;
let second = next_created_name(&mut deltas).await?;
assert_eq!(
vec![first, second],
vec!["tenant-a".to_owned(), "tenant-b".to_owned()]
);
Ok(())
}
#[tokio::test]
async fn mint_without_publisher_creates_record_but_emits_no_delta()
-> Result<(), Box<dyn std::error::Error>> {
let store: Arc<dyn NamespaceStore> = Arc::new(InMemoryStore::default());
let minter = NamespaceMinter::new(Arc::clone(&store), AutoCreate::Open);
minter
.mint_or_gate(&["orders".to_owned()], NamespaceOrigin::WorkerMint)
.await?;
assert!(
store.get_namespace("orders").await?.is_some(),
"the durable record is still minted without a publisher"
);
Ok(())
}
}