ironflow_engine/handler/typed.rs
1//! [`TypedWorkflow`] -- a handler whose input payload has a declared type.
2
3use schemars::JsonSchema;
4use serde::Serialize;
5use serde::de::DeserializeOwned;
6use serde_json::Value;
7
8use super::{WorkflowHandler, input_schema_for};
9
10/// A [`WorkflowHandler`] whose input payload has a declared type.
11///
12/// Implement it on a handler that is called as a sub-workflow:
13/// [`WorkflowContext::workflow`](crate::context::WorkflowContext::workflow)
14/// then only accepts a [`Self::Input`] for it, so a parent cannot pass a
15/// misspelled or incomplete payload. Handlers that are never called as a
16/// sub-workflow do not need it.
17///
18/// [`typed_input_schema`](Self::typed_input_schema) derives the JSON Schema of
19/// the input from the same type, for [`WorkflowHandler::input_schema`].
20///
21/// # Examples
22///
23/// ```no_run
24/// use ironflow_engine::config::ShellConfig;
25/// use ironflow_engine::context::WorkflowContext;
26/// use ironflow_engine::error::EngineError;
27/// use ironflow_engine::handler::{HandlerFuture, TypedWorkflow, WorkflowHandler};
28/// use schemars::JsonSchema;
29/// use serde::{Deserialize, Serialize};
30/// use serde_json::Value;
31///
32/// #[derive(Serialize, Deserialize, JsonSchema)]
33/// struct CollectInput {
34/// host: String,
35/// }
36///
37/// struct Collect;
38///
39/// impl WorkflowHandler for Collect {
40/// fn name(&self) -> &str { "collect" }
41/// fn input_schema(&self) -> Option<Value> { Self::typed_input_schema() }
42/// fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
43/// Box::pin(async move {
44/// let input: CollectInput = ctx.input().await?;
45/// ctx.shell("uptime", ShellConfig::new("ssh \"$HOST\" uptime").env("HOST", &input.host))
46/// .await?;
47/// Ok(())
48/// })
49/// }
50/// }
51///
52/// impl TypedWorkflow for Collect {
53/// type Input = CollectInput;
54/// }
55///
56/// # async fn parent(ctx: &mut WorkflowContext) -> Result<(), EngineError> {
57/// let child = ctx.workflow(&Collect, CollectInput { host: "db-1".to_string() }).await?;
58/// println!("collected by run {}", child.run_id());
59/// # Ok(())
60/// # }
61/// ```
62pub trait TypedWorkflow: WorkflowHandler {
63 /// The payload a run of this workflow is started with.
64 type Input: Serialize + DeserializeOwned + Send;
65
66 /// JSON Schema of [`Self::Input`], ready to return from
67 /// [`WorkflowHandler::input_schema`].
68 ///
69 /// # Examples
70 ///
71 /// ```
72 /// use ironflow_engine::context::WorkflowContext;
73 /// use ironflow_engine::handler::{HandlerFuture, TypedWorkflow, WorkflowHandler};
74 /// use schemars::JsonSchema;
75 /// use serde::{Deserialize, Serialize};
76 ///
77 /// #[derive(Serialize, Deserialize, JsonSchema)]
78 /// struct DeployInput {
79 /// environment: String,
80 /// }
81 ///
82 /// struct Deploy;
83 ///
84 /// impl WorkflowHandler for Deploy {
85 /// fn name(&self) -> &str { "deploy" }
86 /// fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
87 /// Box::pin(async move { Ok(()) })
88 /// }
89 /// }
90 ///
91 /// impl TypedWorkflow for Deploy {
92 /// type Input = DeployInput;
93 /// }
94 ///
95 /// let schema = Deploy::typed_input_schema().expect("a schema");
96 /// assert!(schema["properties"]["environment"].is_object());
97 /// ```
98 fn typed_input_schema() -> Option<Value>
99 where
100 Self: Sized,
101 Self::Input: JsonSchema,
102 {
103 Some(input_schema_for::<Self::Input>())
104 }
105}
106
107/// Names of the given handlers, for [`WorkflowHandler::sub_workflows`].
108///
109/// Pass the handlers themselves rather than their names: a misspelled handler
110/// does not compile, a misspelled string does.
111///
112/// # Examples
113///
114/// ```
115/// use ironflow_engine::context::WorkflowContext;
116/// use ironflow_engine::handler::{HandlerFuture, WorkflowHandler, sub_workflow_names};
117///
118/// struct Collect;
119///
120/// impl WorkflowHandler for Collect {
121/// fn name(&self) -> &str { "collect" }
122/// fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
123/// Box::pin(async move { Ok(()) })
124/// }
125/// }
126///
127/// struct Enrich;
128///
129/// impl WorkflowHandler for Enrich {
130/// fn name(&self) -> &str { "enrich" }
131/// fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
132/// Box::pin(async move { Ok(()) })
133/// }
134/// }
135///
136/// assert_eq!(sub_workflow_names(&[&Collect, &Enrich]), vec!["collect", "enrich"]);
137/// assert!(sub_workflow_names(&[]).is_empty());
138/// ```
139pub fn sub_workflow_names(handlers: &[&dyn WorkflowHandler]) -> Vec<String> {
140 handlers.iter().map(|h| h.name().to_string()).collect()
141}
142
143#[cfg(test)]
144mod tests {
145 use serde::Deserialize;
146
147 use super::*;
148 use crate::context::WorkflowContext;
149 use crate::handler::HandlerFuture;
150
151 #[derive(Serialize, Deserialize, JsonSchema)]
152 struct ProbeInput {
153 count: u32,
154 label: Option<String>,
155 }
156
157 struct Probe;
158
159 impl WorkflowHandler for Probe {
160 fn name(&self) -> &str {
161 "probe"
162 }
163
164 fn input_schema(&self) -> Option<Value> {
165 Self::typed_input_schema()
166 }
167
168 fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
169 Box::pin(async move { Ok(()) })
170 }
171 }
172
173 impl TypedWorkflow for Probe {
174 type Input = ProbeInput;
175 }
176
177 struct Other;
178
179 impl WorkflowHandler for Other {
180 fn name(&self) -> &str {
181 "other-été"
182 }
183
184 fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
185 Box::pin(async move { Ok(()) })
186 }
187 }
188
189 #[test]
190 fn typed_input_schema_describes_the_input_type() {
191 let schema = Probe::typed_input_schema().expect("a schema");
192 assert_eq!(schema, input_schema_for::<ProbeInput>());
193 assert!(schema["properties"]["count"].is_object());
194 assert!(schema["properties"]["label"].is_object());
195 }
196
197 #[test]
198 fn input_schema_can_return_the_typed_schema() {
199 assert_eq!(Probe.input_schema(), Probe::typed_input_schema());
200 assert_eq!(Probe.describe().input_schema, Probe::typed_input_schema());
201 }
202
203 #[test]
204 fn sub_workflow_names_keeps_the_order_of_the_handlers() {
205 assert_eq!(
206 sub_workflow_names(&[&Other, &Probe]),
207 vec!["other-été".to_string(), "probe".to_string()]
208 );
209 }
210
211 #[test]
212 fn sub_workflow_names_of_nothing_is_empty() {
213 assert!(sub_workflow_names(&[]).is_empty());
214 }
215}