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