use std::time::{Duration, Instant};
use crate::{
ActError, ActUserVar, Context, Engine, MessageState, Vars, Workflow, env::Environment,
event::EventAction, utils::consts, utils::test::USES_IRQ,
};
use serde::{Deserialize, Serialize};
use serde_json::json;
use serial_test::serial;
#[test]
fn env_eval_empty() {
let env = Environment::new();
let result = env.eval::<()>("");
assert!(result.is_ok());
}
#[test]
fn env_eval_void() {
let env = Environment::new();
let script = r#"
let v = 5;
console.log(`v=${v}`);
"#;
let result = env.eval::<()>(script);
assert!(result.is_ok());
}
#[test]
fn env_eval_number() {
let env = Environment::new();
let script = r#"
let v = 5;
v
"#;
let result = env.eval::<i64>(script);
assert_eq!(result.unwrap(), 5);
}
#[test]
fn env_eval_throw_error() {
let env = Environment::new();
let script = r#"
throw new Error("err1");
"#;
let result = env.eval::<serde_json::Value>(script);
assert_eq!(
result.err().unwrap(),
ActError::Exception {
ecode: "".to_string(),
message: "err1".to_string()
}
);
}
#[test]
fn env_eval_infinite_loop_timeout() {
let env = Environment::new();
let script = r#"
while (true) {}
"#;
let start = Instant::now();
let result = env.eval::<()>(script);
assert_eq!(
result.err().unwrap(),
ActError::Script("Execution timeout".into())
);
assert!(start.elapsed() < Duration::from_secs(60));
}
#[test]
fn env_eval_expr() {
let env = Environment::new();
let script = r#"
let ret = 10;
ret > 0
"#;
let result = env.eval::<bool>(script);
assert!(result.unwrap());
}
#[test]
fn env_eval_bigint_within_i64() {
let env = Environment::new();
let result = env.eval::<i64>("BigInt(42)");
assert_eq!(result.unwrap(), 42);
}
#[test]
fn env_eval_bigint_outside_u64_is_error() {
let env = Environment::new();
let result = env.eval::<serde_json::Value>(r#"BigInt("1000000000000000000000000000000")"#);
let err = result.unwrap_err();
assert!(
err.to_string().contains("outside the i64/u64 range"),
"unexpected error: {err}"
);
}
#[test]
fn env_eval_bigint_beyond_i64_keeps_u64() {
let env = Environment::new();
let result = env
.eval::<serde_json::Value>(r#"BigInt("18446744073709551615")"#)
.unwrap();
assert_eq!(result, json!(u64::MAX));
}
#[test]
fn env_eval_non_finite_number_is_error() {
let env = Environment::new();
for script in ["0/0", "1e400"] {
let err = env.eval::<serde_json::Value>(script).unwrap_err();
assert!(
err.to_string().contains("non-finite"),
"unexpected error for {script}: {err}"
);
}
}
#[test]
fn env_eval_large_numbers_round_trip_exactly() {
#[derive(Clone)]
struct BigNumberVar;
impl ActUserVar for BigNumberVar {
fn name(&self) -> String {
"bignumbers".to_string()
}
fn default_data(&self) -> Option<Vars> {
Some(
Vars::new()
.with("over_i32", 3_000_000_000i64)
.with("millis", 1_757_318_400_000i64)
.with("snowflake", 1_234_567_890_123_456_789i64)
.with("min", i64::MIN)
.with("over_i64", u64::MAX),
)
}
}
let env = Environment::new();
env.register_var(&BigNumberVar);
assert_eq!(
env.eval::<i64>("bignumbers.over_i32").unwrap(),
3_000_000_000
);
assert_eq!(
env.eval::<i64>("bignumbers.millis").unwrap(),
1_757_318_400_000
);
assert!(
env.eval::<bool>("typeof bignumbers.millis === 'number'")
.unwrap()
);
assert_eq!(
env.eval::<i64>("bignumbers.snowflake").unwrap(),
1_234_567_890_123_456_789
);
assert_eq!(env.eval::<i64>("bignumbers.min").unwrap(), i64::MIN);
assert_eq!(env.eval::<u64>("bignumbers.over_i64").unwrap(), u64::MAX);
assert!(
env.eval::<bool>("bignumbers.snowflake === 1234567890123456789n")
.unwrap()
);
assert_eq!(
env.eval::<i64>("bignumbers.snowflake - 1n").unwrap(),
1_234_567_890_123_456_788
);
}
#[test]
fn env_eval_array() {
let env = Environment::new();
let script = r#"
["u1", "u2"]
"#;
let result = env.eval::<Vec<String>>(script);
assert_eq!(result.unwrap(), ["u1", "u2"]);
}
#[test]
fn env_eval_object() {
let env = Environment::new();
let script = r#"
let ret = { "a": 1, "b": "abc" };
ret
"#;
#[derive(Debug, Deserialize, Serialize, PartialEq, Clone)]
struct Obj {
a: i32,
b: String,
}
let result = env.eval::<Obj>(script);
assert_eq!(
result.unwrap(),
Obj {
a: 1,
b: "abc".to_string()
}
);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn env_eval_sys_env() {
unsafe {
std::env::set_var("TOKEN", "abc");
}
let engine = Engine::builder().start().await.unwrap();
let sig = engine.signal(());
let s1 = sig.clone();
let env = engine.runtime().env().clone();
let workflow = Workflow::new().with_step(|step| step.with_id("step1"));
engine.channel().on_complete(move |_| {
let s1 = s1.clone();
async move { s1.close() }
});
let proc = engine
.runtime()
.start(&workflow, Vars::new())
.await
.unwrap();
sig.recv().await;
let task = proc.root().unwrap();
let script = r#"
$env.TOKEN
"#;
let context = task.create_context();
Context::scope(&context, || {
let result = env.eval::<String>(script);
assert_eq!(result.unwrap(), "abc");
});
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn env_eval_null() {
let engine = Engine::builder().start().await.unwrap();
let sig = engine.signal(());
let s1 = sig.clone();
let env = engine.runtime().env().clone();
let workflow = Workflow::new().with_step(|step| step.with_id("step1"));
engine.channel().on_complete(move |_| {
let s1 = s1.clone();
async move { s1.close() }
});
let proc = engine
.runtime()
.start(&workflow, Vars::new())
.await
.unwrap();
sig.recv().await;
let task = proc.root().unwrap();
let script = r#"
$env.TOKEN2
"#;
let context = task.create_context();
Context::scope(&context, || {
let result = env.eval::<serde_json::Value>(script);
assert_eq!(result.unwrap(), serde_json::Value::Null);
});
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn env_console_module() {
let env = Environment::new();
let script = r#"
let v = 5;
console.log(`v=${v}`);
console.info(`v=${v}`);
console.warn(`v=${v}`);
console.error(`v=${v}`);
// multi-arg support
console.log("hello", 42, true, v);
console.log("joined:", "a", "b", "c");
"#;
let result = env.eval::<()>(script);
assert!(result.is_ok());
}
#[test]
#[serial]
fn env_collection_union() {
let env = Environment::new();
let script = r#"
let a = ["a"];
let b = ["b"];
a.union(b)
"#;
let result = env.eval::<Vec<String>>(script).unwrap();
assert_eq!(result, ["a", "b"]);
}
#[test]
#[serial]
fn env_eval_contexts_do_not_leak_globals() {
let env = Environment::new();
env.eval::<()>(
r#"
var leaked_var = 42;
globalThis.leaked_property = 42;
void 0;
"#,
)
.unwrap();
let leaked = env
.eval::<bool>(
r#"
typeof leaked_var === "undefined" && typeof leaked_property === "undefined"
"#,
)
.unwrap();
assert!(
leaked,
"global declarations must not leak across evaluations"
);
}
#[test]
#[serial]
fn env_eval_contexts_restore_module_globals() {
let env = Environment::new();
env.eval::<()>(
r#"
$get = 42;
os = 42;
void 0;
"#,
)
.unwrap();
let restored = env
.eval::<bool>(
r#"
typeof $get === "function" && typeof os === "string"
"#,
)
.unwrap();
assert!(restored, "static module globals must be restored");
}
#[test]
#[serial]
fn env_eval_isolates_realm_state() {
let env = Environment::new();
env.eval::<()>(
r#"
Array.prototype.__acts_secret = "A";
Object.prototype.__acts_secret = "A";
Math.__acts_secret = "A";
$get.__acts_secret = "A";
console.__acts_secret = "A";
Object.defineProperty($env, "__acts_secret", { value: "A", configurable: true });
globalThis.__acts_secret = "A";
Object.defineProperty(globalThis, "__acts_hidden", { value: "A", configurable: false });
void 0;
"#,
)
.unwrap();
let leaked = env
.eval::<serde_json::Value>(
r#"
({
array: Array.prototype.__acts_secret === "A",
object: Object.prototype.__acts_secret === "A",
namespace: Math.__acts_secret === "A",
module_fn: $get.__acts_secret === "A",
console: console.__acts_secret === "A",
env_target: (Object.getOwnPropertyDescriptor($env, "__acts_secret") || {}).value === "A",
global: globalThis.__acts_secret === "A",
hidden: globalThis.__acts_hidden === "A",
})
"#,
)
.unwrap();
assert_eq!(
leaked,
json!({
"array": false,
"object": false,
"namespace": false,
"module_fn": false,
"console": false,
"env_target": false,
"global": false,
"hidden": false,
}),
"realm state must not survive an evaluation"
);
}
#[test]
#[serial]
fn env_collection_intersect() {
let env = Environment::new();
let script = r#"
let a = ["a", "b"];
let b = ["b", "c"];
a.intersection(b)
"#;
let result = env.eval::<Vec<String>>(script).unwrap();
assert_eq!(result, ["b"]);
}
#[test]
#[serial]
fn env_collection_difference() {
let env = Environment::new();
let script = r#"
let a = ["a", "b"];
let b = ["b"];
a.difference(b)
"#;
let result = env.eval::<Vec<String>>(script).unwrap();
assert_eq!(result, ["a"]);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn env_task_get_value() {
let engine = Engine::builder().start().await.unwrap();
let sig = engine.signal(());
let s1 = sig.clone();
let env = engine.runtime().env().clone();
let workflow = Workflow::new()
.with_var("a", 10)
.with_step(|step| step.with_id("step1"));
engine.channel().on_complete(move |_| {
let s1 = s1.clone();
async move { s1.close() }
});
let proc = engine
.runtime()
.start(&workflow, Vars::new())
.await
.unwrap();
sig.recv().await;
let task = proc.root().unwrap();
let script = r#"
a
"#;
let context = task.create_context();
Context::scope(&context, || {
let result = env.eval::<i64>(script);
assert_eq!(result.unwrap(), 10);
});
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn env_task_get_var_not_exists() {
let engine = Engine::builder().start().await.unwrap();
let sig = engine.signal(());
let s1 = sig.clone();
let env = engine.runtime().env().clone();
let workflow = Workflow::new().with_step(|step| step.with_id("step1"));
engine.channel().on_complete(move |_| {
let s1 = s1.clone();
async move { s1.close() }
});
let proc = engine
.runtime()
.start(&workflow, Vars::new())
.await
.unwrap();
sig.recv().await;
let task = proc.root().unwrap();
let script = r#"
not_exists
"#;
let context = task.create_context();
Context::scope(&context, || {
let result = env.eval::<serde_json::Value>(script);
assert!(result.is_err());
});
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn env_task_get_fn_not_exists() {
let engine = Engine::builder().start().await.unwrap();
let sig = engine.signal(());
let s1 = sig.clone();
let env = engine.runtime().env().clone();
let workflow = Workflow::new().with_step(|step| step.with_id("step1"));
engine.channel().on_complete(move |_| {
let s1 = s1.clone();
async move { s1.close() }
});
let proc = engine
.runtime()
.start(&workflow, Vars::new())
.await
.unwrap();
sig.recv().await;
let task = proc.root().unwrap();
let script = r#"
$get("not_exists")
"#;
let context = task.create_context();
Context::scope(&context, || {
let result = env.eval::<serde_json::Value>(script);
assert_eq!(result.unwrap(), serde_json::Value::Null);
});
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn env_task_set() {
let engine = Engine::builder().start().await.unwrap();
let sig = engine.signal(());
let s1 = sig.clone();
let env = engine.runtime().env().clone();
let workflow = Workflow::new()
.with_var("a", 10)
.with_step(|step| step.with_id("step1"));
engine.channel().on_complete(move |_| {
let s1 = s1.clone();
async move { s1.close() }
});
let proc = engine
.runtime()
.start(&workflow, Vars::new())
.await
.unwrap();
sig.recv().await;
let task = proc.root().unwrap();
let script = r#"
$set("a", 100);
"#;
let context = task.create_context();
Context::scope(&context, || {
env.eval::<()>(script).unwrap();
assert_eq!(proc.data().get::<i64>("a"), Some(100));
});
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn env_task_multi_line() {
let engine = Engine::builder().start().await.unwrap();
let sig = engine.signal(());
let s1 = sig.clone();
let env = engine.runtime().env().clone();
let workflow = Workflow::new().with_step(|step| step.with_id("step1"));
engine.channel().on_complete(move |_| {
let s1 = s1.clone();
async move { s1.close() }
});
let proc = engine
.runtime()
.start(&workflow, Vars::new())
.await
.unwrap();
sig.recv().await;
let task = proc.root().unwrap();
let context = task.create_context();
Context::scope(&context, || {
env.eval::<()>(r#"$set("a", 100)"#).unwrap();
env.eval::<()>(r#"$set("b", 200)"#).unwrap();
let value = env.eval::<bool>(r#"a < b"#).unwrap();
assert!(value);
});
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn env_env_get_local() {
let engine = Engine::builder().start().await.unwrap();
let sig = engine.signal(());
let s1 = sig.clone();
let env = engine.runtime().env().clone();
let workflow = Workflow::new()
.with_env("a", 10)
.with_step(|step| step.with_id("step1"));
engine.channel().on_complete(move |_| {
let s1 = s1.clone();
async move { s1.close() }
});
let proc = engine
.runtime()
.start(&workflow, Vars::new())
.await
.unwrap();
sig.recv().await;
let task = proc.root().unwrap();
let script = r#"
$env.a
"#;
let context = task.create_context();
Context::scope(&context, || {
let result = env.eval::<i64>(script);
assert_eq!(result.unwrap(), 10);
});
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn env_env_set_proc_env() {
let engine = Engine::builder().start().await.unwrap();
let sig = engine.signal(());
let s1 = sig.clone();
let env = engine.runtime().env().clone();
let workflow = Workflow::new()
.with_env("a", 100)
.with_step(|step| step.with_id("step1"));
engine.channel().on_complete(move |_| {
let s1 = s1.clone();
async move { s1.close() }
});
let proc = engine
.runtime()
.start(&workflow, Vars::new())
.await
.unwrap();
sig.recv().await;
let task = proc.root().unwrap();
let script = r#"
$env.a = 200;
"#;
let context = task.create_context();
Context::scope(&context, || {
env.eval::<serde_json::Value>(script).unwrap();
assert_eq!(proc.env().get::<i64>("a"), Some(200));
});
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn env_private_keys_are_engine_only() {
let engine = Engine::builder().start().await.unwrap();
let env = engine.runtime().env().clone();
let workflow = Workflow::new().with_step(|step| step.with_id("step1"));
let proc = engine
.runtime()
.start(&workflow, Vars::new())
.await
.unwrap();
let task = proc.root().unwrap();
let context = task.create_context();
context.set_env(crate::utils::consts::PROC_OWNER, "forged");
context.set_env(crate::utils::consts::PROC_WORKDIR, "/etc");
Context::scope(&context, || {
let read = env
.eval::<serde_json::Value>(&format!("$env.{}", crate::utils::consts::PROC_OWNER))
.unwrap();
assert!(read.is_null(), "private key must not be readable: {read}");
env.eval::<serde_json::Value>(&format!(
"$env.{} = 'escaped';",
crate::utils::consts::PROC_WORKDIR
))
.unwrap();
});
assert_eq!(
proc.env()
.get::<String>(crate::utils::consts::PROC_WORKDIR)
.as_deref(),
Some("/etc")
);
engine.close().await;
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn env_work_dir_names_the_process_directory() {
let engine = Engine::builder().start().await.unwrap();
let env = engine.runtime().env().clone();
let workflow = Workflow::new().with_step(|step| step.with_id("step1"));
let proc = engine
.runtime()
.start(&workflow, Vars::new())
.await
.unwrap();
let task = proc.root().unwrap();
let context = task.create_context();
Context::scope(&context, || {
let read = env.eval::<serde_json::Value>("$env.WORK_DIR").unwrap();
assert!(read.is_null(), "no workdir means no value: {read}");
});
let dir = std::env::temp_dir().join(format!("acts_workdir_{}", crate::utils::longid()));
proc.set_workdir(&dir);
Context::scope(&context, || {
assert_eq!(
env.eval::<String>("$env.WORK_DIR").unwrap(),
dir.display().to_string()
);
assert_eq!(
context.get_env::<std::path::PathBuf>(consts::ENV_WORK_DIR),
Some(dir.clone())
);
env.eval::<serde_json::Value>("$env.WORK_DIR = '/etc'")
.unwrap();
context.set_env(consts::ENV_WORK_DIR, "/etc");
assert_eq!(
env.eval::<String>("$env.WORK_DIR").unwrap(),
dir.display().to_string(),
"a run must not be able to rename its own workdir"
);
});
assert!(
proc.env().get::<String>(consts::ENV_WORK_DIR).is_none(),
"the reserved name must never be stored in the process env"
);
assert_eq!(proc.workdir(), Some(dir));
engine.close().await;
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn env_env_multi_line() {
let engine = Engine::builder().start().await.unwrap();
let sig = engine.signal(());
let s1 = sig.clone();
let env = engine.runtime().env().clone();
let workflow = Workflow::new().with_step(|step| step.with_id("step1"));
engine.channel().on_complete(move |_| {
let s1 = s1.clone();
async move { s1.close() }
});
let proc = engine
.runtime()
.start(&workflow, Vars::new())
.await
.unwrap();
sig.recv().await;
let task = proc.root().unwrap();
let context = task.create_context();
Context::scope(&context, || {
env.eval::<serde_json::Value>(r#"$env.a = 100"#).unwrap();
env.eval::<serde_json::Value>(r#"$env.b = 200"#).unwrap();
let value = env.eval::<bool>(r#"$env.a < $env.b"#).unwrap();
assert!(value);
});
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn env_vars_set_num() {
let engine = Engine::builder().start().await.unwrap();
let sig = engine.signal(());
let s1 = sig.clone();
let workflow = Workflow::new()
.with_env("a", 10)
.with_step(|step| step.with_id("step1"));
engine.channel().on_complete(move |_| {
let s1 = s1.clone();
async move { s1.close() }
});
let proc = engine
.runtime()
.start(&workflow, Vars::new())
.await
.unwrap();
sig.recv().await;
let task = proc.root().unwrap();
let context = task.create_context();
Context::scope(&context, || {
assert_eq!(proc.env().get::<i64>("a"), Some(10));
});
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn env_vars_set_str() {
let engine = Engine::builder().start().await.unwrap();
let sig = engine.signal(());
let s1 = sig.clone();
let workflow = Workflow::new()
.with_env("a", "abc")
.with_step(|step| step.with_id("step1"));
engine.channel().on_complete(move |_| {
let s1 = s1.clone();
async move { s1.close() }
});
let proc = engine
.runtime()
.start(&workflow, Vars::new())
.await
.unwrap();
sig.recv().await;
let task = proc.root().unwrap();
let context = task.create_context();
Context::scope(&context, || {
assert_eq!(proc.env().get::<String>("a"), Some("abc".to_string()));
});
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn env_vars_set_json() {
let engine = Engine::builder().start().await.unwrap();
let sig = engine.signal(());
let s1 = sig.clone();
let workflow = Workflow::new()
.with_env("a", json!({ "count": 1 }))
.with_step(|step| step.with_id("step1"));
engine.channel().on_complete(move |_| {
let s1 = s1.clone();
async move { s1.close() }
});
let proc = engine
.runtime()
.start(&workflow, Vars::new())
.await
.unwrap();
sig.recv().await;
let task = proc.root().unwrap();
let context = task.create_context();
Context::scope(&context, || {
assert_eq!(
proc.env().get::<serde_json::Value>("a"),
Some(json!({ "count": 1 }))
);
});
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn env_vars_update() {
let engine = Engine::builder().start().await.unwrap();
let sig = engine.signal(());
let s1 = sig.clone();
let workflow = Workflow::new()
.with_env("a", 10)
.with_step(|step| step.with_id("step1"));
engine.channel().on_complete(move |_| {
let s1 = s1.clone();
async move { s1.close() }
});
let proc = engine
.runtime()
.start(&workflow, Vars::new())
.await
.unwrap();
sig.recv().await;
let task = proc.root().unwrap();
let context = task.create_context();
context.set_env("a", 100);
Context::scope(&context, || {
assert_eq!(proc.env().get::<i32>("a"), Some(100));
});
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn env_step_get_data_by_id() {
let engine = Engine::builder().start().await.unwrap();
let sig = engine.signal(());
let s1 = sig.clone();
let env = engine.runtime().env().clone();
let workflow = Workflow::new()
.with_step(|step| step.with_id("step1").with_var("a", 10))
.with_step(|step| step.with_id("step2").with_var("b", "abc"));
engine.channel().on_complete(move |_| {
let s1 = s1.clone();
async move { s1.close() }
});
let proc = engine
.runtime()
.start(&workflow, Vars::new())
.await
.unwrap();
sig.recv().await;
let task = proc.root().unwrap();
let script = r#"
step1.a
"#;
proc.print();
let context = task.create_context();
Context::scope(&context, || {
let result = env.eval::<i32>(script);
assert_eq!(result.unwrap(), 10);
});
let script = r#"
step2.b
"#;
let context = task.create_context();
Context::scope(&context, || {
let result = env.eval::<String>(script);
assert_eq!(result.unwrap(), "abc");
});
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn env_step_get_data_null() {
let engine = Engine::builder().start().await.unwrap();
let sig = engine.signal(());
let s1 = sig.clone();
let env = engine.runtime().env().clone();
let workflow = Workflow::new().with_step(|step| step.with_id("step1").with_var("a", 10));
engine.channel().on_complete(move |_| {
let s1 = s1.clone();
async move { s1.close() }
});
let proc = engine
.runtime()
.start(&workflow, Vars::new())
.await
.unwrap();
sig.recv().await;
let task = proc.root().unwrap();
let script = r#"
step1.not_exists
"#;
proc.print();
let context = task.create_context();
Context::scope(&context, || {
let result = env.eval::<serde_json::Value>(script);
assert_eq!(result.unwrap(), serde_json::Value::Null);
});
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn env_step_set_data_err_with_completed_state() {
let engine = Engine::builder().start().await.unwrap();
let sig = engine.signal(());
let s1 = sig.clone();
let env = engine.runtime().env().clone();
let workflow = Workflow::new().with_step(|step| step.with_id("step1").with_var("a", 10));
engine.channel().on_complete(move |_| {
let s1 = s1.clone();
async move { s1.close() }
});
let proc = engine
.runtime()
.start(&workflow, Vars::new())
.await
.unwrap();
sig.recv().await;
let task = proc.root().unwrap();
let script = r#"
step1.a = 100;
"#;
let context = task.create_context();
Context::scope(&context, || {
let result = env.eval::<serde_json::Value>(script);
proc.print();
assert!(result.is_err());
});
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn env_step_set_data_ok_with_running_state() {
let engine = Engine::builder().start().await.unwrap();
let sig = engine.signal(());
let s1 = sig.clone();
let env = engine.runtime().env().clone();
let workflow = Workflow::new().with_step(|step| {
step.with_id("step1")
.with_var("a", 10)
.with_uses(USES_IRQ, Vars::new().with("key", "test"))
});
engine.channel().on_message(move |e| {
let s1 = s1.clone();
async move {
if e.is_irq() {
s1.close()
}
}
});
let proc = engine
.runtime()
.start(&workflow, Vars::new())
.await
.unwrap();
sig.recv().await;
let task = proc.root().unwrap();
let script = r#"
step1.a = 100;
"#;
let context = task.create_context();
Context::scope(&context, || {
env.eval::<serde_json::Value>(script).unwrap();
proc.print();
assert_eq!(
proc.task_by_nid("step1")
.last()
.unwrap()
.data()
.get::<i32>("a")
.unwrap(),
100
);
});
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn env_step_get_data() {
let engine = Engine::builder().start().await.unwrap();
let sig = engine.signal(());
let s1 = sig.clone();
let env = engine.runtime().env().clone();
let workflow = Workflow::new().with_step(|step| {
step.with_id("step1")
.with_var("a", 10)
.with_uses(USES_IRQ, Vars::new().with("key", "test"))
});
engine.channel().on_message(move |e| {
let s1 = s1.clone();
async move {
if e.is_irq() {
s1.close()
}
}
});
let proc = engine
.runtime()
.start(&workflow, Vars::new())
.await
.unwrap();
sig.recv().await;
let task = proc.root().unwrap();
let script = r#"
step1.b = "abc";
step1.data()
"#;
let context = task.create_context();
Context::scope(&context, || {
let result = env.eval::<Vars>(script).unwrap();
proc.print();
assert_eq!(result.get::<String>("b").unwrap(), "abc");
});
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn env_step_get_inputs() {
let engine = Engine::builder().start().await.unwrap();
let sig = engine.signal(());
let s1 = sig.clone();
let env = engine.runtime().env().clone();
let workflow = Workflow::new().with_step(|step| {
step.with_id("step1")
.with_var("a", 10)
.with_uses(USES_IRQ, Vars::new().with("key", "act1"))
});
engine.channel().on_message(move |e| {
let s1 = s1.clone();
async move {
if e.is_irq() {
s1.close()
}
}
});
let proc = engine
.runtime()
.start(&workflow, Vars::new())
.await
.unwrap();
sig.recv().await;
let task = proc.root().unwrap();
let script = r#"
step1.inputs()
"#;
let context = task.create_context();
Context::scope(&context, || {
let result = env.eval::<Vars>(script).unwrap();
proc.print();
assert_eq!(result.get::<i32>("a").unwrap(), 10);
});
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn env_act_get_inputs() {
let engine = Engine::builder().start().await.unwrap();
let sig = engine.signal(());
let s1 = sig.clone();
let env = engine.runtime().env().clone();
let workflow = Workflow::new().with_step(|step| {
step.with_id("step1")
.with_var("a", json!(10))
.with_uses(USES_IRQ, Vars::new().with("key", "act1"))
});
engine.channel().on_message(move |e| {
let s1 = s1.clone();
async move {
if e.is_irq() {
s1.close()
}
}
});
let proc = engine
.runtime()
.start(&workflow, Vars::new())
.await
.unwrap();
sig.recv().await;
let task = proc.task_by_params("key", "act1").last().cloned().unwrap();
let script = r#"
$inputs()
"#;
let context = task.create_context();
Context::scope(&context, || {
let result = env.eval::<Vars>(script).unwrap();
proc.print();
assert_eq!(result.get::<i32>("a").unwrap(), 10);
});
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn env_act_get_data() {
let engine = Engine::builder().start().await.unwrap();
let sig = engine.signal(());
let s1 = sig.clone();
let env = engine.runtime().env().clone();
let workflow = Workflow::new().with_step(|step| {
step.with_id("step1")
.with_uses(USES_IRQ, Vars::new().with("key", "act1"))
});
engine.channel().on_message(move |e| {
let s1 = s1.clone();
async move {
if e.is_irq() {
s1.close()
}
}
});
let proc = engine
.runtime()
.start(&workflow, Vars::new())
.await
.unwrap();
sig.recv().await;
let task = proc.task_by_params("key", "act1").last().cloned().unwrap();
let script = r#"
$set("my_value", 20)
$data()
"#;
let context = task.create_context();
Context::scope(&context, || {
let result = env.eval::<Vars>(script).unwrap();
proc.print();
assert_eq!(result.get::<i32>("my_value").unwrap(), 20);
});
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn env_user_var_get_from_context() {
#[derive(Clone)]
struct MyVarPlugin;
impl ActUserVar for MyVarPlugin {
fn name(&self) -> String {
"test".to_string()
}
}
let engine = Engine::builder().start().await.unwrap();
let sig = engine.signal(());
let s1 = sig.clone();
engine
.executor(&crate::Principal::unrestricted())
.ext()
.register_var(&MyVarPlugin)
.unwrap();
let env = engine.runtime().env().clone();
let workflow = Workflow::new().with_step(|step| {
step.with_id("step1")
.with_uses(USES_IRQ, Vars::new().with("key", "act1"))
});
engine.channel().on_message(move |e| {
let s1 = s1.clone();
async move {
if e.is_irq() {
s1.close()
}
}
});
let proc = engine
.runtime()
.start(
&workflow,
Vars::new().with("test", Vars::new().with("var1", 10)),
)
.await
.unwrap();
sig.recv().await;
let task = proc.task_by_params("key", "act1").last().cloned().unwrap();
let script = r#"
test.var1
"#;
let context = task.create_context();
Context::scope(&context, || {
let result = env.eval::<i32>(script).unwrap();
proc.print();
assert_eq!(result, 10);
});
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn env_user_var_get_default() {
#[derive(Clone)]
struct MyVarPlugin;
impl ActUserVar for MyVarPlugin {
fn name(&self) -> String {
"test".to_string()
}
fn default_data(&self) -> Option<Vars> {
Some(Vars::new().with("var1", 5))
}
}
let engine = Engine::builder().start().await.unwrap();
let sig = engine.signal(());
let s1 = sig.clone();
engine
.executor(&crate::Principal::unrestricted())
.ext()
.register_var(&MyVarPlugin)
.unwrap();
let env = engine.runtime().env().clone();
let workflow = Workflow::new().with_step(|step| {
step.with_id("step1")
.with_uses(USES_IRQ, Vars::new().with("key", "act1"))
});
engine.channel().on_message(move |e| {
let s1 = s1.clone();
async move {
if e.is_irq() {
s1.close()
}
}
});
let proc = engine
.runtime()
.start(&workflow, Vars::new())
.await
.unwrap();
sig.recv().await;
let task = proc.task_by_params("key", "act1").last().cloned().unwrap();
let script = r#"
test.var1
"#;
let context = task.create_context();
Context::scope(&context, || {
let result = env.eval::<i32>(script).unwrap();
proc.print();
assert_eq!(result, 5);
});
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn env_user_var_secrets_get() {
let engine = Engine::builder().start().await.unwrap();
let sig = engine.signal(());
let s1 = sig.clone();
let env = engine.runtime().env().clone();
let workflow = Workflow::new().with_step(|step| {
step.with_id("step1")
.with_uses(USES_IRQ, Vars::new().with("key", "act1"))
});
engine.channel().on_message(move |e| {
let s1 = s1.clone();
async move {
if e.is_irq() {
s1.close()
}
}
});
let proc = engine
.runtime()
.start(
&workflow,
Vars::new().with("secrets", Vars::new().with("TOKEN", "my_token")),
)
.await
.unwrap();
sig.recv().await;
let task = proc.task_by_params("key", "act1").last().cloned().unwrap();
let script = r#"
secrets.TOKEN
"#;
let context = task.create_context();
Context::scope(&context, || {
let result = env.eval::<String>(script).unwrap();
proc.print();
assert_eq!(result, "my_token");
});
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn env_user_var_os_get() {
let env = Environment::new();
let script = r#"
os
"#;
let result = env.eval::<String>(script);
assert!(["linux", "windows", "macos"].contains(&result.unwrap().as_str()));
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn env_act_cost_get() {
let engine = Engine::builder().start().await.unwrap();
let sig = engine.signal(());
let s1 = sig.clone();
let env = engine.runtime().env().clone();
let workflow = Workflow::new().with_step(|step| {
step.with_id("step1")
.with_uses(USES_IRQ, Vars::new().with("key", "act1"))
});
engine.channel().on_message(move |e| {
println!("message: {e:?}");
let s1 = s1.clone();
async move {
if e.is_irq() {
s1.close()
}
}
});
let proc = engine
.runtime()
.start(&workflow, Vars::new())
.await
.unwrap();
sig.recv().await;
let task = proc.task_by_params("key", "act1").last().cloned().unwrap();
let script = r#"
$cost()
"#;
let context = task.create_context();
Context::scope(&context, || {
let result = env.eval::<i32>(script).unwrap();
proc.print();
assert!(result > 0);
});
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn env_act_cost_in_get() {
let engine = Engine::builder().start().await.unwrap();
let sig = engine.signal(());
let s1 = sig.clone();
let env = engine.runtime().env().clone();
let workflow = Workflow::new().with_step(|step| {
step.with_id("step1")
.with_uses(USES_IRQ, Vars::new().with("key", "act1"))
});
engine.channel().on_message(move |e| {
let s1 = s1.clone();
async move {
if e.params().unwrap().get::<String>("key").as_deref() == Some("act1")
&& e.is_state(MessageState::Created)
{
tokio::time::sleep(Duration::from_secs(2)).await;
s1.close()
}
}
});
let proc = engine
.runtime()
.start(&workflow, Vars::new())
.await
.unwrap();
sig.recv().await;
let task = proc.task_by_params("key", "act1").last().cloned().unwrap();
let script = r#"
$cost_in('1s')
"#;
let context = task.create_context();
Context::scope(&context, || {
let result = env.eval::<bool>(script).unwrap();
proc.print();
assert!(result);
});
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn env_act_ecode_get() {
let engine = Engine::builder().start().await.unwrap();
let sig = engine.signal(());
let s1 = sig.clone();
let env = engine.runtime().env().clone();
let workflow = Workflow::new().with_step(|step| {
step.with_id("step1")
.with_uses(USES_IRQ, Vars::new().with("key", "act1"))
});
let runtime = engine.runtime();
engine.channel().on_message(move |e| {
println!("on_message: {:?}", e);
let runtime = runtime.clone();
let s1 = s1.clone();
async move {
if e.is_params_key("act1") && e.is_state(MessageState::Created) {
runtime
.do_action2(
&e.pid,
&e.tid,
EventAction::Error,
Vars::new().with(consts::ACT_ERR_CODE, "err1"),
)
.await
.unwrap();
s1.close()
}
}
});
let proc = engine
.runtime()
.start(&workflow, Vars::new())
.await
.unwrap();
sig.recv().await;
let task = proc.task_by_params("key", "act1").last().cloned().unwrap();
let script = r#"
$ecode()
"#;
let context = task.create_context();
Context::scope(&context, || {
let result = env.eval::<String>(script).unwrap();
proc.print();
assert_eq!(result, "err1");
});
}