frame-core 0.3.0

Component model, lifecycle, process isolation — hosts components as supervised BEAM process trees
Documentation
use std::sync::Arc;
use std::sync::mpsc::{self, Receiver, Sender};
use std::time::{Duration, Instant};

use beamr::atom::Atom;
use beamr::native::native_process::{NativeContext, NativeHandler, NativeOutcome};
use beamr::process::ExitReason;
use beamr::scheduler::Scheduler;
use beamr::term::Term;
use beamr::term::boxed::Tuple;

#[derive(Debug)]
pub(crate) enum ObservationError {
    Spawn(String),
    ReadyDeadline,
    Registration(String),
    TombstoneDeadline,
    Protocol(String),
    CleanupDeadline,
}

pub(crate) struct ExitObservation {
    scheduler: Arc<Scheduler>,
    observer: u64,
    receiver: Receiver<ObserverEvent>,
    timeout: Duration,
}

impl ExitObservation {
    pub(crate) fn register(
        scheduler: &Arc<Scheduler>,
        target: u64,
        timeout: Duration,
    ) -> Result<Self, ObservationError> {
        let (sender, receiver) = mpsc::channel();
        let observer = scheduler
            .spawn_native(Box::new(move || {
                Box::new(TombstoneObserver {
                    target,
                    sender: sender.clone(),
                    ready: false,
                    completion: Some(Err("observer stopped before target DOWN".to_owned())),
                })
            }))
            .map_err(|error| ObservationError::Spawn(error.to_string()))?;
        match receiver.recv_timeout(timeout) {
            Ok(ObserverEvent::Ready) => {}
            Ok(ObserverEvent::Complete(result)) => {
                return Err(ObservationError::Protocol(result.err().unwrap_or_else(
                    || "observer completed before readiness registration".to_owned(),
                )));
            }
            Err(_) => {
                cleanup_observer(scheduler, observer, &receiver, timeout)?;
                return Err(ObservationError::ReadyDeadline);
            }
        }
        if let Err(error) = scheduler.monitor(observer, target) {
            cleanup_observer(scheduler, observer, &receiver, timeout)?;
            return Err(ObservationError::Registration(error.to_string()));
        }
        Ok(Self {
            scheduler: Arc::clone(scheduler),
            observer,
            receiver,
            timeout,
        })
    }

    pub(crate) fn wait(self) -> Result<ExitReason, ObservationError> {
        match self.receiver.recv_timeout(self.timeout) {
            Ok(ObserverEvent::Complete(result)) => result.map_err(ObservationError::Protocol),
            Ok(ObserverEvent::Ready) => {
                self.cancel()?;
                Err(ObservationError::Protocol(
                    "observer emitted readiness more than once".to_owned(),
                ))
            }
            Err(_) => {
                self.cancel()?;
                Err(ObservationError::TombstoneDeadline)
            }
        }
    }

    pub(crate) fn cancel(self) -> Result<(), ObservationError> {
        cleanup_observer(&self.scheduler, self.observer, &self.receiver, self.timeout)
    }
}

fn cleanup_observer(
    scheduler: &Scheduler,
    observer: u64,
    receiver: &Receiver<ObserverEvent>,
    timeout: Duration,
) -> Result<(), ObservationError> {
    scheduler.terminate_process(observer, ExitReason::Kill);
    let deadline = Instant::now() + timeout;
    loop {
        let remaining = deadline.saturating_duration_since(Instant::now());
        match receiver.recv_timeout(remaining) {
            Ok(ObserverEvent::Complete(_)) => return Ok(()),
            // A late Ready is legitimate here: the ReadyDeadline path tears down
            // an observer whose one Ready may still be in flight. Only `wait()`,
            // which runs after register() consumed that Ready, treats a second
            // Ready as a protocol error.
            Ok(ObserverEvent::Ready) => {}
            Err(_) => return Err(ObservationError::CleanupDeadline),
        }
    }
}

enum ObserverEvent {
    Ready,
    Complete(Result<ExitReason, String>),
}

struct TombstoneObserver {
    target: u64,
    sender: Sender<ObserverEvent>,
    ready: bool,
    completion: Option<Result<ExitReason, String>>,
}

impl NativeHandler for TombstoneObserver {
    fn handle(&mut self, context: &mut NativeContext<'_>) -> NativeOutcome {
        if !self.ready {
            self.ready = true;
            if self.sender.send(ObserverEvent::Ready).is_err() {
                return NativeOutcome::Stop(ExitReason::Error);
            }
        }
        let Some(message) = context.recv() else {
            return NativeOutcome::Wait;
        };
        match decode_down(message, self.target) {
            Ok(Some(reason)) => {
                self.completion = Some(Ok(reason));
                NativeOutcome::Stop(ExitReason::Normal)
            }
            Ok(None) => {
                self.completion = Some(Err(
                    "tombstone observer received a non-DOWN message".to_owned()
                ));
                NativeOutcome::Stop(ExitReason::Error)
            }
            Err(detail) => {
                self.completion = Some(Err(detail));
                NativeOutcome::Stop(ExitReason::Error)
            }
        }
    }
}

impl Drop for TombstoneObserver {
    fn drop(&mut self) {
        if let Some(completion) = self.completion.take() {
            let _undelivered = self.sender.send(ObserverEvent::Complete(completion));
        }
    }
}

fn decode_down(message: Term, target: u64) -> Result<Option<ExitReason>, String> {
    let Some(tuple) = Tuple::new(message) else {
        return Ok(None);
    };
    if tuple.arity() != 5 || tuple.get(0).and_then(Term::as_atom) != Some(Atom::DOWN) {
        return Ok(None);
    }
    if tuple.get(2).and_then(Term::as_atom) != Some(Atom::PROCESS)
        || tuple.get(3).and_then(Term::as_pid) != Some(target)
    {
        return Err("DOWN identified the wrong observed process".to_owned());
    }
    tuple
        .get(4)
        .and_then(Term::as_atom)
        .and_then(reason_from_atom)
        .map(Some)
        .ok_or_else(|| "DOWN carried an unknown exit reason".to_owned())
}

fn reason_from_atom(reason: Atom) -> Option<ExitReason> {
    if reason == Atom::NORMAL {
        Some(ExitReason::Normal)
    } else if reason == Atom::KILL {
        Some(ExitReason::Kill)
    } else if reason == Atom::KILLED {
        Some(ExitReason::Killed)
    } else if reason == Atom::ERROR {
        Some(ExitReason::Error)
    } else if reason == Atom::NOCONNECTION {
        Some(ExitReason::NoConnection)
    } else if reason == Atom::NOPROC {
        Some(ExitReason::NoProc)
    } else {
        None
    }
}