use super::*;
use crate::service_protocol::messages::{
signal_notification_message, AwaitingOnMessage, EndMessage, GetLazyStateCommandMessage,
SignalNotificationMessage, SuspensionMessage,
};
use crate::Value;
use test_log::test;
#[test]
fn trigger_suspension_with_get_state() {
let mut output = VMTestCase::new()
.input(start_message(1))
.input(input_entry_message(b"my-data"))
.run_without_closing_input(|vm, _| {
let _ = vm.sys_input().unwrap();
let handle = vm
.sys_state_get("Personaggio".to_owned(), PayloadOptions::default())
.unwrap();
assert_that!(vm.take_notification(handle), ok(none()));
vm.notify_input_closed();
assert_that!(
vm.do_await(UnresolvedFuture::Single(handle)),
err(is_suspended())
);
});
assert_eq!(
output.next_decoded::<GetLazyStateCommandMessage>().unwrap(),
GetLazyStateCommandMessage {
key: Bytes::from_static(b"Personaggio"),
result_completion_id: 1,
..Default::default()
}
);
assert_that!(
output.next_decoded::<SuspensionMessage>().unwrap(),
suspended_waiting_completion(1)
);
assert_eq!(output.next(), None);
}
#[test]
fn trigger_suspension_with_correct_awakeable() {
let mut output = VMTestCase::new()
.input(start_message(1))
.input(input_entry_message(b"my-data"))
.run_without_closing_input(|vm, _| {
vm.sys_input().unwrap();
let _h1 = vm.sys_awakeable().unwrap().handle;
let h2 = vm.sys_awakeable().unwrap().handle;
assert_that!(vm.take_notification(h2), ok(none()));
vm.notify_input_closed();
assert_that!(
vm.do_await(UnresolvedFuture::Single(h2)),
err(is_suspended())
);
});
assert_that!(
output.next_decoded::<SuspensionMessage>().unwrap(),
pat!(SuspensionMessage {
awaiting_on: some(pat!(messages::Future {
waiting_completions: empty(),
waiting_signals: unordered_elements_are![eq(1), eq(18)],
nested_futures: empty(),
waiting_named_signals: empty(),
combinator_type: eq(messages::CombinatorType::FirstCompleted as i32)
}))
})
);
assert_eq!(output.next(), None);
}
#[test]
fn await_many_notifications() {
let mut output = VMTestCase::new()
.input(start_message(1))
.input(input_entry_message(b"my-data"))
.run_without_closing_input(|vm, _| {
vm.sys_input().unwrap();
let h1 = vm.sys_awakeable().unwrap().handle;
let h2 = vm.create_signal_handle("abc".into()).unwrap();
let h3 = vm
.sys_state_get("Personaggio".to_owned(), PayloadOptions::default())
.unwrap();
vm.notify_input_closed();
assert_that!(
vm.do_await(UnresolvedFuture::FirstCompleted(vec![
UnresolvedFuture::Single(h1),
UnresolvedFuture::Single(h2),
UnresolvedFuture::Single(h3)
])),
err(is_suspended())
);
});
assert_eq!(
output.next_decoded::<GetLazyStateCommandMessage>().unwrap(),
GetLazyStateCommandMessage {
key: Bytes::from_static(b"Personaggio"),
result_completion_id: 1,
..Default::default()
}
);
assert_that!(
output.next_decoded::<SuspensionMessage>().unwrap(),
pat!(SuspensionMessage {
awaiting_on: some(pat!(messages::Future {
waiting_completions: empty(),
waiting_signals: eq(&[1]),
waiting_named_signals: empty(),
nested_futures: elements_are![pat!(messages::Future {
waiting_completions: eq(&[1]),
waiting_signals: eq(&[17]),
waiting_named_signals: eq(&["abc".to_owned()]),
nested_futures: empty(),
combinator_type: eq(messages::CombinatorType::FirstCompleted as i32)
})],
combinator_type: eq(messages::CombinatorType::FirstCompleted as i32)
}))
})
);
assert_eq!(output.next(), None);
}
#[test]
fn when_notify_completion_then_notify_await_point_then_notify_input_closed_then_no_suspension() {
let completion = Bytes::from_static(b"completion");
let mut output = VMTestCase::new()
.input(start_message(1))
.input(input_entry_message(b"my-data"))
.run_without_closing_input(|vm, encoder| {
vm.sys_input().unwrap();
let h1 = vm.sys_awakeable().unwrap().handle;
let h2 = vm.sys_awakeable().unwrap().handle;
assert_that!(
vm.do_await(UnresolvedFuture::FirstCompleted(vec![
UnresolvedFuture::Single(h1),
UnresolvedFuture::Single(h2)
])),
ok(eq(AwaitResponse::WaitingExternalProgress {
waiting_input: true,
waiting_run_proposal: false
}))
);
vm.notify_input(encoder.encode(&SignalNotificationMessage {
signal_id: Some(signal_notification_message::SignalId::Idx(18)),
result: Some(signal_notification_message::Result::Value(
completion.clone().into(),
)),
}));
vm.notify_input_closed();
assert_that!(
vm.do_await(UnresolvedFuture::FirstCompleted(vec![
UnresolvedFuture::Single(h1),
UnresolvedFuture::Single(h2)
])),
ok(eq(AwaitResponse::AnyCompleted))
);
assert!(vm.is_completed(h2));
assert_that!(
vm.take_notification(h2),
ok(some(eq(Value::Success(completion.clone()))))
);
vm.sys_write_output(
NonEmptyValue::Success(completion.clone()),
PayloadOptions::default(),
)
.unwrap();
vm.sys_end().unwrap();
});
assert_that!(
output.next_decoded::<AwaitingOnMessage>().unwrap(),
pat!(AwaitingOnMessage {
awaiting_on: some(pat!(messages::Future {
waiting_signals: eq(vec![1]),
nested_futures: elements_are![pat!(messages::Future {
waiting_signals: unordered_elements_are![eq(17), eq(18)],
combinator_type: eq(messages::CombinatorType::FirstCompleted as i32)
})],
combinator_type: eq(messages::CombinatorType::FirstCompleted as i32)
})),
executing_side_effects: eq(false)
})
);
assert_that!(
output.next_decoded::<OutputCommandMessage>().unwrap(),
is_output_with_success(completion)
);
assert_eq!(
output.next_decoded::<EndMessage>().unwrap(),
EndMessage::default()
);
assert_eq!(output.next(), None);
}