Documentation
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));
    }
}