1use std::sync::{Arc, Mutex};
4
5use aion_core::{ActivityError, Payload};
6use beamr::atom::AtomTable;
7use beamr::module::ModuleRegistry;
8use beamr::native::BifRegistryImpl;
9use beamr::process::ExitReason;
10use beamr::scheduler::{Scheduler, SchedulerConfig};
11use beamr::term::Term;
12
13use crate::error::EngineError;
14
15use super::config::{RuntimeConfig, SignalDeliveryConfig};
16#[cfg(test)]
17use super::nif::Mfa;
18use super::nif::NifRegistration;
19use super::payload::payload_to_term;
20
21use self::registration::{nif_registration_error, register_all_bifs};
22
23pub type Pid = u64;
25
26type RetainedHeap = Box<[u64]>;
27type RetainedHeaps = Vec<RetainedHeap>;
28type RetainedSpawnHeaps = Arc<dashmap::DashMap<Pid, Mutex<RetainedHeaps>>>;
29
30#[derive(Debug, Default, Eq, PartialEq)]
36pub struct RuntimeInput {
37 terms: Vec<Term>,
38 heaps: RetainedHeaps,
39}
40
41impl RuntimeInput {
42 pub fn from_payload(payload: &Payload) -> Result<Self, EngineError> {
54 let (term, heaps) = payload_to_term(payload)?.into_parts();
55 Ok(Self {
56 terms: vec![term],
57 heaps,
58 })
59 }
60
61 #[must_use]
63 pub fn arity(&self) -> u8 {
64 u8::try_from(self.terms.len()).unwrap_or(u8::MAX)
65 }
66
67 fn into_spawn_parts(self) -> (Vec<Term>, RetainedHeaps) {
68 (self.terms, self.heaps)
69 }
70}
71
72pub struct RuntimeHandle {
74 pub(super) scheduler: Arc<Scheduler>,
75 pub(super) atom_table: Arc<AtomTable>,
76 pub(super) module_registry: Arc<ModuleRegistry>,
77 pub(super) native_registry: Arc<BifRegistryImpl>,
78 nif_state: Arc<super::nif_state::EngineNifState>,
79 activity_results: Arc<dashmap::DashMap<(Pid, Pid), Payload>>,
80 activity_errors: Arc<dashmap::DashMap<(Pid, Pid), ActivityError>>,
81 activity_delivery_gates: dashmap::DashMap<Pid, Arc<activity_delivery::ActivityDeliveryGate>>,
88 activity_delivery_attempts: Arc<dashmap::DashMap<(Pid, Pid), u32>>,
95 #[cfg(test)]
96 activity_delivery_test_seams: activity_delivery::ActivityDeliveryTestSeams,
97 in_vm_children: Arc<dashmap::DashMap<Pid, std::collections::HashSet<Pid>>>,
108 registered_nif_modules: Arc<dashmap::DashSet<String>>,
109 spawn_heaps: RetainedSpawnHeaps,
110 signal_delivery: SignalDeliveryConfig,
111 outbox_enabled: bool,
114 pub(super) wake_confirmer: super::wake_confirm::WakeConfirmer,
117 pub(super) process_exits: Arc<super::process_exit::ProcessExitRegistry>,
119 pub(super) cleanup_executor: super::cleanup_executor::CleanupExecutor,
121 pub(super) abort_jobs: dashmap::DashMap<Pid, Arc<super::monitor::UnmonitoredProcessAbortJob>>,
123}
124
125impl RuntimeHandle {
126 pub fn new(config: RuntimeConfig) -> Result<Self, EngineError> {
134 let atom_table = Arc::new(AtomTable::with_common_atoms());
135 let module_registry = Arc::new(ModuleRegistry::new());
136 let nif_state = Arc::new(super::nif_state::EngineNifState::default());
139 let scheduler_config = SchedulerConfig {
140 thread_count: config.thread_count,
141 nif_private_data: Some(Arc::clone(&nif_state) as _),
142 ..Default::default()
143 };
144 let native_registry = Arc::new(BifRegistryImpl::new());
145 register_all_bifs(&native_registry, &atom_table, &nif_state)?;
146 let scheduler = Arc::new(
147 Scheduler::with_code_server(
148 scheduler_config,
149 Arc::clone(&module_registry),
150 Arc::clone(&atom_table),
151 Arc::clone(&native_registry),
152 )
153 .map_err(runtime_error_from_display)?,
154 );
155 let wake_confirmer = super::wake_confirm::WakeConfirmer::new(config.signal_delivery)?;
156 let shutdown_timeout = config.signal_delivery.cleanup_shutdown_timeout();
157 let cleanup_executor = super::cleanup_executor::CleanupExecutor::new(
158 config.signal_delivery.max_enqueue_attempts as usize,
159 shutdown_timeout,
160 )?;
161 let process_exits = super::process_exit::ProcessExitRegistry::new(
164 Arc::clone(&scheduler),
165 shutdown_timeout,
166 config.signal_delivery.max_enqueue_attempts as usize,
167 )?;
168 nif_state.set_process_exit_registry(&process_exits)?;
169
170 Ok(Self {
171 scheduler,
172 atom_table,
173 module_registry,
174 native_registry,
175 nif_state,
176 activity_results: Arc::new(dashmap::DashMap::new()),
177 activity_errors: Arc::new(dashmap::DashMap::new()),
178 activity_delivery_gates: dashmap::DashMap::new(),
179 activity_delivery_attempts: Arc::new(dashmap::DashMap::new()),
180 #[cfg(test)]
181 activity_delivery_test_seams: activity_delivery::ActivityDeliveryTestSeams::default(),
182 in_vm_children: Arc::new(dashmap::DashMap::new()),
183 registered_nif_modules: Arc::new(dashmap::DashSet::new()),
184 spawn_heaps: Arc::new(dashmap::DashMap::new()),
185 signal_delivery: config.signal_delivery,
186 outbox_enabled: config.outbox_enabled,
187 wake_confirmer,
188 process_exits,
189 cleanup_executor,
190 abort_jobs: dashmap::DashMap::new(),
191 })
192 }
193
194 pub(crate) fn nif_state(&self) -> &Arc<super::nif_state::EngineNifState> {
196 &self.nif_state
197 }
198
199 pub(crate) fn signal_delivery(&self) -> SignalDeliveryConfig {
201 self.signal_delivery
202 }
203
204 pub(crate) fn outbox_enabled(&self) -> bool {
209 self.outbox_enabled
210 }
211
212 pub fn install_nifs(&self, registration: NifRegistration) -> Result<(), EngineError> {
223 for entry in registration.into_entries() {
224 let mfa = entry.mfa;
225 let module = self.atom_table.intern(&mfa.module);
226 let function = self.atom_table.intern(&mfa.function);
227 let capability = beamr::native::Capability::ExternalIo;
228 let result = if entry.is_dirty {
229 self.native_registry.register_dirty(
230 module,
231 function,
232 mfa.arity,
233 entry.function,
234 beamr::scheduler::dirty::DirtySchedulerKind::Cpu,
235 capability,
236 )
237 } else {
238 self.native_registry.register(
239 module,
240 function,
241 mfa.arity,
242 entry.function,
243 capability,
244 )
245 };
246 result.map_err(|error| nif_registration_error(&mfa, error))?;
247 self.registered_nif_modules.insert(mfa.module);
248 }
249
250 Ok(())
251 }
252
253 #[must_use]
256 pub fn registered_nif_modules(&self) -> Vec<String> {
257 let mut module_names: Vec<_> = self
258 .registered_nif_modules
259 .iter()
260 .map(|module_name| module_name.key().clone())
261 .collect();
262 module_names.sort();
263 module_names
264 }
265
266 pub fn spawn_workflow(
273 &self,
274 deployed_module: &str,
275 function: &str,
276 input: RuntimeInput,
277 ) -> Result<Pid, EngineError> {
278 self.spawn_process(deployed_module, function, input)
279 }
280
281 pub fn spawn_workflow_trapping(
288 &self,
289 deployed_module: &str,
290 function: &str,
291 input: RuntimeInput,
292 ) -> Result<Pid, EngineError> {
293 self.release_dead_spawn_heaps();
294 let module = self.atom_table.intern(deployed_module);
295 let function = self.atom_table.intern(function);
296 let (terms, heaps) = input.into_spawn_parts();
297 let pid = self.spawn_with_exit_ownership(|| {
298 self.scheduler
299 .spawn_trap_exit(module, function, terms)
300 .map_err(runtime_error_from_display)
301 })?;
302 self.retain_spawn_heaps(pid, heaps);
303 Ok(pid)
304 }
305
306 pub fn spawn_activity(
314 &self,
315 parent_pid: Pid,
316 deployed_module: &str,
317 function: &str,
318 input: RuntimeInput,
319 ) -> Result<Pid, EngineError> {
320 self.release_dead_spawn_heaps();
321 self.ensure_live_pid(parent_pid)?;
322 self.wait_for_process_ready(parent_pid)?;
323 let module = self.atom_table.intern(deployed_module);
324 let function_atom = self.atom_table.intern(function);
325 let (terms, heaps) = input.into_spawn_parts();
326 let pid = self.spawn_with_exit_ownership(|| {
327 self.scheduler
328 .spawn_link(parent_pid, module, function_atom, terms)
329 .map_err(runtime_error_from_display)
330 })?;
331 self.retain_spawn_heaps(pid, heaps);
332 Ok(pid)
333 }
334
335 pub fn spawn_activity_closure(
355 &self,
356 parent_pid: Pid,
357 closure_term: Term,
358 ) -> Result<Pid, EngineError> {
359 self.release_dead_spawn_heaps();
360 self.ensure_live_pid(parent_pid)?;
361 let pid = self.spawn_with_exit_ownership(|| {
362 self.scheduler
363 .spawn_link_closure(parent_pid, closure_term)
364 .map_err(runtime_error_from_display)
365 })?;
366 self.in_vm_children
367 .entry(parent_pid)
368 .or_default()
369 .insert(pid);
370 if !self.is_live(parent_pid) {
381 self.kill_in_vm_children(parent_pid);
382 return Err(EngineError::Runtime {
383 reason: format!(
384 "in-vm activity child spawn: parent workflow process {parent_pid} exited during spawn"
385 ),
386 });
387 }
388 Ok(pid)
389 }
390
391 pub(crate) fn deregister_in_vm_child(&self, parent_pid: Pid, child_pid: Pid) {
394 if let Some(mut children) = self.in_vm_children.get_mut(&parent_pid) {
395 children.remove(&child_pid);
396 }
397 self.in_vm_children
398 .remove_if(&parent_pid, |_, children| children.is_empty());
399 }
400
401 pub(crate) fn kill_in_vm_children(&self, workflow_pid: Pid) {
409 let Some((_, children)) = self.in_vm_children.remove(&workflow_pid) else {
410 return;
411 };
412 for child_pid in children {
413 if self.is_live(child_pid) {
414 tracing::debug!(
415 workflow_pid,
416 child_pid,
417 "killing orphaned in-vm activity child on workflow exit"
418 );
419 self.scheduler
420 .terminate_process(child_pid, ExitReason::Kill);
421 }
422 self.release_spawn_heaps(child_pid);
423 }
424 }
425
426 #[must_use]
428 pub fn is_dirty(&self, module: &str, function: &str) -> bool {
429 self.is_dirty_with_arity(module, function, 1)
430 }
431
432 #[must_use]
434 pub fn is_dirty_with_arity(&self, module: &str, function: &str, arity: u8) -> bool {
435 let module = self.atom_table.intern(module);
436 let function = self.atom_table.intern(function);
437 self.native_registry
438 .lookup(module, function, arity)
439 .is_some_and(|entry| entry.dirty_kind.is_some())
440 }
441
442 pub fn cancel_pid(&self, pid: Pid) -> Result<(), EngineError> {
448 self.ensure_live_pid(pid)?;
449 self.scheduler.terminate_process(pid, ExitReason::Kill);
450 self.release_spawn_heaps(pid);
451 Ok(())
452 }
453
454 pub fn set_trap_exit(&self, pid: Pid, value: bool) -> Result<bool, EngineError> {
460 self.scheduler
461 .set_trap_exit(pid, value)
462 .map_err(runtime_error_from_display)
463 }
464
465 #[must_use]
467 pub fn is_live(&self, pid: Pid) -> bool {
468 self.scheduler.process_table().get(pid).is_some()
469 }
470
471 pub fn trap_exit(&self, pid: Pid) -> Result<bool, EngineError> {
477 self.scheduler
478 .trap_exit(pid)
479 .ok_or_else(|| runtime_error(format!("process {pid} is not live")))
480 }
481
482 pub fn is_linked(&self, left: Pid, right: Pid) -> Result<bool, EngineError> {
488 self.ensure_live_pid(left)?;
489 self.ensure_live_pid(right)?;
490 Ok(self.scheduler.is_linked(left, right))
491 }
492
493 pub fn shutdown(&self) -> Result<(), EngineError> {
499 let workflow_pids: Vec<Pid> = self
502 .in_vm_children
503 .iter()
504 .map(|entry| *entry.key())
505 .collect();
506 for workflow_pid in workflow_pids {
507 self.kill_in_vm_children(workflow_pid);
508 }
509 for pid in self.process_exits.begin_shutdown()? {
510 if self.is_live(pid) {
511 self.scheduler.terminate_process(pid, ExitReason::Kill);
512 }
513 }
514 self.cleanup_executor.shutdown()?;
518 self.wake_confirmer.shutdown();
520 self.process_exits.close_and_join_all()?;
521 self.scheduler.shutdown();
522 self.spawn_heaps.clear();
523 Ok(())
524 }
525
526 fn spawn_process(
527 &self,
528 deployed_module: &str,
529 function: &str,
530 input: RuntimeInput,
531 ) -> Result<Pid, EngineError> {
532 self.release_dead_spawn_heaps();
533 let module = self.atom_table.intern(deployed_module);
534 let function = self.atom_table.intern(function);
535 let (terms, heaps) = input.into_spawn_parts();
536 let pid = self.spawn_with_exit_ownership(|| {
537 self.scheduler
538 .spawn(module, function, terms)
539 .map_err(runtime_error_from_display)
540 })?;
541 self.retain_spawn_heaps(pid, heaps);
542 Ok(pid)
543 }
544
545 fn retain_spawn_heaps(&self, pid: Pid, heaps: RetainedHeaps) {
546 if heaps.is_empty() {
547 return;
548 }
549 self.spawn_heaps.insert(pid, Mutex::new(heaps));
550 }
551
552 pub(super) fn release_spawn_heaps(&self, pid: Pid) {
553 self.spawn_heaps.remove(&pid);
554 }
555
556 fn release_dead_spawn_heaps(&self) {
557 let dead_pids: Vec<Pid> = self
558 .spawn_heaps
559 .iter()
560 .filter_map(|entry| {
561 let pid = *entry.key();
562 self.scheduler
563 .process_table()
564 .get(pid)
565 .is_none()
566 .then_some(pid)
567 })
568 .collect();
569 for pid in dead_pids {
570 self.release_spawn_heaps(pid);
571 }
572 }
573
574 pub(super) fn ensure_live_pid(&self, pid: Pid) -> Result<(), EngineError> {
575 if self.scheduler.process_table().get(pid).is_some() {
576 Ok(())
577 } else {
578 Err(runtime_error(format!("process {pid} is not live")))
579 }
580 }
581
582 #[cfg(test)]
583 pub(crate) fn live_processes_for_test(&self) -> usize {
584 self.scheduler.process_table().len()
585 }
586
587 #[cfg(test)]
593 pub fn spawn_test_process(&self) -> Result<Pid, EngineError> {
594 self.spawn_with_exit_ownership(|| Ok(self.scheduler.spawn_test_process(false)))
595 }
596
597 #[cfg(test)]
603 pub fn spawn_test_process_with_trap_exit(&self, trap_exit: bool) -> Result<Pid, EngineError> {
604 self.spawn_with_exit_ownership(|| Ok(self.scheduler.spawn_test_process(trap_exit)))
605 }
606
607 #[cfg(test)]
614 pub fn spawn_linked_test_process(&self, parent_pid: Pid) -> Result<Pid, EngineError> {
615 self.ensure_live_pid(parent_pid)?;
616 self.spawn_with_exit_ownership(|| {
617 self.scheduler
618 .spawn_linked_test_process(parent_pid)
619 .map_err(runtime_error_from_display)
620 })
621 }
622
623 #[cfg(test)]
629 pub fn has_trapped_exit_message(
630 &self,
631 target_pid: Pid,
632 source_pid: Pid,
633 ) -> Result<bool, EngineError> {
634 self.ensure_live_pid(target_pid)?;
635 Ok(self
636 .scheduler
637 .has_trapped_exit_message(target_pid, source_pid)
638 .unwrap_or(false))
639 }
640
641 #[cfg(test)]
650 pub fn wait_for_trapped_exit(
651 &self,
652 target_pid: Pid,
653 source_pid: Pid,
654 ) -> Result<(), EngineError> {
655 let deadline = std::time::Instant::now() + std::time::Duration::from_millis(50);
656 while std::time::Instant::now() < deadline {
657 if self
658 .scheduler
659 .has_trapped_exit_message(target_pid, source_pid)
660 .unwrap_or(false)
661 {
662 return Ok(());
663 }
664 std::thread::sleep(std::time::Duration::from_millis(1));
665 }
666 Err(runtime_error(format!(
667 "trapped exit from {source_pid} to {target_pid} did not arrive"
668 )))
669 }
670
671 #[cfg(test)]
677 pub fn terminate_test_process_with_error(&self, pid: Pid) -> Result<(), EngineError> {
678 self.ensure_live_pid(pid)?;
679 self.scheduler.terminate_process(pid, ExitReason::Error);
680 Ok(())
681 }
682
683 #[cfg(test)]
684 pub(crate) fn lookup_native_for_test(
685 &self,
686 module: &str,
687 function: &str,
688 arity: u8,
689 ) -> Option<beamr::native::NativeEntry> {
690 let module = self.atom_table.intern(module);
691 let function = self.atom_table.intern(function);
692 self.native_registry.lookup(module, function, arity)
693 }
694
695 #[cfg(test)]
696 pub(crate) fn retained_spawn_heap_count_for_test(&self) -> usize {
697 self.release_dead_spawn_heaps();
698 self.spawn_heaps.len()
699 }
700}
701
702fn runtime_error(reason: String) -> EngineError {
703 EngineError::Runtime { reason }
704}
705
706fn runtime_error_from_display(reason: impl std::fmt::Display) -> EngineError {
707 runtime_error(reason.to_string())
708}
709
710mod activity_delivery;
711mod delivery;
712mod process_ownership;
713mod readiness;
714mod registration;
715mod spawn_bifs;
716
717pub(crate) use delivery::InVmChildOutcome;
718
719#[cfg(test)]
720#[path = "handle/test_support.rs"]
721mod test_support;
722
723#[cfg(test)]
724#[path = "handle/tests.rs"]
725mod tests;