use std::collections::BTreeMap;
use std::sync::Arc;
use aion_core::{ActivityId, WorkflowId};
use futures::future::BoxFuture;
#[derive(Clone, Debug)]
pub struct ActivityDispatch {
pub namespace: String,
pub task_queue: String,
pub node: Option<String>,
pub workflow_id: WorkflowId,
pub activity_id: ActivityId,
pub name: String,
pub input: String,
pub config: String,
pub attempt: u32,
pub labels: BTreeMap<String, String>,
}
pub trait ActivityDispatcher: Send + Sync + 'static {
fn dispatch(&self, request: ActivityDispatch) -> Result<String, String>;
fn dispatch_async(
self: Arc<Self>,
request: ActivityDispatch,
) -> BoxFuture<'static, Result<String, String>> {
Box::pin(async move {
let blocking = tokio::task::spawn_blocking(move || self.dispatch(request));
match blocking.await {
Ok(result) => result,
Err(join_error) => Err(format!("activity dispatch task failed: {join_error}")),
}
})
}
}
#[cfg(test)]
mod tests {
use std::sync::Arc;
use std::collections::BTreeMap;
use aion_core::{ActivityId, WorkflowId};
use super::{ActivityDispatch, ActivityDispatcher};
use crate::runtime::EngineNifState;
struct Echo;
impl ActivityDispatcher for Echo {
fn dispatch(&self, request: ActivityDispatch) -> Result<String, String> {
Ok(request.input)
}
}
fn echo_request(input: &str) -> ActivityDispatch {
ActivityDispatch {
namespace: "default".to_owned(),
task_queue: "default".to_owned(),
node: None,
workflow_id: WorkflowId::new_v4(),
activity_id: ActivityId::from_sequence_position(0),
name: "test".to_owned(),
input: input.to_owned(),
config: "{}".to_owned(),
attempt: 1,
labels: BTreeMap::new(),
}
}
#[test]
fn dispatcher_is_accessible_after_install_on_engine_state() {
let state = EngineNifState::default();
state.set_activity_dispatcher(Arc::new(Echo));
let dispatcher = state.activity_dispatcher();
assert!(dispatcher.is_some());
assert_eq!(
dispatcher
.as_ref()
.and_then(|d| d.dispatch(echo_request("hello")).ok()),
Some("hello".to_owned())
);
}
}