use {
super::{
super::{FunctionId, accept::CallReply},
Stats,
},
crate::{
discovery::{Discovery, PeerInfo},
network::LocalNode,
primitives::{Bytes, ShortFmtExt},
tickets::TicketValidator,
},
dashmap::{DashMap, Entry as MapEntry},
futures::future::BoxFuture,
std::sync::Arc,
tokio::sync::Semaphore,
tokio_util::sync::CancellationToken,
};
pub(in crate::functions) type DispatchFn = Arc<
dyn Fn(
Bytes,
crate::functions::CallContext,
) -> BoxFuture<'static, Result<CallReply, DispatchError>>
+ Send
+ Sync,
>;
#[derive(Debug, Clone, Copy)]
pub(in crate::functions) enum DispatchError {
DecodeRequest,
EncodeReply,
}
pub(in crate::functions) struct Entry {
pub dispatch: DispatchFn,
pub require: Box<dyn Fn(&PeerInfo) -> bool + Send + Sync>,
pub caller_auth: Vec<Arc<dyn TicketValidator>>,
pub permits: Arc<Semaphore>,
pub cancel: CancellationToken,
pub stats: Arc<Stats>,
}
pub(in crate::functions) struct Registry {
pub local: LocalNode,
pub discovery: Discovery,
pub active: DashMap<FunctionId, Arc<Entry>>,
}
impl Registry {
pub(in crate::functions) fn new(
local: LocalNode,
discovery: Discovery,
) -> Self {
Self {
local,
discovery,
active: DashMap::new(),
}
}
pub fn create(
&self,
function_id: FunctionId,
entry: Entry,
) -> Result<Arc<Entry>, ()> {
match self.active.entry(function_id) {
MapEntry::Vacant(slot) => {
let entry = slot.insert(Arc::new(entry)).clone();
let labels = [("network", self.local.network_id().short().to_string())];
metrics::gauge!("mosaik.functions.handlers.active", &labels)
.increment(1.0);
self
.discovery
.update_local_entry(move |me| me.add_functions(function_id));
Ok(entry)
}
MapEntry::Occupied(_) => Err(()),
}
}
pub fn open(&self, function_id: FunctionId) -> Option<Arc<Entry>> {
self.active.get(&function_id).map(|entry| entry.clone())
}
pub fn remove(&self, function_id: FunctionId) {
if self.active.remove(&function_id).is_some() {
let labels = [("network", self.local.network_id().short().to_string())];
metrics::gauge!("mosaik.functions.handlers.active", &labels)
.decrement(1.0);
self
.discovery
.update_local_entry(move |me| me.remove_functions(function_id));
}
}
}
impl core::fmt::Debug for Registry {
fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
write!(f, "Registry({} functions)", self.active.len())
}
}