use flo_scene::*;
use flo_scene::programs::*;
use futures::prelude::*;
use futures::future::{select};
use futures::executor;
use futures_timer::*;
use serde::*;
use std::time::{Duration};
use std::sync::*;
use std::sync::atomic::{AtomicUsize, Ordering};
#[test]
fn write_to_filter_target() {
let recv_messages = Arc::new(Mutex::new(vec![]));
let scene = Scene::empty();
let sent_messages = recv_messages.clone();
let usize_to_string = FilterHandle::for_filter(|number_stream: InputStream<usize>| number_stream.map(|num| num.to_string()));
let string_program = SubProgramId::new();
scene.add_subprogram(
string_program,
move |mut strings: InputStream<String>, _| async move {
let next_string = strings.next().await.unwrap();
sent_messages.lock().unwrap().push(next_string);
let next_string = strings.next().await.unwrap();
sent_messages.lock().unwrap().push(next_string);
let next_string = strings.next().await.unwrap();
sent_messages.lock().unwrap().push(next_string);
let next_string = strings.next().await.unwrap();
sent_messages.lock().unwrap().push(next_string);
},
0,
);
let number_program = SubProgramId::new();
scene.add_subprogram(
number_program,
move |_: InputStream<()>, context| async move {
let mut filtered_output = context.send::<usize>(StreamTarget::Filtered(usize_to_string, string_program)).unwrap();
filtered_output.send(1).await.unwrap();
filtered_output.send(2).await.unwrap();
filtered_output.send(3).await.unwrap();
filtered_output.send(4).await.unwrap();
},
0);
let mut has_finished = false;
executor::block_on(select(async {
scene.run_scene().await;
has_finished = true;
}.boxed(), Delay::new(Duration::from_millis(5000))));
let recv_messages = (*recv_messages.lock().unwrap()).clone();
assert!(recv_messages == vec![1.to_string(), 2.to_string(), 3.to_string(), 4.to_string()], "Test program did not send correct numbers (sent {:?})", recv_messages);
assert!(has_finished, "Scene did not finish when the programs terminated");
}
#[test]
fn write_to_conversion_filter() {
let recv_messages = Arc::new(Mutex::new(vec![]));
let scene = Scene::empty();
let sent_messages = recv_messages.clone();
let u32_to_u64 = FilterHandle::conversion_filter::<u32, u64>();
let string_program = SubProgramId::new();
scene.add_subprogram(
string_program,
move |mut numbers: InputStream<u64>, _| async move {
let next_number = numbers.next().await.unwrap();
sent_messages.lock().unwrap().push(next_number.to_string());
let next_number = numbers.next().await.unwrap();
sent_messages.lock().unwrap().push(next_number.to_string());
let next_number = numbers.next().await.unwrap();
sent_messages.lock().unwrap().push(next_number.to_string());
let next_number = numbers.next().await.unwrap();
sent_messages.lock().unwrap().push(next_number.to_string());
},
0,
);
let number_program = SubProgramId::new();
scene.add_subprogram(
number_program,
move |_: InputStream<()>, context| async move {
let mut filtered_output = context.send::<u32>(StreamTarget::Filtered(u32_to_u64, string_program)).unwrap();
filtered_output.send(1u32).await.unwrap();
filtered_output.send(2u32).await.unwrap();
filtered_output.send(3u32).await.unwrap();
filtered_output.send(4u32).await.unwrap();
},
0);
let mut has_finished = false;
executor::block_on(select(async {
scene.run_scene().await;
has_finished = true;
}.boxed(), Delay::new(Duration::from_millis(5000))));
let recv_messages = (*recv_messages.lock().unwrap()).clone();
assert!(recv_messages == vec![1.to_string(), 2.to_string(), 3.to_string(), 4.to_string()], "Test program did not send correct numbers (sent {:?})", recv_messages);
assert!(has_finished, "Scene did not finish when the programs terminated");
}
#[test]
fn apply_filter_to_direct_connection() {
let recv_messages = Arc::new(Mutex::new(vec![]));
let scene = Scene::empty();
let sent_messages = recv_messages.clone();
let usize_to_string = FilterHandle::for_filter(|number_stream: InputStream<usize>| number_stream.map(|num| num.to_string()));
let string_program = SubProgramId::new();
scene.add_subprogram(
string_program,
move |mut strings: InputStream<String>, _| async move {
let next_string = strings.next().await.unwrap();
sent_messages.lock().unwrap().push(next_string);
let next_string = strings.next().await.unwrap();
sent_messages.lock().unwrap().push(next_string);
let next_string = strings.next().await.unwrap();
sent_messages.lock().unwrap().push(next_string);
let next_string = strings.next().await.unwrap();
sent_messages.lock().unwrap().push(next_string);
},
0,
);
scene.connect_programs((), StreamTarget::Filtered(usize_to_string, string_program), StreamId::with_message_type::<usize>().for_target(string_program)).unwrap();
let number_program = SubProgramId::new();
scene.add_subprogram(
number_program,
move |_: InputStream<()>, context| async move {
let mut filtered_output = context.send::<usize>(string_program).unwrap();
filtered_output.send(1).await.unwrap();
filtered_output.send(2).await.unwrap();
filtered_output.send(3).await.unwrap();
filtered_output.send(4).await.unwrap();
},
0);
let mut has_finished = false;
executor::block_on(select(async {
scene.run_scene().await;
has_finished = true;
}.boxed(), Delay::new(Duration::from_millis(5000))));
let recv_messages = (*recv_messages.lock().unwrap()).clone();
assert!(recv_messages == vec![1.to_string(), 2.to_string(), 3.to_string(), 4.to_string()], "Test program did not send correct numbers (sent {:?})", recv_messages);
assert!(has_finished, "Scene did not finish when the programs terminated");
}
#[test]
fn connect_all_to_filter_target() {
let recv_messages = Arc::new(Mutex::new(vec![]));
let scene = Scene::empty();
let sent_messages = recv_messages.clone();
let usize_to_string = FilterHandle::for_filter(|number_stream: InputStream<usize>| number_stream.map(|num| num.to_string()));
let string_program = SubProgramId::new();
scene.add_subprogram(
string_program,
move |mut strings: InputStream<String>, _| async move {
let next_string = strings.next().await.unwrap();
sent_messages.lock().unwrap().push(next_string);
let next_string = strings.next().await.unwrap();
sent_messages.lock().unwrap().push(next_string);
let next_string = strings.next().await.unwrap();
sent_messages.lock().unwrap().push(next_string);
let next_string = strings.next().await.unwrap();
sent_messages.lock().unwrap().push(next_string);
},
0,
);
let number_program = SubProgramId::new();
scene.add_subprogram(
number_program,
move |_: InputStream<()>, context| async move {
let mut filtered_output = context.send::<usize>(StreamTarget::Any).unwrap();
filtered_output.send(1).await.unwrap();
filtered_output.send(2).await.unwrap();
filtered_output.send(3).await.unwrap();
filtered_output.send(4).await.unwrap();
},
0);
scene.connect_programs((), StreamTarget::Filtered(usize_to_string, string_program), StreamId::with_message_type::<usize>()).unwrap();
let mut has_finished = false;
executor::block_on(select(async {
scene.run_scene().await;
has_finished = true;
}.boxed(), Delay::new(Duration::from_millis(5000))));
let recv_messages = (*recv_messages.lock().unwrap()).clone();
assert!(recv_messages == vec![1.to_string(), 2.to_string(), 3.to_string(), 4.to_string()], "Test program did not send correct numbers (sent {:?})", recv_messages);
assert!(has_finished, "Scene did not finish when the programs terminated");
}
#[test]
fn direct_connection_via_general_filter() {
let recv_messages = Arc::new(Mutex::new(vec![]));
let scene = Scene::empty();
let sent_messages = recv_messages.clone();
let usize_to_string = FilterHandle::for_filter(|number_stream: InputStream<usize>| number_stream.map(|num| num.to_string()));
let string_program = SubProgramId::new();
scene.add_subprogram(
string_program,
move |mut strings: InputStream<String>, _| async move {
let next_string = strings.next().await.unwrap();
sent_messages.lock().unwrap().push(next_string);
let next_string = strings.next().await.unwrap();
sent_messages.lock().unwrap().push(next_string);
let next_string = strings.next().await.unwrap();
sent_messages.lock().unwrap().push(next_string);
let next_string = strings.next().await.unwrap();
sent_messages.lock().unwrap().push(next_string);
},
0,
);
let number_program = SubProgramId::new();
scene.add_subprogram(
number_program,
move |_: InputStream<()>, context| async move {
let mut filtered_output = context.send::<usize>(string_program).unwrap();
filtered_output.send(1).await.unwrap();
filtered_output.send(2).await.unwrap();
filtered_output.send(3).await.unwrap();
filtered_output.send(4).await.unwrap();
},
0);
scene.connect_programs((), StreamTarget::Filtered(usize_to_string, string_program), StreamId::with_message_type::<usize>()).unwrap();
let mut has_finished = false;
executor::block_on(select(async {
scene.run_scene().await;
has_finished = true;
}.boxed(), Delay::new(Duration::from_millis(5000))));
let recv_messages = (*recv_messages.lock().unwrap()).clone();
assert!(recv_messages == vec![1.to_string(), 2.to_string(), 3.to_string(), 4.to_string()], "Test program did not send correct numbers (sent {:?})", recv_messages);
assert!(has_finished, "Scene did not finish when the programs terminated");
}
#[test]
fn connect_all_both_filtered_and_unfiltered() {
let recv_messages = Arc::new(Mutex::new(vec![]));
let scene = Scene::empty();
let sent_messages = recv_messages.clone();
let usize_to_string = FilterHandle::for_filter(|number_stream: InputStream<usize>| number_stream.map(|num| num.to_string()));
let string_program = SubProgramId::new();
scene.add_subprogram(
string_program,
move |mut strings: InputStream<String>, _| async move {
for _ in 0..8 {
let next_string = strings.next().await.unwrap();
sent_messages.lock().unwrap().push(next_string);
}
},
3,
);
let string_generator_program = SubProgramId::new();
scene.add_subprogram(
string_generator_program,
move |_: InputStream<()>, context| async move {
let mut filtered_output = context.send::<String>(StreamTarget::Any).unwrap();
filtered_output.send(5.to_string()).await.unwrap();
filtered_output.send(6.to_string()).await.unwrap();
filtered_output.send(7.to_string()).await.unwrap();
filtered_output.send(8.to_string()).await.unwrap();
},
0);
let number_program = SubProgramId::new();
scene.add_subprogram(
number_program,
move |_: InputStream<()>, context| async move {
let mut filtered_output = context.send::<usize>(StreamTarget::Any).unwrap();
filtered_output.send(1).await.unwrap();
filtered_output.send(2).await.unwrap();
filtered_output.send(3).await.unwrap();
filtered_output.send(4).await.unwrap();
},
0);
scene.connect_programs((), string_program, StreamId::with_message_type::<String>()).unwrap();
scene.connect_programs((), StreamTarget::Filtered(usize_to_string, string_program), StreamId::with_message_type::<usize>()).unwrap();
let mut has_finished = false;
executor::block_on(select(async {
scene.run_scene().await;
has_finished = true;
}.boxed(), Delay::new(Duration::from_millis(5000))));
let recv_messages = (*recv_messages.lock().unwrap()).clone();
let one_to_four = recv_messages.iter().filter(|msg| *msg == "1" || *msg == "2" || *msg == "3" || *msg == "4").cloned().collect::<Vec<_>>();
let five_to_eight = recv_messages.iter().filter(|msg| *msg == "5" || *msg == "6" || *msg == "7" || *msg == "8").cloned().collect::<Vec<_>>();
assert!(recv_messages.len() == 8, "Wrong number of messages received: {:?}", recv_messages);
assert!(one_to_four == vec![1.to_string(), 2.to_string(), 3.to_string(), 4.to_string()], "Test program did not send correct numbers (sent {:?})", recv_messages);
assert!(five_to_eight == vec![5.to_string(), 6.to_string(), 7.to_string(), 8.to_string()], "Test program did not send correct numbers (sent {:?})", recv_messages);
assert!(has_finished, "Scene did not finish when the programs terminated");
}
#[test]
fn disconnect_filter_target() {
let recv_messages = Arc::new(Mutex::new(vec![]));
static NUM_DISCONNECTS: AtomicUsize = AtomicUsize::new(0);
struct CountDisconnents {
}
impl CountDisconnents {
fn convert_to_string(&self, num: usize) -> String {
num.to_string()
}
}
impl Drop for CountDisconnents {
fn drop(&mut self) {
println!("Disconnect");
NUM_DISCONNECTS.fetch_add(1, Ordering::Relaxed);
}
}
NUM_DISCONNECTS.store(0, Ordering::Relaxed);
let scene = Arc::new(Scene::empty());
let sent_messages = recv_messages.clone();
println!("Register filter");
let usize_to_string = FilterHandle::for_filter(|number_stream: InputStream<usize>| {
println!("Connecting");
let count_disconnects = CountDisconnents {};
number_stream.map(move |num| count_disconnects.convert_to_string(num))
});
println!("Create string program");
let string_program = SubProgramId::new();
scene.add_subprogram(
string_program,
move |mut strings: InputStream<String>, _| async move {
let next_string = strings.next().await.unwrap();
sent_messages.lock().unwrap().push(next_string);
let next_string = strings.next().await.unwrap();
sent_messages.lock().unwrap().push(next_string);
let next_string = strings.next().await.unwrap();
sent_messages.lock().unwrap().push(next_string);
let next_string = strings.next().await.unwrap();
sent_messages.lock().unwrap().push(next_string);
},
0,
);
println!("Create number program");
let number_program = SubProgramId::new();
let scene2 = scene.clone();
let usize_to_string2 = usize_to_string.clone();
scene.add_subprogram(
number_program,
move |_: InputStream<()>, context| async move {
println!("Create initial stream...");
let mut filtered_output = context.send::<usize>(StreamTarget::Any).unwrap();
println!(" ... created");
filtered_output.send(1).await.unwrap();
filtered_output.send(2).await.unwrap();
println!("Disconnecting initial stream...");
scene2.connect_programs(number_program, StreamTarget::None, StreamId::with_message_type::<usize>()).unwrap();
println!(" ... disconnected");
filtered_output.send(3).await.unwrap();
filtered_output.send(4).await.unwrap();
println!("Reconnecting filter stream...");
scene2.connect_programs(number_program, StreamTarget::Filtered(usize_to_string2, string_program), StreamId::with_message_type::<usize>()).unwrap();
println!(" ... reconnected");
filtered_output.send(5).await.unwrap();
filtered_output.send(6).await.unwrap();
println!("Disconnecting again...");
scene2.connect_programs(number_program, StreamTarget::None, StreamId::with_message_type::<usize>()).unwrap();
println!(" ... disconnected");
},
0);
println!("Connect programs");
scene.connect_programs(number_program, StreamTarget::Filtered(usize_to_string, string_program), StreamId::with_message_type::<usize>()).unwrap();
println!("Start scene");
let mut has_finished = false;
executor::block_on(select(async {
scene.run_scene().await;
has_finished = true;
}.boxed(), Delay::new(Duration::from_millis(5000))));
println!("Scene finished");
let recv_messages = (*recv_messages.lock().unwrap()).clone();
let num_disconnects = NUM_DISCONNECTS.load(Ordering::Relaxed);
assert!(recv_messages == vec![1.to_string(), 2.to_string(), 5.to_string(), 6.to_string()], "Test program did not send correct numbers (sent {:?})", recv_messages);
assert!(num_disconnects == 2, "Filtered stream was not dropped the expected number of times ({} != 2)", num_disconnects);
assert!(has_finished, "Scene did not finish when the programs terminated");
}
#[test]
fn filter_with_send_and_target_filter() {
let scene = Scene::default();
#[derive(Debug, PartialEq, Serialize, Deserialize)]
enum Message1 { Msg(String) }
#[derive(Debug, PartialEq, Serialize, Deserialize)]
enum Message2 { Msg(String) }
impl SceneMessage for Message1 { }
impl SceneMessage for Message2 { }
let str_to_msg1 = FilterHandle::for_filter(|string_value: InputStream<String>| string_value.map(|val| Message1::Msg(val.to_string())));
let msg1_to_msg2 = FilterHandle::for_filter(|msg1: InputStream<Message1>| { msg1.map(|msg| match msg { Message1::Msg(val) => Message2::Msg(val) })});
let message1_sender_program = SubProgramId::new();
let message2_receiver_program = SubProgramId::new();
let test_program = SubProgramId::new();
scene.add_subprogram(message1_sender_program, |_: InputStream<()>, context| async move {
let mut sender = context.send(StreamTarget::Filtered(str_to_msg1, message2_receiver_program)).unwrap();
println!("Sending 1...");
sender.send("Hello".to_string()).await.unwrap();
println!("Sending 2...");
sender.send("Goodbyte".to_string()).await.unwrap();
println!("Done");
}, 0);
scene.add_subprogram(message2_receiver_program, |mut input: InputStream<Message2>, context| async move {
let mut test_program = context.send(()).unwrap();
while let Some(input) = input.next().await {
test_program.send(input).await.unwrap();
}
}, 0);
scene.connect_programs((), StreamTarget::Filtered(msg1_to_msg2, message2_receiver_program), StreamId::with_message_type::<Message1>()).unwrap();
TestBuilder::new()
.expect_message(|msg2: Message2| if msg2 != Message2::Msg("Hello".to_string()) { Err(format!("Expected 'Hello'")) } else { Ok(()) })
.expect_message(|msg2: Message2| if msg2 != Message2::Msg("Goodbyte".to_string()) { Err(format!("Expected 'Goodbyte'")) } else { Ok(()) })
.run_in_scene(&scene, test_program);
}
#[test]
fn filter_at_target() {
let scene = Scene::default();
#[derive(Debug, PartialEq, Serialize, Deserialize)]
enum Message1 { Msg(String) }
#[derive(Debug, PartialEq, Serialize, Deserialize)]
enum Message2 { Msg(String) }
impl SceneMessage for Message1 { }
impl SceneMessage for Message2 { }
let msg1_to_msg2 = FilterHandle::for_filter(|msg1: InputStream<Message1>| { msg1.map(|msg| match msg { Message1::Msg(val) => Message2::Msg(val) })});
let message1_sender_program = SubProgramId::new();
let message2_receiver_program = SubProgramId::new();
let test_program = SubProgramId::new();
scene.add_subprogram(message1_sender_program, |_: InputStream<()>, context| async move {
let mut sender = context.send(()).unwrap();
println!("Sending 1...");
sender.send(Message1::Msg("Hello".to_string())).await.unwrap();
println!("Sending 2...");
sender.send(Message1::Msg("Goodbyte".to_string())).await.unwrap();
println!("Done");
}, 0);
scene.add_subprogram(message2_receiver_program, |mut input: InputStream<Message2>, context| async move {
let mut test_program = context.send(()).unwrap();
while let Some(input) = input.next().await {
test_program.send(input).await.unwrap();
}
}, 0);
scene.connect_programs((), StreamTarget::Filtered(msg1_to_msg2, message2_receiver_program), StreamId::with_message_type::<Message1>()).unwrap();
TestBuilder::new()
.expect_message(|msg2: Message2| if msg2 != Message2::Msg("Hello".to_string()) { Err(format!("Expected 'Hello'")) } else { Ok(()) })
.expect_message(|msg2: Message2| if msg2 != Message2::Msg("Goodbyte".to_string()) { Err(format!("Expected 'Goodbyte'")) } else { Ok(()) })
.run_in_scene(&scene, test_program);
}
#[test]
fn filter_target_using_source_filter() {
let scene = Scene::default();
#[derive(Debug, PartialEq, Serialize, Deserialize)]
enum Message1 { Msg(String) }
#[derive(Debug, PartialEq, Serialize, Deserialize)]
enum Message2 { Msg(String) }
impl SceneMessage for Message1 { }
impl SceneMessage for Message2 { }
let msg1_to_msg2 = FilterHandle::for_filter(|msg1: InputStream<Message1>| { msg1.map(|msg| match msg { Message1::Msg(val) => Message2::Msg(val) })});
let message1_sender_program = SubProgramId::new();
let message2_receiver_program = SubProgramId::new();
let test_program = SubProgramId::new();
scene.add_subprogram(message1_sender_program, |_: InputStream<()>, context| async move {
let mut sender = context.send(()).unwrap();
println!("Sending 1...");
sender.send(Message1::Msg("Hello".to_string())).await.unwrap();
println!("Sending 2...");
sender.send(Message1::Msg("Goodbyte".to_string())).await.unwrap();
println!("Done");
}, 0);
scene.add_subprogram(message2_receiver_program, |mut input: InputStream<Message2>, context| async move {
let mut test_program = context.send(()).unwrap();
while let Some(input) = input.next().await {
test_program.send(input).await.unwrap();
}
}, 0);
scene.connect_programs(StreamSource::Filtered(msg1_to_msg2), message2_receiver_program, StreamId::with_message_type::<Message1>()).unwrap();
TestBuilder::new()
.expect_message(|msg2: Message2| if msg2 != Message2::Msg("Hello".to_string()) { Err(format!("Expected 'Hello'")) } else { Ok(()) })
.expect_message(|msg2: Message2| if msg2 != Message2::Msg("Goodbyte".to_string()) { Err(format!("Expected 'Goodbyte'")) } else { Ok(()) })
.run_in_scene(&scene, test_program);
}
#[test]
fn filter_at_source_with_specific_target() {
let scene = Scene::default();
#[derive(Debug, Serialize, Deserialize)]
enum Message1 { Msg(String) }
#[derive(Debug, PartialEq, Serialize, Deserialize)]
enum Message2 { Msg(String) }
impl SceneMessage for Message1 { }
impl SceneMessage for Message2 { }
let msg1_to_msg2 = FilterHandle::for_filter(|msg1: InputStream<Message1>| { msg1.map(|msg| match msg { Message1::Msg(val) => Message2::Msg(val) })});
scene.connect_programs(msg1_to_msg2, (), StreamId::with_message_type::<Message1>()).unwrap();
let message1_sender_program = SubProgramId::new();
let message2_receiver_program = SubProgramId::new();
let test_program = SubProgramId::new();
scene.add_subprogram(message1_sender_program, |_: InputStream<()>, context| async move {
let mut sender = context.send(()).unwrap();
println!("Sending 1...");
sender.send(Message1::Msg("Hello".to_string())).await.unwrap();
println!("Sending 2...");
sender.send(Message1::Msg("Goodbyte".to_string())).await.unwrap();
println!("Done");
}, 0);
scene.add_subprogram(message2_receiver_program, |mut input: InputStream<Message2>, context| async move {
let mut test_program = context.send(()).unwrap();
while let Some(input) = input.next().await {
test_program.send(input).await.unwrap();
}
}, 0);
scene.connect_programs(message1_sender_program, message2_receiver_program, StreamId::with_message_type::<Message1>()).unwrap();
TestBuilder::new()
.expect_message(|msg2: Message2| if msg2 != Message2::Msg("Hello".to_string()) { Err(format!("Expected 'Hello'")) } else { Ok(()) })
.expect_message(|msg2: Message2| if msg2 != Message2::Msg("Goodbyte".to_string()) { Err(format!("Expected 'Goodbyte'")) } else { Ok(()) })
.run_in_scene(&scene, test_program);
}
#[test]
fn filter_at_source_with_any_target() {
let scene = Scene::default();
#[derive(Debug, Serialize, Deserialize)]
enum Message1 { Msg(String) }
#[derive(Debug, PartialEq, Serialize, Deserialize)]
enum Message2 { Msg(String) }
impl SceneMessage for Message1 { }
impl SceneMessage for Message2 { }
let msg1_to_msg2 = FilterHandle::for_filter(|msg1: InputStream<Message1>| { msg1.map(|msg| match msg { Message1::Msg(val) => Message2::Msg(val) })});
scene.connect_programs(msg1_to_msg2, (), StreamId::with_message_type::<Message1>()).unwrap();
let message1_sender_program = SubProgramId::new();
let message2_receiver_program = SubProgramId::new();
let test_program = SubProgramId::new();
scene.add_subprogram(message1_sender_program, |_: InputStream<()>, context| async move {
let mut sender = context.send(()).unwrap();
println!("Sending 1...");
sender.send(Message1::Msg("Hello".to_string())).await.unwrap();
println!("Sending 2...");
sender.send(Message1::Msg("Goodbyte".to_string())).await.unwrap();
println!("Done");
}, 0);
scene.add_subprogram(message2_receiver_program, |mut input: InputStream<Message2>, context| async move {
let mut test_program = context.send(()).unwrap();
while let Some(input) = input.next().await {
match input {
Message2::Msg(msg) => test_program.send(msg).await.unwrap()
}
}
}, 0);
scene.connect_programs((), message2_receiver_program, StreamId::with_message_type::<Message2>()).unwrap();
TestBuilder::new()
.expect_message(|msg2: String| if msg2 != "Hello".to_string() { Err(format!("Expected 'Hello'")) } else { Ok(()) })
.expect_message(|msg2: String| if msg2 != "Goodbyte".to_string() { Err(format!("Expected 'Goodbyte'")) } else { Ok(()) })
.run_in_scene(&scene, test_program);
}
#[test]
fn filter_at_source_with_direct_target() {
let scene = Scene::default();
#[derive(Debug, Serialize, Deserialize)]
enum Message1 { Msg(String) }
#[derive(Debug, PartialEq, Serialize, Deserialize)]
enum Message2 { Msg(String) }
impl SceneMessage for Message1 { }
impl SceneMessage for Message2 { }
let msg1_to_msg2 = FilterHandle::for_filter(|msg1: InputStream<Message1>| { msg1.map(|msg| match msg { Message1::Msg(val) => Message2::Msg(val) })});
scene.connect_programs(msg1_to_msg2, (), StreamId::with_message_type::<Message1>()).unwrap();
let message1_sender_program = SubProgramId::new();
let message2_receiver_program = SubProgramId::new();
let test_program = SubProgramId::new();
scene.add_subprogram(message1_sender_program, |_: InputStream<()>, context| async move {
let mut sender = context.send(message2_receiver_program).unwrap();
println!("Sending 1...");
sender.send(Message1::Msg("Hello".to_string())).await.unwrap();
println!("Sending 2...");
sender.send(Message1::Msg("Goodbyte".to_string())).await.unwrap();
println!("Done");
}, 0);
scene.add_subprogram(message2_receiver_program, |mut input: InputStream<Message2>, context| async move {
let mut test_program = context.send(()).unwrap();
while let Some(input) = input.next().await {
match input {
Message2::Msg(msg) => test_program.send(msg).await.unwrap()
}
}
}, 0);
TestBuilder::new()
.expect_message(|msg2: String| if msg2 != "Hello".to_string() { Err(format!("Expected 'Hello'")) } else { Ok(()) })
.expect_message(|msg2: String| if msg2 != "Goodbyte".to_string() { Err(format!("Expected 'Goodbyte'")) } else { Ok(()) })
.run_in_scene(&scene, test_program);
}
#[test]
fn chain_filters_with_target_filter() {
let scene = Scene::default();
#[derive(Debug, Serialize, Deserialize)]
enum Message1 { Msg(String) }
#[derive(Debug, Serialize, Deserialize)]
enum Message2 { Msg(String) }
#[derive(Debug, Serialize, Deserialize)]
enum Message3 { Msg(String) }
impl SceneMessage for Message1 { }
impl SceneMessage for Message2 { }
impl SceneMessage for Message3 { }
let msg1_to_msg2 = FilterHandle::for_filter(|msg1: InputStream<Message1>| { msg1.map(|msg| match msg { Message1::Msg(val) => Message2::Msg(val) })});
let msg2_to_msg3 = FilterHandle::for_filter(|msg1: InputStream<Message2>| { msg1.map(|msg| match msg { Message2::Msg(val) => Message3::Msg(val) })});
scene.connect_programs(msg1_to_msg2, (), StreamId::with_message_type::<Message1>()).unwrap();
let message1_sender_program = SubProgramId::new();
let message2_receiver_program = SubProgramId::new();
let message3_receiver_program = SubProgramId::new();
let test_program = SubProgramId::new();
scene.add_subprogram(message1_sender_program, |_: InputStream<()>, context| async move {
let mut sender = context.send(()).unwrap();
sender.send(Message1::Msg("Hello".to_string())).await.unwrap();
sender.send(Message1::Msg("Goodbyte".to_string())).await.unwrap();
}, 0);
scene.add_subprogram(message2_receiver_program, |mut input: InputStream<Message2>, _| async move {
while let Some(_input) = input.next().await {
assert!(false, "Should receive only Message3");
}
}, 0);
scene.add_subprogram(message3_receiver_program, |mut input, context| async move {
let mut test_program = context.send(()).unwrap();
while let Some(input) = input.next().await {
match input {
Message3::Msg(msg) => test_program.send(msg).await.unwrap()
}
}
}, 0);
scene.connect_programs((), message2_receiver_program, StreamId::with_message_type::<Message2>()).unwrap();
scene.connect_programs((), message3_receiver_program, StreamId::with_message_type::<Message3>()).unwrap();
scene.connect_programs((), StreamTarget::Filtered(msg2_to_msg3, message3_receiver_program), StreamId::with_message_type::<Message2>()).unwrap();
TestBuilder::new()
.expect_message(|msg2: String| if msg2 != "Hello".to_string() { Err(format!("Expected 'Hello'")) } else { Ok(()) })
.expect_message(|msg2: String| if msg2 != "Goodbyte".to_string() { Err(format!("Expected 'Goodbyte'")) } else { Ok(()) })
.run_in_scene(&scene, test_program);
}
#[test]
fn chain_two_filters_with_target_filter_1() {
let scene = Scene::default();
#[derive(Debug, Serialize, Deserialize)]
enum Message1 { Msg(String) }
#[derive(Debug, Serialize, Deserialize)]
enum Message2 { Msg(String) }
#[derive(Debug, Serialize, Deserialize)]
enum Message3 { Msg(String) }
#[derive(Debug, Serialize, Deserialize)]
enum Message4 { Msg(String) }
impl SceneMessage for Message1 { }
impl SceneMessage for Message2 { }
impl SceneMessage for Message3 { }
impl SceneMessage for Message4 { }
let msg1_to_msg2 = FilterHandle::for_filter(|msg1: InputStream<Message1>| { msg1.map(|msg| match msg { Message1::Msg(val) => Message2::Msg(val) })});
let msg1_to_msg4 = FilterHandle::for_filter(|msg1: InputStream<Message1>| { msg1.map(|msg| match msg { Message1::Msg(val) => Message4::Msg(val) })});
let msg2_to_msg3 = FilterHandle::for_filter(|msg1: InputStream<Message2>| { msg1.map(|msg| match msg { Message2::Msg(val) => Message3::Msg(val) })});
scene.connect_programs(msg1_to_msg2, (), StreamId::with_message_type::<Message1>()).unwrap();
scene.connect_programs(msg1_to_msg4, (), StreamId::with_message_type::<Message1>()).unwrap();
let message1_sender_program = SubProgramId::new();
let message2_receiver_program = SubProgramId::new();
let message3_receiver_program = SubProgramId::new();
let test_program = SubProgramId::new();
scene.add_subprogram(message1_sender_program, |_: InputStream<()>, context| async move {
let mut sender = context.send(()).unwrap();
sender.send(Message1::Msg("Hello".to_string())).await.unwrap();
sender.send(Message1::Msg("Goodbyte".to_string())).await.unwrap();
}, 0);
scene.add_subprogram(message2_receiver_program, |mut input: InputStream<Message2>, _| async move {
while let Some(_input) = input.next().await {
assert!(false, "Should receive only Message3");
}
}, 0);
scene.add_subprogram(message3_receiver_program, |mut input, context| async move {
let mut test_program = context.send(()).unwrap();
while let Some(input) = input.next().await {
match input {
Message3::Msg(msg) => test_program.send(msg).await.unwrap()
}
}
}, 0);
scene.connect_programs((), message2_receiver_program, StreamId::with_message_type::<Message2>()).unwrap();
scene.connect_programs((), message3_receiver_program, StreamId::with_message_type::<Message3>()).unwrap();
scene.connect_programs((), StreamTarget::Filtered(msg2_to_msg3, message3_receiver_program), StreamId::with_message_type::<Message2>()).unwrap();
TestBuilder::new()
.expect_message(|msg2: String| if msg2 != "Hello".to_string() { Err(format!("Expected 'Hello'")) } else { Ok(()) })
.expect_message(|msg2: String| if msg2 != "Goodbyte".to_string() { Err(format!("Expected 'Goodbyte'")) } else { Ok(()) })
.run_in_scene(&scene, test_program);
}
#[test]
fn chain_two_filters_with_target_filter_2() {
let scene = Scene::default();
#[derive(Debug, Serialize, Deserialize)]
enum Message1 { Msg(String) }
#[derive(Debug, Serialize, Deserialize)]
enum Message2 { Msg(String) }
#[derive(Debug, Serialize, Deserialize)]
enum Message3 { Msg(String) }
#[derive(Debug, Serialize, Deserialize)]
enum Message4 { Msg(String) }
impl SceneMessage for Message1 { }
impl SceneMessage for Message2 { }
impl SceneMessage for Message3 { }
impl SceneMessage for Message4 { }
let msg1_to_msg2 = FilterHandle::for_filter(|msg1: InputStream<Message1>| { msg1.map(|msg| match msg { Message1::Msg(val) => Message2::Msg(val) })});
let msg1_to_msg4 = FilterHandle::for_filter(|msg1: InputStream<Message1>| { msg1.map(|msg| match msg { Message1::Msg(val) => Message4::Msg(val) })});
let msg2_to_msg3 = FilterHandle::for_filter(|msg1: InputStream<Message2>| { msg1.map(|msg| match msg { Message2::Msg(val) => Message3::Msg(val) })});
scene.connect_programs(msg1_to_msg4, (), StreamId::with_message_type::<Message1>()).unwrap();
scene.connect_programs(msg1_to_msg2, (), StreamId::with_message_type::<Message1>()).unwrap();
let message1_sender_program = SubProgramId::new();
let message2_receiver_program = SubProgramId::new();
let message3_receiver_program = SubProgramId::new();
let test_program = SubProgramId::new();
scene.add_subprogram(message1_sender_program, |_: InputStream<()>, context| async move {
let mut sender = context.send(()).unwrap();
sender.send(Message1::Msg("Hello".to_string())).await.unwrap();
sender.send(Message1::Msg("Goodbyte".to_string())).await.unwrap();
}, 0);
scene.add_subprogram(message2_receiver_program, |mut input: InputStream<Message2>, _| async move {
while let Some(_input) = input.next().await {
assert!(false, "Should receive only Message3");
}
}, 0);
scene.add_subprogram(message3_receiver_program, |mut input, context| async move {
let mut test_program = context.send(()).unwrap();
while let Some(input) = input.next().await {
match input {
Message3::Msg(msg) => test_program.send(msg).await.unwrap()
}
}
}, 0);
scene.connect_programs((), message2_receiver_program, StreamId::with_message_type::<Message2>()).unwrap();
scene.connect_programs((), message3_receiver_program, StreamId::with_message_type::<Message3>()).unwrap();
scene.connect_programs((), StreamTarget::Filtered(msg2_to_msg3, message3_receiver_program), StreamId::with_message_type::<Message2>()).unwrap();
TestBuilder::new()
.expect_message(|msg2: String| if msg2 != "Hello".to_string() { Err(format!("Expected 'Hello'")) } else { Ok(()) })
.expect_message(|msg2: String| if msg2 != "Goodbyte".to_string() { Err(format!("Expected 'Goodbyte'")) } else { Ok(()) })
.run_in_scene(&scene, test_program);
}
#[test]
fn chain_with_direct_target() {
let scene = Scene::default();
#[derive(Debug, Serialize, Deserialize)]
enum Message1 { Msg(String) }
#[derive(Debug, PartialEq, Serialize, Deserialize)]
enum Message2 { Msg(String) }
#[derive(Debug, Serialize, Deserialize)]
enum Message3 { Msg(String) }
impl SceneMessage for Message1 { }
impl SceneMessage for Message2 { }
impl SceneMessage for Message3 { }
let msg1_to_msg2 = FilterHandle::for_filter(|msg1: InputStream<Message1>| { msg1.map(|msg| match msg { Message1::Msg(val) => Message2::Msg(val) })});
let msg2_to_msg3 = FilterHandle::for_filter(|msg1: InputStream<Message2>| { msg1.map(|msg| match msg { Message2::Msg(val) => Message3::Msg(val) })});
scene.connect_programs(msg1_to_msg2, (), StreamId::with_message_type::<Message1>()).unwrap();
let message1_sender_program = SubProgramId::new();
let message3_receiver_program = SubProgramId::new();
let test_program = SubProgramId::new();
scene.add_subprogram(message1_sender_program, |_: InputStream<()>, context| async move {
let mut sender = context.send(message3_receiver_program).unwrap();
println!("Sending 1...");
sender.send(Message1::Msg("Hello".to_string())).await.unwrap();
println!("Sending 2...");
sender.send(Message1::Msg("Goodbyte".to_string())).await.unwrap();
println!("Done");
}, 0);
scene.add_subprogram(message3_receiver_program, |mut input, context| async move {
let mut test_program = context.send(()).unwrap();
while let Some(input) = input.next().await {
match input {
Message3::Msg(msg) => test_program.send(msg).await.unwrap()
}
}
}, 0);
scene.connect_programs((), StreamTarget::Filtered(msg2_to_msg3, message3_receiver_program), StreamId::with_message_type::<Message2>()).unwrap();
TestBuilder::new()
.expect_message(|msg2: String| if msg2 != "Hello".to_string() { Err(format!("Expected 'Hello'")) } else { Ok(()) })
.expect_message(|msg2: String| if msg2 != "Goodbyte".to_string() { Err(format!("Expected 'Goodbyte'")) } else { Ok(()) })
.run_in_scene(&scene, test_program);
}