use std::sync::{Arc, Mutex};
use super::*;
use crate::atom::{Atom, AtomTable};
use crate::module::ModuleRegistry;
use crate::native::BifRegistryImpl;
use crate::native::native_process::{NativeContext, NativeHandler, NativeOutcome};
use crate::process::ExitReason;
use crate::term::Term;
struct Echo {
reply_to: u64,
}
impl NativeHandler for Echo {
fn handle(&mut self, ctx: &mut NativeContext<'_>) -> NativeOutcome {
match ctx.recv() {
Some(message) => {
ctx.send(self.reply_to, message);
NativeOutcome::Stop(ExitReason::Normal)
}
None => NativeOutcome::Wait,
}
}
}
struct Collector {
sink: Arc<Mutex<Option<i64>>>,
}
impl NativeHandler for Collector {
fn handle(&mut self, ctx: &mut NativeContext<'_>) -> NativeOutcome {
while let Some(message) = ctx.recv() {
if let Some(value) = message.as_small_int()
&& let Ok(mut guard) = self.sink.lock()
{
*guard = Some(value);
}
}
NativeOutcome::Wait
}
}
fn scheduler() -> WasmScheduler {
let atom_table = Arc::new(AtomTable::with_common_atoms());
let modules = Arc::new(ModuleRegistry::new());
let bifs = Arc::new(BifRegistryImpl::new());
WasmScheduler::new(atom_table, modules, bifs)
}
fn drain_until_exit(scheduler: &mut WasmScheduler, pid: u64, max_turns: usize) -> bool {
for _ in 0..max_turns {
let exited = scheduler.run_native_until_idle();
if exited.contains(&pid) {
return true;
}
}
false
}
#[test]
fn native_actor_spawns_receives_one_message_and_replies_with_captured_result() {
let mut scheduler = scheduler();
let sink = Arc::new(Mutex::new(None));
let collector = scheduler.spawn_native_root({
let sink = Arc::clone(&sink);
Box::new(move || {
Box::new(Collector {
sink: Arc::clone(&sink),
})
})
});
let echo = scheduler.spawn_native_root(Box::new(move || {
Box::new(Echo {
reply_to: collector,
})
}));
let exited = scheduler.run_native_until_idle();
assert!(exited.is_empty(), "nothing exits before a message arrives");
assert_eq!(
*sink.lock().expect("sink lock"),
None,
"collector has received nothing yet"
);
scheduler
.send_owned(echo, &crate::ets::OwnedTerm::immediate(Term::small_int(42)))
.expect("message delivers to the parked echo actor");
assert!(
drain_until_exit(&mut scheduler, echo, 4),
"the echo actor exits after handling its one message"
);
assert_eq!(
scheduler.native_exit_reason(echo),
Some(ExitReason::Normal),
"the echo actor stopped normally"
);
for _ in 0..4 {
let _exited = scheduler.run_native_until_idle();
if sink.lock().expect("sink lock").is_some() {
break;
}
}
assert_eq!(
*sink.lock().expect("sink lock"),
Some(42),
"the result the native actor produced is observable end-to-end"
);
}
struct Parent {
reply_to: u64,
}
impl NativeHandler for Parent {
fn handle(&mut self, ctx: &mut NativeContext<'_>) -> NativeOutcome {
let Some(_trigger) = ctx.recv() else {
return NativeOutcome::Wait;
};
let reply_to = self.reply_to;
let child = ctx
.spawn_native(Box::new(move || Box::new(Echo { reply_to })), None)
.expect("cooperative spawn_native succeeds");
ctx.send(child, Term::small_int(7));
NativeOutcome::Stop(ExitReason::Normal)
}
}
#[test]
fn handler_spawns_child_and_sends_it_a_message_cooperatively() {
let mut scheduler = scheduler();
let sink = Arc::new(Mutex::new(None));
let collector = scheduler.spawn_native_root({
let sink = Arc::clone(&sink);
Box::new(move || {
Box::new(Collector {
sink: Arc::clone(&sink),
})
})
});
let parent = scheduler.spawn_native_root(Box::new(move || {
Box::new(Parent {
reply_to: collector,
})
}));
let _first = scheduler.run_native_until_idle();
scheduler
.send_owned(
parent,
&crate::ets::OwnedTerm::immediate(Term::atom(Atom::OK)),
)
.expect("trigger delivers to the parent");
assert!(
drain_until_exit(&mut scheduler, parent, 8),
"the parent exits after spawning and sending"
);
assert_eq!(
scheduler.native_exit_reason(parent),
Some(ExitReason::Normal),
"parent stopped after spawning and sending"
);
for _ in 0..8 {
let _exited = scheduler.run_native_until_idle();
if sink.lock().expect("sink lock").is_some() {
break;
}
}
assert_eq!(
*sink.lock().expect("sink lock"),
Some(7),
"the cooperatively-spawned child received and forwarded the message"
);
}