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}