use flo_scene::*;
use flo_scene::programs::*;
use flo_scene::commands::*;
use futures::prelude::*;
use serde::*;
#[test]
fn simple_command() {
let scene = Scene::default();
let test_command = FnCommand::<(), usize>::new(|_input, context| async move {
let mut output = context.send::<usize>(()).unwrap();
output.send(1).await.unwrap();
output.send(2).await.unwrap();
output.send(3).await.unwrap();
output.send(4).await.unwrap();
});
let test_program = SubProgramId::new();
TestBuilder::new()
.run_command(test_command.clone(), vec![], |output| if &output != &vec![1, 2, 3, 4] { Err(format!("Unexpected command output: {:?}", output)) } else { Ok(()) })
.run_in_scene_with_threads(&scene, test_program, 5);
}
#[test]
fn pipe_command() {
let scene = Scene::empty();
let test_command = FnCommand::<(), usize>::new(|_input, context| async move {
let mut output = context.send::<usize>(()).unwrap();
println!("send(1)");
output.send(1).await.unwrap();
println!("send(2)");
output.send(2).await.unwrap();
println!("send(3)");
output.send(3).await.unwrap();
println!("send(4)");
output.send(4).await.unwrap();
println!("done.");
});
let add_one_command = FnCommand::<usize, usize>::new(|input, context| async move {
let mut input = input;
let mut output = context.send::<usize>(()).unwrap();
println!("+1 start");
while let Some(next) = input.next().await {
println!("+1: {:?}", next);
output.send(next+1).await.unwrap();
println!(" = {:?}", next+1);
}
println!("+1 done");
});
let combined_command = test_command.pipe_to(add_one_command);
let test_program = SubProgramId::new();
TestBuilder::new()
.run_command(combined_command.clone(), vec![], |output| if &output != &vec![2, 3, 4, 5] { Err(format!("Unexpected command output: {:?}", output)) } else { Ok(()) })
.run_in_scene_with_threads(&scene, test_program, 5);
}
#[test]
fn query_command() {
let scene = Scene::default();
let test_program = SubProgramId::new();
TestBuilder::new()
.send_message(IdleRequest::WhenIdle(test_program))
.expect_message(|_: IdleNotification| { Ok(()) })
.run_query(ReadCommand::default(), Query::<SceneUpdate>::with_no_target(), *SCENE_CONTROL_PROGRAM, |output| if output.len() == 0 { Err(format!("Unexpected command output: {:?}", output)) } else { Ok(()) })
.run_in_scene_with_threads(&scene, test_program, 5);
}
#[test]
fn connect_filter_source_in_command() {
let scene = Scene::default();
#[derive(Serialize, Deserialize)]
struct TestMessage { }
impl SceneMessage for TestMessage { }
let filter_handle = FilterHandle::for_filter::<String, _>(|src| src.map(|_| TestMessage { }));
let cmd_scene = scene.clone();
let test_command = FnCommand::<(), usize>::new(move |_input, context| {
let cmd_scene = cmd_scene.clone();
let filter_handle = filter_handle.clone();
async move {
let mut output = context.send::<usize>(()).unwrap();
println!("Sending initial values");
output.send(1).await.unwrap();
output.send(2).await.unwrap();
println!("Reconnecting");
cmd_scene.connect_programs(StreamSource::Filtered(filter_handle), StreamTarget::Any, StreamId::with_message_type::<String>()).unwrap();
println!("Sending remaining values");
output.send(3).await.unwrap();
output.send(4).await.unwrap();
println!("Done");
}
});
let test_program = SubProgramId::new();
TestBuilder::new()
.run_command(test_command.clone(), vec![], |output| if &output != &vec![1, 2, 3, 4] { Err(format!("Unexpected command output: {:?}", output)) } else { Ok(()) })
.run_in_scene_with_threads(&scene, test_program, 5);
}