use crate::{BoxFuture, Container, Injectable, ProviderDef, Result};
use serde_json::Value;
use std::{collections::BTreeMap, sync::Arc, time::SystemTime};
use tokio_util::sync::CancellationToken;
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum MessageHandlerKind {
Command,
Event,
}
#[derive(Clone, Debug, Default)]
pub struct MessageDelivery {
pub stream: Option<String>,
pub consumer: Option<String>,
pub attempt: u64,
}
#[derive(Clone, Debug)]
pub struct MessageContext {
subject: String,
headers: BTreeMap<String, String>,
correlation_id: Option<String>,
deadline: Option<SystemTime>,
cancellation: CancellationToken,
delivery: Option<MessageDelivery>,
event_id: Option<String>,
}
impl MessageContext {
pub fn new(
subject: impl Into<String>,
headers: BTreeMap<String, String>,
correlation_id: Option<String>,
deadline: Option<SystemTime>,
cancellation: CancellationToken,
delivery: Option<MessageDelivery>,
event_id: Option<String>,
) -> Self {
Self {
subject: subject.into(),
headers,
correlation_id,
deadline,
cancellation,
delivery,
event_id,
}
}
pub fn subject(&self) -> &str {
&self.subject
}
pub fn headers(&self) -> &BTreeMap<String, String> {
&self.headers
}
pub fn correlation_id(&self) -> Option<&str> {
self.correlation_id.as_deref()
}
pub fn deadline(&self) -> Option<SystemTime> {
self.deadline
}
pub fn cancellation_token(&self) -> &CancellationToken {
&self.cancellation
}
pub fn delivery(&self) -> Option<&MessageDelivery> {
self.delivery.as_ref()
}
pub fn delivery_attempt(&self) -> u64 {
self.delivery
.as_ref()
.map_or(0, |delivery| delivery.attempt)
}
pub fn event_id(&self) -> Option<&str> {
self.event_id.as_deref()
}
}
type InvokeMessageFn = Arc<
dyn for<'a> Fn(&'a Container, MessageContext, Value) -> BoxFuture<'a, Result<Option<Value>>>
+ Send
+ Sync,
>;
#[derive(Clone)]
pub struct MessageHandlerDef {
pub kind: MessageHandlerKind,
pub pattern: &'static str,
invoke: InvokeMessageFn,
}
impl MessageHandlerDef {
pub fn new(
kind: MessageHandlerKind,
pattern: &'static str,
invoke: impl for<'a> Fn(
&'a Container,
MessageContext,
Value,
) -> BoxFuture<'a, Result<Option<Value>>>
+ Send
+ Sync
+ 'static,
) -> Self {
Self {
kind,
pattern,
invoke: Arc::new(invoke),
}
}
pub fn invoke<'a>(
&self,
container: &'a Container,
context: MessageContext,
payload: Value,
) -> BoxFuture<'a, Result<Option<Value>>> {
(self.invoke)(container, context, payload)
}
}
pub struct MicroserviceDef {
pub provider: ProviderDef,
handlers_fn: fn() -> Vec<MessageHandlerDef>,
}
impl MicroserviceDef {
pub fn of<T: Injectable>(handlers_fn: fn() -> Vec<MessageHandlerDef>) -> Self {
Self {
provider: ProviderDef::of::<T>(),
handlers_fn,
}
}
pub fn handlers(&self) -> Vec<MessageHandlerDef> {
(self.handlers_fn)()
}
}
pub trait Microservice: Injectable {
fn definition() -> MicroserviceDef;
}
#[doc(hidden)]
pub fn _assert_microservice_send_sync<T: Send + Sync + 'static>() {
let _ = std::marker::PhantomData::<Arc<T>>;
}