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;

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;

pub(crate) struct ExitWitness {
    receiver: Receiver<Result<(), String>>,
    deadline: Duration,
}

impl ExitWitness {
    pub(crate) fn register(
        scheduler: &Arc<Scheduler>,
        target: u64,
        deadline: Duration,
    ) -> Result<Self, Box<dyn std::error::Error>> {
        let (ready_tx, ready_rx) = mpsc::channel();
        let (done_tx, done_rx) = mpsc::channel();
        let watcher = scheduler.spawn_native(Box::new(move || {
            Box::new(WitnessHandler {
                target,
                ready: Some(ready_tx.clone()),
                done: done_tx.clone(),
                result: Some(Err("witness stopped before DOWN".to_owned())),
            })
        }))?;
        ready_rx.recv_timeout(deadline)?;
        let _registration = scheduler.monitor_with_result(watcher, target)?;
        Ok(Self {
            receiver: done_rx,
            deadline,
        })
    }

    pub(crate) fn wait(self) -> Result<(), Box<dyn std::error::Error>> {
        self.receiver
            .recv_timeout(self.deadline)?
            .map_err(Into::into)
    }
}

struct WitnessHandler {
    target: u64,
    ready: Option<Sender<()>>,
    done: Sender<Result<(), String>>,
    result: Option<Result<(), String>>,
}

impl NativeHandler for WitnessHandler {
    fn handle(&mut self, context: &mut NativeContext<'_>) -> NativeOutcome {
        if let Some(ready) = self.ready.take()
            && ready.send(()).is_err()
        {
            return NativeOutcome::Stop(ExitReason::Error);
        }
        let Some(message) = context.recv() else {
            return NativeOutcome::Wait;
        };
        self.result = Some(decode_down(message, self.target));
        NativeOutcome::Stop(ExitReason::Normal)
    }
}

impl Drop for WitnessHandler {
    fn drop(&mut self) {
        if let Some(result) = self.result.take() {
            let _undelivered = self.done.send(result);
        }
    }
}

fn decode_down(message: Term, target: u64) -> Result<(), String> {
    let tuple = Tuple::new(message).ok_or_else(|| "DOWN was not a tuple".to_owned())?;
    if tuple.arity() == 5
        && tuple.get(0).and_then(Term::as_atom) == Some(Atom::DOWN)
        && tuple.get(2).and_then(Term::as_atom) == Some(Atom::PROCESS)
        && tuple.get(3).and_then(Term::as_pid) == Some(target)
    {
        Ok(())
    } else {
        Err("witness received malformed DOWN".to_owned())
    }
}