Skip to main content

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}