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(()),
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
}
}