use crate::{
Event, Message, Result, Vars, event::ActWorkflowMessageHandle, scheduler::Runtime, utils,
};
use std::pin::Pin;
use std::sync::Arc;
use tracing::{debug, error, info};
type GlobSet = (
globset::GlobMatcher,
globset::GlobMatcher,
globset::GlobMatcher,
Vec<(String, globset::GlobMatcher)>,
);
#[derive(Debug, Clone)]
pub struct ChannelOptions {
pub id: String,
pub ack: bool,
pub r#type: String,
pub state: String,
pub uses: String,
pub options: Vars,
}
impl Default for ChannelOptions {
fn default() -> Self {
Self {
id: utils::shortid(),
ack: false,
r#type: "*".to_string(),
state: "*".to_string(),
uses: "*".to_string(),
options: Vars::new(),
}
}
}
impl ChannelOptions {
pub fn pattern(&self) -> String {
let mut options = Vars::new()
.with("ack", self.ack)
.with("type", self.r#type.clone())
.with("state", self.state.clone())
.with("uses", self.uses.clone());
for (key, value) in self.options.iter() {
options.set(key, value);
}
options.to_string()
}
}
pub struct Channel {
runtime: Arc<Runtime>,
ack: bool,
chan_id: String,
pattern: String,
glob: GlobSet,
}
impl Channel {
pub fn new(rt: &Arc<Runtime>) -> Self {
Self::channel(rt, &ChannelOptions::default())
}
#[allow(clippy::self_named_constructors)]
pub fn channel(rt: &Arc<Runtime>, options: &ChannelOptions) -> Self {
debug!("channel created");
let pat_type = globset::Glob::new(&options.r#type)
.unwrap()
.compile_matcher();
let pat_state = globset::Glob::new(&options.state)
.unwrap()
.compile_matcher();
let pat_uses = globset::Glob::new(&options.uses).unwrap().compile_matcher();
let opt_globs: Vec<(String, globset::GlobMatcher)> = options
.options
.iter()
.filter_map(|(k, v)| {
v.as_str()
.and_then(|pattern| globset::Glob::new(pattern).ok())
.map(|g| (k.clone(), g.compile_matcher()))
})
.collect();
Self {
runtime: rt.clone(),
ack: options.ack,
chan_id: options.id.clone(),
pattern: options.pattern(),
glob: (pat_type, pat_state, pat_uses, opt_globs),
}
}
pub fn on_message<F, Fut>(self: &Arc<Self>, f: F)
where
F: Fn(Event<Message>) -> Fut + Send + Sync + 'static,
Fut: Future<Output = ()> + Send + 'static,
{
let glob = self.glob.clone();
let runtime = self.runtime.clone();
let ack = self.ack;
let chan_id = self.chan_id.clone();
let pattern = self.pattern.clone();
let handle: ActWorkflowMessageHandle =
Arc::new(move |e| -> Pin<Box<dyn Future<Output = ()> + Send>> { Box::pin(f(e)) });
let chan = chan_id.clone();
self.runtime.emitter().on_message(&self.chan_id, move |e| {
debug!(chan = %chan, "on message");
let glob = glob.clone();
let runtime = runtime.clone();
let handle = handle.clone();
let chan_id = chan_id.clone();
let pattern = pattern.clone();
async move {
deliver(&glob, &runtime, ack, &chan_id, &pattern, &handle, e).await;
}
});
}
pub fn on_start<F, Fut>(self: &Arc<Self>, f: F)
where
F: Fn(Event<Message>) -> Fut + Send + Sync + 'static,
Fut: Future<Output = ()> + Send + 'static,
{
let glob = self.glob.clone();
let runtime = self.runtime.clone();
let ack = self.ack;
let chan_id = self.chan_id.clone();
let pattern = self.pattern.clone();
let handle: ActWorkflowMessageHandle =
Arc::new(move |e| -> Pin<Box<dyn Future<Output = ()> + Send>> { Box::pin(f(e)) });
self.runtime.emitter().on_start(&self.chan_id, move |e| {
let glob = glob.clone();
let runtime = runtime.clone();
let handle = handle.clone();
let chan_id = chan_id.clone();
let pattern = pattern.clone();
async move {
deliver(&glob, &runtime, ack, &chan_id, &pattern, &handle, e).await;
}
});
}
pub fn on_complete<F, Fut>(self: &Arc<Self>, f: F)
where
F: Fn(Event<Message>) -> Fut + Send + Sync + 'static,
Fut: Future<Output = ()> + Send + 'static,
{
let glob = self.glob.clone();
let runtime = self.runtime.clone();
let ack = self.ack;
let chan_id = self.chan_id.clone();
let pattern = self.pattern.clone();
let handle: ActWorkflowMessageHandle =
Arc::new(move |e| -> Pin<Box<dyn Future<Output = ()> + Send>> { Box::pin(f(e)) });
let chan = chan_id.clone();
self.runtime.emitter().on_complete(&self.chan_id, move |e| {
debug!(chan = %chan, "on complete");
let glob = glob.clone();
let runtime = runtime.clone();
let handle = handle.clone();
let chan_id = chan_id.clone();
let pattern = pattern.clone();
async move {
deliver(&glob, &runtime, ack, &chan_id, &pattern, &handle, e).await;
}
});
}
pub fn on_error<F, Fut>(self: &Arc<Self>, f: F)
where
F: Fn(Event<Message>) -> Fut + Send + Sync + 'static,
Fut: Future<Output = ()> + Send + 'static,
{
let glob = self.glob.clone();
let runtime = self.runtime.clone();
let ack = self.ack;
let chan_id = self.chan_id.clone();
let pattern = self.pattern.clone();
let handle: ActWorkflowMessageHandle =
Arc::new(move |e| -> Pin<Box<dyn Future<Output = ()> + Send>> { Box::pin(f(e)) });
self.runtime.emitter().on_error(&self.chan_id, move |e| {
let glob = glob.clone();
let runtime = runtime.clone();
let handle = handle.clone();
let chan_id = chan_id.clone();
let pattern = pattern.clone();
async move {
deliver(&glob, &runtime, ack, &chan_id, &pattern, &handle, e).await;
}
});
}
pub fn close(&self) {
self.runtime.emitter().remove(&self.chan_id);
}
}
async fn deliver(
glob: &GlobSet,
runtime: &Arc<Runtime>,
ack: bool,
chan_id: &str,
pattern: &str,
f: &ActWorkflowMessageHandle,
e: Event<Message>,
) {
if !is_match(glob, &e) {
return;
}
match store_if(runtime, ack, chan_id, pattern, &e).await {
Ok(Some(delivery_id)) => {
let mut msg = e.inner().clone();
msg.delivery_id = Some(delivery_id.clone());
let event = Event::from_inner(msg);
f(event).await;
if let Err(err) = runtime.cache().store().mark_delivered(&delivery_id).await {
error!(error = %err, delivery_id = %delivery_id, "mark delivery succeeded failed");
}
}
Ok(None) => {
f(e).await;
}
Err(err) => error!(error = %err, chan = %chan_id, "delivery store failed, message dropped"),
}
}
async fn store_if(
runtime: &Arc<Runtime>,
ack: bool,
chan_id: &str,
pattern: &str,
message: &Message,
) -> Result<Option<String>> {
if ack && !chan_id.is_empty() && message.delivery_id.is_none() {
info!(r#type = message.r#type, pid = %message.pid, tid = %message.tid, mid = %message.mid, state = %message.state, "delivery stored");
let store = runtime.cache().store();
if !store.messages().exists(&message.id).await? {
store.messages().create(&message.into_message()).await?;
}
let delivery = message.into_delivery(chan_id, pattern);
match store.deliveries().create(&delivery).await {
Ok(_) => Ok(Some(delivery.id)),
Err(err) => {
error!(error = %err, "channel store failure");
Err(err)
}
}
} else {
Ok(None)
}
}
fn is_match(glob: &GlobSet, e: &Event<Message>) -> bool {
let (pat_type, pat_state, pat_uses, pat_options) = glob;
if !pat_type.is_match(&e.r#type)
|| !pat_state.is_match(e.state.as_ref())
|| !pat_uses.is_match(e.uses.as_deref().unwrap_or_default())
{
return false;
}
let msg_options = e.options();
for (key, pat) in pat_options {
let value = msg_options
.as_ref()
.and_then(|o| o.get::<String>(key))
.unwrap_or_default();
if !pat.is_match(&value) {
return false;
}
}
true
}