aion/activity/
dispatch.rs1use aion_core::{ActivityError, Payload};
22
23use crate::{EngineError, Pid, RuntimeHandle, RuntimeInput};
24
25pub fn dispatch_activity(
36 runtime: &RuntimeHandle,
37 parent_pid: Pid,
38 deployed_module: &str,
39 function: &str,
40 input: &Payload,
41) -> Result<Pid, EngineError> {
42 runtime.spawn_activity(
43 parent_pid,
44 deployed_module,
45 function,
46 RuntimeInput::from_payload(input)?,
47 )
48}
49
50pub fn propagate_activity_outcome(
62 runtime: &RuntimeHandle,
63 parent_pid: Pid,
64 activity_pid: Pid,
65) -> Result<(), EngineError> {
66 runtime.propagate_activity_outcome(parent_pid, activity_pid)
67}
68
69pub fn surface_activity_error(
78 runtime: &RuntimeHandle,
79 parent_pid: Pid,
80 activity_pid: Pid,
81 error: ActivityError,
82) -> Result<(), EngineError> {
83 runtime.deliver_activity_error(parent_pid, activity_pid, error)
84}
85
86#[cfg(test)]
87mod tests {
88 use aion_core::{ActivityErrorKind, ContentType};
89 use serde_json::json;
90
91 use super::{dispatch_activity, propagate_activity_outcome, surface_activity_error};
92 use crate::runtime::RuntimeConfig;
93 use crate::{EngineError, RuntimeHandle};
94
95 fn runtime() -> Result<RuntimeHandle, EngineError> {
96 RuntimeHandle::new(RuntimeConfig::new(Some(1)))
97 }
98
99 fn payload() -> Result<aion_core::Payload, aion_core::PayloadError> {
100 aion_core::Payload::from_json(&json!(null))
101 }
102
103 fn fixture_workflow_beam() -> &'static [u8] {
104 include_bytes!("../../tests/fixtures/aion_fixture_workflow.beam")
105 }
106
107 #[test]
108 fn dispatch_spawns_linked_child_and_uses_dirty_registration()
109 -> Result<(), Box<dyn std::error::Error>> {
110 let runtime = runtime()?;
111 runtime.install_test_activity_nif("activity_host", "answer", true, true)?;
112 runtime.register_native_call_module_for_test(
113 "activity_mod",
114 "run",
115 "activity_host",
116 "answer",
117 );
118 let workflow = runtime.spawn_test_process_with_trap_exit(true)?;
119
120 let activity = dispatch_activity(&runtime, workflow, "activity_mod", "run", &payload()?)?;
121
122 if let Ok(linked) = runtime.is_linked(workflow, activity) {
123 assert!(linked);
124 }
125 assert!(runtime.is_dirty("activity_host", "answer"));
126 runtime.shutdown()?;
127 Ok(())
128 }
129
130 #[test]
131 fn successful_activity_result_is_delivered_to_workflow()
132 -> Result<(), Box<dyn std::error::Error>> {
133 let runtime = runtime()?;
134 runtime.install_test_activity_nif("activity_host", "answer", false, true)?;
135 runtime.register_native_call_module_for_test(
136 "activity_ok",
137 "run",
138 "activity_host",
139 "answer",
140 );
141 let workflow = runtime.spawn_test_process_with_trap_exit(true)?;
142 let activity = dispatch_activity(&runtime, workflow, "activity_ok", "run", &payload()?)?;
143
144 propagate_activity_outcome(&runtime, workflow, activity)?;
145
146 let result = runtime.activity_result(workflow, activity);
147 assert_eq!(result, Some(aion_core::Payload::from_json(&json!(42))?));
148 runtime.shutdown()?;
149 Ok(())
150 }
151
152 #[test]
153 fn failing_activity_surfaces_typed_error_with_trapped_exit()
154 -> Result<(), Box<dyn std::error::Error>> {
155 let runtime = runtime()?;
156 runtime.register_module("aion_fixture_workflow", fixture_workflow_beam())?;
157 let workflow = runtime.spawn_test_process_with_trap_exit(true)?;
158 let activity = dispatch_activity(
159 &runtime,
160 workflow,
161 "aion_fixture_workflow",
162 "activity",
163 &payload()?,
164 )?;
165 assert!(runtime.is_linked(workflow, activity)?);
166 let details = aion_core::Payload::new(ContentType::Json, br#"{"code":"boom"}"#.to_vec());
167 let error = aion_core::ActivityError {
168 kind: ActivityErrorKind::Retryable,
169 message: String::from("boom"),
170 details: Some(details),
171 };
172 surface_activity_error(&runtime, workflow, activity, error.clone())?;
173 runtime.wait_for_process_ready(workflow)?;
174 runtime.wait_for_process_ready(activity)?;
175 runtime.terminate_test_process_with_error(activity)?;
176
177 propagate_activity_outcome(&runtime, workflow, activity)?;
178 runtime.wait_for_trapped_exit(workflow, activity)?;
179
180 assert!(runtime.has_trapped_exit_message(workflow, activity)?);
181 assert_eq!(runtime.activity_error(workflow, activity), Some(error));
182 runtime.shutdown()?;
183 Ok(())
184 }
185}