1use aion_core::{ActivityError, ActivityErrorKind, ContentType, Payload};
8use beamr::atom::Atom;
9use beamr::process::ExitReason;
10
11use crate::error::EngineError;
12
13use super::{Pid, RuntimeHandle, runtime_error};
14use crate::runtime::payload::term_to_payload;
15
16impl RuntimeHandle {
17 pub fn propagate_activity_outcome(
28 &self,
29 parent_pid: Pid,
30 activity_pid: Pid,
31 ) -> Result<(), EngineError> {
32 self.ensure_live_pid(parent_pid)?;
33 let (reason, owned_result) = self.scheduler.run_until_exit(activity_pid);
34 self.release_spawn_heaps(activity_pid);
35 if reason == ExitReason::Normal {
36 let payload = term_to_payload(owned_result.root(), &self.atom_table)?;
37 self.deliver_activity_result(parent_pid, activity_pid, payload)
38 } else {
39 let error = self
40 .activity_errors
41 .get(&(parent_pid, activity_pid))
42 .map_or_else(
43 || ActivityError {
44 kind: ActivityErrorKind::Terminal,
45 message: format!("activity process {activity_pid} exited: {reason:?}"),
46 details: None,
47 },
48 |entry| entry.clone(),
49 );
50 self.deliver_activity_error(parent_pid, activity_pid, error)
51 }
52 }
53
54 pub fn deliver_signal_received(&self, workflow_pid: Pid) -> Result<(), EngineError> {
70 self.ensure_live_pid(workflow_pid)?;
71 self.wait_for_process_ready(workflow_pid)?;
72 let marker = self.atom_table.intern("aion_signal_received");
73 self.enqueue_signal_marker_with_retry(workflow_pid, marker)
74 }
75
76 pub(crate) async fn deliver_signal_received_async(
85 &self,
86 workflow_pid: Pid,
87 ) -> Result<(), EngineError> {
88 self.ensure_live_pid(workflow_pid)?;
89 self.wait_for_process_ready_async(workflow_pid).await?;
90 let marker = self.atom_table.intern("aion_signal_received");
91 self.enqueue_signal_marker_with_retry_async(workflow_pid, marker)
92 .await
93 }
94
95 pub(crate) fn deliver_query_request(&self, workflow_pid: Pid) -> Result<(), EngineError> {
107 self.ensure_live_pid(workflow_pid)?;
108 self.wait_for_process_ready(workflow_pid)?;
109 let marker = self.atom_table.intern("aion_query");
110 self.enqueue_signal_marker_with_retry(workflow_pid, marker)
111 }
112
113 pub(crate) async fn deliver_child_terminal(
132 &self,
133 workflow_pid: Pid,
134 ) -> Result<(), EngineError> {
135 self.ensure_live_pid(workflow_pid)?;
136 self.wait_for_process_ready_async(workflow_pid).await?;
137 let marker = self.atom_table.intern("aion_child_terminal");
138 self.enqueue_signal_marker_with_retry_async(workflow_pid, marker)
139 .await
140 }
141
142 pub(crate) fn deliver_activity_completion_message(
154 &self,
155 workflow_pid: Pid,
156 correlation_id: &str,
157 result: String,
158 ) -> Result<(), EngineError> {
159 self.ensure_live_pid(workflow_pid)?;
160 let activity_id = correlation_to_activity_pid(correlation_id)?;
161 self.activity_results.insert(
162 (workflow_pid, activity_id),
163 Payload::new(ContentType::Json, result.into_bytes()),
164 );
165 let marker = self.atom_table.intern("activity_complete");
166 self.enqueue_activity_marker(workflow_pid, marker, correlation_id)
167 }
168
169 pub(crate) fn deliver_activity_failure_message(
176 &self,
177 workflow_pid: Pid,
178 correlation_id: &str,
179 reason: String,
180 ) -> Result<(), EngineError> {
181 self.ensure_live_pid(workflow_pid)?;
182 let activity_id = correlation_to_activity_pid(correlation_id)?;
183 self.activity_errors
184 .insert((workflow_pid, activity_id), activity_failure(reason));
185 let marker = self.atom_table.intern("activity_failed");
186 self.enqueue_activity_marker(workflow_pid, marker, correlation_id)
187 }
188
189 pub fn deliver_activity_result(
196 &self,
197 parent_pid: Pid,
198 activity_pid: Pid,
199 payload: Payload,
200 ) -> Result<(), EngineError> {
201 self.ensure_live_pid(parent_pid)?;
202 self.activity_results
203 .insert((parent_pid, activity_pid), payload);
204 let marker = self.atom_table.intern("aion_activity_result");
205 if self.scheduler.enqueue_atom_message(parent_pid, marker) {
206 self.confirm_marker_wake(parent_pid);
207 Ok(())
208 } else {
209 Err(runtime_error(format!(
210 "failed to deliver activity result from {activity_pid} to {parent_pid}"
211 )))
212 }
213 }
214
215 pub(crate) fn wake_workflow(&self, workflow_pid: Pid) -> Result<(), EngineError> {
224 self.ensure_live_pid(workflow_pid)?;
225 let marker = self.atom_table.intern("aion_timer_fired");
226 self.enqueue_signal_marker_with_retry(workflow_pid, marker)
230 }
231
232 fn enqueue_activity_marker(
233 &self,
234 workflow_pid: Pid,
235 marker: Atom,
236 correlation_id: &str,
237 ) -> Result<(), EngineError> {
238 if self.scheduler.enqueue_atom_message(workflow_pid, marker) {
239 self.confirm_marker_wake(workflow_pid);
240 tracing::debug!(
241 workflow_pid,
242 correlation_id,
243 "delivered activity completion marker to workflow mailbox via scheduler queue"
244 );
245 Ok(())
246 } else {
247 Err(runtime_error(format!(
248 "failed to deliver activity completion marker {correlation_id} to {workflow_pid}"
249 )))
250 }
251 }
252
253 pub fn deliver_activity_error(
259 &self,
260 parent_pid: Pid,
261 activity_pid: Pid,
262 error: ActivityError,
263 ) -> Result<(), EngineError> {
264 self.ensure_live_pid(parent_pid)?;
265 self.activity_errors
266 .insert((parent_pid, activity_pid), error);
267 Ok(())
268 }
269
270 #[must_use]
272 pub fn activity_result(&self, parent_pid: Pid, activity_pid: Pid) -> Option<Payload> {
273 self.activity_results
274 .get(&(parent_pid, activity_pid))
275 .map(|entry| entry.clone())
276 }
277
278 #[must_use]
280 pub fn activity_error(&self, parent_pid: Pid, activity_pid: Pid) -> Option<ActivityError> {
281 self.activity_errors
282 .get(&(parent_pid, activity_pid))
283 .map(|entry| entry.clone())
284 }
285
286 pub(crate) fn take_activity_result(
287 &self,
288 parent_pid: Pid,
289 activity_sequence: Pid,
290 ) -> Option<Payload> {
291 self.activity_results
292 .remove(&(parent_pid, activity_sequence))
293 .map(|(_, payload)| payload)
294 }
295
296 pub(crate) fn take_activity_error(
297 &self,
298 parent_pid: Pid,
299 activity_sequence: Pid,
300 ) -> Option<ActivityError> {
301 self.activity_errors
302 .remove(&(parent_pid, activity_sequence))
303 .map(|(_, error)| error)
304 }
305
306 pub(crate) fn drain_activity_completions(&self, workflow_pid: Pid) {
313 self.activity_results
314 .retain(|(parent, _), _| *parent != workflow_pid);
315 self.activity_errors
316 .retain(|(parent, _), _| *parent != workflow_pid);
317 }
318
319 #[must_use]
326 pub fn retained_activity_completions(&self) -> usize {
327 self.activity_results.len() + self.activity_errors.len()
328 }
329
330 pub(crate) fn activity_complete_atom(&self) -> Atom {
331 self.atom_table.intern("activity_complete")
332 }
333
334 pub(crate) fn activity_failed_atom(&self) -> Atom {
335 self.atom_table.intern("activity_failed")
336 }
337
338 pub(crate) fn activity_result_atom(&self) -> Atom {
339 self.atom_table.intern("aion_activity_result")
340 }
341
342 pub(crate) fn signal_received_atom(&self) -> Atom {
343 self.atom_table.intern("aion_signal_received")
344 }
345
346 pub(crate) fn timer_fired_atom(&self) -> Atom {
347 self.atom_table.intern("aion_timer_fired")
348 }
349
350 pub(crate) fn query_marker_atom(&self) -> Atom {
351 self.atom_table.intern("aion_query")
352 }
353
354 pub(crate) fn child_terminal_atom(&self) -> Atom {
355 self.atom_table.intern("aion_child_terminal")
356 }
357
358 pub(crate) fn wait_for_process_ready(&self, pid: Pid) -> Result<(), EngineError> {
359 let deadline = std::time::Instant::now() + self.signal_delivery.ready_timeout;
360 while std::time::Instant::now() < deadline {
361 if self.scheduler.trap_exit(pid).is_some() {
362 return Ok(());
363 }
364 sleep_signal_delivery_backoff(self.signal_delivery.initial_backoff);
365 }
366 self.scheduler
367 .trap_exit(pid)
368 .map(|_| ())
369 .ok_or_else(|| runtime_error(format!("process {pid} is not ready")))
370 }
371
372 pub(crate) async fn wait_for_process_ready_async(&self, pid: Pid) -> Result<(), EngineError> {
377 let deadline = std::time::Instant::now() + self.signal_delivery.ready_timeout;
378 while std::time::Instant::now() < deadline {
379 if self.scheduler.trap_exit(pid).is_some() {
380 return Ok(());
381 }
382 yield_signal_delivery_backoff(self.signal_delivery.initial_backoff).await;
383 }
384 self.scheduler
385 .trap_exit(pid)
386 .map(|_| ())
387 .ok_or_else(|| runtime_error(format!("process {pid} is not ready")))
388 }
389
390 fn enqueue_signal_marker_with_retry(
391 &self,
392 workflow_pid: Pid,
393 marker: Atom,
394 ) -> Result<(), EngineError> {
395 let attempts = self.signal_delivery.max_enqueue_attempts.max(1);
396 let mut backoff = self.signal_delivery.initial_backoff;
397 for attempt in 1..=attempts {
398 if self.scheduler.enqueue_atom_message(workflow_pid, marker) {
399 self.confirm_marker_wake(workflow_pid);
400 return Ok(());
401 }
402
403 if self.scheduler.process_table().get(workflow_pid).is_none() {
404 return Err(runtime_error(format!(
405 "failed to deliver signal to workflow process {workflow_pid}: process is not live"
406 )));
407 }
408
409 if attempt < attempts {
410 sleep_signal_delivery_backoff(backoff);
417 backoff = next_signal_delivery_backoff(backoff, self.signal_delivery.max_backoff);
418 }
419 }
420
421 Err(runtime_error(format!(
422 "failed to deliver signal to workflow process {workflow_pid} after {attempts} attempts"
423 )))
424 }
425
426 async fn enqueue_signal_marker_with_retry_async(
430 &self,
431 workflow_pid: Pid,
432 marker: Atom,
433 ) -> Result<(), EngineError> {
434 let attempts = self.signal_delivery.max_enqueue_attempts.max(1);
435 let mut backoff = self.signal_delivery.initial_backoff;
436 for attempt in 1..=attempts {
437 if self.scheduler.enqueue_atom_message(workflow_pid, marker) {
438 self.confirm_marker_wake(workflow_pid);
439 return Ok(());
440 }
441
442 if self.scheduler.process_table().get(workflow_pid).is_none() {
443 return Err(runtime_error(format!(
444 "failed to deliver signal to workflow process {workflow_pid}: process is not live"
445 )));
446 }
447
448 if attempt < attempts {
449 yield_signal_delivery_backoff(backoff).await;
451 backoff = next_signal_delivery_backoff(backoff, self.signal_delivery.max_backoff);
452 }
453 }
454
455 Err(runtime_error(format!(
456 "failed to deliver signal to workflow process {workflow_pid} after {attempts} attempts"
457 )))
458 }
459
460 fn confirm_marker_wake(&self, workflow_pid: Pid) {
472 let state = std::sync::Arc::clone(self.nif_state());
473 let snapshot = state.wake_observation_epoch(workflow_pid);
474 self.wake_confirmer
475 .confirm(self.scheduler.wake_notifier(workflow_pid), move || {
476 state.wake_ladder_done(workflow_pid, snapshot)
477 });
478 }
479}
480
481fn activity_failure(message: String) -> ActivityError {
482 ActivityError {
483 kind: ActivityErrorKind::Terminal,
484 message,
485 details: None,
486 }
487}
488
489fn correlation_to_activity_pid(correlation_id: &str) -> Result<Pid, EngineError> {
490 let Some(raw) = correlation_id.strip_prefix("activity:") else {
491 return Err(runtime_error(format!(
492 "invalid activity correlation id {correlation_id}"
493 )));
494 };
495 raw.parse::<Pid>().map_err(|error| {
496 runtime_error(format!(
497 "invalid activity correlation sequence {correlation_id}: {error}"
498 ))
499 })
500}
501
502fn next_signal_delivery_backoff(
503 current: std::time::Duration,
504 max: std::time::Duration,
505) -> std::time::Duration {
506 let doubled = current.saturating_mul(2);
507 if doubled > max { max } else { doubled }
508}
509
510fn sleep_signal_delivery_backoff(duration: std::time::Duration) {
511 if duration.is_zero() {
512 std::thread::yield_now();
513 } else {
514 std::thread::sleep(duration);
515 }
516}
517
518async fn yield_signal_delivery_backoff(duration: std::time::Duration) {
519 if duration.is_zero() {
520 tokio::task::yield_now().await;
521 } else {
522 tokio::time::sleep(duration).await;
523 }
524}