use flo_scene::*;
use flo_scene::programs::*;
use futures::prelude::*;
use serde::*;
#[test]
fn auto_start_on_first_message() {
#[derive(Serialize, Deserialize)]
struct AutoStartMessage;
impl SceneMessage for AutoStartMessage {
fn initialise(scene: &impl SceneInitialisationContext) {
scene.add_subprogram(SubProgramId::called("AutoStart"),
|mut input_stream: InputStream<AutoStartMessage>, _context| async move {
while let Some(_) = input_stream.next().await { }
}, 0);
scene.connect_programs((), SubProgramId::called("AutoStart"), StreamId::with_message_type::<AutoStartMessage>()).unwrap();
}
}
let scene = Scene::default();
TestBuilder::new()
.send_message(AutoStartMessage)
.send_message(AutoStartMessage)
.run_in_scene(&scene, SubProgramId::new());
}
#[test]
fn auto_start_on_connect() {
#[derive(Serialize, Deserialize)]
struct AutoStartMessage;
impl SceneMessage for AutoStartMessage {
fn initialise(scene: &impl SceneInitialisationContext) {
scene.add_subprogram(SubProgramId::called("AutoStart"),
|mut input_stream: InputStream<AutoStartMessage>, _context| async move {
while let Some(_) = input_stream.next().await { }
}, 0);
}
}
let scene = Scene::default();
scene.connect_programs((), SubProgramId::called("AutoStart"), StreamId::with_message_type::<AutoStartMessage>()).unwrap();
TestBuilder::new()
.send_message(AutoStartMessage)
.send_message(AutoStartMessage)
.run_in_scene(&scene, SubProgramId::new());
}
#[test]
fn connect_before_program_start() {
use futures::task::{Poll};
#[derive(Serialize, Deserialize, Debug)]
struct TestMessage;
#[derive(Serialize, Deserialize, Debug)]
struct ReadyMessage;
impl SceneMessage for TestMessage { }
impl SceneMessage for ReadyMessage { }
let test_program = SubProgramId::called("test_program");
let sending_program = SubProgramId::called("sending_program");
let receiving_program = SubProgramId::called("receiving_program");
let scene = Scene::default();
scene.connect_programs((), receiving_program, StreamId::with_message_type::<TestMessage>()).unwrap();
let scene2 = scene.clone();
scene.add_subprogram(sending_program, move |_: InputStream<()>, context| async move {
println!("Connecting to 'TestMessage'");
let mut channel = context.send(()).unwrap();
let mut send_message = channel.send(TestMessage);
println!("Waiting for send to block");
future::poll_fn(|context| {
match send_message.poll_unpin(context) {
Poll::Ready(Ok(_)) => panic!("Message finished sending unexpectedly"),
Poll::Ready(Err(err)) => panic!("Error: {:?}", err),
Poll::Pending => Poll::Ready(()),
}
}).await;
println!("Starting receiving subprogram");
scene2.add_subprogram(receiving_program, |input, context| async move {
println!("Receiving subprogram started");
let mut input = input;
while let Some(message) = input.next().await {
let _message: TestMessage = message;
println!("Received message, sending 'Ready'");
context.send(test_program).unwrap()
.send(ReadyMessage).await.unwrap();
}
}, 20);
println!("Waiting for message to finish sending");
send_message.await.unwrap();
println!("Message has finished sending");
}, 0);
println!("Running test");
TestBuilder::new()
.expect_message::<ReadyMessage>(|_| Ok(()))
.run_in_scene(&scene, test_program);
}
#[test]
fn send_before_connect_before_program_start() {
use futures::task::{Poll};
#[derive(Serialize, Deserialize, Debug)]
struct TestMessage;
#[derive(Serialize, Deserialize, Debug)]
struct ReadyMessage;
impl SceneMessage for TestMessage { }
impl SceneMessage for ReadyMessage { }
let test_program = SubProgramId::called("test_program");
let sending_program = SubProgramId::called("sending_program");
let receiving_program = SubProgramId::called("receiving_program");
let scene = Scene::default();
let scene2 = scene.clone();
scene.add_subprogram(sending_program, move |_: InputStream<()>, context| async move {
println!("Connecting to 'TestMessage'");
let mut channel = context.send(()).unwrap();
let mut send_message = channel.send(TestMessage);
println!("Waiting for send to block");
future::poll_fn(|context| {
match send_message.poll_unpin(context) {
Poll::Ready(Ok(_)) => panic!("Message finished sending unexpectedly"),
Poll::Ready(Err(err)) => panic!("Error: {:?}", err),
Poll::Pending => Poll::Ready(()),
}
}).await;
scene2.connect_programs((), receiving_program, StreamId::with_message_type::<TestMessage>()).unwrap();
println!("Starting receiving subprogram");
scene2.add_subprogram(receiving_program, |input, context| async move {
println!("Receiving subprogram started");
let mut input = input;
while let Some(message) = input.next().await {
let _message: TestMessage = message;
println!("Received message, sending 'Ready'");
context.send(test_program).unwrap()
.send(ReadyMessage).await.unwrap();
}
}, 20);
println!("Waiting for message to finish sending");
send_message.await.unwrap();
println!("Message has finished sending");
}, 0);
println!("Running test");
TestBuilder::new()
.expect_message::<ReadyMessage>(|_| Ok(()))
.run_in_scene(&scene, test_program);
}
#[test]
fn connect_before_program_start_without_waiting() {
#[derive(Serialize, Deserialize, Debug)]
struct TestMessage;
#[derive(Serialize, Deserialize, Debug)]
struct ReadyMessage;
impl SceneMessage for TestMessage { }
impl SceneMessage for ReadyMessage { }
let test_program = SubProgramId::called("test_program");
let sending_program = SubProgramId::called("sending_program");
let receiving_program = SubProgramId::called("receiving_program");
let scene = Scene::default();
scene.connect_programs((), receiving_program, StreamId::with_message_type::<TestMessage>()).unwrap();
let scene2 = scene.clone();
scene.add_subprogram(sending_program, move |_: InputStream<()>, context| async move {
let mut channel = context.send(()).unwrap();
let send_message = channel.send(TestMessage);
scene2.add_subprogram(receiving_program, |input, context| async move {
let mut input = input;
while let Some(message) = input.next().await {
let _message: TestMessage = message;
context.send(test_program).unwrap()
.send(ReadyMessage).await.unwrap();
}
}, 20);
send_message.await.unwrap();
}, 0);
TestBuilder::new()
.expect_message::<ReadyMessage>(|_| Ok(()))
.run_in_scene(&scene, test_program);
}