Skip to main content

ironflow_engine/
run_creator.rs

1//! [`RunCreator`] trait and [`CreateRunOpts`] builder -- centralised run creation.
2//!
3//! [`RunCreator`] is a minimal trait with a single method: create a run from
4//! a [`NewRun`]. It is intentionally thinner than [`RunStore`] so that any
5//! store implementation can be used as a run creator through the blanket impl.
6//!
7//! [`CreateRunOpts`] is a builder for the optional fields of a run creation
8//! request. Combined with [`WorkflowHandler::create_run`](crate::handler::WorkflowHandler::create_run), it assembles a
9//! [`NewRun`] from the handler's own metadata, removing duplication across
10//! call sites.
11//!
12//! # Examples
13//!
14//! ```no_run
15//! use ironflow_engine::run_creator::{CreateRunOpts, RunCreator};
16//! use ironflow_store::entities::TriggerKind;
17//! use ironflow_store::memory::InMemoryStore;
18//!
19//! # async fn example() -> Result<(), ironflow_engine::error::EngineError> {
20//! let store = InMemoryStore::new();
21//! let creator: &dyn RunCreator = &store;
22//!
23//! let new_run = CreateRunOpts::new()
24//!     .trigger(TriggerKind::Api)
25//!     .build("deploy", Some("1.0.0"), None);
26//!
27//! let creation = creator.create_run(new_run).await?;
28//! # Ok(())
29//! # }
30//! ```
31
32use std::collections::HashMap;
33use std::future::Future;
34use std::pin::Pin;
35
36use chrono::{DateTime, Utc};
37use rust_decimal::Decimal;
38use serde_json::Value;
39
40use ironflow_store::entities::{NewRun, RunActor, RunCreation, TriggerKind};
41use ironflow_store::store::RunStore;
42
43use crate::error::EngineError;
44
45/// Future returned by [`RunCreator::create_run`].
46pub type RunCreatorFuture<'a> =
47    Pin<Box<dyn Future<Output = Result<RunCreation, EngineError>> + Send + 'a>>;
48
49/// Minimal trait for creating workflow runs.
50///
51/// Automatically implemented for every [`RunStore`] via a blanket impl,
52/// so any store (InMemory, Postgres, ApiRunStore) is a valid [`RunCreator`].
53///
54/// # Examples
55///
56/// ```no_run
57/// use ironflow_engine::run_creator::RunCreator;
58/// use ironflow_store::entities::{NewRun, TriggerKind};
59///
60/// # async fn example(creator: &dyn RunCreator) -> Result<(), ironflow_engine::error::EngineError> {
61/// let new_run = NewRun {
62///     workflow_name: "deploy".to_string(),
63///     trigger: TriggerKind::Manual,
64///     payload: serde_json::json!({}),
65///     max_retries: 0,
66///     handler_version: None,
67///     labels: Default::default(),
68///     scheduled_at: None,
69///     created_by: None,
70///     idempotency_key: None,
71///     concurrency_key: None,
72///     max_cost_usd: None,
73/// };
74/// let creation = creator.create_run(new_run).await?;
75/// # Ok(())
76/// # }
77/// ```
78pub trait RunCreator: Send + Sync {
79    /// Create a new workflow run.
80    ///
81    /// # Errors
82    ///
83    /// Returns [`EngineError`] if the run could not be created.
84    fn create_run(&self, req: NewRun) -> RunCreatorFuture<'_>;
85}
86
87impl<T: RunStore + ?Sized> RunCreator for T {
88    fn create_run(&self, req: NewRun) -> RunCreatorFuture<'_> {
89        Box::pin(async move {
90            RunStore::create_run(self, req)
91                .await
92                .map_err(EngineError::from)
93        })
94    }
95}
96
97/// Builder for optional run creation fields.
98///
99/// Builds a [`NewRun`] from handler metadata plus user-supplied overrides.
100/// Use with [`WorkflowHandler::create_run`] to avoid duplicating workflow
101/// name, version, and cost cap at every call site.
102///
103/// [`WorkflowHandler::create_run`]: crate::handler::WorkflowHandler::create_run
104///
105/// # Examples
106///
107/// ```
108/// use ironflow_engine::run_creator::CreateRunOpts;
109/// use ironflow_store::entities::TriggerKind;
110/// use serde_json::json;
111///
112/// let new_run = CreateRunOpts::new()
113///     .trigger(TriggerKind::Webhook { path: "/hooks/gh".to_string() })
114///     .payload(json!({"ref": "main"}))
115///     .max_retries(3)
116///     .build("deploy", Some("2.0.0"), None);
117///
118/// assert_eq!(new_run.workflow_name, "deploy");
119/// assert_eq!(new_run.max_retries, 3);
120/// assert_eq!(new_run.handler_version, Some("2.0.0".to_string()));
121/// ```
122#[derive(Debug, Clone, Default)]
123pub struct CreateRunOpts {
124    trigger: Option<TriggerKind>,
125    payload: Option<Value>,
126    max_retries: Option<u32>,
127    scheduled_at: Option<DateTime<Utc>>,
128    created_by: Option<RunActor>,
129    idempotency_key: Option<String>,
130    concurrency_key: Option<String>,
131    labels: Option<HashMap<String, String>>,
132    max_cost_usd: Option<Decimal>,
133}
134
135impl CreateRunOpts {
136    /// Create a new builder with all fields unset.
137    ///
138    /// # Examples
139    ///
140    /// ```
141    /// use ironflow_engine::run_creator::CreateRunOpts;
142    ///
143    /// let opts = CreateRunOpts::new();
144    /// let new_run = opts.build("my-workflow", None, None);
145    /// assert_eq!(new_run.workflow_name, "my-workflow");
146    /// ```
147    pub fn new() -> Self {
148        Self::default()
149    }
150
151    /// Set how the run was triggered.
152    ///
153    /// Defaults to [`TriggerKind::Manual`] if not set.
154    pub fn trigger(mut self, trigger: TriggerKind) -> Self {
155        self.trigger = Some(trigger);
156        self
157    }
158
159    /// Set the trigger-specific payload.
160    ///
161    /// Defaults to `json!({})` if not set.
162    pub fn payload(mut self, payload: Value) -> Self {
163        self.payload = Some(payload);
164        self
165    }
166
167    /// Set the maximum retry attempts.
168    ///
169    /// Defaults to `0` if not set.
170    pub fn max_retries(mut self, max_retries: u32) -> Self {
171        self.max_retries = Some(max_retries);
172        self
173    }
174
175    /// Schedule the run for later execution.
176    pub fn scheduled_at(mut self, at: DateTime<Utc>) -> Self {
177        self.scheduled_at = Some(at);
178        self
179    }
180
181    /// Set the authenticated principal creating this run.
182    pub fn created_by(mut self, actor: RunActor) -> Self {
183        self.created_by = Some(actor);
184        self
185    }
186
187    /// Set an idempotency key to prevent duplicate runs.
188    pub fn idempotency_key(mut self, key: impl Into<String>) -> Self {
189        self.idempotency_key = Some(key.into());
190        self
191    }
192
193    /// Set a concurrency key: the store refuses the run while another
194    /// non-terminal run holds the same key.
195    ///
196    /// # Examples
197    ///
198    /// ```
199    /// use ironflow_engine::run_creator::CreateRunOpts;
200    ///
201    /// let new_run = CreateRunOpts::new()
202    ///     .concurrency_key("issue:12")
203    ///     .build("deploy", None, None);
204    /// assert_eq!(new_run.concurrency_key.as_deref(), Some("issue:12"));
205    /// ```
206    pub fn concurrency_key(mut self, key: impl Into<String>) -> Self {
207        self.concurrency_key = Some(key.into());
208        self
209    }
210
211    /// Set user-defined labels for categorization.
212    pub fn labels(mut self, labels: HashMap<String, String>) -> Self {
213        self.labels = Some(labels);
214        self
215    }
216
217    /// Set the maximum cumulative cost allowed for this run.
218    pub fn max_cost_usd(mut self, cap: Decimal) -> Self {
219        self.max_cost_usd = Some(cap);
220        self
221    }
222
223    /// Assemble a [`NewRun`] from these options and handler metadata.
224    ///
225    /// * `workflow_name` -- typically from [`WorkflowHandler::name`].
226    /// * `handler_version` -- typically from [`WorkflowHandler::version`].
227    /// * `default_max_cost_usd` -- typically from [`WorkflowHandler::default_max_cost_usd`].
228    ///   Applied only when [`max_cost_usd`](Self::max_cost_usd) was not set.
229    ///
230    /// [`WorkflowHandler::name`]: crate::handler::WorkflowHandler::name
231    /// [`WorkflowHandler::version`]: crate::handler::WorkflowHandler::version
232    /// [`WorkflowHandler::default_max_cost_usd`]: crate::handler::WorkflowHandler::default_max_cost_usd
233    ///
234    /// # Examples
235    ///
236    /// ```
237    /// use ironflow_engine::run_creator::CreateRunOpts;
238    /// use rust_decimal::Decimal;
239    ///
240    /// let new_run = CreateRunOpts::new()
241    ///     .build("my-handler", Some("3.0.0"), Some(Decimal::new(1000, 2)));
242    ///
243    /// assert_eq!(new_run.workflow_name, "my-handler");
244    /// assert_eq!(new_run.handler_version, Some("3.0.0".to_string()));
245    /// assert_eq!(new_run.max_cost_usd, Some(Decimal::new(1000, 2)));
246    /// ```
247    pub fn build(
248        self,
249        workflow_name: &str,
250        handler_version: Option<&str>,
251        default_max_cost_usd: Option<Decimal>,
252    ) -> NewRun {
253        NewRun {
254            workflow_name: workflow_name.to_string(),
255            trigger: self.trigger.unwrap_or(TriggerKind::Manual),
256            payload: self.payload.unwrap_or_else(|| serde_json::json!({})),
257            max_retries: self.max_retries.unwrap_or(0),
258            handler_version: handler_version.map(str::to_string),
259            labels: self.labels.unwrap_or_default(),
260            scheduled_at: self.scheduled_at,
261            created_by: self.created_by,
262            idempotency_key: self.idempotency_key,
263            concurrency_key: self.concurrency_key,
264            max_cost_usd: self.max_cost_usd.or(default_max_cost_usd),
265        }
266    }
267}
268
269#[cfg(test)]
270mod tests {
271    use super::*;
272    use serde_json::json;
273
274    #[test]
275    fn create_run_opts_default_produces_correct_defaults() {
276        let opts = CreateRunOpts::new();
277        let new_run = opts.build("test-workflow", None, None);
278
279        assert_eq!(new_run.workflow_name, "test-workflow");
280        assert_eq!(new_run.trigger, TriggerKind::Manual);
281        assert_eq!(new_run.payload, json!({}));
282        assert_eq!(new_run.max_retries, 0);
283        assert_eq!(new_run.handler_version, None);
284        assert!(new_run.labels.is_empty());
285        assert_eq!(new_run.scheduled_at, None);
286        assert_eq!(new_run.created_by, None);
287        assert_eq!(new_run.idempotency_key, None);
288        assert_eq!(new_run.concurrency_key, None);
289        assert_eq!(new_run.max_cost_usd, None);
290    }
291
292    #[test]
293    fn create_run_opts_builder_sets_all_fields() {
294        let labels = HashMap::from([("env".to_string(), "prod".to_string())]);
295        let scheduled = Utc::now();
296
297        let new_run = CreateRunOpts::new()
298            .trigger(TriggerKind::Webhook {
299                path: "/hooks/gh".to_string(),
300            })
301            .payload(json!({"ref": "main"}))
302            .max_retries(3)
303            .scheduled_at(scheduled)
304            .idempotency_key("key-123")
305            .labels(labels.clone())
306            .max_cost_usd(Decimal::new(500, 2))
307            .build("deploy", Some("2.0.0"), None);
308
309        assert_eq!(new_run.workflow_name, "deploy");
310        assert_eq!(
311            new_run.trigger,
312            TriggerKind::Webhook {
313                path: "/hooks/gh".to_string()
314            }
315        );
316        assert_eq!(new_run.payload, json!({"ref": "main"}));
317        assert_eq!(new_run.max_retries, 3);
318        assert_eq!(new_run.handler_version, Some("2.0.0".to_string()));
319        assert_eq!(new_run.scheduled_at, Some(scheduled));
320        assert_eq!(new_run.idempotency_key, Some("key-123".to_string()));
321        assert_eq!(new_run.labels, labels);
322        assert_eq!(new_run.max_cost_usd, Some(Decimal::new(500, 2)));
323    }
324
325    #[test]
326    fn create_run_opts_build_carries_the_concurrency_key() {
327        let new_run = CreateRunOpts::new()
328            .concurrency_key("issue:12")
329            .build("deploy", None, None);
330
331        assert_eq!(new_run.concurrency_key.as_deref(), Some("issue:12"));
332        assert_eq!(new_run.idempotency_key, None);
333    }
334
335    #[test]
336    fn create_run_opts_build_uses_handler_metadata() {
337        let new_run =
338            CreateRunOpts::new().build("my-handler", Some("3.0.0"), Some(Decimal::new(1000, 2)));
339
340        assert_eq!(new_run.workflow_name, "my-handler");
341        assert_eq!(new_run.handler_version, Some("3.0.0".to_string()));
342        assert_eq!(new_run.max_cost_usd, Some(Decimal::new(1000, 2)));
343    }
344
345    #[test]
346    fn create_run_opts_explicit_max_cost_overrides_handler_default() {
347        let new_run = CreateRunOpts::new()
348            .max_cost_usd(Decimal::new(200, 2))
349            .build("handler", Some("1"), Some(Decimal::new(1000, 2)));
350
351        assert_eq!(new_run.max_cost_usd, Some(Decimal::new(200, 2)));
352    }
353
354    #[tokio::test]
355    async fn run_creator_blanket_impl_with_in_memory_store() {
356        use ironflow_store::memory::InMemoryStore;
357
358        let store = InMemoryStore::new();
359        let creator: &dyn RunCreator = &store;
360
361        let new_run =
362            CreateRunOpts::new()
363                .trigger(TriggerKind::Api)
364                .build("blanket-test", None, None);
365
366        let creation = creator.create_run(new_run).await.expect("create_run");
367        let run = creation.into_run();
368        assert_eq!(run.workflow_name, "blanket-test");
369    }
370
371    #[tokio::test]
372    async fn create_run_with_reused_idempotency_key_returns_existing() {
373        use ironflow_store::memory::InMemoryStore;
374
375        let store = InMemoryStore::new();
376        let creator: &dyn RunCreator = &store;
377
378        let first = creator
379            .create_run(CreateRunOpts::new().idempotency_key("dedup-1").build(
380                "idem-test",
381                None,
382                None,
383            ))
384            .await
385            .expect("first create_run");
386        assert!(first.is_created());
387
388        let second = creator
389            .create_run(CreateRunOpts::new().idempotency_key("dedup-1").build(
390                "idem-test",
391                None,
392                None,
393            ))
394            .await
395            .expect("second create_run");
396        assert!(!second.is_created());
397        assert_eq!(first.into_run().id, second.into_run().id);
398    }
399}