use std::net::SocketAddr;
use std::sync::Arc;
use async_trait::async_trait;
use aion_proto::WireError;
use aion_proto::generated;
use aion_store::NamespaceOrigin;
use crate::error::ServerError;
pub trait NamespaceShardResolver: Send + Sync {
fn shard_for_namespace(&self, name: &str) -> usize;
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub enum MintOwner {
Local,
Remote(SocketAddr),
Unknown,
}
pub trait MintShardOwners: Send + Sync {
fn owner_of(&self, shard: usize) -> MintOwner;
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub enum MintRoute {
Local {
shard: usize,
},
Remote {
shard: usize,
target: SocketAddr,
},
Unknown {
shard: usize,
},
}
#[derive(Clone, Default)]
pub struct MintCredentials {
metadata: tonic::metadata::MetadataMap,
}
impl MintCredentials {
#[must_use]
pub fn from_grpc_metadata(metadata: &tonic::metadata::MetadataMap) -> Self {
Self {
metadata: metadata.clone(),
}
}
#[must_use]
pub fn to_grpc_metadata(&self) -> tonic::metadata::MetadataMap {
self.metadata.clone()
}
}
impl std::fmt::Debug for MintCredentials {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter
.debug_struct("MintCredentials")
.field("entries", &self.metadata.len())
.finish()
}
}
#[derive(Debug)]
pub enum ForwardMintError {
Refused(WireError),
Unreachable(String),
}
#[async_trait]
pub trait MintForwarder: Send + Sync {
async fn forward_mint(
&self,
target: SocketAddr,
credentials: &MintCredentials,
namespaces: &[String],
origin: NamespaceOrigin,
) -> Result<(), ForwardMintError>;
}
#[derive(Clone)]
pub struct NamespaceRouting {
shards: Arc<dyn NamespaceShardResolver>,
owners: Arc<dyn MintShardOwners>,
forwarder: Arc<dyn MintForwarder>,
credentials: MintCredentials,
}
impl std::fmt::Debug for NamespaceRouting {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter
.debug_struct("NamespaceRouting")
.field("credentials", &self.credentials)
.finish_non_exhaustive()
}
}
impl NamespaceRouting {
#[must_use]
pub fn new(
shards: Arc<dyn NamespaceShardResolver>,
owners: Arc<dyn MintShardOwners>,
forwarder: Arc<dyn MintForwarder>,
) -> Self {
Self {
shards,
owners,
forwarder,
credentials: MintCredentials::default(),
}
}
#[must_use]
pub fn with_credentials(mut self, credentials: MintCredentials) -> Self {
self.credentials = credentials;
self
}
#[must_use]
pub fn route_for(&self, name: &str) -> MintRoute {
let shard = self.shards.shard_for_namespace(name);
match self.owners.owner_of(shard) {
MintOwner::Local => MintRoute::Local { shard },
MintOwner::Remote(target) => MintRoute::Remote { shard, target },
MintOwner::Unknown => MintRoute::Unknown { shard },
}
}
pub async fn forward(
&self,
name: &str,
target: SocketAddr,
shard: usize,
origin: NamespaceOrigin,
) -> Result<(), ServerError> {
let namespaces = [name.to_owned()];
match self
.forwarder
.forward_mint(target, &self.credentials, &namespaces, origin)
.await
{
Ok(()) => Ok(()),
Err(ForwardMintError::Refused(wire)) => Err(ServerError::Wire { wire }),
Err(ForwardMintError::Unreachable(detail)) => Err(ServerError::Wire {
wire: WireError::not_owner(format!(
"namespace `{name}` is registered on shard {shard}, owned by \
cluster node {target}, which did not answer the mint: {detail}"
))
.with_error_type("NotOwner"),
}),
}
}
}
#[cfg(feature = "haematite-backend")]
impl NamespaceShardResolver for aion_store_haematite::HaematiteStore {
fn shard_for_namespace(&self, name: &str) -> usize {
Self::shard_for_namespace(self, name)
}
}
#[cfg(feature = "haematite-backend")]
impl MintShardOwners for crate::routing::StaticShardDirectory {
fn owner_of(&self, shard: usize) -> MintOwner {
use crate::routing::{OwnerView, ShardDirectory};
match ShardDirectory::owner_of(self, shard) {
OwnerView::Local => MintOwner::Local,
OwnerView::Remote(node) => node.grpc_addr.map_or(MintOwner::Unknown, MintOwner::Remote),
OwnerView::Unknown => MintOwner::Unknown,
}
}
}
#[must_use]
pub const fn encode_mint_origin(origin: NamespaceOrigin) -> i32 {
let code = match origin {
NamespaceOrigin::WorkerMint => generated::NamespaceMintOrigin::Worker,
NamespaceOrigin::StartMint => generated::NamespaceMintOrigin::Start,
NamespaceOrigin::Explicit => generated::NamespaceMintOrigin::Explicit,
NamespaceOrigin::InferredFromState => generated::NamespaceMintOrigin::InferredFromState,
};
code as i32
}
#[must_use]
pub fn decode_mint_origin(code: i32) -> Option<NamespaceOrigin> {
match generated::NamespaceMintOrigin::try_from(code).ok()? {
generated::NamespaceMintOrigin::Unspecified => None,
generated::NamespaceMintOrigin::Worker => Some(NamespaceOrigin::WorkerMint),
generated::NamespaceMintOrigin::Start => Some(NamespaceOrigin::StartMint),
generated::NamespaceMintOrigin::Explicit => Some(NamespaceOrigin::Explicit),
generated::NamespaceMintOrigin::InferredFromState => {
Some(NamespaceOrigin::InferredFromState)
}
}
}
#[cfg(test)]
#[path = "route_tests.rs"]
mod tests;