use crate::{
ActSchema, ChannelOptions, Engine, Message, Signal, Variant, VariantTypes, Vars, Workflow,
data::{self, Package},
event::MessageState,
scheduler::TaskState,
store::query::*,
utils::{
self, consts,
test::{USES_IRQ, USES_SET, auto_complete},
},
};
use parking_lot::Mutex;
use serde_json::json;
use std::sync::Arc;
use serial_test::serial;
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_manager_publish_ok() {
let engine = Engine::new().start().await.unwrap();
let manager = engine.executor();
let pack = data::Package {
id: "pack1".to_string(),
desc: "desc".to_string(),
icon: "icon".to_string(),
doc: "doc".to_string(),
version: "0.1.0".to_string(),
schema: "{}".to_string(),
run_as: crate::ActRunAs::Func,
resources: "[]".to_string(),
catalog: crate::package::ActPackageCatalog::Core,
..Default::default()
};
let result = manager.pack().publish(&pack).await;
assert!(result.is_ok());
assert!(manager.pack().publish(&pack).await.is_ok());
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_manager_deploy_ok() {
let engine = Engine::new().start().await.unwrap();
let manager = engine.executor();
let model = Workflow::new()
.with_id(&utils::longid())
.with_step(|step| step.with_uses(USES_IRQ, Vars::new().with("key", "test")));
let result = manager.model().deploy(&model, None).await;
assert!(result.is_ok());
assert!(manager.model().get(&model.id, "text").await.is_ok());
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_manager_deploy_many_times() {
let engine = Engine::new().start().await.unwrap();
let manager = engine.executor();
let model = Workflow::new()
.with_id(&utils::longid())
.with_step(|step| step.with_id("step1"));
let mut result = true;
for _ in 0..10 {
let state = manager.model().deploy(&model, None).await;
result &= state.is_ok();
}
assert!(result);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_manager_deploy_no_model_id_error() {
let engine = Engine::new().start().await.unwrap();
let manager = engine.executor();
let model = Workflow::new().with_step(|step| step.with_id("step1"));
let result = manager.model().deploy(&model, None).await;
assert!(result.is_err());
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_manager_deploy_dup_id_error() {
let engine = Engine::new().start().await.unwrap();
let executor = engine.executor();
let model = Workflow::new()
.with_id(&utils::longid())
.with_step(|step| step.with_id("step1"))
.with_step(|step| step.with_id("step1"));
let result = executor.model().deploy(&model, None).await;
assert!(result.is_err());
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn engine_executor_start_no_pid() {
let engine = Engine::new().start().await.unwrap();
let executor = engine.executor();
let mid = utils::longid();
let workflow = Workflow::new()
.with_id(&mid)
.with_step(|step| step.with_uses(USES_IRQ, Vars::new().with("key", "test")));
engine
.executor()
.model()
.deploy(&workflow, None)
.await
.unwrap();
let options = Vars::new();
let result = executor.proc().start(&workflow.id, options).await;
assert!(result.is_ok());
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn engine_executor_start_with_pid() {
let engine = Engine::new().start().await.unwrap();
let executor = engine.executor();
let mid = utils::longid();
let workflow = Workflow::new()
.with_id(&mid)
.with_step(|step| step.with_uses(USES_IRQ, Vars::new().with("key", "test")));
engine
.executor()
.model()
.deploy(&workflow, None)
.await
.unwrap();
let mut options = Vars::new();
options.insert("pid".to_string(), "123".into());
let result = executor.proc().start(&workflow.id, options).await;
assert!(result.is_ok());
assert_eq!(result.unwrap(), "123");
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_executor_start_empty_pid() {
let engine = Engine::new().start().await.unwrap();
let executor = engine.executor();
let mid = utils::longid();
let workflow = Workflow::new()
.with_id(&mid)
.with_step(|step| step.with_uses(USES_IRQ, Vars::new().with("key", "test")));
engine
.executor()
.model()
.deploy(&workflow, None)
.await
.unwrap();
let mut options = Vars::new();
options.insert("pid".to_string(), "".into());
let result = executor.proc().start(&workflow.id, options).await;
assert_eq!(
result.err().unwrap().to_string(),
"external process id cannot be empty"
);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_executor_start_dup_pid_error() {
let engine = Engine::new().start().await.unwrap();
let executor = engine.executor();
let pid = utils::longid();
let mid = utils::longid();
let model = Workflow::new()
.with_id(&mid)
.with_step(|step| step.with_uses(USES_IRQ, Vars::new().with("key", "test")));
let store = engine.runtime().cache().store();
let proc = data::Proc {
id: pid.clone(),
name: model.name.clone(),
mid: model.id.clone(),
state: TaskState::None.to_string(),
start_time: 0,
end_time: 0,
timestamp: 0,
model: model.to_json().unwrap(),
env: "{}".to_string(),
err: None,
removable: false,
v: 0,
};
store.procs().create(&proc).await.expect("create process");
engine
.executor()
.model()
.deploy(&model, None)
.await
.expect("fail to deploy workflow");
let mut options = Vars::new();
options.insert("pid".to_string(), json!(pid.to_string()));
let result = executor.proc().start(&model.id, options).await;
assert!(result.is_err());
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_executor_start_from_yaml() {
let engine = Engine::new().start().await.unwrap();
let executor = engine.executor();
let mid = utils::longid();
let model = Workflow::new()
.with_id(&mid)
.with_step(|step| step.with_uses(USES_IRQ, Vars::new().with("key", "test")))
.to_yml()
.unwrap();
let result = executor
.proc()
.start_from_model(&model, "yaml", Vars::new())
.await;
assert!(result.is_ok());
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_executor_start_from_json() {
let engine = Engine::new().start().await.unwrap();
let executor = engine.executor();
let mid = utils::longid();
let model = Workflow::new()
.with_id(&mid)
.with_step(|step| step.with_uses(USES_IRQ, Vars::new().with("key", "test")))
.to_json()
.unwrap();
let result = executor
.proc()
.start_from_model(&model, "json", Vars::new())
.await;
assert!(result.is_ok());
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_executor_start_with_inputs_schema_ok() {
let engine = Engine::new().start().await.unwrap();
let executor = engine.executor();
let mid = utils::longid();
let workflow = Workflow::new()
.with_id(&mid)
.with_inputs(ActSchema::Multiple(vec![Variant::create(
"a",
json!("string"),
)]))
.with_step(|step| step.with_uses(USES_IRQ, Vars::new().with("key", "test")));
engine
.executor()
.model()
.deploy(&workflow, None)
.await
.unwrap();
let result = executor
.proc()
.start(&mid, Vars::new().with("a", "abc"))
.await;
assert!(result.is_ok());
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_executor_start_with_inputs_schema_err() {
let engine = Engine::new().start().await.unwrap();
let executor = engine.executor();
let mid = utils::longid();
let workflow = Workflow::new()
.with_id(&mid)
.with_inputs(ActSchema::Multiple(vec![Variant::create(
"a",
json!("string"),
)]))
.with_step(|step| step.with_uses(USES_IRQ, Vars::new().with("key", "test")));
engine
.executor()
.model()
.deploy(&workflow, None)
.await
.unwrap();
let result = executor
.proc()
.start(&mid, Vars::new().with("a", 100))
.await;
assert!(result.is_err());
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_executor_start_with_outputs_schema_ok() {
let engine = Engine::new().start().await.unwrap();
let executor = engine.executor();
let mid = utils::longid();
let workflow = Workflow::new()
.with_id(&mid)
.with_expose(Variant::new().name("a").r#type(VariantTypes::Number));
engine
.executor()
.model()
.deploy(&workflow, None)
.await
.unwrap();
let (s1, s2) = Signal::new(0).double();
engine.channel().on_complete(move |e| {
let s2 = s2.clone();
async move {
s2.send(e.outputs.get::<i32>("a").unwrap());
}
});
executor
.proc()
.start(&mid, Vars::new().with("a", 100))
.await
.unwrap();
let ret = s1.recv().await;
assert_eq!(ret, 100);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_executor_start_with_outputs_schema_err() {
let engine = Engine::new().start().await.unwrap();
let executor = engine.executor();
let mid = utils::longid();
let workflow = Workflow::new()
.with_id(&mid)
.with_expose(Variant::new().name("a").r#type(VariantTypes::Number));
engine
.executor()
.model()
.deploy(&workflow, None)
.await
.unwrap();
let (s1, s2) = Signal::new(None).double();
engine.channel().on_error(move |e| {
let s2 = s2.clone();
async move {
s2.send(e.inputs.get::<String>("message"));
}
});
executor
.proc()
.start(&mid, Vars::new().with("a", "abc"))
.await
.unwrap();
let ret = s1.recv().await;
assert!(ret.is_some());
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_executor_start_from_empty_fmt() {
let engine = Engine::new().start().await.unwrap();
let executor = engine.executor();
let mid = utils::longid();
let model = Workflow::new()
.with_id(&mid)
.with_step(|step| step.with_uses(USES_IRQ, Vars::new().with("key", "test")))
.to_json()
.unwrap();
let result = executor
.proc()
.start_from_model(&model, "", Vars::new())
.await;
assert!(result.is_err());
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_executor_start_from_error_fmt() {
let engine = Engine::new().start().await.unwrap();
let executor = engine.executor();
let mid = utils::longid();
let model = Workflow::new()
.with_id(&mid)
.with_step(|step| step.with_uses(USES_IRQ, Vars::new().with("key", "test")))
.to_json()
.unwrap();
let result = executor
.proc()
.start_from_model(&model, "xml", Vars::new())
.await;
assert!(result.is_err());
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_manager_models_get_count() {
let engine = Engine::new().start().await.unwrap();
let manager = engine.executor();
let mut model = Workflow::new().with_step(|step| step.with_id("step1"));
for _ in 0..5 {
model.set_id(&utils::longid());
manager.model().deploy(&model, None).await.unwrap();
}
let result = manager
.model()
.list(&Query::new().offset(0).limit(10))
.await
.unwrap();
assert_eq!(result.count, 5);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_manager_models_order() {
let engine = Engine::new().start().await.unwrap();
let manager = engine.executor();
let mut model = Workflow::new().with_step(|step| step.with_id("step1"));
for i in 0..5 {
model.set_id(&utils::longid());
model.name = format!("model-{}", i + 1);
manager.model().deploy(&model, None).await.unwrap();
}
let result = manager
.model()
.list(&Query::new().order("timestamp", Sort::Desc))
.await
.unwrap();
assert_eq!(result.rows.first().unwrap().name, "model-5");
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_manager_models_get_rows() {
let engine = Engine::new().start().await.unwrap();
let manager = engine.executor();
let mut model = Workflow::new().with_step(|step| step.with_id("step1"));
for _ in 0..5 {
model.set_id(&utils::longid());
manager.model().deploy(&model, None).await.unwrap();
}
let result = manager
.model()
.list(&Query::new().offset(4).limit(10))
.await
.unwrap();
assert_eq!(result.rows.len(), 1);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_manager_models_query() {
let engine = Engine::new().start().await.unwrap();
let manager = engine.executor();
let mut model = Workflow::new().with_step(|step| step.with_id("step1"));
for i in 0..5 {
model.set_id(&utils::longid());
model.name = format!("model-{}", i + 1);
manager.model().deploy(&model, None).await.unwrap();
}
let result = manager
.model()
.list(&Query::new().filter(Filter::and().expr(Expr::eq("name", "model-3"))))
.await
.unwrap();
assert_eq!(result.rows.len(), 1);
assert_eq!(result.rows.first().unwrap().name, "model-3");
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_manager_model_get_text() {
let engine = Engine::new().start().await.unwrap();
let manager = engine.executor();
let mut model = Workflow::new().with_step(|step| step.with_id("step1"));
model.set_id(&utils::longid());
manager.model().deploy(&model, None).await.unwrap();
let result = manager.model().get(&model.id, "text").await.unwrap();
assert_eq!(result.id, model.id);
assert!(!result.data.is_empty());
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_manager_model_get_tree() {
let engine = Engine::new().start().await.unwrap();
let manager = engine.executor();
let mut model = Workflow::new().with_step(|step| step.with_id("step1"));
model.set_id(&utils::longid());
manager.model().deploy(&model, None).await.unwrap();
let result = manager.model().get(&model.id, "tree").await.unwrap();
assert_eq!(result.id, model.id);
assert!(!result.data.is_empty());
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_manager_model_remove() {
let engine = Engine::new().start().await.unwrap();
let manager = engine.executor();
let mut model = Workflow::new().with_step(|step| step.with_id("step1"));
model.set_id(&utils::longid());
manager.model().deploy(&model, None).await.unwrap();
manager.model().rm(&model.id).await.unwrap();
assert_eq!(
manager
.model()
.list(&Query::new().offset(0).limit(10))
.await
.unwrap()
.count,
0
);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_manager_model_remove_with_events() {
let engine = Engine::new().start().await.unwrap();
let manager = engine.executor();
let mut model = Workflow::new()
.with_trigger(|t| t.with_id("event1").with_kind("manual"))
.with_trigger(|t| t.with_id("event2").with_kind("manual"))
.with_step(|step| step.with_id("step1"));
model.set_id(&utils::longid());
manager.model().deploy(&model, None).await.unwrap();
assert_eq!(
manager
.evt()
.list(&Query::new().filter(Filter::and().expr(Expr::eq(consts::MODEL_ID, &model.id))))
.await
.unwrap()
.count,
2
);
manager.model().rm(&model.id).await.unwrap();
assert_eq!(
manager
.model()
.list(&Query::new().offset(0).limit(10))
.await
.unwrap()
.count,
0
);
assert_eq!(
manager
.evt()
.list(&Query::new().filter(Filter::and().expr(Expr::eq(consts::MODEL_ID, &model.id))))
.await
.unwrap()
.count,
0
);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_manager_procs_one() {
let engine = Engine::new().start().await.unwrap();
let manager = engine.executor();
let model = Workflow::new().with_step(|step| {
step.with_id("step1")
.with_uses(USES_IRQ, Vars::new().with("key", "act1"))
});
let rt = engine.runtime();
let sig = engine.signal(());
let s1 = sig.clone();
let proc = rt.create_proc(&utils::longid(), &model);
engine.channel().on_start(move |_| {
let s1 = s1.clone();
async move {
s1.close();
}
});
rt.launch(&proc).await.unwrap();
sig.recv().await;
assert_eq!(
manager
.proc()
.list(&Query::new().offset(0).limit(10))
.await
.unwrap()
.count,
1
);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_manager_procs_count() {
let engine = Engine::new().start().await.unwrap();
let manager = engine.executor();
let model = Workflow::new().with_step(|step| {
step.with_id("step1")
.with_uses(USES_IRQ, Vars::new().with("key", "act1"))
});
let rt = engine.runtime();
let sig = engine.signal(());
let s1 = sig.clone();
let count = Arc::new(Mutex::new(0));
engine.channel().on_start(move |_e| {
let count = count.clone();
let s1 = s1.clone();
async move {
println!("message:{_e:?}");
let mut count = count.lock();
*count += 1;
if *count == 5 {
s1.close();
}
}
});
for _ in 0..5 {
let proc = rt.create_proc(&utils::longid(), &model);
rt.launch(&proc).await.unwrap();
}
sig.recv().await;
assert_eq!(
manager
.proc()
.list(&Query::new().offset(0).limit(10))
.await
.unwrap()
.count,
5
);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_manager_procs_offset_in_range() {
let engine = Engine::new().start().await.unwrap();
let manager = engine.executor();
let model = Workflow::new().with_step(|step| {
step.with_id("step1")
.with_uses(USES_IRQ, Vars::new().with("key", "act1"))
});
let rt = engine.runtime();
let sig = engine.signal(());
let s1 = sig.clone();
let count = Arc::new(Mutex::new(0));
engine.channel().on_start(move |_e| {
let count = count.clone();
let s1 = s1.clone();
async move {
println!("message:{_e:?}");
let mut count = count.lock();
*count += 1;
if *count == 5 {
s1.close();
}
}
});
for _ in 0..5 {
let proc = rt.create_proc(&utils::longid(), &model);
rt.launch(&proc).await.unwrap();
}
sig.recv().await;
assert_eq!(
manager
.proc()
.list(&Query::new().offset(4).limit(10))
.await
.unwrap()
.rows
.len(),
1
);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_manager_procs_offset_out_range() {
let engine = Engine::new().start().await.unwrap();
let manager = engine.executor();
let model = Workflow::new().with_step(|step| {
step.with_id("step1")
.with_uses(USES_IRQ, Vars::new().with("key", "act1"))
});
let rt = engine.runtime();
let sig = engine.signal(());
let s1 = sig.clone();
let count = Arc::new(Mutex::new(0));
engine.channel().on_start(move |_e| {
let count = count.clone();
let s1 = s1.clone();
async move {
println!("message:{_e:?}");
let mut count = count.lock();
*count += 1;
if *count == 5 {
s1.close();
}
}
});
for _ in 0..5 {
let proc = rt.create_proc(&utils::longid(), &model);
rt.launch(&proc).await.unwrap();
}
sig.recv().await;
assert_eq!(
manager
.proc()
.list(&Query::new().offset(1000).limit(10))
.await
.unwrap()
.rows
.len(),
0
);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_manager_procs_query() {
let engine = Engine::new().start().await.unwrap();
let manager = engine.executor();
let model = Workflow::new().with_step(|step| {
step.with_id("step1")
.with_uses(USES_IRQ, Vars::new().with("key", "act1"))
});
let rt = engine.runtime();
let sig = engine.signal(());
let s1 = sig.clone();
let count = Arc::new(Mutex::new(0));
engine.channel().on_start(move |_e| {
let count = count.clone();
let s1 = s1.clone();
async move {
println!("message:{_e:?}");
let mut count = count.lock();
*count += 1;
if *count == 5 {
s1.close();
}
}
});
let pid = utils::longid();
for i in 0..5 {
let proc = rt.create_proc(&format!("{pid}-{}", i + 1), &model);
rt.launch(&proc).await.unwrap();
}
sig.recv().await;
let rows = manager
.proc()
.list(&Query::new().filter(Filter::and().expr(Expr::eq("id", format!("{pid}-3")))))
.await
.unwrap()
.rows;
assert_eq!(rows.len(), 1);
assert_eq!(rows.first().unwrap().id, format!("{pid}-3"));
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_manager_procs_order() {
let engine = Engine::new().start().await.unwrap();
let manager = engine.executor();
let model = Workflow::new().with_step(|step| {
step.with_id("step1")
.with_uses(USES_IRQ, Vars::new().with("key", "act1"))
});
let rt = engine.runtime();
let sig = engine.signal(());
let s1 = sig.clone();
let count = Arc::new(Mutex::new(0));
engine.channel().on_start(move |_e| {
let count = count.clone();
let s1 = s1.clone();
async move {
println!("message:{_e:?}");
let mut count = count.lock();
*count += 1;
if *count == 5 {
s1.close();
}
}
});
let pid = utils::longid();
for i in 0..5 {
let proc = rt.create_proc(&format!("{pid}-{}", i + 1), &model);
rt.launch(&proc).await.unwrap();
}
sig.recv().await;
let rows = manager
.proc()
.list(&Query::new().order("timestamp", Sort::Desc))
.await
.unwrap()
.rows;
assert_eq!(rows.first().unwrap().id, format!("{pid}-5"));
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_manager_proc_get() {
let engine = Engine::new().start().await.unwrap();
let manager = engine.executor();
let model = Workflow::new().with_step(|step| {
step.with_id("step1")
.with_uses(USES_IRQ, Vars::new().with("key", "act1"))
});
let rt = engine.runtime();
let sig = engine.signal(());
let s1 = sig.clone();
engine.channel().on_start(move |_| {
let s1 = s1.clone();
async move {
s1.close();
}
});
let pid = utils::longid();
let proc = rt.create_proc(&pid, &model);
rt.launch(&proc).await.unwrap();
sig.recv().await;
let info = manager.proc().get(&pid).await.unwrap();
assert_eq!(info.id, pid);
assert!(!info.tasks.is_empty());
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_manager_tasks_count() {
let engine = Engine::new().start().await.unwrap();
let manager = engine.executor();
let model = Workflow::new().with_step(|step| {
step.with_id("step1")
.with_uses(USES_IRQ, Vars::new().with("key", "act1"))
});
let rt = engine.runtime();
let sig = engine.signal(());
let s1 = sig.clone();
engine.channel().on_message(move |e| {
let s1 = s1.clone();
async move {
if e.is_params_key("act1") {
s1.close()
}
}
});
let pid = utils::longid();
let mut vars = Vars::new();
vars.insert("uid".to_string(), json!("u1"));
vars.insert("pid".to_string(), json!(pid));
rt.start(&model, vars).await.unwrap();
sig.recv().await;
engine.runtime().cache().flush().await.unwrap();
let tasks = manager
.task()
.list(
&Query::new()
.filter(Filter::and().expr(Expr::eq("pid", &pid)))
.offset(0)
.limit(10),
)
.await
.unwrap();
assert_eq!(tasks.count, 3); }
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_manager_tasks_offset_in_range() {
let engine = Engine::new().start().await.unwrap();
let manager = engine.executor();
let model = Workflow::new().with_step(|step| {
step.with_id("step1")
.with_uses(USES_IRQ, Vars::new().with("key", "act1"))
});
let rt = engine.runtime();
let sig = engine.signal(());
let s1 = sig.clone();
engine.channel().on_message(move |e| {
let s1 = s1.clone();
async move {
if e.is_params_key("act1") {
s1.close()
}
}
});
let pid = utils::longid();
let mut vars = Vars::new();
vars.insert("uid".to_string(), json!("u1"));
vars.insert("pid".to_string(), json!(pid));
rt.start(&model, vars).await.unwrap();
sig.recv().await;
engine.runtime().cache().flush().await.unwrap();
let tasks = manager
.task()
.list(
&Query::new()
.filter(Filter::and().expr(Expr::eq("pid", &pid)))
.offset(2)
.limit(10),
)
.await
.unwrap();
assert_eq!(tasks.rows.len(), 1); }
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_manager_tasks_offset_out_range() {
let engine = Engine::new().start().await.unwrap();
let manager = engine.executor();
let model = Workflow::new().with_step(|step| {
step.with_id("step1")
.with_uses(USES_IRQ, Vars::new().with("key", "act1"))
});
let rt = engine.runtime();
let sig = engine.signal(());
let s1 = sig.clone();
engine.channel().on_message(move |e| {
let s1 = s1.clone();
async move {
if e.is_params_key("act1") {
s1.close()
}
}
});
let pid = utils::longid();
let mut vars = Vars::new();
vars.insert("uid".to_string(), json!("u1"));
vars.insert("pid".to_string(), json!(pid));
rt.start(&model, vars).await.unwrap();
sig.recv().await;
engine.runtime().cache().flush().await.unwrap();
let tasks = manager
.task()
.list(
&Query::new()
.filter(Filter::and().expr(Expr::eq("pid", &pid)))
.offset(1000)
.limit(10),
)
.await
.unwrap();
assert_eq!(tasks.rows.len(), 0); }
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_manager_tasks_query() {
let engine = Engine::new().start().await.unwrap();
let manager = engine.executor();
let model = Workflow::new().with_step(|step| {
step.with_id("step1")
.with_uses(USES_IRQ, Vars::new().with("key", "act1"))
});
let rt = engine.runtime();
let sig = engine.signal(());
let s1 = sig.clone();
engine.channel().on_message(move |e| {
let s1 = s1.clone();
async move {
if e.is_params_key("act1") {
s1.close()
}
}
});
let pid = utils::longid();
let mut vars = Vars::new();
vars.insert("uid".to_string(), json!("u1"));
vars.insert("pid".to_string(), json!(pid));
rt.start(&model, vars).await.unwrap();
sig.recv().await;
engine.runtime().cache().flush().await.unwrap();
let tasks = manager
.task()
.list(
&Query::new().filter(
Filter::and()
.expr(Expr::eq("pid", &pid))
.expr(Expr::eq("state", "interrupted")),
),
)
.await
.unwrap();
assert_eq!(tasks.rows.first().unwrap().r#type, "act");
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_manager_tasks_order() {
let engine = Engine::new().start().await.unwrap();
let manager = engine.executor();
let model = Workflow::new().with_step(|step| {
step.with_id("step1")
.with_uses(USES_IRQ, Vars::new().with("key", "act1"))
});
let rt = engine.runtime();
let sig = engine.signal(());
let s1 = sig.clone();
engine.channel().on_message(move |e| {
let s1 = s1.clone();
async move {
if e.is_params_key("act1") {
s1.close()
}
}
});
let pid = utils::longid();
let mut vars = Vars::new();
vars.insert("uid".to_string(), json!("u1"));
vars.insert("pid".to_string(), json!(pid));
rt.start(&model, vars).await.unwrap();
sig.recv().await;
engine.runtime().cache().flush().await.unwrap();
let tasks = manager
.task()
.list(
&Query::new()
.filter(Filter::and().expr(Expr::eq("pid", &pid)))
.order("timestamp", Sort::Desc),
)
.await
.unwrap();
assert_eq!(tasks.rows.first().unwrap().r#type, "act");
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_manager_task_get() {
let engine = Engine::new().start().await.unwrap();
let manager = engine.executor();
let model = Workflow::new().with_step(|step| {
step.with_id("step1")
.with_uses(USES_IRQ, Vars::new().with("key", "act1"))
});
let rt = engine.runtime();
let sig = engine.signal(());
let s1 = sig.clone();
engine.channel().on_message(move |e| {
let s1 = s1.clone();
async move {
if e.is_params_key("act1") {
s1.close()
}
}
});
let pid = utils::longid();
let mut vars = Vars::new();
vars.insert("uid".to_string(), json!("u1"));
vars.insert("pid".to_string(), json!(pid));
rt.start(&model, vars).await.unwrap();
sig.recv().await;
let tasks = manager
.task()
.list(
&Query::new()
.filter(Filter::and().expr(Expr::eq("pid", &pid)))
.offset(0)
.limit(10),
)
.await
.unwrap();
let mut result = true;
for task in tasks.rows.iter() {
result &= manager.task().get(&pid, &task.id).await.is_ok();
}
assert!(result);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_manager_messages_all() {
let engine = Engine::new().start().await.unwrap();
let manager = engine.executor();
let model = Workflow::new().with_step(|step| {
step.with_id("step1")
.with_uses(USES_IRQ, Vars::new().with("key", "act1"))
});
let rt = engine.runtime();
let sig = engine.signal(());
let chan = engine.channel_with_options(&ChannelOptions {
ack: true,
..Default::default()
});
chan.on_message(move |e| async move {
println!("message:{e:?}");
});
let pid = utils::longid();
let proc = rt.create_proc(&pid, &model);
rt.launch(&proc).await.unwrap();
sig.timeout(100).await;
assert_eq!(
manager
.msg()
.list(&Query::new().offset(0).limit(1000))
.await
.unwrap()
.count,
3
);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_manager_messages_query() {
let engine = Engine::new().start().await.unwrap();
let manager = engine.executor();
let model = Workflow::new().with_step(|step| {
step.with_id("step1")
.with_uses(USES_IRQ, Vars::new().with("key", "act1"))
});
let rt: Arc<crate::scheduler::Runtime> = engine.runtime();
let (sig, rx) = engine.signal(()).double();
auto_complete(&engine, &rx);
let chan = engine.channel_with_options(&ChannelOptions {
ack: true,
..Default::default()
});
chan.on_message(move |e| async move {
println!("message:{e:?}");
});
let pid = utils::longid();
let proc = rt.create_proc(&pid, &model);
rt.launch(&proc).await.unwrap();
sig.timeout(100).await;
assert_eq!(
manager
.msg()
.list(
&Query::new()
.filter(Filter::and().expr(Expr::eq("pid", &pid)))
.offset(0)
.limit(1000)
)
.await
.unwrap()
.count,
3
);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_manager_messages_order() {
let engine = Engine::new().start().await.unwrap();
let manager = engine.executor();
let model = Workflow::new().with_step(|step| {
step.with_id("step1")
.with_uses(USES_IRQ, Vars::new().with("key", "act1"))
});
let rt = engine.runtime();
let sig = engine.signal(());
let chan = engine.channel_with_options(&ChannelOptions {
ack: true,
..Default::default()
});
chan.on_message(move |e| async move {
println!("message:{e:?}");
});
let pid = utils::longid();
let proc = rt.create_proc(&pid, &model);
rt.launch(&proc).await.unwrap();
sig.timeout(100).await;
assert_eq!(
manager
.msg()
.list(
&Query::new()
.filter(Filter::and().expr(Expr::eq("pid", &pid)))
.order("timestamp", Sort::Desc)
)
.await
.unwrap()
.rows
.first()
.unwrap()
.r#type,
"act"
);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_manager_messages_count() {
let engine = Engine::new().start().await.unwrap();
let manager = engine.executor();
let model = Workflow::new().with_step(|step| {
step.with_id("step1")
.with_uses(USES_IRQ, Vars::new().with("key", "act1"))
});
let rt = engine.runtime();
let (sig, s1) = engine.signal(0).double();
let count = Arc::new(Mutex::new(0));
let chan = engine.channel_with_options(&ChannelOptions {
ack: true,
..Default::default()
});
chan.on_message(move |e| {
let count = count.clone();
let s1 = s1.clone();
async move {
println!("message:{e:?}");
let mut count = count.lock();
*count += 1;
if e.is_params_key("act1") && e.is_state(MessageState::Created) {
s1.send(*count);
}
}
});
let pid = utils::longid();
let proc = rt.create_proc(&pid, &model);
rt.launch(&proc).await.unwrap();
let _ = sig.recv().await;
assert_eq!(
manager
.msg()
.list(
&Query::new()
.filter(Filter::and().expr(Expr::eq("pid", &pid)))
.offset(0)
.limit(1)
)
.await
.unwrap()
.rows
.len(),
1
);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_manager_messages_offset_in_range() {
let engine = Engine::new().start().await.unwrap();
let manager = engine.executor();
let model = Workflow::new().with_step(|step| {
step.with_id("step1")
.with_uses(USES_IRQ, Vars::new().with("key", "act1"))
});
let rt = engine.runtime();
let sig = engine.signal(());
let chan = engine.channel_with_options(&ChannelOptions {
ack: true,
..Default::default()
});
chan.on_message(move |e| async move {
println!("message:{e:?}");
});
let pid = utils::longid();
let proc = rt.create_proc(&pid, &model);
rt.launch(&proc).await.unwrap();
sig.timeout(100).await;
assert_eq!(
manager
.msg()
.list(
&Query::new()
.filter(Filter::and().expr(Expr::eq("pid", &pid)))
.offset(1)
.limit(2)
)
.await
.unwrap()
.rows
.len(),
2
);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_manager_messages_offset_out_range() {
let engine = Engine::new().start().await.unwrap();
let manager = engine.executor();
let model = Workflow::new().with_step(|step| {
step.with_id("step1")
.with_uses(USES_IRQ, Vars::new().with("key", "act1"))
});
let rt = engine.runtime();
let (sig, s1) = engine.signal(0).double();
let count = Arc::new(Mutex::new(0));
let chan = engine.channel_with_options(&ChannelOptions {
ack: true,
..Default::default()
});
chan.on_message(move |e| {
let count = count.clone();
let s1 = s1.clone();
async move {
println!("message:{e:?}");
let mut count = count.lock();
*count += 1;
if e.is_params_key("act1") && e.is_state(MessageState::Created) {
s1.send(*count);
}
}
});
let pid = utils::longid();
let proc = rt.create_proc(&pid, &model);
rt.launch(&proc).await.unwrap();
let _ = sig.recv().await;
assert_eq!(
manager
.msg()
.list(
&Query::new()
.filter(Filter::and().expr(Expr::eq("pid", &pid)))
.offset(1000)
.limit(100)
)
.await
.unwrap()
.rows
.len(),
0
);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_manager_message_get() {
let engine = Engine::new().start().await.unwrap();
let manager = engine.executor();
let model = Workflow::new().with_step(|step| {
step.with_id("step1")
.with_uses(USES_IRQ, Vars::new().with("key", "act1"))
});
let rt = engine.runtime();
let (sig, s1) = engine.signal(0).double();
let count = Arc::new(Mutex::new(0));
let chan = engine.channel_with_options(&ChannelOptions {
ack: true,
..Default::default()
});
chan.on_message(move |e| {
let count = count.clone();
let s1 = s1.clone();
async move {
println!("message:{e:?}");
let mut count = count.lock();
*count += 1;
if e.is_params_key("act1") && e.is_state(MessageState::Created) {
s1.send(*count);
}
}
});
let pid = utils::longid();
let proc = rt.create_proc(&pid, &model);
rt.launch(&proc).await.unwrap();
sig.recv().await;
let messages = manager
.msg()
.list(
&Query::new()
.filter(Filter::and().expr(Expr::eq("pid", &pid)))
.offset(0)
.limit(1000),
)
.await
.unwrap();
let message = messages.rows[0].clone();
let m = manager.msg().get(&message.id).await.unwrap();
assert_eq!(m.id, message.id);
assert_eq!(m.name, message.name);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_manager_message_rm() {
let engine = Engine::new().start().await.unwrap();
let manager = engine.executor();
let model = Workflow::new().with_step(|step| {
step.with_id("step1")
.with_uses(USES_IRQ, Vars::new().with("key", "act1"))
});
let rt = engine.runtime();
let (sig, s1) = engine.signal(0).double();
let count = Arc::new(Mutex::new(0));
let chan = engine.channel_with_options(&ChannelOptions {
ack: true,
..Default::default()
});
chan.on_message(move |e| {
let count = count.clone();
let s1 = s1.clone();
async move {
println!("message:{e:?}");
let mut count = count.lock();
*count += 1;
if e.is_type("act") && e.is_params_key("act1") && e.is_state(MessageState::Created) {
s1.send(*count);
}
}
});
let pid = utils::longid();
let proc = rt.create_proc(&pid, &model);
rt.launch(&proc).await.unwrap();
sig.recv().await;
let messages = manager
.msg()
.list(
&Query::new()
.filter(Filter::and().expr(Expr::eq("pid", &pid)))
.offset(0)
.limit(1),
)
.await
.unwrap();
let message = messages.rows[0].clone();
let ret = manager.msg().rm(&message.id).await.unwrap();
assert!(ret);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_manager_packages_count() {
let engine = Engine::new().start().await.unwrap();
let manager = engine.executor();
let count = 5;
for i in 0..count {
let package = Package {
id: utils::longid(),
name: format!("test-{}", i + 1),
desc: i.to_string(),
icon: "icon".to_string(),
doc: "doc".to_string(),
version: "0.1.0".to_string(),
schema: "{}".to_string(),
options: None,
run_as: crate::ActRunAs::Func,
resources: "[]".to_string(),
catalog: crate::package::ActPackageCatalog::Core,
create_time: utils::time::time_millis(),
update_time: 0,
timestamp: utils::time::timestamp(),
built_in: false,
v: 0,
};
manager.pack().publish(&package).await.unwrap();
}
assert_eq!(
manager
.pack()
.list(
&Query::new()
.filter(Filter::and().expr(Expr::eq("built_in", false)))
.offset(0)
.limit(1000)
)
.await
.unwrap()
.count,
count
);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_manager_packages_order() {
let engine = Engine::new().start().await.unwrap();
let manager = engine.executor();
let count = 5;
for i in 0..count {
let package = Package {
id: utils::longid(),
name: format!("test-{}", i + 1),
desc: format!("test-{}", i + 1),
icon: "icon".to_string(),
doc: "doc".to_string(),
version: "0.1.0".to_string(),
schema: "{}".to_string(),
options: None,
run_as: crate::ActRunAs::Func,
resources: "[]".to_string(),
catalog: crate::package::ActPackageCatalog::Core,
create_time: utils::time::time_millis(),
update_time: 0,
timestamp: utils::time::timestamp(),
built_in: false,
v: 0,
};
manager.pack().publish(&package).await.unwrap();
}
assert_eq!(
manager
.pack()
.list(&Query::new().order("timestamp", Sort::Desc))
.await
.unwrap()
.rows
.first()
.unwrap()
.desc,
"test-5"
);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_manager_packages_query() {
let engine = Engine::new().start().await.unwrap();
let manager = engine.executor();
let count = 5;
for i in 0..count {
let package = Package {
id: utils::longid(),
name: format!("test-{}", i + 1),
desc: format!("test-{}", i + 1),
icon: "icon".to_string(),
doc: "doc".to_string(),
version: "0.1.0".to_string(),
schema: "{}".to_string(),
options: None,
run_as: crate::ActRunAs::Func,
resources: "[]".to_string(),
catalog: crate::package::ActPackageCatalog::Core,
create_time: utils::time::time_millis(),
update_time: 0,
timestamp: utils::time::timestamp(),
built_in: false,
v: 0,
};
manager.pack().publish(&package).await.unwrap();
}
let rows = manager
.pack()
.list(&Query::new().filter(Filter::and().expr(Expr::eq("desc", "test-3"))))
.await
.unwrap()
.rows;
assert_eq!(rows.len(), 1);
assert_eq!(rows.first().unwrap().desc, "test-3");
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_manager_packages_offset_in_range() {
let engine = Engine::new().start().await.unwrap();
let manager = engine.executor();
let count = 5;
for i in 0..count {
let package = Package {
id: utils::longid(),
name: format!("test-{}", i + 1),
desc: i.to_string(),
icon: "icon".to_string(),
doc: "doc".to_string(),
version: "0.1.0".to_string(),
schema: "{}".to_string(),
options: None,
run_as: crate::ActRunAs::Func,
resources: "[]".to_string(),
catalog: crate::package::ActPackageCatalog::Core,
create_time: utils::time::time_millis(),
update_time: 0,
timestamp: utils::time::timestamp(),
built_in: false,
v: 0,
};
manager.pack().publish(&package).await.unwrap();
}
assert_eq!(
manager
.pack()
.list(
&Query::new()
.filter(Filter::and().expr(Expr::eq("built_in", false)))
.offset(4)
.limit(1000)
)
.await
.unwrap()
.rows
.len(),
1
);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_manager_packages_offset_out_range() {
let engine = Engine::new().start().await.unwrap();
let manager = engine.executor();
let count = 5;
for i in 0..count {
let package = Package {
id: utils::longid(),
name: format!("test-{}", i + 1),
desc: "desc".to_string(),
icon: "icon".to_string(),
doc: "doc".to_string(),
version: "0.1.0".to_string(),
schema: "{}".to_string(),
options: None,
run_as: crate::ActRunAs::Func,
resources: "[]".to_string(),
catalog: crate::package::ActPackageCatalog::Core,
create_time: utils::time::time_millis(),
update_time: 0,
timestamp: utils::time::timestamp(),
built_in: false,
v: 0,
};
manager.pack().publish(&package).await.unwrap();
}
assert_eq!(
manager
.pack()
.list(&Query::new().offset(1000).limit(1000))
.await
.unwrap()
.rows
.len(),
0
);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_manager_package_rm() {
let engine = Engine::new().start().await.unwrap();
let manager = engine.executor();
let package = Package {
id: utils::longid(),
name: "test name".to_string(),
desc: "desc".to_string(),
icon: "icon".to_string(),
doc: "doc".to_string(),
version: "0.1.0".to_string(),
schema: "{}".to_string(),
options: None,
run_as: crate::ActRunAs::Func,
resources: "[]".to_string(),
catalog: crate::package::ActPackageCatalog::Core,
create_time: utils::time::time_millis(),
update_time: 0,
timestamp: utils::time::timestamp(),
built_in: false,
v: 0,
};
manager.pack().publish(&package).await.unwrap();
assert!(manager.pack().rm(&package.id).await.unwrap());
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_executor_start() {
let engine = Engine::new().start().await.unwrap();
let model = Workflow::new()
.with_id(&utils::longid())
.with_step(|step| step.with_id("step1"));
let sig = engine.signal(());
let s1 = sig.clone();
engine.channel().on_complete(move |_| {
let s1 = s1.clone();
async move {
s1.close();
}
});
engine
.executor()
.model()
.deploy(&model, None)
.await
.unwrap();
let pid = utils::longid();
let mut vars = Vars::new();
vars.insert("uid".to_string(), json!("u1"));
vars.insert("pid".to_string(), json!(pid));
let result = engine.executor().proc().start(&model.id, vars).await;
sig.recv().await;
assert!(result.is_ok());
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_executor_start_not_found_model() {
let engine = Engine::new().start().await.unwrap();
let sig = engine.signal(());
let s1 = sig.clone();
engine.channel().on_complete(move |_| {
let s1 = s1.clone();
async move {
s1.close();
}
});
let pid = utils::longid();
let mut vars = Vars::new();
vars.insert("uid".to_string(), json!("u1"));
vars.insert("pid".to_string(), json!(pid));
let result = engine.executor().proc().start("not_exists", vars).await;
assert!(result.is_err());
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_executor_complete_normal() {
let engine = Engine::new().start().await.unwrap();
let model = Workflow::new().with_step(|step| {
step.with_id("step1")
.with_uses(USES_IRQ, Vars::new().with("key", "act1"))
});
let rt = engine.runtime();
let sig = engine.signal(false);
let s1 = sig.clone();
let executor = engine.executor();
engine.channel().on_message(move |e| {
let executor = executor.clone();
let s1 = s1.clone();
async move {
if e.is_params_key("act1") && e.is_state(MessageState::Created) {
let mut vars = Vars::new();
vars.insert("uid".to_string(), json!("u1"));
let ret = executor.act().complete(&e.pid, &e.tid, vars).await;
s1.send(ret.is_ok());
}
}
});
let mut vars = Vars::new();
vars.insert("uid".to_string(), json!("u1"));
rt.start(&model, vars).await.unwrap();
let ret = sig.recv().await;
assert!(ret);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_executor_complete_no_uid() {
let engine = Engine::new().start().await.unwrap();
let model = Workflow::new().with_step(|step| {
step.with_id("step1")
.with_uses(USES_IRQ, Vars::new().with("key", "act1"))
});
let rt = engine.runtime();
let sig = engine.signal(false);
let s1 = sig.clone();
let executor = engine.executor();
engine.channel().on_message(move |e| {
let executor = executor.clone();
let s1 = s1.clone();
async move {
if e.is_params_key("act1") && e.is_state(MessageState::Created) {
let vars = Vars::new();
let ret = executor.act().complete(&e.pid, &e.tid, vars).await;
s1.send(ret.is_ok());
}
}
});
rt.start(&model, Vars::new()).await.unwrap();
let ret = sig.recv().await;
assert!(ret);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_executor_submit() {
let engine = Engine::new().start().await.unwrap();
let model = Workflow::new().with_step(|step| {
step.with_id("step1")
.with_uses(USES_IRQ, Vars::new().with("key", "act1"))
});
let rt = engine.runtime();
let sig = engine.signal(false);
let s1 = sig.clone();
let executor = engine.executor();
engine.channel().on_message(move |e| {
let executor = executor.clone();
let s1 = s1.clone();
async move {
if e.is_params_key("act1") && e.is_state(MessageState::Created) {
let mut vars = Vars::new();
vars.insert("uid".to_string(), json!("u1"));
let ret = executor.act().submit(&e.pid, &e.tid, vars).await;
s1.send(ret.is_ok());
}
}
});
let mut vars = Vars::new();
vars.insert("uid".to_string(), json!("u1"));
rt.start(&model, vars).await.unwrap();
let ret = sig.recv().await;
assert!(ret);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_executor_skip() {
let engine = Engine::new().start().await.unwrap();
let model = Workflow::new().with_step(|step| {
step.with_id("step1")
.with_uses(USES_IRQ, Vars::new().with("key", "act1"))
});
let rt = engine.runtime();
let sig = engine.signal(false);
let s1 = sig.clone();
let executor = engine.executor();
engine.channel().on_message(move |e| {
let executor = executor.clone();
let s1 = s1.clone();
async move {
if e.params().unwrap().get::<String>("key").as_deref() == Some("act1")
&& e.is_state(MessageState::Created)
{
let mut vars = Vars::new();
vars.insert("uid".to_string(), json!("u1"));
let ret = executor.act().skip(&e.pid, &e.tid, vars).await;
s1.send(ret.is_ok());
}
}
});
let mut vars = Vars::new();
vars.insert("uid".to_string(), json!("u1"));
rt.start(&model, vars).await.unwrap();
let ret = sig.recv().await;
assert!(ret);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_executor_error() {
let engine = Engine::new().start().await.unwrap();
let model = Workflow::new().with_step(|step| {
step.with_id("step1")
.with_uses(USES_IRQ, Vars::new().with("key", "act1"))
});
let rt = engine.runtime();
let sig = engine.signal(false);
let s1 = sig.clone();
let executor = engine.executor();
engine.channel().on_message(move |e| {
let executor = executor.clone();
let s1 = s1.clone();
async move {
if e.params().unwrap().get::<String>("key").as_deref() == Some("act1")
&& e.is_state(MessageState::Created)
{
let mut vars = Vars::new();
vars.insert("uid".to_string(), json!("u1"));
vars.insert("ecode".to_string(), json!("code_1"));
let ret = executor.act().fail(&e.pid, &e.tid, vars).await;
s1.send(ret.is_ok());
}
}
});
let mut vars = Vars::new();
vars.insert("uid".to_string(), json!("u1"));
rt.start(&model, vars).await.unwrap();
let ret = sig.recv().await;
assert!(ret);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_executor_abort() {
let engine = Engine::new().start().await.unwrap();
let model = Workflow::new().with_step(|step| {
step.with_id("step1")
.with_uses(USES_IRQ, Vars::new().with("key", "act1"))
});
let rt = engine.runtime();
let sig = engine.signal(false);
let s1 = sig.clone();
let executor = engine.executor();
engine.channel().on_message(move |e| {
let executor = executor.clone();
let s1 = s1.clone();
async move {
println!("message: {:?}", e.inner());
if e.is_params_key("act1") && e.is_state(MessageState::Created) {
let mut vars = Vars::new();
vars.insert("uid".to_string(), json!("u1"));
let ret = executor.act().abort(&e.pid, &e.tid, vars).await;
s1.send(ret.is_ok());
}
}
});
let mut vars = Vars::new();
vars.insert("uid".to_string(), json!("u1"));
rt.start(&model, vars).await.unwrap();
let ret = sig.recv().await;
assert!(ret);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_executor_back() {
let engine = Engine::new().start().await.unwrap();
let model = Workflow::new()
.with_step(|step| {
step.with_id("step1")
.with_uses(USES_IRQ, Vars::new().with("key", "act1"))
})
.with_step(|step| {
step.with_id("step2")
.with_uses(USES_IRQ, Vars::new().with("key", "act2"))
});
let rt = engine.runtime();
let sig = engine.signal(false);
let s1 = sig.clone();
let executor = engine.executor();
let count = Arc::new(Mutex::new(0));
engine.channel().on_message(move |e| {
let executor = executor.clone();
let count = count.clone();
let s1 = s1.clone();
async move {
println!("message: {e:?}");
if e.is_params_key("act1") && e.is_state(MessageState::Created) {
let mut vars = Vars::new();
vars.insert("uid".to_string(), json!("u1"));
let count = {
let mut count = count.lock();
*count += 1;
*count
};
if count == 2 {
s1.close();
}
executor.act().complete(&e.pid, &e.tid, vars).await.unwrap();
}
if e.is_params_key("act2") && e.is_state(MessageState::Created) {
let mut vars = Vars::new();
vars.insert("uid".to_string(), json!("u1"));
vars.insert("to".to_string(), json!("step1"));
let ret = executor.act().back(&e.pid, &e.tid, vars).await;
s1.update(|data| *data = ret.is_ok());
}
}
});
let mut vars = Vars::new();
vars.insert("uid".to_string(), json!("u1"));
rt.start(&model, vars).await.unwrap();
let ret = sig.recv().await;
assert!(ret);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_executor_cancel() {
let engine = Engine::new().start().await.unwrap();
let model = Workflow::new()
.with_step(|step| {
step.with_id("step1")
.with_uses(USES_IRQ, Vars::new().with("key", "act1"))
})
.with_step(|step| {
step.with_id("step2")
.with_uses(USES_IRQ, Vars::new().with("key", "act2"))
});
let rt = engine.runtime();
let sig = engine.signal(false);
let s1 = sig.clone();
let executor = engine.executor();
let count = Arc::new(Mutex::new(0));
let tid = Arc::new(Mutex::new("".to_string()));
engine.channel().on_message(move |e| {
let executor = executor.clone();
let count = count.clone();
let tid = tid.clone();
let s1 = s1.clone();
async move {
println!("message: {:?}", e.inner());
if e.is_params_key("act1") && e.is_state(MessageState::Created) {
let mut vars = Vars::new();
vars.insert("uid".to_string(), json!("u1"));
let count = {
let mut count = count.lock();
*count += 1;
*count
};
if count == 2 {
s1.close();
}
executor.act().complete(&e.pid, &e.tid, vars).await.unwrap();
*tid.lock() = e.tid.clone();
}
if e.is_params_key("act2") && e.is_state(MessageState::Created) {
let mut vars = Vars::new();
vars.insert("uid".to_string(), json!("u1"));
let tid_val = tid.lock().clone();
let ret = executor.act().cancel(&e.pid, &tid_val, vars).await;
s1.update(|data| *data = ret.is_ok());
}
}
});
let mut vars = Vars::new();
vars.insert("uid".to_string(), json!("u1"));
rt.start(&model, vars).await.unwrap();
let ret = sig.recv().await;
assert!(ret);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_executor_push() {
let engine = Engine::new().start().await.unwrap();
let model = Workflow::new().with_step(|step| {
step.with_id("step1")
.with_uses(USES_IRQ, Vars::new().with("key", "act1"))
});
let rt = engine.runtime();
let sig = engine.signal(false);
let s1 = sig.clone();
let executor = engine.executor();
engine.channel().on_message(move |e| {
let executor = executor.clone();
let s1 = s1.clone();
async move {
println!("message: {e:?}");
if e.is_type("step") && e.is_state(MessageState::Created) {
let vars = Vars::new()
.with("uses", "acts.core.irq")
.with("params", Vars::new().with("key", "act2"));
executor.act().push(&e.pid, &e.tid, vars).await.unwrap();
}
if e.is_type("act") && e.is_params_key("act2") && e.is_state(MessageState::Created) {
s1.send(true);
}
}
});
let mut vars = Vars::new();
vars.insert("uid".to_string(), json!("u1"));
rt.start(&model, vars).await.unwrap();
let ret = sig.recv().await;
assert!(ret);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_executor_push_no_key_error() {
let engine = Engine::new().start().await.unwrap();
let model = Workflow::new().with_step(|step| {
step.with_id("step1")
.with_uses(USES_IRQ, Vars::new().with("key", "act1"))
});
let rt = engine.runtime();
let sig = engine.signal(false);
let s1 = sig.clone();
let executor = engine.executor();
engine.channel().on_message(move |e| {
let executor = executor.clone();
let s1 = s1.clone();
async move {
println!("message: {e:?}");
if e.is_type("step") && e.is_state(MessageState::Created) {
s1.send(
executor
.act()
.push(&e.pid, &e.tid, Vars::new())
.await
.is_err(),
);
}
}
});
let mut vars = Vars::new();
vars.insert("uid".to_string(), json!("u1"));
rt.start(&model, vars).await.unwrap();
let ret = sig.recv().await;
assert!(ret);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_executor_push_not_step_id_error() {
let engine = Engine::new().start().await.unwrap();
let model = Workflow::new().with_step(|step| {
step.with_id("step1")
.with_uses(USES_IRQ, Vars::new().with("key", "act1"))
});
let rt = engine.runtime();
let sig = engine.signal(false);
let s1 = sig.clone();
let executor = engine.executor();
engine.channel().on_message(move |e| {
let executor = executor.clone();
let s1 = s1.clone();
async move {
println!("message: {e:?}");
if e.is_type("act")
&& e.params().unwrap().get::<String>("key").as_deref() == Some("act1")
&& e.is_state(MessageState::Created)
{
let vars = Vars::new();
s1.send(executor.act().push(&e.pid, &e.tid, vars).await.is_err());
}
}
});
let mut vars = Vars::new();
vars.insert("uid".to_string(), json!("u1"));
rt.start(&model, vars).await.unwrap();
let ret = sig.recv().await;
assert!(ret);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_executor_remove() {
let engine = Engine::new().start().await.unwrap();
let model = Workflow::new().with_step(|step| {
step.with_id("step1")
.with_uses(USES_IRQ, Vars::new().with("key", "act1"))
});
let rt = engine.runtime();
let sig = engine.signal(false);
let s1 = sig.clone();
let executor = engine.executor();
engine.channel().on_message(move |e| {
let executor = executor.clone();
let s1 = s1.clone();
async move {
println!("message: {e:?}");
if e.params().unwrap().get::<String>("key").as_deref() == Some("act1")
&& e.is_state(MessageState::Created)
{
s1.send(
executor
.act()
.remove(&e.pid, &e.tid, Vars::new())
.await
.is_ok(),
);
}
}
});
let mut vars = Vars::new();
vars.insert("uid".to_string(), json!("u1"));
rt.start(&model, vars).await.unwrap();
let ret = sig.recv().await;
assert!(ret);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_extender_set_process_var() {
let engine = Engine::new().start().await.unwrap();
let model = Workflow::new().with_step(|step| {
step.with_id("step1")
.with_uses(USES_IRQ, Vars::new().with("key", "act1"))
});
let rt = engine.runtime();
let sig = engine.signal(false);
let s1 = sig.clone();
let executor = engine.executor();
engine.channel().on_message(move |e| {
let executor = executor.clone();
let s1 = s1.clone();
async move {
println!("message: {e:?}");
if e.params().unwrap().get::<String>("key").as_deref() == Some("act1")
&& e.is_state(MessageState::Created)
{
s1.send(
executor
.act()
.set_process_vars(&e.pid, &e.tid, Vars::new().with("var1", 1))
.await
.is_ok(),
);
}
}
});
let mut vars = Vars::new();
vars.set("var1", 0);
rt.start(&model, vars).await.unwrap();
let ret = sig.recv().await;
assert!(ret);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_extender_register_module() {
let engine = Engine::new().start().await.unwrap();
let extender = engine.extender();
let before_count = engine.runtime().env().user_env_count();
let module = test_module::TestModule;
extender.register_var(&module);
let count = engine.runtime().env().user_env_count();
assert_eq!(count, before_count + 1);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_emitter_default() {
let engine = Engine::new().start().await.unwrap();
let emitter = engine.channel();
let sig = engine.signal::<Vec<Message>>(Vec::new());
let s = sig.clone();
emitter.on_message(move |e| {
let s = s.clone();
async move {
s.update(|data| data.push(e.inner().clone()));
s.close();
}
});
let msg = Message::default();
engine.runtime().emitter().emit_message(&msg);
let ret = sig.recv().await;
assert_eq!(ret.len(), 1);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_emitter_type_match() {
let engine = Engine::new().start().await.unwrap();
let emitter = engine.channel_with_options(&ChannelOptions {
r#type: "a*".to_string(),
..Default::default()
});
let sig = engine.signal::<Vec<Message>>(Vec::new());
let s = sig.clone();
emitter.on_message(move |e| {
let s = s.clone();
async move {
s.update(|data| data.push(e.inner().clone()));
s.close();
}
});
let msg = Message {
r#type: "abc".to_string(),
..Message::default()
};
engine.runtime().emitter().emit_message(&msg);
let ret = sig.recv().await;
assert_eq!(ret.len(), 1);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_emitter_type_not_match() {
let engine = Engine::new().start().await.unwrap();
let emitter = engine.channel_with_options(&ChannelOptions {
r#type: "a*".to_string(),
..Default::default()
});
let sig = engine.signal::<Vec<Message>>(Vec::new());
let s = sig.clone();
emitter.on_message(move |e| {
let s = s.clone();
async move {
s.update(|data| data.push(e.inner().clone()));
s.close();
}
});
let msg = Message {
r#type: "bac".to_string(),
..Message::default()
};
engine.runtime().emitter().emit_message(&msg);
let ret = sig.timeout(100).await;
assert_eq!(ret.len(), 0);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_emitter_state_match() {
let engine = Engine::new().start().await.unwrap();
let emitter = engine.channel_with_options(&ChannelOptions {
state: "completed".to_string(),
..Default::default()
});
let sig = engine.signal::<Vec<Message>>(Vec::new());
let s = sig.clone();
emitter.on_message(move |e| {
let s = s.clone();
async move {
s.update(|data| data.push(e.inner().clone()));
s.close();
}
});
let msg = Message {
state: MessageState::Completed,
..Message::default()
};
engine.runtime().emitter().emit_message(&msg);
let ret = sig.recv().await;
assert_eq!(ret.len(), 1);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_emitter_state_not_match() {
let engine = Engine::new().start().await.unwrap();
let emitter = engine.channel_with_options(&ChannelOptions {
r#type: "error".to_string(),
..Default::default()
});
let sig = engine.signal::<Vec<Message>>(Vec::new());
let s = sig.clone();
emitter.on_message(move |e| {
let s = s.clone();
async move {
s.update(|data| data.push(e.inner().clone()));
s.close();
}
});
let msg = Message {
state: MessageState::Completed,
..Message::default()
};
engine.runtime().emitter().emit_message(&msg);
let ret = sig.timeout(100).await;
assert_eq!(ret.len(), 0);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_emitter_tag_match() {
let engine = Engine::new().start().await.unwrap();
let emitter = engine.channel_with_options(&ChannelOptions {
options: Vars::new().with("tag", "tag*"),
..Default::default()
});
let sig = engine.signal::<Vec<Message>>(Vec::new());
let s = sig.clone();
emitter.on_message(move |e| {
let s = s.clone();
async move {
s.update(|data| data.push(e.inner().clone()));
}
});
let msg = Message {
inputs: Vars::new().with("options", Vars::new().with("tag", "tag1")),
..Message::default()
};
engine.runtime().emitter().emit_message(&msg);
let msg = Message {
inputs: Vars::new().with("options", Vars::new().with("tag", "aaaa")),
..Message::default()
};
engine.runtime().emitter().emit_message(&msg);
let ret = sig.timeout(100).await;
assert_eq!(ret.len(), 1);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_emitter_tag_not_match() {
let engine = Engine::new().start().await.unwrap();
let emitter = engine.channel_with_options(&ChannelOptions {
options: Vars::new().with("tag", "tag*"),
..Default::default()
});
let sig = engine.signal::<Vec<Message>>(Vec::new());
let s = sig.clone();
emitter.on_message(move |e| {
let s = s.clone();
async move {
s.update(|data| data.push(e.inner().clone()));
s.close();
}
});
let msg = Message {
inputs: Vars::new().with("options", Vars::new().with("tag", "aaaa")),
..Message::default()
};
engine.runtime().emitter().emit_message(&msg);
let ret = sig.timeout(100).await;
assert_eq!(ret.len(), 0);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_emitter_on_message_with_dup_id() {
let engine = Engine::new().start().await.unwrap();
let emitter = engine.channel_with_options(&ChannelOptions {
id: "dup_id".to_string(),
..Default::default()
});
let sig = engine.signal::<Vec<Message>>(Vec::new());
let s = sig.clone();
emitter.on_message(move |e| {
let s = s.clone();
async move {
s.update(|data| data.push(e.inner().clone()));
s.close();
}
});
let s2 = sig.clone();
let emitter2 = engine.channel_with_options(&ChannelOptions {
id: "dup_id".to_string(),
..Default::default()
});
emitter2.on_message(move |e| {
let s2 = s2.clone();
async move {
println!("message: {e:?}");
s2.update(|data| data.push(e.inner().clone()));
s2.close();
}
});
let msg = Message {
inputs: Vars::new().with("options", Vars::new().with("tag", "aaaa")),
..Message::default()
};
engine.runtime().emitter().emit_message(&msg);
let ret = sig.recv().await;
assert_eq!(ret.len(), 1);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_message_store_with_emit_id() {
let engine = Engine::new().start().await.unwrap();
let emitter = engine.channel_with_options(&ChannelOptions {
id: "my_emit_id".to_string(),
ack: true,
..Default::default()
});
let (s1, s2) = engine.signal::<Message>(Message::default()).double();
emitter.on_message(move |e| {
let s1 = s1.clone();
async move {
s1.send(e.inner().clone());
}
});
let msg = Message {
id: "1".to_string(),
..Message::default()
};
engine.runtime().emitter().emit_message(&msg);
let ret = s2.recv().await;
assert_eq!(ret.id, "1");
assert!(
engine
.runtime()
.cache()
.store()
.messages()
.exists("1")
.await
.unwrap()
);
let delivery = engine
.runtime()
.cache()
.store()
.deliveries()
.find(&ret.delivery_id.clone().unwrap())
.await
.unwrap();
assert_eq!(delivery.msg_id, "1");
assert_eq!(delivery.chan_id, "my_emit_id");
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_message_store_with_emit_id_and_options() {
let engine = Engine::new().start().await.unwrap();
let emitter = engine.channel_with_options(&ChannelOptions {
id: "my_emit_id".to_string(),
options: Vars::new().with("tag", "tag*"),
ack: true,
..Default::default()
});
let (s1, s2) = engine.signal::<Message>(Message::default()).double();
emitter.on_message(move |e| {
let s1 = s1.clone();
async move {
s1.send(e.inner().clone());
}
});
let msg = Message {
id: utils::longid(),
inputs: Vars::new().with("options", Vars::new().with("tag", "tagaaaa")),
..Message::default()
};
engine.runtime().emitter().emit_message(&msg);
let ret = s2.recv().await;
assert_eq!(ret.id, msg.id);
let delivery = engine
.runtime()
.cache()
.store()
.deliveries()
.find(&ret.delivery_id.clone().unwrap())
.await
.unwrap();
let pattern = serde_json::from_str::<Vars>(&delivery.chan_pattern).unwrap();
assert_eq!(delivery.msg_id, msg.id);
assert_eq!(delivery.chan_id, "my_emit_id");
assert_eq!(pattern.get::<String>("tag").unwrap(), "tag*");
assert!(pattern.get::<bool>("ack").unwrap());
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_message_multi_channels_share_message_single_delivery_each() {
let engine = Engine::new().start().await.unwrap();
let (s1, r1) = engine.signal::<Message>(Message::default()).double();
let emitter1 = engine.channel_with_options(&ChannelOptions {
id: "chan_a".to_string(),
ack: true,
..Default::default()
});
let s1a = s1.clone();
emitter1.on_message(move |e| {
let s1a = s1a.clone();
async move {
s1a.send(e.inner().clone());
}
});
let (s2, r2) = engine.signal::<Message>(Message::default()).double();
let emitter2 = engine.channel_with_options(&ChannelOptions {
id: "chan_b".to_string(),
ack: true,
..Default::default()
});
let s2a = s2.clone();
emitter2.on_message(move |e| {
let s2a = s2a.clone();
async move {
s2a.send(e.inner().clone());
}
});
let msg = Message {
id: utils::longid(),
..Message::default()
};
engine.runtime().emitter().emit_message(&msg);
let (received_a, received_b) = tokio::join!(r1.recv(), r2.recv());
assert_eq!(received_a.id, msg.id);
assert_eq!(received_b.id, msg.id);
let store = engine.runtime().cache().store();
let messages = store
.messages()
.query(&Query::new().filter(Filter::and().expr(Expr::eq("id", msg.id.clone()))))
.await
.unwrap();
assert_eq!(messages.count, 1);
let q = Query::new().filter(Filter::and().expr(Expr::eq("msg_id", msg.id.clone())));
let deliveries = store.deliveries().query(&q).await.unwrap();
assert_eq!(deliveries.count, 2);
let mut chan_ids: Vec<String> = deliveries.rows.iter().map(|d| d.chan_id.clone()).collect();
chan_ids.sort();
assert_eq!(chan_ids, vec!["chan_a".to_string(), "chan_b".to_string()]);
assert_ne!(
deliveries.rows[0].id, deliveries.rows[1].id,
"each delivery must have its own delivery id"
);
let delivery_a = received_a.delivery_id.clone().unwrap();
let delivery_b = received_b.delivery_id.clone().unwrap();
assert_ne!(delivery_a, delivery_b);
engine.executor().msg().ack(&delivery_a).await.unwrap();
let row_a = store.deliveries().find(&delivery_a).await.unwrap();
assert_eq!(row_a.status, data::DeliveryStatus::Acked);
let row_b = store.deliveries().find(&delivery_b).await.unwrap();
assert_eq!(row_b.status, data::DeliveryStatus::Delivered);
}
#[serial]
#[tokio::test(flavor = "multi_thread")]
async fn export_message_not_store_without_match() {
let engine = Engine::new().start().await.unwrap();
let emitter = engine.channel_with_options(&ChannelOptions {
id: "my_emit_id".to_string(),
options: Vars::new().with("tag", "tag*"),
..Default::default()
});
let (s1, s2) = engine.signal::<Message>(Message::default()).double();
emitter.on_message(move |e| {
let s1 = s1.clone();
async move {
s1.send(e.inner().clone());
}
});
let msg = Message {
id: utils::longid(),
inputs: Vars::new().with("options", Vars::new().with("tag", "not_match_tag")),
..Message::default()
};
engine.runtime().emitter().emit_message(&msg);
s2.timeout(20).await;
assert!(
!engine
.runtime()
.cache()
.store()
.messages()
.exists(&msg.id)
.await
.unwrap()
);
}
#[serial]
#[tokio::test(flavor = "multi_thread")]
async fn export_message_not_store_with_empty_emit_id_and_not_match_option() {
let engine = Engine::new().start().await.unwrap();
let emitter = engine.channel_with_options(&ChannelOptions {
id: "".to_string(),
..Default::default()
});
let (s1, s2) = engine.signal::<Message>(Message::default()).double();
emitter.on_message(move |e| {
let s1 = s1.clone();
async move {
s1.send(e.inner().clone());
}
});
let msg = Message {
id: "1".to_string(),
..Message::default()
};
engine.runtime().emitter().emit_message(&msg);
s2.timeout(20).await;
assert!(
!engine
.runtime()
.cache()
.store()
.messages()
.exists("1")
.await
.unwrap()
);
}
#[serial]
#[tokio::test(flavor = "multi_thread")]
async fn export_message_clear_error_messages_by_none() {
let engine = Engine::new().start().await.unwrap();
let delivery = data::Delivery {
id: utils::longid(),
msg_id: utils::longid(),
status: data::DeliveryStatus::Error,
..data::Delivery::default()
};
engine
.runtime()
.cache()
.store()
.deliveries()
.create(&delivery)
.await
.unwrap();
let ret = engine
.runtime()
.cache()
.store()
.deliveries()
.find(&delivery.id)
.await
.unwrap();
assert_eq!(ret.status, data::DeliveryStatus::Error);
engine.executor().msg().clear(None).await.unwrap();
assert!(
!engine
.runtime()
.cache()
.store()
.deliveries()
.exists(&delivery.id)
.await
.unwrap()
);
}
#[serial]
#[tokio::test(flavor = "multi_thread")]
async fn export_message_clear_error_messages_by_pid() {
let engine = Engine::new().start().await.unwrap();
let pid = utils::longid();
engine
.runtime()
.cache()
.store()
.deliveries()
.create(&data::Delivery {
id: utils::longid(),
msg_id: utils::longid(),
pid: pid.clone(),
status: data::DeliveryStatus::Error,
..data::Delivery::default()
})
.await
.unwrap();
engine
.runtime()
.cache()
.store()
.deliveries()
.create(&data::Delivery {
id: utils::longid(),
msg_id: utils::longid(),
pid: pid.clone(),
status: data::DeliveryStatus::Error,
..data::Delivery::default()
})
.await
.unwrap();
engine
.executor()
.msg()
.clear(Some(pid.clone()))
.await
.unwrap();
let deliveries = engine
.runtime()
.cache()
.store()
.deliveries()
.query(
&Query::new().filter(
Filter::and()
.expr(Expr::eq("pid", pid))
.expr(Expr::eq("status", data::DeliveryStatus::Error)),
),
)
.await
.unwrap();
assert_eq!(deliveries.rows.len(), 0);
}
#[serial]
#[tokio::test(flavor = "multi_thread")]
async fn export_message_resend_error_messages() {
let engine = Engine::new().start().await.unwrap();
let delivery = data::Delivery {
id: utils::longid(),
msg_id: utils::longid(),
status: data::DeliveryStatus::Error,
..data::Delivery::default()
};
engine
.runtime()
.cache()
.store()
.deliveries()
.create(&delivery)
.await
.unwrap();
let ret = engine
.runtime()
.cache()
.store()
.deliveries()
.find(&delivery.id)
.await
.unwrap();
assert_eq!(ret.status, data::DeliveryStatus::Error);
engine.executor().msg().redo().await.unwrap();
let ret = engine
.runtime()
.cache()
.store()
.deliveries()
.find(&delivery.id)
.await
.unwrap();
assert_eq!(ret.status, data::DeliveryStatus::Created);
assert_eq!(ret.retry_times, 0);
}
mod test_module {
use crate::{ActUserVar, Vars};
#[derive(Clone)]
pub struct TestModule;
impl ActUserVar for TestModule {
fn name(&self) -> String {
"test_env".to_string()
}
fn default_data(&self) -> Option<crate::Vars> {
Some(Vars::new().with("a", 10))
}
}
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_emitter_options_multi_key() {
let engine = Engine::new().start().await.unwrap();
let emitter = engine.channel_with_options(&ChannelOptions {
options: Vars::new().with("tag", "tag*").with("rn", "a:*"),
..Default::default()
});
let sig = engine.signal::<Vec<Message>>(Vec::new());
let s = sig.clone();
emitter.on_message(move |e| {
let s = s.clone();
async move {
s.update(|data| data.push(e.inner().clone()));
}
});
let msg = Message {
inputs: Vars::new().with(
"options",
Vars::new().with("tag", "tag1").with("rn", "a:b:c"),
),
..Message::default()
};
engine.runtime().emitter().emit_message(&msg);
let msg = Message {
inputs: Vars::new().with(
"options",
Vars::new().with("tag", "tag2").with("rn", "b:x:y"),
),
..Message::default()
};
engine.runtime().emitter().emit_message(&msg);
let ret = sig.timeout(100).await;
assert_eq!(ret.len(), 1);
let tag = ret[0]
.options()
.and_then(|o| o.get::<String>("tag"))
.unwrap_or_default();
assert_eq!(tag, "tag1");
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_emitter_options_custom_key() {
let engine = Engine::new().start().await.unwrap();
let emitter = engine.channel_with_options(&ChannelOptions {
options: Vars::new().with("priority", "high"),
..Default::default()
});
let sig = engine.signal::<Vec<Message>>(Vec::new());
let s = sig.clone();
emitter.on_message(move |e| {
let s = s.clone();
async move {
s.update(|data| data.push(e.inner().clone()));
}
});
let msg = Message {
inputs: Vars::new().with("options", Vars::new().with("priority", "high")),
..Message::default()
};
engine.runtime().emitter().emit_message(&msg);
let msg = Message {
inputs: Vars::new().with("options", Vars::new().with("priority", "low")),
..Message::default()
};
engine.runtime().emitter().emit_message(&msg);
let ret = sig.timeout(100).await;
assert_eq!(ret.len(), 1);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_emitter_options_missing_key() {
let engine = Engine::new().start().await.unwrap();
let emitter = engine.channel_with_options(&ChannelOptions {
options: Vars::new().with("nonexistent", "value"),
..Default::default()
});
let sig = engine.signal::<Vec<Message>>(Vec::new());
let s = sig.clone();
emitter.on_message(move |e| {
let s = s.clone();
async move {
s.update(|data| data.push(e.inner().clone()));
}
});
let msg = Message {
inputs: Vars::new().with("options", Vars::new().with("tag", "tag1")),
..Message::default()
};
engine.runtime().emitter().emit_message(&msg);
let ret = sig.timeout(100).await;
assert_eq!(ret.len(), 0);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_trigger_deploy_all_kinds() {
let engine = Engine::new().start().await.unwrap();
let manager = engine.executor();
let model = Workflow::new()
.with_id("trigger-model")
.with_trigger(|t| t.with_id("e-manual").with_kind("manual"))
.with_trigger(|t| t.with_id("e-chat").with_kind("chat"))
.with_trigger(|t| t.with_id("e-hook").with_kind("hook"))
.with_trigger(|t| {
t.with_id("e-schedule")
.with_kind("schedule")
.with_schedule("* * * * * *")
})
.with_step(|step| step.with_id("step1"));
manager.model().deploy(&model, None).await.unwrap();
let rows = manager
.evt()
.list(&Query::new().filter(Filter::and().expr(Expr::eq(consts::MODEL_ID, "trigger-model"))))
.await
.unwrap();
assert_eq!(rows.count, 4);
let manual = manager.evt().get("trigger-model:e-manual").await.unwrap();
assert_eq!(manual.kind, "manual");
let schedule = manager.evt().get("trigger-model:e-schedule").await.unwrap();
assert_eq!(schedule.kind, "schedule");
assert!(schedule.next_run > 0, "schedule must be armed on deploy");
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_trigger_redeploy_removes_stale_rows() {
let engine = Engine::new().start().await.unwrap();
let manager = engine.executor();
let mut model = Workflow::new()
.with_id("trigger-reconcile")
.with_trigger(|t| t.with_id("e1").with_kind("manual"))
.with_trigger(|t| t.with_id("e2").with_kind("manual"))
.with_step(|step| step.with_id("step1"));
manager.model().deploy(&model, None).await.unwrap();
assert_eq!(
manager
.evt()
.get("trigger-reconcile:e2")
.await
.unwrap()
.kind,
"manual"
);
model.on = vec![crate::Trigger {
id: "e1".to_string(),
kind: "manual".to_string(),
..Default::default()
}];
manager.model().deploy(&model, None).await.unwrap();
assert!(manager.evt().get("trigger-reconcile:e2").await.is_err());
assert!(manager.evt().get("trigger-reconcile:e1").await.is_ok());
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_trigger_redeploy_keeps_schedule_state() {
let engine = Engine::new().start().await.unwrap();
let manager = engine.executor();
let mut model = Workflow::new()
.with_id("trigger-reconcile-state")
.with_trigger(|t| {
t.with_id("e1")
.with_kind("schedule")
.with_schedule("0 0 12 * * *")
.with_params_vars(|vars| vars.with("a", 1))
})
.with_step(|step| step.with_id("step1"));
manager.model().deploy(&model, None).await.unwrap();
let first = manager
.evt()
.get("trigger-reconcile-state:e1")
.await
.unwrap();
assert_eq!(first.kind, "schedule");
assert!(first.next_run > 0, "schedule must be armed on deploy");
model.on = vec![crate::Trigger {
id: "e1".to_string(),
kind: "schedule".to_string(),
schedule: Some("0 0 12 * * *".to_string()),
params: serde_json::json!({"a": 2}),
..Default::default()
}];
let before = manager
.evt()
.get("trigger-reconcile-state:e1")
.await
.unwrap();
manager.model().deploy(&model, None).await.unwrap();
let after = manager
.evt()
.get("trigger-reconcile-state:e1")
.await
.unwrap();
assert_eq!(after.params, "{\"a\":2}");
assert_eq!(after.next_run, before.next_run, "schedule state preserved");
assert_eq!(after.last_run, before.last_run);
let armed_at = utils::time::time_millis();
model.on = vec![crate::Trigger {
id: "e1".to_string(),
kind: "schedule".to_string(),
schedule: Some("0 0 6 * * *".to_string()),
params: serde_json::json!({"a": 2}),
..Default::default()
}];
manager.model().deploy(&model, None).await.unwrap();
let rearmed = manager
.evt()
.get("trigger-reconcile-state:e1")
.await
.unwrap();
assert_eq!(
rearmed.last_run, after.last_run,
"last_run kept on cron change"
);
assert!(
rearmed.next_run >= armed_at,
"cron change must re-arm next_run"
);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_trigger_manual_start() {
let engine = Engine::new().start().await.unwrap();
let manager = engine.executor();
let model = Workflow::new()
.with_id("trigger-manual")
.with_trigger(|t| {
t.with_id("event1")
.with_kind("manual")
.with_params_vars(|vars| vars.with("test", 10))
})
.with_step(|step| step.with_id("step1"));
manager.model().deploy(&model, None).await.unwrap();
let ret = manager
.evt()
.start("trigger-manual:event1", &Vars::new().into())
.await
.unwrap();
assert!(ret.unwrap().get::<String>("pid").is_some());
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_trigger_chat_start() {
let engine = Engine::new().start().await.unwrap();
let manager = engine.executor();
let model = Workflow::new()
.with_id("trigger-chat")
.with_trigger(|t| t.with_id("event1").with_kind("chat"))
.with_step(|step| step.with_id("step1"));
manager.model().deploy(&model, None).await.unwrap();
let ret = manager
.evt()
.start("trigger-chat:event1", &"hello".into())
.await
.unwrap();
assert!(ret.unwrap().get::<String>("pid").is_some());
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_trigger_hook_start() {
let engine = Engine::new().start().await.unwrap();
let manager = engine.executor();
let model = Workflow::new()
.with_var("ret", 0)
.with_expose(Variant::create("ret", json!(null)))
.with_id("trigger-hook")
.with_trigger(|t| {
t.with_id("event1")
.with_kind("hook")
.with_params_vars(|vars| vars.with("var1", 10))
})
.with_step(|step| {
step.with_id("step1")
.with_uses(USES_SET, Vars::new().with("ret", 100))
});
manager.model().deploy(&model, None).await.unwrap();
let ret = manager
.evt()
.start("trigger-hook:event1", &Vars::new().into())
.await
.unwrap();
assert_eq!(ret.unwrap().get::<i32>("ret").unwrap(), 100);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_trigger_schedule_cannot_start_manually() {
let engine = Engine::new().start().await.unwrap();
let manager = engine.executor();
let model = Workflow::new()
.with_id("trigger-schedule-blocked")
.with_trigger(|t| {
t.with_id("event1")
.with_kind("schedule")
.with_schedule("* * * * * *")
})
.with_step(|step| step.with_id("step1"));
manager.model().deploy(&model, None).await.unwrap();
assert!(
manager
.evt()
.start("trigger-schedule-blocked:event1", &Vars::new().into())
.await
.is_err()
);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_trigger_schedule_auto_fire() {
let engine = Engine::new().start().await.unwrap();
let manager = engine.executor();
let model = Workflow::new()
.with_id("trigger-schedule")
.with_trigger(|t| {
t.with_id("every-sec")
.with_kind("schedule")
.with_schedule("* * * * * *")
})
.with_step(|step| step.with_id("step1"));
manager.model().deploy(&model, None).await.unwrap();
let evt = manager
.evt()
.get("trigger-schedule:every-sec")
.await
.unwrap();
assert_eq!(evt.kind, "schedule");
let mut fired = false;
for _ in 0..40 {
std::thread::sleep(std::time::Duration::from_millis(300));
let evt = manager
.evt()
.get("trigger-schedule:every-sec")
.await
.unwrap();
if evt.last_run > 0 && evt.next_run > evt.last_run {
fired = true;
break;
}
}
assert!(fired, "schedule trigger never fired");
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_trigger_deploy_invalid() {
let engine = Engine::new().start().await.unwrap();
let manager = engine.executor();
let workflow = r#"
id: "trigger-invalid"
on:
- id: event1
kind: manual
- id: event1
kind: manual
steps:
- id: step1
"#;
let workflow = Workflow::from_yml(workflow).unwrap();
assert!(manager.model().deploy(&workflow, None).await.is_err());
let workflow = r#"
id: "trigger-invalid"
on:
- kind: manual
steps:
- id: step1
"#;
let workflow = Workflow::from_yml(workflow).unwrap();
assert!(manager.model().deploy(&workflow, None).await.is_err());
let workflow = r#"
id: "trigger-invalid"
on:
- id: event1
kind: schedule
steps:
- id: step1
"#;
let workflow = Workflow::from_yml(workflow).unwrap();
assert!(manager.model().deploy(&workflow, None).await.is_err());
let workflow = r#"
id: "trigger-invalid"
on:
- id: event1
kind: schedule
schedule: "not a cron"
steps:
- id: step1
"#;
let workflow = Workflow::from_yml(workflow).unwrap();
assert!(manager.model().deploy(&workflow, None).await.is_err());
}
#[derive(Clone, serde::Deserialize)]
struct TriggerTestPackage;
#[async_trait::async_trait]
impl crate::package::ActPackage for TriggerTestPackage {
fn new(_config: &crate::Config) -> crate::Result<Self> {
Ok(Self)
}
fn definition() -> crate::package::ActPackageDefinition {
crate::package::ActPackageDefinition {
id: "test.trigger.pkg",
name: "Test Trigger",
desc: "custom trigger kind test package",
icon: "",
doc: "",
version: "0.1.0",
schema: json!({}),
options: None,
run_as: crate::ActRunAs::Func,
resources: vec![],
catalog: crate::package::ActPackageCatalog::Event,
}
}
async fn start(
&self,
rt: &Arc<crate::scheduler::Runtime>,
params: &serde_json::Value,
options: &Vars,
) -> crate::Result<Option<Vars>> {
let mid = options
.get::<String>(consts::MODEL_ID)
.ok_or(crate::ActError::Runtime(format!(
"cannot find '{}' in options",
consts::MODEL_ID
)))?;
let model: crate::ModelInfo = rt.cache().store().models().find(&mid).await?.into();
let workflow = model.workflow()?;
let start_params = serde_json::from_value::<Vars>(params.clone()).map_err(|e| {
crate::ActError::Package(format!("invalid trigger package params: {e}"))
})?;
let ret = rt.start(&workflow, start_params).await?;
Ok(Some(Vars::new().with(consts::PROCESS_ID, ret.id())))
}
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn export_trigger_custom_kind_package() {
let engine = Engine::new()
.add_package::<TriggerTestPackage>()
.start()
.await
.unwrap();
let manager = engine.executor();
let model = Workflow::new()
.with_id("trigger-custom")
.with_trigger(|t| t.with_id("event1").with_kind("test.trigger.pkg"))
.with_step(|step| step.with_id("step1"));
manager.model().deploy(&model, None).await.unwrap();
let evt = manager.evt().get("trigger-custom:event1").await.unwrap();
assert_eq!(evt.kind, "test.trigger.pkg");
let ret = manager
.evt()
.start("trigger-custom:event1", &json!({"var1": 10}))
.await
.unwrap();
assert!(ret.unwrap().get::<String>("pid").is_some());
}