use flo_scene::*;
use flo_scene::programs::*;
use futures::prelude::*;
use futures::future::{select, join};
use futures::executor;
use futures_timer::*;
use std::time::{Duration};
use std::sync::*;
#[test]
fn run_subprogram_and_stop_when_scene_is_empty() {
let has_run = Arc::new(Mutex::new(false));
let scene = Scene::empty();
let run_flag = has_run.clone();
scene.add_subprogram(
SubProgramId::new(),
move |_: InputStream<()>, _| async move {
*run_flag.lock().unwrap() = true;
},
0,
);
let mut has_stopped = false;
executor::block_on(select(async {
scene.run_scene().await;
has_stopped = true;
}.boxed(), Delay::new(Duration::from_millis(5000))));
assert!(*has_run.lock().unwrap() == true, "Test program did not run");
assert!(has_stopped, "Scene did not stop when all the subprograms finished");
}
#[test]
fn send_output_to_subprogram_directly() {
let sent_message = Arc::new(Mutex::new(None));
let scene = Scene::empty();
let program_1 = SubProgramId::new();
let program_2 = SubProgramId::new();
let recv_message = sent_message.clone();
scene.add_subprogram(program_1.clone(),
move |mut input: InputStream<usize>, _| async move {
let message = input.next().await.unwrap();
*recv_message.lock().unwrap() = Some(message);
},
0);
scene.add_subprogram(program_2,
move |_: InputStream<()>, context| async move {
let mut send_usize = context.send::<usize>(program_1).unwrap();
send_usize.send(42).await.unwrap();
},
0);
executor::block_on(select(async {
scene.run_scene().await;
}.boxed(), Delay::new(Duration::from_millis(5000))));
assert!(*sent_message.lock().unwrap() == Some(42), "Message was not sent");
}
#[test]
fn connect_before_starting() {
let sent_message = Arc::new(Mutex::new(None));
let scene = Scene::empty();
let program_1 = SubProgramId::new();
let program_2 = SubProgramId::new();
let recv_message = sent_message.clone();
let scene_ref = &scene;
scene.add_subprogram(program_1.clone(),
move |mut input: InputStream<usize>, _| {
scene_ref.connect_programs((), program_1, StreamId::with_message_type::<usize>()).unwrap();
async move {
let message = input.next().await.unwrap();
*recv_message.lock().unwrap() = Some(message);
}
},
0);
scene.add_subprogram(program_2,
move |_: InputStream<()>, context| async move {
let mut send_usize = context.send::<usize>(()).unwrap();
send_usize.send(42).await.unwrap();
},
0);
executor::block_on(select(async {
scene.run_scene().await;
}.boxed(), Delay::new(Duration::from_millis(5000))));
assert!(*sent_message.lock().unwrap() == Some(42), "Message was not sent");
}
#[test]
fn send_output_to_subprogram_via_all_connection() {
let sent_message = Arc::new(Mutex::new(None));
let scene = Scene::empty();
let program_1 = SubProgramId::new();
let program_2 = SubProgramId::new();
let recv_message = sent_message.clone();
scene.add_subprogram(program_1.clone(),
move |mut input: InputStream<usize>, _| async move {
let message = input.next().await.unwrap();
*recv_message.lock().unwrap() = Some(message);
},
0);
scene.add_subprogram(program_2.clone(),
move |_: InputStream<()>, context| async move {
let mut send_usize = context.send::<usize>(StreamTarget::Any).unwrap();
send_usize.send(42).await.unwrap();
},
0);
scene.connect_programs(StreamSource::All, program_1, StreamId::with_message_type::<usize>()).unwrap();
executor::block_on(select(async {
scene.run_scene().await;
}.boxed(), Delay::new(Duration::from_millis(5000))));
assert!(*sent_message.lock().unwrap() == Some(42), "Message was not sent");
}
#[test]
fn send_output_to_subprogram_via_specific_connection() {
let sent_message = Arc::new(Mutex::new(None));
let scene = Scene::empty();
let program_1 = SubProgramId::new();
let program_2 = SubProgramId::new();
let recv_message = sent_message.clone();
scene.add_subprogram(program_1.clone(),
move |mut input: InputStream<usize>, _| async move {
let message = input.next().await.unwrap();
*recv_message.lock().unwrap() = Some(message);
},
0);
scene.add_subprogram(program_2.clone(),
move |_: InputStream<()>, context| async move {
let mut send_usize = context.send::<usize>(StreamTarget::Any).unwrap();
send_usize.send(42).await.unwrap();
},
0);
scene.connect_programs(program_2, program_1, StreamId::with_message_type::<usize>()).unwrap();
executor::block_on(select(async {
scene.run_scene().await;
}.boxed(), Delay::new(Duration::from_millis(5000))));
assert!(*sent_message.lock().unwrap() == Some(42), "Message was not sent");
}
#[test]
fn retrieve_subprogram_id() {
let received_messages = Arc::new(Mutex::new(vec![]));
let scene = Scene::empty();
let program_1 = SubProgramId::new();
let program_2 = SubProgramId::new();
let program_3 = SubProgramId::new();
let stored_messages = received_messages.clone();
scene.add_subprogram(program_1,
move |input: InputStream<String>, _context| async move {
let mut input = input.messages_with_sources();
for _ in 0..4 {
let next = input.next().await.unwrap();
stored_messages.lock().unwrap().push(next);
}
}, 0);
scene.add_subprogram(program_2,
move |_: InputStream<()>, context| async move {
let mut target = context.send::<String>(program_1).unwrap();
target.send("Program 2 message 1".into()).await.unwrap();
target.send("Program 2 message 2".into()).await.unwrap();
}, 0);
scene.add_subprogram(program_3,
move |_: InputStream<()>, context| async move {
let mut target = context.send::<String>(program_1).unwrap();
target.send("Program 3 message 1".into()).await.unwrap();
target.send("Program 3 message 2".into()).await.unwrap();
}, 0);
executor::block_on(select(async {
scene.run_scene().await;
}.boxed(), Delay::new(Duration::from_millis(5000))));
let received_messages = received_messages.lock().unwrap();
assert!(received_messages.len() == 4, "Expected 4 messages to be sent, {:?}", received_messages);
assert!(received_messages.contains(&(program_2, "Program 2 message 1".into())), "Expected program 2 message 1, {:?}", received_messages);
assert!(received_messages.contains(&(program_2, "Program 2 message 2".into())), "Expected program 2 message 2, {:?}", received_messages);
assert!(received_messages.contains(&(program_3, "Program 3 message 1".into())), "Expected program 3 message 1, {:?}", received_messages);
assert!(received_messages.contains(&(program_3, "Program 3 message 2".into())), "Expected program 3 message 2, {:?}", received_messages);
}
#[test]
fn connect_multiple_prorgams_via_any_connection() {
let received_messages = Arc::new(Mutex::new(vec![]));
let scene = Scene::empty();
let program_1 = SubProgramId::new();
let program_2 = SubProgramId::new();
let program_3 = SubProgramId::new();
let stored_messages = received_messages.clone();
scene.add_subprogram(program_1,
move |input: InputStream<String>, _context| async move {
let mut input = input.messages_with_sources();
for _ in 0..4 {
let next = input.next().await.unwrap();
stored_messages.lock().unwrap().push(next);
}
}, 0);
scene.add_subprogram(program_2,
move |_: InputStream<()>, context| async move {
let mut target = context.send::<String>(StreamTarget::Any).unwrap();
target.send("Program 2 message 1".into()).await.unwrap();
target.send("Program 2 message 2".into()).await.unwrap();
}, 0);
scene.add_subprogram(program_3,
move |_: InputStream<()>, context| async move {
let mut target = context.send::<String>(StreamTarget::Any).unwrap();
target.send("Program 3 message 1".into()).await.unwrap();
target.send("Program 3 message 2".into()).await.unwrap();
}, 0);
scene.connect_programs(StreamSource::All, program_1, StreamId::with_message_type::<String>()).unwrap();
executor::block_on(select(async {
scene.run_scene().await;
}.boxed(), Delay::new(Duration::from_millis(5000))));
let received_messages = received_messages.lock().unwrap();
assert!(received_messages.len() == 4, "Expected 4 messages to be sent, {:?}", received_messages);
assert!(received_messages.contains(&(program_2, "Program 2 message 1".into())), "Expected program 2 message 1, {:?}", received_messages);
assert!(received_messages.contains(&(program_2, "Program 2 message 2".into())), "Expected program 2 message 2, {:?}", received_messages);
assert!(received_messages.contains(&(program_3, "Program 3 message 1".into())), "Expected program 3 message 1, {:?}", received_messages);
assert!(received_messages.contains(&(program_3, "Program 3 message 2".into())), "Expected program 3 message 2, {:?}", received_messages);
}
#[test]
fn send_output_via_thread_context() {
let sent_message = Arc::new(Mutex::new(None));
let scene = Scene::empty();
let program_1 = SubProgramId::new();
let program_2 = SubProgramId::new();
let recv_message = sent_message.clone();
scene.add_subprogram(program_1.clone(),
move |mut input: InputStream<usize>, _| async move {
let message = input.next().await.unwrap();
*recv_message.lock().unwrap() = Some(message);
},
0);
scene.add_subprogram(program_2,
move |_: InputStream<()>, _| async move {
let context = scene_context().unwrap();
let mut send_usize = context.send::<usize>(program_1).unwrap();
send_usize.send(42).await.unwrap();
},
0);
executor::block_on(select(async {
scene.run_scene().await;
}.boxed(), Delay::new(Duration::from_millis(5000))));
assert!(*sent_message.lock().unwrap() == Some(42), "Message was not sent");
assert!(scene_context().is_none(), "Scene context should be none outside of the scene");
}
#[test]
fn send_output_from_outside() {
let received_messages = Arc::new(Mutex::new(vec![]));
let scene = Scene::default();
let program_1 = SubProgramId::new();
let stored_messages = received_messages.clone();
scene.add_subprogram(program_1,
move |input: InputStream<String>, _context| async move {
let mut input = input.messages_with_sources();
for _ in 0..4 {
let next = input.next().await.unwrap();
stored_messages.lock().unwrap().push(next);
}
scene_context().unwrap().send_message(SceneControl::StopScene).await.unwrap();
}, 0);
let mut write_messages = scene.send_to_scene::<String>(program_1).unwrap();
executor::block_on(select(async {
join(scene.run_scene(), async {
write_messages.send("One".to_string()).await.unwrap();
write_messages.send("Two".to_string()).await.unwrap();
write_messages.send("Three".to_string()).await.unwrap();
write_messages.send("Four".to_string()).await.unwrap();
}).await;
}.boxed(), Delay::new(Duration::from_millis(5000))));
let received_messages = received_messages.lock().unwrap();
assert!(received_messages.len() == 4, "Expected 4 messages to be sent, {:?}", received_messages);
assert!(received_messages.contains(&(*OUTSIDE_SCENE_PROGRAM, "One".into())), "Expected one, {:?}", received_messages);
assert!(received_messages.contains(&(*OUTSIDE_SCENE_PROGRAM, "Two".into())), "Expected two, {:?}", received_messages);
assert!(received_messages.contains(&(*OUTSIDE_SCENE_PROGRAM, "Three".into())), "Expected three, {:?}", received_messages);
assert!(received_messages.contains(&(*OUTSIDE_SCENE_PROGRAM, "Four".into())), "Expected four, {:?}", received_messages);
}
#[test]
fn send_message_with_thread_stealing() {
let scene = Scene::default();
let received_message = Arc::new(Mutex::new(0));
let received_immediate = Arc::new(Mutex::new(0));
let receiver_program = SubProgramId::new();
let receiver_program_counter = Arc::clone(&received_message);
scene.add_subprogram(receiver_program,
move |messages: InputStream<()>, context| {
messages.allow_thread_stealing(true);
async move {
let mut messages = messages;
while let Some(_msg) = messages.next().await {
println!("Recv");
assert!(context.current_program_id() == Some(receiver_program), "Context program is {:?}, should be {:?}", context.current_program_id(), receiver_program);
assert!(scene_context().unwrap().current_program_id() == Some(receiver_program), "Thread program is {:?}, should be {:?}", scene_context().unwrap().current_program_id(), receiver_program);
*receiver_program_counter.lock().unwrap() += 1;
}
}
}, 20);
let sender_program = SubProgramId::new();
let receiver_program_counter = Arc::clone(&received_message);
let output_counter = Arc::clone(&received_immediate);
scene.add_subprogram(sender_program,
move |_: InputStream<()>, context| {
let mut message_sender = context.send::<()>(receiver_program).unwrap();
async move {
println!("Send");
message_sender.send(()).await.unwrap();
println!("Send");
message_sender.send(()).await.unwrap();
println!("Send");
message_sender.send(()).await.unwrap();
*output_counter.lock().unwrap() = *receiver_program_counter.lock().unwrap();
println!("Send Done");
context.send_message(SceneControl::StopScene).await.unwrap();
}
}, 0);
let mut finished = false;
executor::block_on(select(async {
scene.run_scene().await;
finished = true;
}.boxed(), Delay::new(Duration::from_millis(5000))));
assert!(*received_immediate.lock().unwrap() == 3, "Expected to have processed 3 messages immediated (processed: {:?})", *received_immediate.lock().unwrap());
assert!(finished, "Scene did not finish");
}