use flo_scene::*;
use flo_scene::programs::*;
use futures::prelude::*;
use serde::*;
#[derive(Debug, Clone, Serialize, Deserialize)]
#[allow(dead_code)]
struct TestResult(String);
impl SceneMessage for TestResult { }
#[test]
pub fn connect_two_subprograms() {
let scene = Scene::default();
let test_program = SubProgramId::new();
let program_1 = SubProgramId::new();
let program_2 = SubProgramId::new();
scene.add_subprogram(program_1, |_: InputStream<()>, context| {
async move {
let mut send_strings = context.send(()).unwrap();
send_strings.send("Test".to_string()).await.unwrap();
}
}, 0);
scene.add_subprogram(program_2, move |input: InputStream<String>, context| {
async move {
let mut test_program = context.send(test_program).unwrap();
let mut input = input;
while let Some(input) = input.next().await {
test_program.send(TestResult(input)).await.unwrap();
}
}
}, 0);
scene.connect_programs(program_1, program_2, StreamId::with_message_type::<String>()).unwrap();
TestBuilder::new()
.expect_message(|_msg: TestResult| { Ok(()) })
.run_in_scene(&scene, test_program);
}
#[test]
pub fn connect_two_subprograms_using_filter() {
let scene = Scene::default();
let test_program = SubProgramId::new();
let program_1 = SubProgramId::new();
let program_2 = SubProgramId::new();
#[derive(Debug, Serialize, Deserialize)]
struct TestMessage;
impl SceneMessage for TestMessage { }
let test_string_filter = FilterHandle::for_filter(|messages| messages.map(|_: TestMessage| "Test".to_string()));
scene.add_subprogram(program_1, |_: InputStream<()>, context| {
async move {
let mut send_strings = context.send(()).unwrap();
send_strings.send(TestMessage).await.unwrap();
}
}, 0);
scene.add_subprogram(program_2, move |input: InputStream<String>, context| {
async move {
let mut test_program = context.send(test_program).unwrap();
let mut input = input;
while let Some(input) = input.next().await {
test_program.send(TestResult(input)).await.unwrap();
}
}
}, 0);
scene.connect_programs((), program_2, StreamId::with_message_type::<String>()).unwrap();
scene.connect_programs((), StreamTarget::Filtered(test_string_filter, program_2), StreamId::with_message_type::<TestMessage>()).unwrap();
TestBuilder::new()
.expect_message(|_msg: TestResult| { Ok(()) })
.run_in_scene(&scene, test_program);
}
#[test]
pub fn connect_two_subprograms_using_filter_then_all_no_delay() {
let scene = Scene::default();
let test_program = SubProgramId::called("test_program");
let program_1 = SubProgramId::called("program_1");
let program_2 = SubProgramId::called("program_2");
#[derive(Debug, Serialize, Deserialize)]
struct TestMessage;
impl SceneMessage for TestMessage { }
let test_string_filter = FilterHandle::for_filter(|messages| messages.map(|_: TestMessage| "Test".to_string()));
scene.add_subprogram(program_1, |_: InputStream<()>, context| {
async move {
let mut send_strings = context.send(()).unwrap();
send_strings.send(TestMessage).await.unwrap();
}
}, 0);
scene.add_subprogram(program_2, move |input: InputStream<String>, context| {
async move {
let mut test_program = context.send(test_program).unwrap();
let mut input = input;
while let Some(input) = input.next().await {
test_program.send(TestResult(input)).await.unwrap();
}
}
}, 0);
scene.connect_programs((), program_2, StreamId::with_message_type::<String>()).unwrap();
scene.connect_programs((), program_2, StreamId::with_message_type::<TestMessage>()).unwrap();
scene.connect_programs((), StreamTarget::Filtered(test_string_filter, program_2), StreamId::with_message_type::<TestMessage>()).unwrap();
TestBuilder::new()
.expect_message(|_msg: TestResult| { Ok(()) })
.run_in_scene(&scene, test_program);
}
#[test]
pub fn connect_default_using_source_filter_after_all() {
let scene = Scene::default();
let test_program = SubProgramId::called("test_program");
let program_1 = SubProgramId::called("program_1");
let program_2 = SubProgramId::called("program_2");
#[derive(Debug, Serialize, Deserialize)]
struct TestMessage;
impl SceneMessage for TestMessage { }
let test_string_filter = FilterHandle::for_filter(|messages| messages.map(|_: TestMessage| "Test".to_string()));
scene.add_subprogram(program_1, |_: InputStream<()>, context| {
async move {
let mut send_strings = context.send(()).unwrap();
context.send_message(SceneControl::connect(StreamSource::Filtered(test_string_filter), (), StreamId::with_message_type::<TestMessage>())).await.unwrap();
send_strings.send(TestMessage).await.unwrap();
}
}, 0);
scene.add_subprogram(program_2, move |input: InputStream<String>, context| {
async move {
let mut test_program = context.send(test_program).unwrap();
let mut input = input;
while let Some(input) = input.next().await {
test_program.send(TestResult(input)).await.unwrap();
}
}
}, 0);
scene.connect_programs((), program_2, StreamId::with_message_type::<TestMessage>()).unwrap();
scene.connect_programs((), program_2, StreamId::with_message_type::<String>()).unwrap();
TestBuilder::new()
.expect_message(|_msg: TestResult| { Ok(()) })
.run_in_scene(&scene, test_program);
}
#[test]
pub fn connect_default_using_source_filter_after_all_no_delay() {
let scene = Scene::default();
let test_program = SubProgramId::called("test_program");
let program_1 = SubProgramId::called("program_1");
let program_2 = SubProgramId::called("program_2");
#[derive(Debug, Serialize, Deserialize)]
struct TestMessage;
impl SceneMessage for TestMessage { }
let test_string_filter = FilterHandle::for_filter(|messages| messages.map(|_: TestMessage| "Test".to_string()));
scene.add_subprogram(program_1, |_: InputStream<()>, context| {
async move {
let mut send_strings = context.send(()).unwrap();
send_strings.send(TestMessage).await.unwrap();
}
}, 0);
scene.add_subprogram(program_2, move |input: InputStream<String>, context| {
async move {
let mut test_program = context.send(test_program).unwrap();
let mut input = input;
while let Some(input) = input.next().await {
test_program.send(TestResult(input)).await.unwrap();
}
}
}, 0);
scene.connect_programs((), program_2, StreamId::with_message_type::<TestMessage>()).unwrap();
scene.connect_programs((), program_2, StreamId::with_message_type::<String>()).unwrap();
scene.connect_programs(StreamSource::Filtered(test_string_filter), (), StreamId::with_message_type::<TestMessage>()).unwrap();
TestBuilder::new()
.expect_message(|_msg: TestResult| { Ok(()) })
.run_in_scene(&scene, test_program);
}
#[test]
pub fn connect_default_using_chained_filter() {
let scene = Scene::default();
let test_program = SubProgramId::called("test_program");
let program_1 = SubProgramId::called("program_1");
let program_2 = SubProgramId::called("program_2");
#[derive(Debug, Serialize, Deserialize)]
struct TestMessage1;
impl SceneMessage for TestMessage1 { }
#[derive(Debug, Serialize, Deserialize)]
struct TestMessage2;
impl SceneMessage for TestMessage2 { }
let test_message_filter = FilterHandle::for_filter(|messages| messages.map(|_: TestMessage1| TestMessage2));
let test_string_filter = FilterHandle::for_filter(|messages| messages.map(|_: TestMessage2| "Test".to_string()));
scene.add_subprogram(program_1, |_: InputStream<()>, context| {
async move {
let mut send_strings = context.send(()).unwrap();
send_strings.send(TestMessage1).await.unwrap();
}
}, 0);
scene.add_subprogram(program_2, move |input: InputStream<String>, context| {
async move {
let mut test_program = context.send(test_program).unwrap();
let mut input = input;
while let Some(input) = input.next().await {
test_program.send(TestResult(input)).await.unwrap();
}
}
}, 0);
scene.connect_programs((), StreamTarget::Filtered(test_string_filter, program_2), StreamId::with_message_type::<TestMessage2>()).unwrap();
scene.connect_programs((), program_2, StreamId::with_message_type::<String>()).unwrap();
scene.connect_programs(StreamSource::Filtered(test_message_filter), (), StreamId::with_message_type::<TestMessage1>()).unwrap();
TestBuilder::new()
.expect_message(|_msg: TestResult| { Ok(()) })
.run_in_scene(&scene, test_program);
}
#[test]
pub fn connect_default_using_chained_filter_later_1() {
let scene = Scene::default();
let test_program = SubProgramId::called("test_program");
let program_1 = SubProgramId::called("program_1");
let program_2 = SubProgramId::called("program_2");
#[derive(Debug, Serialize, Deserialize)]
struct TestMessage1;
impl SceneMessage for TestMessage1 { }
#[derive(Debug, Serialize, Deserialize)]
struct TestMessage2;
impl SceneMessage for TestMessage2 { }
let test_message_filter = FilterHandle::for_filter(|messages| messages.map(|_: TestMessage1| TestMessage2));
let test_string_filter = FilterHandle::for_filter(|messages| messages.map(|_: TestMessage2| "Test".to_string()));
scene.add_subprogram(program_1, |_: InputStream<()>, context| {
async move {
let mut send_strings = context.send(()).unwrap();
context.send_message(SceneControl::connect((), StreamTarget::Filtered(test_string_filter, program_2), StreamId::with_message_type::<TestMessage2>())).await.unwrap();
send_strings.send(TestMessage1).await.unwrap();
}
}, 0);
scene.add_subprogram(program_2, move |input: InputStream<String>, context| {
async move {
let mut test_program = context.send(test_program).unwrap();
let mut input = input;
while let Some(input) = input.next().await {
test_program.send(TestResult(input)).await.unwrap();
}
}
}, 0);
scene.connect_programs((), program_2, StreamId::with_message_type::<String>()).unwrap();
scene.connect_programs(StreamSource::Filtered(test_message_filter), (), StreamId::with_message_type::<TestMessage1>()).unwrap();
TestBuilder::new()
.expect_message(|_msg: TestResult| { Ok(()) })
.run_in_scene(&scene, test_program);
}
#[test]
pub fn connect_default_using_chained_filter_later_2() {
let scene = Scene::default();
let test_program = SubProgramId::called("test_program");
let program_1 = SubProgramId::called("program_1");
let program_2 = SubProgramId::called("program_2");
#[derive(Debug, Serialize, Deserialize)]
struct TestMessage1;
impl SceneMessage for TestMessage1 { }
#[derive(Debug, Serialize, Deserialize)]
struct TestMessage2;
impl SceneMessage for TestMessage2 { }
let test_message_filter = FilterHandle::for_filter(|messages| messages.map(|_: TestMessage1| TestMessage2));
let test_string_filter = FilterHandle::for_filter(|messages| messages.map(|_: TestMessage2| "Test".to_string()));
scene.add_subprogram(program_1, |_: InputStream<()>, context| {
async move {
let mut send_strings = context.send(()).unwrap();
context.send_message(SceneControl::connect(StreamSource::Filtered(test_message_filter), (), StreamId::with_message_type::<TestMessage1>())).await.unwrap();
send_strings.send(TestMessage1).await.unwrap();
}
}, 0);
scene.add_subprogram(program_2, move |input: InputStream<String>, context| {
async move {
let mut test_program = context.send(test_program).unwrap();
let mut input = input;
while let Some(input) = input.next().await {
test_program.send(TestResult(input)).await.unwrap();
}
}
}, 0);
scene.connect_programs((), StreamTarget::Filtered(test_string_filter, program_2), StreamId::with_message_type::<TestMessage2>()).unwrap();
scene.connect_programs((), program_2, StreamId::with_message_type::<String>()).unwrap();
TestBuilder::new()
.expect_message(|_msg: TestResult| { Ok(()) })
.run_in_scene(&scene, test_program);
}
#[test]
pub fn connect_default_using_chained_filter_later_3() {
let scene = Scene::default();
let test_program = SubProgramId::called("test_program");
let program_1 = SubProgramId::called("program_1");
let program_2 = SubProgramId::called("program_2");
#[derive(Debug, Serialize, Deserialize)]
struct TestMessage1;
impl SceneMessage for TestMessage1 { }
#[derive(Debug, Serialize, Deserialize)]
struct TestMessage2;
impl SceneMessage for TestMessage2 { }
let test_message_filter = FilterHandle::for_filter(|messages| messages.map(|_: TestMessage1| TestMessage2));
let test_string_filter = FilterHandle::for_filter(|messages| messages.map(|_: TestMessage2| "Test".to_string()));
scene.add_subprogram(program_1, |_: InputStream<()>, context| {
async move {
let mut send_strings = context.send(()).unwrap();
context.send_message(SceneControl::connect(StreamSource::Filtered(test_message_filter), (), StreamId::with_message_type::<TestMessage1>())).await.unwrap();
context.send_message(SceneControl::connect((), StreamTarget::Filtered(test_string_filter, program_2), StreamId::with_message_type::<TestMessage2>())).await.unwrap();
send_strings.send(TestMessage1).await.unwrap();
}
}, 0);
scene.add_subprogram(program_2, move |input: InputStream<String>, context| {
async move {
let mut test_program = context.send(test_program).unwrap();
let mut input = input;
while let Some(input) = input.next().await {
test_program.send(TestResult(input)).await.unwrap();
}
}
}, 0);
scene.connect_programs((), program_2, StreamId::with_message_type::<String>()).unwrap();
TestBuilder::new()
.expect_message(|_msg: TestResult| { Ok(()) })
.run_in_scene(&scene, test_program);
}
#[test]
pub fn connect_default_using_filter_added_later() {
let scene = Scene::default();
let test_program = SubProgramId::called("test_program");
let program_1 = SubProgramId::called("program_1");
let program_2 = SubProgramId::called("program_2");
#[derive(Debug, Serialize, Deserialize)]
struct TestMessage2;
impl SceneMessage for TestMessage2 { }
let test_string_filter = FilterHandle::for_filter(|messages| messages.map(|_: TestMessage2| "Test".to_string()));
scene.add_subprogram(program_1, |_: InputStream<()>, context| {
async move {
let mut send_strings = context.send(()).unwrap();
context.send_message(SceneControl::connect((), StreamTarget::Filtered(test_string_filter, program_2), StreamId::with_message_type::<TestMessage2>())).await.unwrap();
send_strings.send(TestMessage2).await.unwrap();
}
}, 0);
scene.add_subprogram(program_2, move |input: InputStream<String>, context| {
async move {
let mut test_program = context.send(test_program).unwrap();
let mut input = input;
while let Some(input) = input.next().await {
test_program.send(TestResult(input)).await.unwrap();
}
}
}, 0);
scene.connect_programs((), program_2, StreamId::with_message_type::<String>()).unwrap();
TestBuilder::new()
.expect_message(|_msg: TestResult| { Ok(()) })
.run_in_scene(&scene, test_program);
}
#[test]
pub fn connect_default_using_source_filter_added_later_1() {
let scene = Scene::default();
let test_program = SubProgramId::called("test_program");
let program_1 = SubProgramId::called("program_1");
let program_2 = SubProgramId::called("program_2");
#[derive(Debug, Serialize, Deserialize)]
struct TestMessage2;
impl SceneMessage for TestMessage2 { }
let test_string_filter = FilterHandle::for_filter(|messages| messages.map(|_: TestMessage2| "Test".to_string()));
scene.add_subprogram(program_1, |_: InputStream<()>, context| {
async move {
let mut send_strings = context.send(()).unwrap();
context.send_message(SceneControl::connect(StreamSource::Filtered(test_string_filter), (), StreamId::with_message_type::<TestMessage2>())).await.unwrap();
send_strings.send(TestMessage2).await.unwrap();
}
}, 0);
scene.add_subprogram(program_2, move |input: InputStream<String>, context| {
async move {
let mut test_program = context.send(test_program).unwrap();
let mut input = input;
while let Some(input) = input.next().await {
test_program.send(TestResult(input)).await.unwrap();
}
}
}, 0);
scene.connect_programs((), program_2, StreamId::with_message_type::<String>()).unwrap();
TestBuilder::new()
.expect_message(|_msg: TestResult| { Ok(()) })
.run_in_scene(&scene, test_program);
}
#[test]
pub fn connect_default_using_source_filter_added_later_2() {
let scene = Scene::default();
let test_program = SubProgramId::called("test_program");
let program_1 = SubProgramId::called("program_1");
let program_2 = SubProgramId::called("program_2");
#[derive(Debug, Serialize, Deserialize)]
struct TestMessage2;
impl SceneMessage for TestMessage2 { }
let test_string_filter = FilterHandle::for_filter(|messages| messages.map(|_: TestMessage2| "Test".to_string()));
scene.add_subprogram(program_1, |_: InputStream<()>, context| {
async move {
let mut send_strings = context.send(()).unwrap();
context.send_message(SceneControl::connect((), program_2, StreamId::with_message_type::<String>())).await.unwrap();
send_strings.send(TestMessage2).await.unwrap();
}
}, 0);
scene.add_subprogram(program_2, move |input: InputStream<String>, context| {
async move {
let mut test_program = context.send(test_program).unwrap();
let mut input = input;
while let Some(input) = input.next().await {
test_program.send(TestResult(input)).await.unwrap();
}
}
}, 0);
scene.connect_programs(StreamSource::Filtered(test_string_filter), (), StreamId::with_message_type::<TestMessage2>()).unwrap();
TestBuilder::new()
.expect_message(|_msg: TestResult| { Ok(()) })
.run_in_scene(&scene, test_program);
}
#[test]
pub fn connect_two_subprograms_using_source_filter() {
let scene = Scene::default();
let test_program = SubProgramId::new();
let program_1 = SubProgramId::new();
let program_2 = SubProgramId::new();
#[derive(Debug, Serialize, Deserialize)]
struct TestMessage;
impl SceneMessage for TestMessage { }
let test_string_filter = FilterHandle::for_filter(|messages| messages.map(|_: TestMessage| "Test".to_string()));
scene.add_subprogram(program_1, |_: InputStream<()>, context| {
async move {
let mut send_strings = context.send(()).unwrap();
send_strings.send(TestMessage).await.unwrap();
}
}, 0);
scene.add_subprogram(program_2, move |input: InputStream<String>, context| {
async move {
let mut test_program = context.send(test_program).unwrap();
let mut input = input;
while let Some(input) = input.next().await {
test_program.send(TestResult(input)).await.unwrap();
}
}
}, 0);
scene.connect_programs(StreamSource::Filtered(test_string_filter), (), StreamId::with_message_type::<TestMessage>()).unwrap();
scene.connect_programs(program_1, program_2, StreamId::with_message_type::<TestMessage>()).unwrap();
TestBuilder::new()
.expect_message(|_msg: TestResult| { Ok(()) })
.run_in_scene(&scene, test_program);
}
#[test]
pub fn connect_two_subprograms_using_source_filter_later() {
let scene = Scene::default();
let test_program = SubProgramId::new();
let program_1 = SubProgramId::new();
let program_2 = SubProgramId::new();
#[derive(Debug, Serialize, Deserialize)]
struct TestMessage;
impl SceneMessage for TestMessage { }
let test_string_filter = FilterHandle::for_filter(|messages| messages.map(|_: TestMessage| "Test".to_string()));
scene.add_subprogram(program_1, |_: InputStream<()>, context| {
async move {
let mut send_strings = context.send(()).unwrap();
send_strings.send(TestMessage).await.unwrap();
}
}, 0);
scene.add_subprogram(program_2, move |input: InputStream<String>, context| {
async move {
let mut test_program = context.send(test_program).unwrap();
let mut input = input;
while let Some(input) = input.next().await {
test_program.send(TestResult(input)).await.unwrap();
}
}
}, 0);
scene.connect_programs(program_1, program_2, StreamId::with_message_type::<TestMessage>()).unwrap();
scene.connect_programs(StreamSource::Filtered(test_string_filter), (), StreamId::with_message_type::<TestMessage>()).unwrap();
TestBuilder::new()
.expect_message(|_msg: TestResult| { Ok(()) })
.run_in_scene(&scene, test_program);
}
#[test]
pub fn connect_two_subprograms_using_string_type_then_source_filter() {
let scene = Scene::default();
let test_program = SubProgramId::new();
let program_1 = SubProgramId::new();
let program_2 = SubProgramId::new();
#[derive(Debug, Serialize, Deserialize)]
struct TestMessage;
impl SceneMessage for TestMessage { }
let test_string_filter = FilterHandle::for_filter(|messages| messages.map(|_: TestMessage| "Test".to_string()));
scene.add_subprogram(program_1, |_: InputStream<()>, context| {
async move {
let mut send_strings = context.send(()).unwrap();
send_strings.send("Test".to_string()).await.unwrap();
}
}, 0);
scene.add_subprogram(program_2, move |input: InputStream<String>, context| {
async move {
let mut test_program = context.send(test_program).unwrap();
let mut input = input;
while let Some(input) = input.next().await {
test_program.send(TestResult(input)).await.unwrap();
}
}
}, 0);
scene.connect_programs(program_1, program_2, StreamId::with_message_type::<String>()).unwrap();
scene.connect_programs(StreamSource::Filtered(test_string_filter), (), StreamId::with_message_type::<TestMessage>()).unwrap();
TestBuilder::new()
.expect_message(|_msg: TestResult| { Ok(()) })
.run_in_scene(&scene, test_program);
}
#[test]
pub fn connect_two_subprograms_after_creating_stream() {
let scene = Scene::default();
let test_program = SubProgramId::new();
let program_1 = SubProgramId::new();
let program_2 = SubProgramId::new();
scene.add_subprogram(program_1, |_: InputStream<()>, context| {
async move {
let mut send_strings = context.send(()).unwrap();
context.send_message(SceneControl::connect(program_1, program_2, StreamId::with_message_type::<String>())).await.unwrap();
send_strings.send("Test".to_string()).await.unwrap();
}
}, 0);
scene.add_subprogram(program_2, move |input: InputStream<String>, context| {
async move {
let mut test_program = context.send(test_program).unwrap();
let mut input = input;
while let Some(input) = input.next().await {
test_program.send(TestResult(input)).await.unwrap();
}
}
}, 0);
TestBuilder::new()
.expect_message(|_msg: TestResult| { Ok(()) })
.run_in_scene(&scene, test_program);
}
#[test]
pub fn connect_default_subprogram_after_creating_stream() {
let scene = Scene::default();
let test_program = SubProgramId::new();
let program_1 = SubProgramId::new();
let program_2 = SubProgramId::new();
scene.add_subprogram(program_1, |_: InputStream<()>, context| {
async move {
let mut send_strings = context.send(()).unwrap();
context.send_message(SceneControl::connect((), program_2, StreamId::with_message_type::<String>())).await.unwrap();
send_strings.send("Test".to_string()).await.unwrap();
}
}, 0);
scene.add_subprogram(program_2, move |input: InputStream<String>, context| {
async move {
let mut test_program = context.send(test_program).unwrap();
let mut input = input;
while let Some(input) = input.next().await {
test_program.send(TestResult(input)).await.unwrap();
}
}
}, 0);
TestBuilder::new()
.expect_message(|_msg: TestResult| { Ok(()) })
.run_in_scene(&scene, test_program);
}
#[test]
pub fn connect_two_subprograms_before_creating() {
let scene = Scene::default();
let test_program = SubProgramId::new();
let program_1 = SubProgramId::new();
let program_2 = SubProgramId::new();
scene.connect_programs(program_1, program_2, StreamId::with_message_type::<String>()).unwrap();
scene.add_subprogram(program_1, |_: InputStream<()>, context| {
async move {
context.send_message("Test".to_string()).await.unwrap();
}
}, 0);
scene.add_subprogram(program_2, move |input: InputStream<String>, context| {
async move {
let mut test_program = context.send(test_program).unwrap();
let mut input = input;
while let Some(input) = input.next().await {
test_program.send(TestResult(input)).await.unwrap();
}
}
}, 0);
TestBuilder::new()
.expect_message(|_msg: TestResult| { Ok(()) })
.run_in_scene(&scene, test_program);
}
#[test]
pub fn connect_default_subprograms_before_creating() {
let scene = Scene::default();
let test_program = SubProgramId::new();
let program_1 = SubProgramId::new();
let program_2 = SubProgramId::new();
scene.connect_programs((), program_2, StreamId::with_message_type::<String>()).unwrap();
scene.add_subprogram(program_1, |_: InputStream<()>, context| {
async move {
context.send_message("Test".to_string()).await.unwrap();
}
}, 0);
scene.add_subprogram(program_2, move |input: InputStream<String>, context| {
async move {
let mut test_program = context.send(test_program).unwrap();
let mut input = input;
while let Some(input) = input.next().await {
test_program.send(TestResult(input)).await.unwrap();
}
}
}, 0);
TestBuilder::new()
.expect_message(|_msg: TestResult| { Ok(()) })
.run_in_scene(&scene, test_program);
}
#[test]
pub fn connect_default_subprograms_before_launching() {
let scene = Scene::default();
let test_program = SubProgramId::new();
let program_1 = SubProgramId::new();
let program_2 = SubProgramId::new();
scene.connect_programs((), program_2, StreamId::with_message_type::<String>()).unwrap();
scene.add_subprogram(program_1, |_: InputStream<()>, context| {
async move {
let mut send_strings = context.send(()).unwrap();
context.wait_for_idle(100).await;
context.send_message(SceneControl::start_program(program_2, move |input: InputStream<String>, context| {
async move {
let mut test_program = context.send(test_program).unwrap();
let mut input = input;
while let Some(input) = input.next().await {
test_program.send(TestResult(input)).await.unwrap();
}
}
}, 0)).await.unwrap();
send_strings.send("Test".to_string()).await.unwrap();
}
}, 0);
TestBuilder::new()
.expect_message(|_msg: TestResult| { Ok(()) })
.run_in_scene(&scene, test_program);
}
#[test]
pub fn connect_two_subprograms_after_creating_stream_using_filter_target() {
let scene = Scene::default();
let test_program = SubProgramId::new();
let program_1 = SubProgramId::new();
let program_2 = SubProgramId::new();
#[derive(Debug, Serialize, Deserialize)]
struct TestMessage;
impl SceneMessage for TestMessage { }
let test_string_filter = FilterHandle::for_filter(|messages| messages.map(|_: TestMessage| "Test".to_string()));
scene.add_subprogram(program_1, |_: InputStream<()>, context| {
async move {
let mut send_strings = context.send(program_2).unwrap();
context.send_message(SceneControl::connect(program_1, StreamTarget::Filtered(test_string_filter, program_2), StreamId::with_message_type::<TestMessage>())).await.unwrap();
send_strings.send(TestMessage).await.unwrap();
}
}, 0);
scene.add_subprogram(program_2, move |input: InputStream<String>, context| {
async move {
let mut test_program = context.send(test_program).unwrap();
let mut input = input;
while let Some(input) = input.next().await {
test_program.send(TestResult(input)).await.unwrap();
}
}
}, 0);
TestBuilder::new()
.expect_message(|_msg: TestResult| { Ok(()) })
.run_in_scene(&scene, test_program);
}
#[test]
pub fn connect_default_after_creating_stream_using_filter_target() {
let scene = Scene::default();
let test_program = SubProgramId::new();
let program_1 = SubProgramId::new();
let program_2 = SubProgramId::new();
#[derive(Debug, Serialize, Deserialize)]
struct TestMessage;
impl SceneMessage for TestMessage { }
let test_string_filter = FilterHandle::for_filter(|messages| messages.map(|_: TestMessage| "Test".to_string()));
scene.add_subprogram(program_1, |_: InputStream<()>, context| {
async move {
let mut send_strings = context.send(()).unwrap();
context.send_message(SceneControl::connect((), StreamTarget::Filtered(test_string_filter, program_2), StreamId::with_message_type::<TestMessage>())).await.unwrap();
send_strings.send(TestMessage).await.unwrap();
}
}, 0);
scene.add_subprogram(program_2, move |input: InputStream<String>, context| {
async move {
let mut test_program = context.send(test_program).unwrap();
let mut input = input;
while let Some(input) = input.next().await {
test_program.send(TestResult(input)).await.unwrap();
}
}
}, 0);
TestBuilder::new()
.expect_message(|_msg: TestResult| { Ok(()) })
.run_in_scene(&scene, test_program);
}
#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
pub struct SimpleTestMessage {
value: String,
}
#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
pub struct SimpleResponseMessage {
value: String,
}
impl SceneMessage for SimpleTestMessage {
fn message_type_name() -> String {
"flo_scene_tests::guest_subprogram_tests::SimpleTestMessage".into()
}
}
impl SceneMessage for SimpleResponseMessage {
fn message_type_name() -> String {
"flo_scene_tests::guest_subprogram_tests::SimpleResponseMessage".into()
}
}
#[test]
fn connect_single_source_to_single_target() {
let scene = Scene::default();
let guest_subprogram_id = SubProgramId::called("Guest subprogram");
let sender_subprogram_id = SubProgramId::called("Sender subprogram");
let test_subprogram_id = SubProgramId::called("Test subprogram");
scene.add_subprogram(guest_subprogram_id, move |input_stream: InputStream<SimpleTestMessage>, context| async move {
let mut response = context.send::<SimpleResponseMessage>(()).unwrap();
let mut input_stream = input_stream;
while let Some(msg) = input_stream.next().await {
println!("Received message: {:?}", msg);
response.send(SimpleResponseMessage { value: msg.value }).await.unwrap();
println!("Sent message");
}
}, 10);
scene.add_subprogram(sender_subprogram_id, move |_input: InputStream<()>, context| async move {
let mut test_messages = context.send(guest_subprogram_id).unwrap();
test_messages.send(SimpleTestMessage { value: "Hello".into() }).await.unwrap();
test_messages.send(SimpleTestMessage { value: "Goodbyte".into() }).await.unwrap();
}, 0);
scene.connect_programs(guest_subprogram_id, test_subprogram_id, StreamId::with_message_type::<SimpleResponseMessage>()).unwrap();
TestBuilder::new()
.expect_message(|_: SimpleResponseMessage| { Ok(()) })
.expect_message(|_: SimpleResponseMessage| { Ok(()) })
.run_in_scene(&scene, test_subprogram_id);
}
#[test]
fn connect_single_source_to_single_target_before_creation() {
let scene = Scene::default();
let guest_subprogram_id = SubProgramId::called("Guest subprogram");
let sender_subprogram_id = SubProgramId::called("Sender subprogram");
let test_subprogram_id = SubProgramId::called("Test subprogram");
scene.connect_programs(guest_subprogram_id, test_subprogram_id, StreamId::with_message_type::<SimpleResponseMessage>()).unwrap();
scene.add_subprogram(guest_subprogram_id, move |input_stream: InputStream<SimpleTestMessage>, context| async move {
let mut response = context.send::<SimpleResponseMessage>(()).unwrap();
let mut input_stream = input_stream;
while let Some(msg) = input_stream.next().await {
println!("Received message: {:?}", msg);
response.send(SimpleResponseMessage { value: msg.value }).await.unwrap();
println!("Sent message");
}
}, 10);
scene.add_subprogram(sender_subprogram_id, move |_input: InputStream<()>, context| async move {
let mut test_messages = context.send(guest_subprogram_id).unwrap();
test_messages.send(SimpleTestMessage { value: "Hello".into() }).await.unwrap();
test_messages.send(SimpleTestMessage { value: "Goodbyte".into() }).await.unwrap();
}, 0);
TestBuilder::new()
.expect_message(|_: SimpleResponseMessage| { Ok(()) })
.expect_message(|_: SimpleResponseMessage| { Ok(()) })
.run_in_scene(&scene, test_subprogram_id);
}