use std::sync::Arc;
use crate::durability::{Command, CorrelationKey, ResolveOutcome};
use crate::runtime::nif_activity::{
ScheduledActivity, context_error_term, correlation_id, decode_string_arg, error_result_term,
json_payload, ok_result_term, record_started, runtime_context,
};
use crate::runtime::nif_activity_dispatch::FIRST_DELIVERY_ATTEMPT;
use crate::runtime::nif_context::NifContext;
use aion_core::ActivityId;
use beamr::native::ProcessContext;
use beamr::term::Term;
use beamr::term::boxed::Closure;
pub(super) fn dispatch_activity_in_vm_impl(
args: &[Term],
ctx: &mut ProcessContext,
) -> Result<Term, Term> {
if args.len() > 255 {
return Err(Term::NIL);
}
let (name, input, config) = match decode_in_vm_args(args) {
Ok(parts) => parts,
Err(reason) => return Ok(error_result_term(ctx, &reason).unwrap_or(Term::NIL)),
};
let thunk = args[3];
let Some(pid) = ctx.pid() else {
return Ok(
error_result_term(ctx, "dispatch_activity_in_vm: missing calling process pid")
.unwrap_or(Term::NIL),
);
};
let state = match super::nif_state::engine_nif_state(ctx) {
Ok(state) => state,
Err(error) => return Ok(error_result_term(ctx, &error).unwrap_or(Term::NIL)),
};
if let Err(error) =
super::nif_query_pump::ensure_not_servicing_query(&state, pid, "dispatch_activity_in_vm")
{
return Ok(error_result_term(ctx, &error).unwrap_or(Term::NIL));
}
let runtime = match runtime_context(&state) {
Ok(runtime) => runtime,
Err(error) => return Ok(context_error_term(ctx, &error)),
};
let context = match NifContext::new(
pid,
runtime.registry.as_ref(),
runtime.tokio_handle.clone(),
runtime.runtime.signal_delivery(),
) {
Ok(context) => context,
Err(error) => return Ok(context_error_term(ctx, &error)),
};
dispatch_in_vm_with_context(
ctx,
context,
&runtime.runtime,
&runtime.tokio_handle,
(name, input, config),
thunk,
)
}
fn decode_in_vm_args(args: &[Term]) -> Result<(String, String, String), String> {
if args.len() != 4 {
return Err(format!(
"dispatch_activity_in_vm: expected 4 arguments, got {}",
args.len()
));
}
let name = decode_string_arg(args[0])
.map_err(|error| format!("dispatch_activity_in_vm name: {error}"))?;
let input = decode_string_arg(args[1])
.map_err(|error| format!("dispatch_activity_in_vm input: {error}"))?;
let config = decode_string_arg(args[2])
.map_err(|error| format!("dispatch_activity_in_vm config: {error}"))?;
let Some(thunk) = Closure::new(args[3]) else {
return Err("dispatch_activity_in_vm: thunk argument is not a closure".to_owned());
};
if thunk.arity() != 0 {
return Err(format!(
"dispatch_activity_in_vm: thunk arity is {}, expected 0",
thunk.arity()
));
}
Ok((name, input, config))
}
fn dispatch_in_vm_with_context(
ctx: &mut ProcessContext,
mut context: NifContext,
runtime: &Arc<crate::RuntimeHandle>,
tokio_handle: &tokio::runtime::Handle,
(name, input, config): (String, String, String),
thunk: Term,
) -> Result<Term, Term> {
let input_payload = json_payload(ctx, &input, "dispatch_activity_in_vm", "input")?;
let ordinal = context.next_activity_ordinal();
let key = CorrelationKey::Activity(ordinal);
let activity_id = ActivityId::from_sequence_position(ordinal);
let correlation = correlation_id(ordinal);
match context
.resolve_command(Command::RunActivity {
key,
activity_type: name.clone(),
input: input_payload.clone(),
})
.map_err(|error| context_error_term(ctx, &error))?
{
ResolveOutcome::Recorded(_) => {
Ok(ok_result_term(ctx, correlation.as_bytes()).unwrap_or(Term::NIL))
}
ResolveOutcome::ResumeLive => {
let start_time_task_queue = context.start_time_task_queue();
let task_queue =
super::nif_activity::resolve_task_queue(&config, start_time_task_queue.as_deref());
let node = super::nif_activity::resolve_node(&config);
record_started(
ctx,
&context,
activity_id,
ScheduledActivity {
activity_type: name,
input: input_payload,
task_queue,
node,
attempt: FIRST_DELIVERY_ATTEMPT,
},
)?;
match runtime.spawn_activity_closure(context.pid(), thunk) {
Ok(child_pid) => {
spawn_in_vm_completion_watcher(
tokio_handle,
Arc::clone(runtime),
context.pid(),
child_pid,
correlation.clone(),
);
}
Err(error) => retain_spawn_failure(runtime, context.pid(), &correlation, &error),
}
Ok(ok_result_term(ctx, correlation.as_bytes()).unwrap_or(Term::NIL))
}
}
}
pub(super) fn retain_spawn_failure(
runtime: &crate::RuntimeHandle,
workflow_pid: u64,
correlation: &str,
error: &impl std::fmt::Display,
) {
let reason = format!("terminal:in-vm activity child spawn failed: {error}");
if let Err(delivery_error) =
runtime.deliver_activity_failure_message(workflow_pid, correlation, reason)
{
tracing::debug!(
%delivery_error,
workflow_pid,
correlation_id = correlation,
"in-vm spawn-failure marker not queued; retained entry settles the await"
);
}
}
fn spawn_in_vm_completion_watcher(
tokio_handle: &tokio::runtime::Handle,
runtime: Arc<crate::RuntimeHandle>,
workflow_pid: u64,
child_pid: u64,
correlation_id: String,
) {
tokio_handle.spawn_blocking(move || {
let outcome = runtime.in_vm_child_outcome(child_pid);
runtime.deregister_in_vm_child(workflow_pid, child_pid);
let delivered = match outcome {
Ok(crate::runtime::handle::InVmChildOutcome::Completed(payload)) => {
runtime.deliver_activity_completion_message(workflow_pid, &correlation_id, payload)
}
Ok(crate::runtime::handle::InVmChildOutcome::Failed(reason)) => {
runtime.deliver_activity_failure_message(workflow_pid, &correlation_id, reason)
}
Err(error) => Err(error),
};
if let Err(error) = delivered {
tracing::warn!(
%error,
workflow_pid,
child_pid,
correlation_id,
"in-vm activity outcome delivery failed"
);
}
});
}