use std::cell::RefCell;
use std::future::Future;
use std::pin::Pin;
use std::rc::Rc;
use std::sync::Arc;
use std::task::{Context, Poll, RawWaker, RawWakerVTable, Waker};
use std::time::{Duration, Instant};
use super::super::{Actor, ActorContext, ActorError};
use super::{CallFuture, CoopSenderHandle, spawn_actor_cooperative};
use crate::atom::AtomTable;
use crate::module::ModuleRegistry;
use crate::native::BifRegistryImpl;
use crate::scheduler::WasmScheduler;
struct Adder;
impl Actor for Adder {
type Call = i64;
type Reply = i64;
type Cast = i64;
fn handle_call(&mut self, request: i64, _ctx: &mut ActorContext<'_, '_>) -> i64 {
request + 1
}
fn handle_cast(&mut self, _request: i64, _ctx: &mut ActorContext<'_, '_>) {}
}
fn scheduler() -> Rc<RefCell<WasmScheduler>> {
let atom_table = Arc::new(AtomTable::with_common_atoms());
let modules = Arc::new(ModuleRegistry::new());
let bifs = Arc::new(BifRegistryImpl::new());
Rc::new(RefCell::new(WasmScheduler::new(atom_table, modules, bifs)))
}
fn noop_waker() -> Waker {
fn clone(_: *const ()) -> RawWaker {
RawWaker::new(std::ptr::null(), &VTABLE)
}
fn noop(_: *const ()) {}
static VTABLE: RawWakerVTable = RawWakerVTable::new(clone, noop, noop, noop);
unsafe { Waker::from_raw(RawWaker::new(std::ptr::null(), &VTABLE)) }
}
fn poll_once<R>(future: &mut Pin<Box<CallFuture<R>>>) -> Poll<Result<R, ActorError>> {
let waker = noop_waker();
let mut cx = Context::from_waker(&waker);
future.as_mut().poll(&mut cx)
}
#[test]
fn call_async_resolves_with_the_actors_reply() {
let scheduler = scheduler();
let actor = spawn_actor_cooperative::<Adder, _>(&scheduler, || Adder);
let mut future = Box::pin(actor.sender.call_async(41));
assert!(
matches!(poll_once(&mut future), Poll::Pending),
"future is pending before any turn runs"
);
let mut resolved = None;
for _ in 0..8 {
scheduler.borrow_mut().run_until_idle();
if let Poll::Ready(result) = poll_once(&mut future) {
resolved = Some(result);
break;
}
}
assert_eq!(
resolved,
Some(Ok(42)),
"the call_async future resolved with the actor's reply (41 + 1)"
);
}
#[test]
fn concurrent_call_asyncs_never_cross_replies() {
let scheduler = scheduler();
let actor = spawn_actor_cooperative::<Adder, _>(&scheduler, || Adder);
let mut first = Box::pin(actor.sender.call_async(10));
let mut second = Box::pin(actor.sender.call_async(20));
let mut first_result = None;
let mut second_result = None;
for _ in 0..16 {
scheduler.borrow_mut().run_until_idle();
if first_result.is_none()
&& let Poll::Ready(result) = poll_once(&mut first)
{
first_result = Some(result);
}
if second_result.is_none()
&& let Poll::Ready(result) = poll_once(&mut second)
{
second_result = Some(result);
}
if first_result.is_some() && second_result.is_some() {
break;
}
}
assert_eq!(first_result, Some(Ok(11)), "first call got its own reply");
assert_eq!(second_result, Some(Ok(21)), "second call got its own reply");
}
#[test]
fn call_async_rejects_on_timeout_when_no_reply_arrives() {
let scheduler = scheduler();
let target = scheduler
.borrow_mut()
.spawn_native_root(Box::new(|| Box::new(ParkForever)));
let handle = CoopSenderHandle::<Adder>::attach(&scheduler, target);
let delay = Duration::from_secs(10);
let mut future = Box::pin(handle.call_async_timeout(7, delay));
scheduler.borrow_mut().run_until_idle();
assert!(
matches!(poll_once(&mut future), Poll::Pending),
"future is pending while the request is outstanding and the timeout is armed"
);
let start = Instant::now();
let woken_early = scheduler
.borrow_mut()
.tick_native_timers_at(start + Duration::from_secs(5));
let _ = woken_early;
scheduler.borrow_mut().run_until_idle();
assert!(
matches!(poll_once(&mut future), Poll::Pending),
"future stays pending before the timeout deadline"
);
let _woken = scheduler
.borrow_mut()
.tick_native_timers_at(start + delay + Duration::from_secs(5));
let mut rejected = None;
for _ in 0..8 {
scheduler.borrow_mut().run_until_idle();
if let Poll::Ready(result) = poll_once(&mut future) {
rejected = Some(result);
break;
}
}
assert_eq!(
rejected,
Some(Err(ActorError::Timeout)),
"the call_async future rejected with Timeout when no reply arrived"
);
}
struct ParkForever;
impl crate::native::native_process::NativeHandler for ParkForever {
fn handle(
&mut self,
_ctx: &mut crate::native::native_process::NativeContext<'_>,
) -> crate::native::native_process::NativeOutcome {
crate::native::native_process::NativeOutcome::Wait
}
}