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}