use {
super::FunctionId,
crate::primitives::Datum,
core::{
marker::PhantomData,
sync::atomic::{AtomicUsize, Ordering},
},
std::sync::Arc,
tokio_util::sync::CancellationToken,
};
mod builder;
pub(in crate::functions) mod registry;
pub use builder::{Builder, BuilderError, HandlerConfig};
pub(in crate::functions) use registry::Registry;
pub struct Handler<Req: Datum, Res: Datum, E: Datum = ()> {
function_id: FunctionId,
stats: Arc<Stats>,
registry: Arc<Registry>,
cancel: CancellationToken,
_marker: PhantomData<fn(&Req, &Res, &E)>,
}
impl<Req: Datum, Res: Datum, E: Datum> Handler<Req, Res, E> {
pub(in crate::functions) fn new(
function_id: FunctionId,
stats: Arc<Stats>,
registry: Arc<Registry>,
cancel: CancellationToken,
) -> Self {
Self {
function_id,
stats,
registry,
cancel,
_marker: PhantomData,
}
}
pub const fn function_id(&self) -> &FunctionId {
&self.function_id
}
pub fn stats(&self) -> &Stats {
&self.stats
}
}
impl<Req: Datum, Res: Datum, E: Datum> Drop for Handler<Req, Res, E> {
fn drop(&mut self) {
self.cancel.cancel();
self.registry.remove(self.function_id);
}
}
impl<Req: Datum, Res: Datum, E: Datum> core::fmt::Debug
for Handler<Req, Res, E>
{
fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
f.debug_struct("Handler")
.field("function_id", &self.function_id)
.finish_non_exhaustive()
}
}
#[derive(Debug, Default)]
pub struct Stats {
served: AtomicUsize,
failed: AtomicUsize,
inflight: AtomicUsize,
}
impl Stats {
pub fn served(&self) -> usize {
self.served.load(Ordering::Relaxed)
}
pub fn failed(&self) -> usize {
self.failed.load(Ordering::Relaxed)
}
pub fn inflight(&self) -> usize {
self.inflight.load(Ordering::Relaxed)
}
pub(in crate::functions) fn record_started(&self) {
self.inflight.fetch_add(1, Ordering::Relaxed);
}
pub(in crate::functions) fn record_served(&self) {
self.inflight.fetch_sub(1, Ordering::Relaxed);
self.served.fetch_add(1, Ordering::Relaxed);
}
pub(in crate::functions) fn record_failed(&self) {
self.inflight.fetch_sub(1, Ordering::Relaxed);
self.failed.fetch_add(1, Ordering::Relaxed);
}
}