1use aion_core::{Payload, WorkflowError};
4use beamr::atom::{Atom, AtomTable};
5use beamr::process::ExitReason;
6use beamr::term::Term;
7use beamr::term::boxed::Tuple;
8
9use crate::{EngineError, Pid};
10
11use super::payload::term_to_payload;
12use super::process_exit::OwnedProcessExitOutcome;
13
14#[derive(Clone)]
16pub enum WorkflowProcessOutcome {
17 Completed(Payload),
19 Failed(WorkflowError),
21}
22
23pub(super) fn workflow_outcome(
24 atoms: &AtomTable,
25 pid: Pid,
26 owned: &OwnedProcessExitOutcome,
27) -> Result<Result<Payload, WorkflowError>, EngineError> {
28 match workflow_outcome_from_owned_exit(atoms, pid, owned)? {
29 WorkflowProcessOutcome::Completed(payload) => Ok(Ok(payload)),
30 WorkflowProcessOutcome::Failed(error) => Ok(Err(error)),
31 }
32}
33
34pub(super) fn workflow_outcome_from_owned_exit(
35 atoms: &AtomTable,
36 pid: Pid,
37 owned: &OwnedProcessExitOutcome,
38) -> Result<WorkflowProcessOutcome, EngineError> {
39 let observed = match owned {
40 OwnedProcessExitOutcome::Observed(observed) => observed,
41 OwnedProcessExitOutcome::ObservationFailed {
42 process_id,
43 failure,
44 } => return Err(failure.into_engine_error(*process_id)),
45 };
46 if observed.reason != ExitReason::Normal {
47 if let Some(error) = &observed.execution_error {
50 let formatted = error.format_with_atoms(atoms);
51 let residue = observed
52 .exception
53 .as_ref()
54 .map_or_else(String::new, |exception| {
55 format!(
56 " (residual exception: {})",
57 exception.format_with_atoms(atoms)
58 )
59 });
60 return Ok(WorkflowProcessOutcome::Failed(WorkflowError {
61 message: format!(
62 "workflow process {pid} exited: {:?}: VM execution error: {formatted}{residue}",
63 observed.reason
64 ),
65 details: None,
66 }));
67 }
68 if let Some(exception) = &observed.exception {
69 let formatted = exception.format_with_atoms(atoms);
70 let view = exception.view();
71 let details = beamr::ets::copy_term_to_ets(view.reason)
72 .ok()
73 .and_then(|reason| {
74 term_to_payload(reason.root(), atoms, reason.borrow_terms()).ok()
75 });
76 return Ok(WorkflowProcessOutcome::Failed(WorkflowError {
77 message: format!("workflow process {pid} exited: {formatted}"),
78 details,
79 }));
80 }
81 }
82 convert_process_outcome(
83 atoms,
84 pid,
85 observed.reason,
86 observed.result.root(),
87 observed.result.borrow_terms(),
88 )
89}
90
91pub(super) fn convert_process_outcome(
92 atoms: &AtomTable,
93 pid: Pid,
94 reason: ExitReason,
95 result: Term,
96 heap: beamr::term::heap_borrow::HeapBorrow<'_>,
97) -> Result<WorkflowProcessOutcome, EngineError> {
98 if reason == ExitReason::Normal {
99 unwrap_gleam_result(result, atoms, pid, heap)
100 } else {
101 let formatted = beamr::term::format::format_term(result, atoms);
102 let details = term_to_payload(result, atoms, heap).ok();
103 Ok(WorkflowProcessOutcome::Failed(WorkflowError {
104 message: format!("workflow process {pid} exited: {reason:?}: {formatted}"),
105 details,
106 }))
107 }
108}
109
110fn unwrap_gleam_result(
111 result: Term,
112 atoms: &AtomTable,
113 pid: Pid,
114 heap: beamr::term::heap_borrow::HeapBorrow<'_>,
115) -> Result<WorkflowProcessOutcome, EngineError> {
116 if let Some(tuple) = Tuple::new(result)
117 && tuple.arity() == 2
118 && let (Some(tag), Some(value)) = (tuple.get(0), tuple.get(1))
119 && let Some(atom) = tag.as_atom()
120 {
121 if atom == Atom::OK {
122 return Ok(WorkflowProcessOutcome::Completed(term_to_payload(
123 value, atoms, heap,
124 )?));
125 }
126 if atom == Atom::ERROR {
127 let details = term_to_payload(value, atoms, heap).ok();
128 return Ok(WorkflowProcessOutcome::Failed(WorkflowError {
129 message: format!("workflow {pid} returned error"),
130 details,
131 }));
132 }
133 }
134 Ok(WorkflowProcessOutcome::Completed(term_to_payload(
135 result, atoms, heap,
136 )?))
137}