temporalio_workflow/runtime/
entry.rs1use crate::{
4 SyncWorkflowContext, WorkflowContext, WorkflowContextView,
5 runtime::{
6 mark_intercepted_handler_ready, model::WorkflowTermination,
7 types::WorkflowDefinitionDescriptor,
8 },
9 workflow_interceptors::WorkflowOutputValue,
10};
11use futures_util::future::{FutureExt, LocalBoxFuture};
12use std::any::Any;
13use temporalio_common_wasm::{
14 QueryDefinition, SignalDefinition, UpdateDefinition, WorkflowDefinition,
15 data_converters::{
16 GenericPayloadConverter, PayloadConversionError, PayloadConverter, SerializationContext,
17 SerializationContextData, TemporalSerializable,
18 },
19 protos::temporal::api::{
20 common::v1::{Payload, Payloads},
21 failure::v1::Failure,
22 },
23};
24
25#[derive(Debug, thiserror::Error)]
27pub enum WorkflowError {
28 #[error("Payload conversion error: {0}")]
30 PayloadConversion(#[from] PayloadConversionError),
31
32 #[error("Workflow execution error: {0}")]
34 Execution(#[from] Box<dyn std::error::Error + Send + Sync>),
35}
36
37impl From<WorkflowError> for Failure {
38 fn from(err: WorkflowError) -> Self {
39 Failure {
40 message: err.to_string(),
41 ..Default::default()
42 }
43 }
44}
45
46fn downcast_handler_input<T: Any>(input: Box<dyn Any>, handler_kind: &'static str) -> T {
47 *input.downcast::<T>().unwrap_or_else(|_| {
48 panic!("typed {handler_kind} dispatch received input with wrong concrete type")
49 })
50}
51
52pub trait WorkflowImplementation: Sized + 'static {
57 type Run: WorkflowDefinition;
59
60 const HAS_INIT: bool;
63
64 const INIT_TAKES_INPUT: bool;
67
68 fn name() -> &'static str;
70
71 fn definition() -> WorkflowDefinitionDescriptor;
73
74 fn init(
76 ctx: WorkflowContextView,
77 input: Option<<Self::Run as WorkflowDefinition>::Input>,
78 ) -> Self;
79
80 fn run(
82 ctx: WorkflowContext<Self>,
83 input: Option<<Self::Run as WorkflowDefinition>::Input>,
84 ) -> LocalBoxFuture<'static, Result<Box<dyn WorkflowOutputValue>, WorkflowTermination>>;
85
86 fn decode_signal_input(
88 _name: &str,
89 _payloads: Payloads,
90 _converter: &PayloadConverter,
91 ) -> Result<Option<Box<dyn Any>>, WorkflowError> {
92 Ok(None)
93 }
94
95 fn dispatch_signal(
97 ctx: WorkflowContext<Self>,
98 name: &str,
99 input: Box<dyn Any>,
100 ) -> LocalBoxFuture<'static, Result<(), WorkflowError>>;
101
102 fn decode_query_input(
104 _name: &str,
105 _payloads: &Payloads,
106 _converter: &PayloadConverter,
107 ) -> Result<Option<Box<dyn Any>>, WorkflowError> {
108 Ok(None)
109 }
110
111 fn dispatch_query(
113 &self,
114 ctx: WorkflowContextView,
115 name: &str,
116 input: Box<dyn Any>,
117 ) -> Result<Box<dyn WorkflowOutputValue>, WorkflowError>;
118
119 fn decode_update_input(
121 _name: &str,
122 _payloads: Payloads,
123 _converter: &PayloadConverter,
124 ) -> Result<Option<Box<dyn Any>>, WorkflowError> {
125 Ok(None)
126 }
127
128 fn dispatch_update(
130 ctx: WorkflowContext<Self>,
131 name: &str,
132 input: Box<dyn Any>,
133 ) -> LocalBoxFuture<'static, Result<Box<dyn WorkflowOutputValue>, WorkflowError>>;
134
135 fn validate_update(
137 &self,
138 ctx: WorkflowContextView,
139 name: &str,
140 input: Box<dyn Any>,
141 ) -> Result<(), WorkflowError>;
142}
143
144pub trait ExecutableSyncSignal<S: SignalDefinition>: WorkflowImplementation {
146 fn handle(&mut self, ctx: &mut SyncWorkflowContext<Self>, input: S::Input);
148
149 fn dispatch(
151 ctx: WorkflowContext<Self>,
152 input: Box<dyn Any>,
153 ) -> LocalBoxFuture<'static, Result<(), WorkflowError>> {
154 let input = downcast_handler_input::<S::Input>(input, "signal");
155 let mut sync_ctx = ctx.sync_context();
156 ctx.state_mut(|wf| Self::handle(wf, &mut sync_ctx, input));
157 mark_intercepted_handler_ready();
158 std::future::ready(Ok(())).boxed_local()
159 }
160}
161
162pub trait ExecutableAsyncSignal<S: SignalDefinition>: WorkflowImplementation {
164 fn handle(ctx: WorkflowContext<Self>, input: S::Input) -> LocalBoxFuture<'static, ()>;
166
167 fn dispatch(
169 ctx: WorkflowContext<Self>,
170 input: Box<dyn Any>,
171 ) -> LocalBoxFuture<'static, Result<(), WorkflowError>> {
172 let input = downcast_handler_input::<S::Input>(input, "signal");
173 Self::handle(ctx, input).map(|()| Ok(())).boxed_local()
174 }
175}
176
177pub trait ExecutableQuery<Q: QueryDefinition>: WorkflowImplementation {
179 fn handle(
181 &self,
182 ctx: &WorkflowContextView,
183 input: Q::Input,
184 ) -> Result<Q::Output, Box<dyn std::error::Error + Send + Sync>>;
185
186 fn dispatch(
188 &self,
189 ctx: &WorkflowContextView,
190 input: Box<dyn Any>,
191 ) -> Result<Box<dyn WorkflowOutputValue>, WorkflowError> {
192 let input = downcast_handler_input::<Q::Input>(input, "query");
193 let output = self.handle(ctx, input).map_err(WorkflowError::Execution)?;
194 Ok(Box::new(output))
195 }
196}
197
198pub trait ExecutableSyncUpdate<U: UpdateDefinition>: WorkflowImplementation {
200 fn handle(
202 &mut self,
203 ctx: &mut SyncWorkflowContext<Self>,
204 input: U::Input,
205 ) -> Result<U::Output, Box<dyn std::error::Error + Send + Sync>>;
206
207 fn validate(
209 &self,
210 _ctx: &WorkflowContextView,
211 _input: &U::Input,
212 ) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
213 Ok(())
214 }
215
216 fn dispatch(
218 ctx: WorkflowContext<Self>,
219 input: Box<dyn Any>,
220 ) -> LocalBoxFuture<'static, Result<Box<dyn WorkflowOutputValue>, WorkflowError>> {
221 let input = downcast_handler_input::<U::Input>(input, "update");
222 let mut sync_ctx = ctx.sync_context();
223 let result = ctx.state_mut(|wf| Self::handle(wf, &mut sync_ctx, input));
224 mark_intercepted_handler_ready();
225 match result {
226 Ok(output) => std::future::ready(Ok(Box::new(output) as Box<dyn WorkflowOutputValue>))
227 .boxed_local(),
228 Err(e) => std::future::ready(Err(WorkflowError::Execution(e))).boxed_local(),
229 }
230 }
231
232 fn dispatch_validate(
234 &self,
235 ctx: &WorkflowContextView,
236 input: Box<dyn Any>,
237 ) -> Result<(), WorkflowError> {
238 let input = downcast_handler_input::<U::Input>(input, "update validation");
239 self.validate(ctx, &input).map_err(WorkflowError::Execution)
240 }
241}
242
243pub trait ExecutableAsyncUpdate<U: UpdateDefinition>: WorkflowImplementation {
245 fn handle(
247 ctx: WorkflowContext<Self>,
248 input: U::Input,
249 ) -> LocalBoxFuture<'static, Result<U::Output, Box<dyn std::error::Error + Send + Sync>>>;
250
251 fn validate(
253 &self,
254 _ctx: &WorkflowContextView,
255 _input: &U::Input,
256 ) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
257 Ok(())
258 }
259
260 fn dispatch(
262 ctx: WorkflowContext<Self>,
263 input: Box<dyn Any>,
264 ) -> LocalBoxFuture<'static, Result<Box<dyn WorkflowOutputValue>, WorkflowError>> {
265 let input = downcast_handler_input::<U::Input>(input, "update");
266 let future = async move {
267 let output = Self::handle(ctx, input)
268 .await
269 .map_err(WorkflowError::Execution)?;
270 Ok(Box::new(output) as Box<dyn WorkflowOutputValue>)
271 };
272 future.boxed_local()
273 }
274
275 fn dispatch_validate(
277 &self,
278 ctx: &WorkflowContextView,
279 input: Box<dyn Any>,
280 ) -> Result<(), WorkflowError> {
281 let input = downcast_handler_input::<U::Input>(input, "update validation");
282 self.validate(ctx, &input).map_err(WorkflowError::Execution)
283 }
284}
285
286pub(crate) fn serialize_output<O: TemporalSerializable + 'static>(
288 output: &O,
289 converter: &PayloadConverter,
290) -> Result<Payload, WorkflowError> {
291 let ctx = SerializationContext {
292 data: &SerializationContextData::Workflow,
293 converter,
294 };
295 converter.to_payload(&ctx, output).map_err(Into::into)
296}
297
298pub fn serialize_result<T: TemporalSerializable + 'static>(
300 result: T,
301 converter: &PayloadConverter,
302) -> Result<Payload, WorkflowError> {
303 serialize_output(&result, converter)
304}