mod context;
mod define;
mod engine;
mod event;
mod schedule;
pub use context::*;
pub use define::*;
pub use engine::*;
pub use event::*;
pub use schedule::*;
#[cfg(test)]
mod test {
use super::*;
use std::sync::{
Arc, Mutex,
atomic::{AtomicUsize, Ordering},
};
#[tokio::test]
async fn event_engine_registers_and_executes_lambda_events() {
let engine = EventEngine::default();
let observed = Arc::new(Mutex::new(Vec::new()));
let first_observed = Arc::clone(&observed);
engine
.register_event("math", move |ctx: &mut Context, input: i32| {
first_observed.lock().unwrap().push(input);
ctx.set_input(input + 1);
async { Ok(()) }
})
.await;
let second_observed = Arc::clone(&observed);
engine
.register_event("math", move |ctx: &mut Context, input: i32| {
second_observed.lock().unwrap().push(input);
ctx.set_input(input * 2);
async { Ok(()) }
})
.await;
let ctx = engine.invoke("math", 10_i32).await.unwrap();
assert_eq!(*observed.lock().unwrap(), vec![10, 11]);
assert_eq!(ctx.deref_fn::<i32, _>(|input| input.copied()), Some(22));
}
#[tokio::test]
async fn event_engine_registers_and_executes_once_lambda_event_once() {
let engine = EventEngine::default();
let calls = Arc::new(AtomicUsize::new(0));
let event_calls = Arc::clone(&calls);
engine.register_event_once("once", move |ctx: &mut Context, input: String| {
event_calls.fetch_add(1, Ordering::SeqCst);
ctx.set_input(format!("{input}->handled"));
async { Ok(()) }
});
let first = engine
.invoke("once", String::from("payload"))
.await
.unwrap();
assert_eq!(
first.deref_fn::<String, _>(|input| input.cloned()),
Some(String::from("payload->handled"))
);
let second = engine
.invoke("once", String::from("payload"))
.await
.unwrap();
assert_eq!(
second.deref_fn::<String, _>(|input| input.cloned()),
Some(String::from("payload"))
);
assert_eq!(calls.load(Ordering::SeqCst), 1);
}
#[tokio::test]
async fn global_agent_engine_exposes_register_and_invoke() {
register_event(
"global_agent_engine_exposes_register_and_invoke",
|ctx: &mut Context, input: i32| {
ctx.set_input(input + 10);
async { Ok(()) }
},
)
.await;
let ctx = invoke("global_agent_engine_exposes_register_and_invoke", 2_i32)
.await
.unwrap();
assert_eq!(ctx.deref_fn::<i32, _>(|input| input.copied()), Some(12));
}
#[tokio::test]
async fn event_engine_launch_fn_invokes_completion_callback() {
let engine = EventEngine::default();
let observed = Arc::new(Mutex::new(None));
engine
.register_event("launch_fn", |ctx: &mut Context, input: i32| {
ctx.set_input(input + 5);
async { Ok(()) }
})
.await;
let callback_observed = Arc::clone(&observed);
engine.launch_fn("launch_fn", 7_i32, move |ctx| {
let value = ctx.deref_fn::<i32, _>(|input| input.copied());
*callback_observed.lock().unwrap() = value;
});
for _ in 0..16 {
if observed.lock().unwrap().is_some() {
break;
}
tokio::task::yield_now().await;
}
assert_eq!(*observed.lock().unwrap(), Some(12));
}
#[tokio::test]
async fn event_engine_launch_fut_awaits_completion_future() {
let engine = EventEngine::default();
let observed = Arc::new(Mutex::new(None));
engine
.register_event("launch_fut", |ctx: &mut Context, input: i32| {
ctx.set_input(input * 3);
async { Ok(()) }
})
.await;
let callback_observed = Arc::clone(&observed);
engine.launch_fut("launch_fut", 4_i32, move |ctx| {
let callback_observed = Arc::clone(&callback_observed);
async move {
tokio::task::yield_now().await;
let value = ctx.deref_fn::<i32, _>(|input| input.copied());
*callback_observed.lock().unwrap() = value;
}
});
for _ in 0..16 {
if observed.lock().unwrap().is_some() {
break;
}
tokio::task::yield_now().await;
}
assert_eq!(*observed.lock().unwrap(), Some(12));
}
}