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, normalize_worker_tags};
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///     priority: 0,
73///     concurrency_limits: Vec::new(),
74///     max_cost_usd: None,
75///     worker_tags: Vec::new(),
76/// };
77/// let creation = creator.create_run(new_run).await?;
78/// # Ok(())
79/// # }
80/// ```
81pub trait RunCreator: Send + Sync {
82    /// Create a new workflow run.
83    ///
84    /// # Errors
85    ///
86    /// Returns [`EngineError`] if the run could not be created.
87    fn create_run(&self, req: NewRun) -> RunCreatorFuture<'_>;
88}
89
90impl<T: RunStore + ?Sized> RunCreator for T {
91    fn create_run(&self, req: NewRun) -> RunCreatorFuture<'_> {
92        Box::pin(async move {
93            RunStore::create_run(self, req)
94                .await
95                .map_err(EngineError::from)
96        })
97    }
98}
99
100/// Builder for optional run creation fields.
101///
102/// Builds a [`NewRun`] from handler metadata plus user-supplied overrides.
103/// Use with [`WorkflowHandler::create_run`] to avoid duplicating workflow
104/// name, version, and cost cap at every call site.
105///
106/// [`WorkflowHandler::create_run`]: crate::handler::WorkflowHandler::create_run
107///
108/// # Examples
109///
110/// ```
111/// use ironflow_engine::run_creator::CreateRunOpts;
112/// use ironflow_store::entities::TriggerKind;
113/// use serde_json::json;
114///
115/// let new_run = CreateRunOpts::new()
116///     .trigger(TriggerKind::Webhook { path: "/hooks/gh".to_string() })
117///     .payload(json!({"ref": "main"}))
118///     .max_retries(3)
119///     .build("deploy", Some("2.0.0"), None);
120///
121/// assert_eq!(new_run.workflow_name, "deploy");
122/// assert_eq!(new_run.max_retries, 3);
123/// assert_eq!(new_run.handler_version, Some("2.0.0".to_string()));
124/// ```
125#[derive(Debug, Clone, Default)]
126pub struct CreateRunOpts {
127    trigger: Option<TriggerKind>,
128    payload: Option<Value>,
129    max_retries: Option<u32>,
130    scheduled_at: Option<DateTime<Utc>>,
131    created_by: Option<RunActor>,
132    idempotency_key: Option<String>,
133    concurrency_key: Option<String>,
134    labels: Option<HashMap<String, String>>,
135    max_cost_usd: Option<Decimal>,
136    priority: Option<i16>,
137    worker_tags: Vec<String>,
138}
139
140impl CreateRunOpts {
141    /// Create a new builder with all fields unset.
142    ///
143    /// # Examples
144    ///
145    /// ```
146    /// use ironflow_engine::run_creator::CreateRunOpts;
147    ///
148    /// let opts = CreateRunOpts::new();
149    /// let new_run = opts.build("my-workflow", None, None);
150    /// assert_eq!(new_run.workflow_name, "my-workflow");
151    /// ```
152    pub fn new() -> Self {
153        Self::default()
154    }
155
156    /// Set how the run was triggered.
157    ///
158    /// Defaults to [`TriggerKind::Manual`] if not set.
159    pub fn trigger(mut self, trigger: TriggerKind) -> Self {
160        self.trigger = Some(trigger);
161        self
162    }
163
164    /// Set the trigger-specific payload.
165    ///
166    /// Defaults to `json!({})` if not set.
167    pub fn payload(mut self, payload: Value) -> Self {
168        self.payload = Some(payload);
169        self
170    }
171
172    /// Set the maximum retry attempts.
173    ///
174    /// Defaults to `0` if not set.
175    pub fn max_retries(mut self, max_retries: u32) -> Self {
176        self.max_retries = Some(max_retries);
177        self
178    }
179
180    /// Schedule the run for later execution.
181    pub fn scheduled_at(mut self, at: DateTime<Utc>) -> Self {
182        self.scheduled_at = Some(at);
183        self
184    }
185
186    /// Set the authenticated principal creating this run.
187    pub fn created_by(mut self, actor: RunActor) -> Self {
188        self.created_by = Some(actor);
189        self
190    }
191
192    /// Set an idempotency key to prevent duplicate runs.
193    pub fn idempotency_key(mut self, key: impl Into<String>) -> Self {
194        self.idempotency_key = Some(key.into());
195        self
196    }
197
198    /// Set a concurrency key: the store refuses the run while another
199    /// non-terminal run holds the same key.
200    ///
201    /// # Examples
202    ///
203    /// ```
204    /// use ironflow_engine::run_creator::CreateRunOpts;
205    ///
206    /// let new_run = CreateRunOpts::new()
207    ///     .concurrency_key("issue:12")
208    ///     .build("deploy", None, None);
209    /// assert_eq!(new_run.concurrency_key.as_deref(), Some("issue:12"));
210    /// ```
211    pub fn concurrency_key(mut self, key: impl Into<String>) -> Self {
212        self.concurrency_key = Some(key.into());
213        self
214    }
215
216    /// Set user-defined labels for categorization.
217    pub fn labels(mut self, labels: HashMap<String, String>) -> Self {
218        self.labels = Some(labels);
219        self
220    }
221
222    /// Set the maximum cumulative cost allowed for this run.
223    pub fn max_cost_usd(mut self, cap: Decimal) -> Self {
224        self.max_cost_usd = Some(cap);
225        self
226    }
227
228    /// Set the queue priority of the run.
229    ///
230    /// Workers pick the pending run with the highest priority first, then
231    /// the oldest among equal priorities. The value must lie in
232    /// [`MIN_PRIORITY`]`..=`[`MAX_PRIORITY`]; the store rejects anything
233    /// else. Defaults to `0` if not set.
234    ///
235    /// [`MIN_PRIORITY`]: ironflow_store::entities::MIN_PRIORITY
236    /// [`MAX_PRIORITY`]: ironflow_store::entities::MAX_PRIORITY
237    ///
238    /// # Examples
239    ///
240    /// ```
241    /// use ironflow_engine::run_creator::CreateRunOpts;
242    ///
243    /// let new_run = CreateRunOpts::new().priority(10).build("deploy", None, None);
244    /// assert_eq!(new_run.priority, 10);
245    /// ```
246    pub fn priority(mut self, priority: i16) -> Self {
247        self.priority = Some(priority);
248        self
249    }
250
251    /// Set the queue priority only when [`priority`](Self::priority) was not
252    /// called, typically with [`WorkflowHandler::priority`].
253    ///
254    /// [`WorkflowHandler::priority`]: crate::handler::WorkflowHandler::priority
255    ///
256    /// # Examples
257    ///
258    /// ```
259    /// use ironflow_engine::run_creator::CreateRunOpts;
260    ///
261    /// let handler_default = CreateRunOpts::new().default_priority(5).build("a", None, None);
262    /// assert_eq!(handler_default.priority, 5);
263    ///
264    /// let explicit = CreateRunOpts::new()
265    ///     .priority(-3)
266    ///     .default_priority(5)
267    ///     .build("a", None, None);
268    /// assert_eq!(explicit.priority, -3);
269    /// ```
270    pub fn default_priority(mut self, priority: i16) -> Self {
271        self.priority.get_or_insert(priority);
272        self
273    }
274
275    /// Add worker tags the run requires. Extends the tags already set.
276    ///
277    /// Tags are trimmed, sorted and deduplicated by [`build`](Self::build).
278    /// The store refuses invalid ones when the run is created.
279    ///
280    /// # Examples
281    ///
282    /// ```
283    /// use ironflow_engine::run_creator::CreateRunOpts;
284    ///
285    /// let new_run = CreateRunOpts::new()
286    ///     .worker_tags(["region:eu", "gpu"])
287    ///     .worker_tags(["gpu"])
288    ///     .build("transcode", None, None);
289    /// assert_eq!(new_run.worker_tags, vec!["gpu".to_string(), "region:eu".to_string()]);
290    /// ```
291    pub fn worker_tags<I, S>(mut self, tags: I) -> Self
292    where
293        I: IntoIterator<Item = S>,
294        S: Into<String>,
295    {
296        self.worker_tags.extend(tags.into_iter().map(Into::into));
297        self
298    }
299
300    /// Assemble a [`NewRun`] from these options and handler metadata.
301    ///
302    /// * `workflow_name` -- typically from [`WorkflowHandler::name`].
303    /// * `handler_version` -- typically from [`WorkflowHandler::version`].
304    /// * `default_max_cost_usd` -- typically from [`WorkflowHandler::default_max_cost_usd`].
305    ///   Applied only when [`max_cost_usd`](Self::max_cost_usd) was not set.
306    ///
307    /// [`WorkflowHandler::name`]: crate::handler::WorkflowHandler::name
308    /// [`WorkflowHandler::version`]: crate::handler::WorkflowHandler::version
309    /// [`WorkflowHandler::default_max_cost_usd`]: crate::handler::WorkflowHandler::default_max_cost_usd
310    ///
311    /// # Examples
312    ///
313    /// ```
314    /// use ironflow_engine::run_creator::CreateRunOpts;
315    /// use rust_decimal::Decimal;
316    ///
317    /// let new_run = CreateRunOpts::new()
318    ///     .build("my-handler", Some("3.0.0"), Some(Decimal::new(1000, 2)));
319    ///
320    /// assert_eq!(new_run.workflow_name, "my-handler");
321    /// assert_eq!(new_run.handler_version, Some("3.0.0".to_string()));
322    /// assert_eq!(new_run.max_cost_usd, Some(Decimal::new(1000, 2)));
323    /// ```
324    pub fn build(
325        self,
326        workflow_name: &str,
327        handler_version: Option<&str>,
328        default_max_cost_usd: Option<Decimal>,
329    ) -> NewRun {
330        NewRun {
331            workflow_name: workflow_name.to_string(),
332            trigger: self.trigger.unwrap_or(TriggerKind::Manual),
333            payload: self.payload.unwrap_or_else(|| serde_json::json!({})),
334            max_retries: self.max_retries.unwrap_or(0),
335            handler_version: handler_version.map(str::to_string),
336            labels: self.labels.unwrap_or_default(),
337            scheduled_at: self.scheduled_at,
338            created_by: self.created_by,
339            idempotency_key: self.idempotency_key,
340            concurrency_key: self.concurrency_key,
341            priority: self.priority.unwrap_or(0),
342            concurrency_limits: Vec::new(),
343            max_cost_usd: self.max_cost_usd.or(default_max_cost_usd),
344            worker_tags: normalize_worker_tags(self.worker_tags),
345        }
346    }
347}
348
349#[cfg(test)]
350mod tests {
351    use super::*;
352    use serde_json::json;
353
354    #[test]
355    fn create_run_opts_default_produces_correct_defaults() {
356        let opts = CreateRunOpts::new();
357        let new_run = opts.build("test-workflow", None, None);
358
359        assert_eq!(new_run.workflow_name, "test-workflow");
360        assert_eq!(new_run.trigger, TriggerKind::Manual);
361        assert_eq!(new_run.payload, json!({}));
362        assert_eq!(new_run.max_retries, 0);
363        assert_eq!(new_run.handler_version, None);
364        assert!(new_run.labels.is_empty());
365        assert_eq!(new_run.scheduled_at, None);
366        assert_eq!(new_run.created_by, None);
367        assert_eq!(new_run.idempotency_key, None);
368        assert_eq!(new_run.concurrency_key, None);
369        assert_eq!(new_run.max_cost_usd, None);
370    }
371
372    #[test]
373    fn create_run_opts_without_worker_tags_requires_none() {
374        let new_run = CreateRunOpts::new().build("test-workflow", None, None);
375        assert!(new_run.worker_tags.is_empty());
376    }
377
378    #[test]
379    fn create_run_opts_worker_tags_are_merged_and_normalized() {
380        let new_run = CreateRunOpts::new()
381            .worker_tags(["region:eu", " gpu "])
382            .worker_tags(vec!["gpu".to_string()])
383            .build("test-workflow", None, None);
384        assert_eq!(
385            new_run.worker_tags,
386            vec!["gpu".to_string(), "region:eu".to_string()]
387        );
388    }
389
390    #[test]
391    fn create_run_opts_builder_sets_all_fields() {
392        let labels = HashMap::from([("env".to_string(), "prod".to_string())]);
393        let scheduled = Utc::now();
394
395        let new_run = CreateRunOpts::new()
396            .trigger(TriggerKind::Webhook {
397                path: "/hooks/gh".to_string(),
398            })
399            .payload(json!({"ref": "main"}))
400            .max_retries(3)
401            .scheduled_at(scheduled)
402            .idempotency_key("key-123")
403            .labels(labels.clone())
404            .max_cost_usd(Decimal::new(500, 2))
405            .build("deploy", Some("2.0.0"), None);
406
407        assert_eq!(new_run.workflow_name, "deploy");
408        assert_eq!(
409            new_run.trigger,
410            TriggerKind::Webhook {
411                path: "/hooks/gh".to_string()
412            }
413        );
414        assert_eq!(new_run.payload, json!({"ref": "main"}));
415        assert_eq!(new_run.max_retries, 3);
416        assert_eq!(new_run.handler_version, Some("2.0.0".to_string()));
417        assert_eq!(new_run.scheduled_at, Some(scheduled));
418        assert_eq!(new_run.idempotency_key, Some("key-123".to_string()));
419        assert_eq!(new_run.labels, labels);
420        assert_eq!(new_run.max_cost_usd, Some(Decimal::new(500, 2)));
421    }
422
423    #[test]
424    fn create_run_opts_build_carries_the_concurrency_key() {
425        let new_run = CreateRunOpts::new()
426            .concurrency_key("issue:12")
427            .build("deploy", None, None);
428
429        assert_eq!(new_run.concurrency_key.as_deref(), Some("issue:12"));
430        assert_eq!(new_run.idempotency_key, None);
431    }
432
433    #[test]
434    fn create_run_opts_build_uses_handler_metadata() {
435        let new_run =
436            CreateRunOpts::new().build("my-handler", Some("3.0.0"), Some(Decimal::new(1000, 2)));
437
438        assert_eq!(new_run.workflow_name, "my-handler");
439        assert_eq!(new_run.handler_version, Some("3.0.0".to_string()));
440        assert_eq!(new_run.max_cost_usd, Some(Decimal::new(1000, 2)));
441    }
442
443    #[test]
444    fn create_run_opts_explicit_max_cost_overrides_handler_default() {
445        let new_run = CreateRunOpts::new()
446            .max_cost_usd(Decimal::new(200, 2))
447            .build("handler", Some("1"), Some(Decimal::new(1000, 2)));
448
449        assert_eq!(new_run.max_cost_usd, Some(Decimal::new(200, 2)));
450    }
451
452    #[test]
453    fn create_run_opts_priority_defaults_to_zero() {
454        let new_run = CreateRunOpts::new().build("handler", None, None);
455        assert_eq!(new_run.priority, 0);
456    }
457
458    #[test]
459    fn create_run_opts_priority_is_carried() {
460        let new_run = CreateRunOpts::new()
461            .priority(-20)
462            .build("handler", None, None);
463        assert_eq!(new_run.priority, -20);
464    }
465
466    #[test]
467    fn create_run_opts_default_priority_applies_only_when_unset() {
468        let defaulted = CreateRunOpts::new()
469            .default_priority(40)
470            .build("handler", None, None);
471        assert_eq!(defaulted.priority, 40);
472
473        let explicit = CreateRunOpts::new()
474            .priority(0)
475            .default_priority(40)
476            .build("handler", None, None);
477        assert_eq!(explicit.priority, 0);
478    }
479
480    #[tokio::test]
481    async fn run_creator_blanket_impl_with_in_memory_store() {
482        use ironflow_store::memory::InMemoryStore;
483
484        let store = InMemoryStore::new();
485        let creator: &dyn RunCreator = &store;
486
487        let new_run =
488            CreateRunOpts::new()
489                .trigger(TriggerKind::Api)
490                .build("blanket-test", None, None);
491
492        let creation = creator.create_run(new_run).await.expect("create_run");
493        let run = creation.into_run();
494        assert_eq!(run.workflow_name, "blanket-test");
495    }
496
497    #[tokio::test]
498    async fn create_run_with_reused_idempotency_key_returns_existing() {
499        use ironflow_store::memory::InMemoryStore;
500
501        let store = InMemoryStore::new();
502        let creator: &dyn RunCreator = &store;
503
504        let first = creator
505            .create_run(CreateRunOpts::new().idempotency_key("dedup-1").build(
506                "idem-test",
507                None,
508                None,
509            ))
510            .await
511            .expect("first create_run");
512        assert!(first.is_created());
513
514        let second = creator
515            .create_run(CreateRunOpts::new().idempotency_key("dedup-1").build(
516                "idem-test",
517                None,
518                None,
519            ))
520            .await
521            .expect("second create_run");
522        assert!(!second.is_created());
523        assert_eq!(first.into_run().id, second.into_run().id);
524    }
525}