use std::sync::Arc;
use crate::ets::copy::OwnedTerm;
use crate::native::native_process::{NativeContext, NativeOutcome};
use crate::process::heap::DEFAULT_HEAP_SIZE;
use crate::process::{ExitReason, Process, ProcessStatus};
use crate::scheduler::{SharedState, supervision_integration};
use crate::term::Term;
use super::core::SliceOutcome;
pub(in crate::scheduler) fn run_native_slice(
shared: &Arc<SharedState>,
process: &mut Process,
) -> SliceOutcome {
let pid = process.pid();
if let Some(reason) = shared.exit_tombstones.get(&pid) {
return SliceOutcome::Exited(reason, OwnedTerm::immediate(Term::NIL));
}
if transition_to_running(process).is_err() {
return SliceOutcome::Exited(ExitReason::Error, OwnedTerm::immediate(Term::NIL));
}
let mut handler = match process.native_body_mut() {
Some(body) => match body.handler.take() {
Some(handler) => handler,
None => (body.factory)(),
},
None => {
return SliceOutcome::Exited(ExitReason::Normal, OwnedTerm::immediate(Term::NIL));
}
};
let services = supervision_integration::build_native_services(shared, process.namespace_id());
let teardown_admission_facility = services.teardown_admission_facility.clone();
#[cfg(feature = "readiness")]
let readiness_facility = services.readiness_facility.clone();
let (Some(local_send), Some(spawn)) = (services.local_send, services.spawn_facility) else {
if let Some(body) = process.native_body_mut() {
body.handler = Some(handler);
}
return SliceOutcome::Exited(ExitReason::Error, OwnedTerm::immediate(Term::NIL));
};
let replay_driver = shared.replay_driver.clone();
let timers = Some(shared.timers.clone());
let mut context = NativeContext::new(process, local_send, spawn, replay_driver, timers);
context.set_teardown_admission_facility(teardown_admission_facility);
#[cfg(feature = "readiness")]
context.set_readiness_facility(readiness_facility);
let outcome = handler.handle(&mut context);
let replay_error = context.take_replay_error();
drop(context);
if let Some(body) = process.native_body_mut() {
body.handler = Some(handler);
}
if let Some(error) = replay_error {
shared.exit_errors.insert(pid, error);
let result = crate::scheduler::exit_capture::capture_term(process.x_reg(0));
process.terminate(ExitReason::Error);
return SliceOutcome::Exited(ExitReason::Error, result);
}
match outcome {
NativeOutcome::Continue => {
let _transition = process.transition_to(ProcessStatus::Yielded);
SliceOutcome::Requeue(take_process(process))
}
NativeOutcome::Wait => {
let _transition = process.transition_to(ProcessStatus::Waiting);
SliceOutcome::Wait(take_process(process))
}
NativeOutcome::Stop(reason) => {
let result = crate::scheduler::exit_capture::capture_term(process.x_reg(0));
process.terminate(reason);
SliceOutcome::Exited(reason, result)
}
}
}
fn transition_to_running(process: &mut Process) -> Result<(), ()> {
match process.status() {
ProcessStatus::Running => Ok(()),
ProcessStatus::New | ProcessStatus::Yielded | ProcessStatus::Waiting => process
.transition_to(ProcessStatus::Running)
.map_err(|_| ()),
_ => Err(()),
}
}
fn take_process(process: &mut Process) -> Process {
std::mem::replace(process, Process::new(u64::MAX, DEFAULT_HEAP_SIZE))
}
#[cfg(test)]
mod tests {
use std::sync::atomic::{AtomicBool, Ordering};
use super::{SliceOutcome, run_native_slice};
use crate::atom::Atom;
use crate::module::ModuleRegistry;
use crate::namespace::NamespaceId;
use crate::native::native_process::{NativeBody, NativeContext, NativeHandler, NativeOutcome};
use crate::process::heap::DEFAULT_HEAP_SIZE;
use crate::process::{ExitReason, Process};
use crate::replay::{RecordedDeliveryKind, RecordedMessageDelivery, ReplayEvent, ReplayLog};
use crate::scheduler::{Scheduler, SchedulerConfig, supervision_integration};
use crate::term::Term;
use std::sync::Arc;
struct FlagHandler {
invoked: Arc<AtomicBool>,
}
impl NativeHandler for FlagHandler {
fn handle(&mut self, _ctx: &mut NativeContext<'_>) -> NativeOutcome {
self.invoked.store(true, Ordering::SeqCst);
NativeOutcome::Wait
}
}
fn single_thread_scheduler() -> Scheduler {
let config = SchedulerConfig {
thread_count: Some(1),
..Default::default()
};
Scheduler::new(config, Arc::new(ModuleRegistry::new())).expect("scheduler starts")
}
#[test]
fn tombstone_is_checked_before_handler_runs() {
let scheduler = single_thread_scheduler();
let invoked = Arc::new(AtomicBool::new(false));
let invoked_for_handler = Arc::clone(&invoked);
let pid = 4242;
let mut process = Process::new(pid, DEFAULT_HEAP_SIZE);
process.set_native_body(NativeBody::new(Box::new(move || {
Box::new(FlagHandler {
invoked: Arc::clone(&invoked_for_handler),
})
})));
scheduler
.shared
.exit_tombstones
.insert(pid, ExitReason::Kill);
let outcome = run_native_slice(&scheduler.shared, &mut process);
assert!(matches!(outcome, SliceOutcome::Exited(ExitReason::Kill, _)));
assert!(
!invoked.load(Ordering::SeqCst),
"handler must not run for a tombstoned pid"
);
scheduler.shutdown();
}
#[test]
fn native_send_validates_recorded_delivery_under_replay() {
let message = Term::atom(Atom::OK);
let log = ReplayLog::new(vec![ReplayEvent::MessageDelivery(
RecordedMessageDelivery {
order: 0,
kind: RecordedDeliveryKind::Message,
sender_pid: Some(2),
receiver_pid: 1,
sender_clock: 1,
receiver_clock: 2,
message,
},
)]);
let scheduler =
Scheduler::new_replay(SchedulerConfig::default(), log).expect("replay scheduler");
let receiver_pid = scheduler.spawn_test_process(false);
assert_eq!(receiver_pid, 1, "first spawned pid is 1");
let services =
supervision_integration::build_native_services(&scheduler.shared, NamespaceId::DEFAULT);
let local_send = services.local_send.clone().expect("local send facility");
let spawn = services.spawn_facility.clone().expect("spawn facility");
let replay_driver = scheduler.shared.replay_driver.clone();
let mut sender = Process::new(2, DEFAULT_HEAP_SIZE);
{
let mut ctx = NativeContext::new(&mut sender, local_send, spawn, replay_driver, None);
ctx.send(receiver_pid, message);
assert!(
ctx.take_replay_error().is_none(),
"a matching recorded delivery must not error"
);
}
assert_eq!(sender.logical_clock(), 1, "sender clock ticked once");
assert_eq!(scheduler.has_message(receiver_pid, message), Some(true));
assert!(
scheduler
.shared
.replay_driver
.as_ref()
.expect("driver")
.lock()
.expect("driver lock")
.is_complete(),
"the recorded delivery was consumed"
);
scheduler.shutdown();
}
#[test]
fn native_send_rolls_back_on_replay_mismatch() {
let message = Term::atom(Atom::OK);
let log = ReplayLog::new(vec![ReplayEvent::MessageDelivery(
RecordedMessageDelivery {
order: 0,
kind: RecordedDeliveryKind::Message,
sender_pid: Some(2),
receiver_pid: 1,
sender_clock: 99,
receiver_clock: 2,
message,
},
)]);
let scheduler =
Scheduler::new_replay(SchedulerConfig::default(), log).expect("replay scheduler");
let receiver_pid = scheduler.spawn_test_process(false);
let services =
supervision_integration::build_native_services(&scheduler.shared, NamespaceId::DEFAULT);
let local_send = services.local_send.clone().expect("local send facility");
let spawn = services.spawn_facility.clone().expect("spawn facility");
let replay_driver = scheduler.shared.replay_driver.clone();
let mut sender = Process::new(2, DEFAULT_HEAP_SIZE);
{
let mut ctx = NativeContext::new(&mut sender, local_send, spawn, replay_driver, None);
ctx.send(receiver_pid, message);
assert!(
ctx.take_replay_error().is_some(),
"a clock mismatch must surface a replay error"
);
}
assert_eq!(sender.logical_clock(), 0, "sender clock rolled back");
assert_eq!(
scheduler.has_message(receiver_pid, message),
Some(false),
"no message is delivered on a replay mismatch"
);
scheduler.shutdown();
}
}