ironflow_engine/handler.rs
1//! [`WorkflowHandler`] trait — dynamic workflows with context chaining.
2//!
3//! Implement this trait to define workflows where steps can reference
4//! outputs from previous steps. The handler receives a [`WorkflowContext`]
5//! that provides step execution methods with automatic persistence.
6//!
7//! # Examples
8//!
9//! ```no_run
10//! use ironflow_engine::handler::WorkflowHandler;
11//! use ironflow_engine::context::WorkflowContext;
12//! use ironflow_engine::config::{ShellConfig, AgentStepConfig};
13//! use ironflow_engine::error::EngineError;
14//! use std::future::Future;
15//! use std::pin::Pin;
16//!
17//! struct DeployWorkflow;
18//!
19//! impl WorkflowHandler for DeployWorkflow {
20//! fn name(&self) -> &str {
21//! "deploy"
22//! }
23//!
24//! fn execute<'a>(
25//! &'a self,
26//! ctx: &'a mut WorkflowContext,
27//! ) -> Pin<Box<dyn Future<Output = Result<(), EngineError>> + Send + 'a>> {
28//! Box::pin(async move {
29//! let build = ctx.shell("build", ShellConfig::new("cargo build --release")).await?;
30//! let tests = ctx.shell("test", ShellConfig::new("cargo test")).await?;
31//!
32//! let review = ctx.agent("review", AgentStepConfig::new(
33//! &format!("Build:\n{}\nTests:\n{}\nReview.",
34//! build.stdout(), tests.stdout())
35//! )).await?;
36//!
37//! if review.text().contains("LGTM") {
38//! ctx.shell("deploy", ShellConfig::new("./deploy.sh")).await?;
39//! }
40//!
41//! Ok(())
42//! })
43//! }
44//! }
45//! ```
46
47use std::collections::HashMap;
48use std::future::Future;
49use std::pin::Pin;
50
51use ironflow_store::entities::{MAX_PRIORITY, MIN_PRIORITY};
52use rust_decimal::Decimal;
53use schemars::JsonSchema;
54use serde::Serialize;
55use serde_json::Value;
56
57use crate::context::WorkflowContext;
58use crate::error::EngineError;
59use crate::guard::WorkflowGuardConfig;
60use crate::run_creator::{CreateRunOpts, RunCreator, RunCreatorFuture};
61use crate::schedule::CronSchedule;
62
63mod typed;
64
65pub use typed::{TypedWorkflow, sub_workflow_names};
66
67/// Generate a JSON Schema [`Value`] from a type that derives [`JsonSchema`].
68///
69/// Use this in [`WorkflowHandler::input_schema`] to automatically derive the
70/// schema from your input struct instead of writing JSON by hand.
71///
72/// # Examples
73///
74/// ```
75/// use schemars::JsonSchema;
76/// use serde::Deserialize;
77/// use ironflow_engine::handler::input_schema_for;
78///
79/// #[derive(Deserialize, JsonSchema)]
80/// struct DeployInput {
81/// environment: String,
82/// dry_run: Option<bool>,
83/// }
84///
85/// let schema = input_schema_for::<DeployInput>();
86/// assert_eq!(schema["type"], "object");
87/// assert!(schema["properties"]["environment"].is_object());
88/// ```
89pub fn input_schema_for<T: JsonSchema>() -> Value {
90 let schema = schemars::schema_for!(T);
91 serde_json::to_value(schema).expect("schema serialization cannot fail")
92}
93
94/// Boxed future returned by [`WorkflowHandler::execute`].
95pub type HandlerFuture<'a> = Pin<Box<dyn Future<Output = Result<(), EngineError>> + Send + 'a>>;
96
97/// Metadata about a workflow, returned by [`WorkflowHandler::describe`].
98///
99/// Contains a human-readable description and optional Rust source code
100/// for display in the dashboard.
101///
102/// Most handlers never build this struct by hand: override
103/// [`WorkflowHandler::description`] and [`WorkflowHandler::source_code`] and
104/// the default [`WorkflowHandler::describe`] assembles it from the other
105/// trait methods. The builder below exists for handlers that override
106/// `describe` entirely.
107///
108/// # Examples
109///
110/// ```
111/// use ironflow_engine::handler::WorkflowInfo;
112///
113/// let info = WorkflowInfo::new("Deploy to production")
114/// .with_category("ops")
115/// .with_version("2.0.0")
116/// .with_sub_workflows(["build"]);
117///
118/// assert_eq!(info.description, "Deploy to production");
119/// assert_eq!(info.category.as_deref(), Some("ops"));
120/// assert_eq!(info.sub_workflows, vec!["build".to_string()]);
121/// ```
122#[derive(Debug, Clone, Default, Serialize)]
123pub struct WorkflowInfo {
124 /// Human-readable description of what the workflow does.
125 pub description: String,
126 /// Optional Rust source code of the handler (for UI display).
127 pub source_code: Option<String>,
128 /// Names of sub-workflows invoked by this handler.
129 #[serde(default, skip_serializing_if = "Vec::is_empty")]
130 pub sub_workflows: Vec<String>,
131 /// Optional `/`-separated category path used to group workflows in the UI tree.
132 ///
133 /// A value like `"data/etl"` places the workflow under `data` → `etl`.
134 /// `None` means the workflow is uncategorized.
135 #[serde(default, skip_serializing_if = "Option::is_none")]
136 pub category: Option<String>,
137 /// Handler version string, used to trace which code produced a given run.
138 #[serde(default, skip_serializing_if = "Option::is_none")]
139 pub version: Option<String>,
140 /// Versions accepted for replay without `force`.
141 #[serde(default, skip_serializing_if = "Vec::is_empty")]
142 pub compatible_versions: Vec<String>,
143 /// JSON Schema describing the expected input payload.
144 ///
145 /// When present, the dashboard renders a dynamic form from this schema
146 /// and the engine validates the payload before creating a run.
147 #[serde(default, skip_serializing_if = "Option::is_none")]
148 pub input_schema: Option<Value>,
149 /// Labels automatically applied to every run of this workflow.
150 #[serde(default, skip_serializing_if = "HashMap::is_empty")]
151 pub default_labels: HashMap<String, String>,
152 /// Optional cron schedule for automatic execution.
153 #[serde(default, skip_serializing_if = "Option::is_none")]
154 pub schedule: Option<CronSchedule>,
155 /// Default cumulative cost cap applied to runs of this workflow, in USD.
156 ///
157 /// Overridden by a cap supplied at run creation, and takes precedence over
158 /// the server-wide default. `None` means the handler declares no default.
159 #[serde(default, skip_serializing_if = "Option::is_none")]
160 pub default_max_cost_usd: Option<Decimal>,
161 /// Queue priority applied to runs of this workflow when the creation
162 /// request does not supply one. See [`WorkflowHandler::priority`].
163 #[serde(default, skip_serializing_if = "is_default_priority")]
164 pub priority: i16,
165}
166
167/// Whether `priority` is the default `0`, so it is left out of the JSON.
168fn is_default_priority(priority: &i16) -> bool {
169 *priority == 0
170}
171
172/// Bring a handler priority into the range the store accepts, so a
173/// misconfigured handler never makes every run creation fail.
174pub(crate) fn clamp_priority(priority: i16) -> i16 {
175 priority.clamp(MIN_PRIORITY, MAX_PRIORITY)
176}
177
178impl WorkflowInfo {
179 /// Create metadata with a description and every other field at its default.
180 ///
181 /// # Examples
182 ///
183 /// ```
184 /// use ironflow_engine::handler::WorkflowInfo;
185 ///
186 /// let info = WorkflowInfo::new("Nightly backup");
187 /// assert_eq!(info.description, "Nightly backup");
188 /// assert!(info.source_code.is_none());
189 /// assert!(info.sub_workflows.is_empty());
190 /// ```
191 pub fn new(description: impl Into<String>) -> Self {
192 Self {
193 description: description.into(),
194 ..Self::default()
195 }
196 }
197
198 /// Attach the handler source code, typically via `include_str!`.
199 ///
200 /// # Examples
201 ///
202 /// ```
203 /// use ironflow_engine::handler::WorkflowInfo;
204 ///
205 /// let info = WorkflowInfo::new("Demo").with_source_code("struct Demo;");
206 /// assert_eq!(info.source_code.as_deref(), Some("struct Demo;"));
207 /// ```
208 pub fn with_source_code(mut self, source: impl Into<String>) -> Self {
209 self.source_code = Some(source.into());
210 self
211 }
212
213 /// Declare the sub-workflows this handler invokes.
214 ///
215 /// # Examples
216 ///
217 /// ```
218 /// use ironflow_engine::handler::WorkflowInfo;
219 ///
220 /// let info = WorkflowInfo::new("Report").with_sub_workflows(["collect", "enrich"]);
221 /// assert_eq!(info.sub_workflows, vec!["collect".to_string(), "enrich".to_string()]);
222 /// ```
223 pub fn with_sub_workflows<I, S>(mut self, names: I) -> Self
224 where
225 I: IntoIterator<Item = S>,
226 S: Into<String>,
227 {
228 self.sub_workflows = names.into_iter().map(Into::into).collect();
229 self
230 }
231
232 /// Set the `/`-separated category path.
233 ///
234 /// # Examples
235 ///
236 /// ```
237 /// use ironflow_engine::handler::WorkflowInfo;
238 ///
239 /// let info = WorkflowInfo::new("ETL").with_category("data/etl");
240 /// assert_eq!(info.category.as_deref(), Some("data/etl"));
241 /// ```
242 pub fn with_category(mut self, category: impl Into<String>) -> Self {
243 self.category = Some(category.into());
244 self
245 }
246
247 /// Set the handler version.
248 ///
249 /// # Examples
250 ///
251 /// ```
252 /// use ironflow_engine::handler::WorkflowInfo;
253 ///
254 /// let info = WorkflowInfo::new("Deploy").with_version("1.2.0");
255 /// assert_eq!(info.version.as_deref(), Some("1.2.0"));
256 /// ```
257 pub fn with_version(mut self, version: impl Into<String>) -> Self {
258 self.version = Some(version.into());
259 self
260 }
261
262 /// Set the versions accepted for replay without `force`.
263 ///
264 /// # Examples
265 ///
266 /// ```
267 /// use ironflow_engine::handler::WorkflowInfo;
268 ///
269 /// let info = WorkflowInfo::new("Deploy").with_compatible_versions(["1.0.0"]);
270 /// assert_eq!(info.compatible_versions, vec!["1.0.0".to_string()]);
271 /// ```
272 pub fn with_compatible_versions<I, S>(mut self, versions: I) -> Self
273 where
274 I: IntoIterator<Item = S>,
275 S: Into<String>,
276 {
277 self.compatible_versions = versions.into_iter().map(Into::into).collect();
278 self
279 }
280
281 /// Set the JSON Schema of the expected input payload.
282 ///
283 /// # Examples
284 ///
285 /// ```
286 /// use ironflow_engine::handler::WorkflowInfo;
287 /// use serde_json::json;
288 ///
289 /// let info = WorkflowInfo::new("Greet").with_input_schema(json!({"type": "object"}));
290 /// assert_eq!(info.input_schema.unwrap()["type"], "object");
291 /// ```
292 pub fn with_input_schema(mut self, schema: Value) -> Self {
293 self.input_schema = Some(schema);
294 self
295 }
296
297 /// Set the labels applied to every run of this workflow.
298 ///
299 /// # Examples
300 ///
301 /// ```
302 /// use std::collections::HashMap;
303 /// use ironflow_engine::handler::WorkflowInfo;
304 ///
305 /// let labels = HashMap::from([("team".to_string(), "core".to_string())]);
306 /// let info = WorkflowInfo::new("Sync").with_default_labels(labels);
307 /// assert_eq!(info.default_labels["team"], "core");
308 /// ```
309 pub fn with_default_labels(mut self, labels: HashMap<String, String>) -> Self {
310 self.default_labels = labels;
311 self
312 }
313
314 /// Set the cron schedule.
315 ///
316 /// # Examples
317 ///
318 /// ```
319 /// use ironflow_engine::handler::WorkflowInfo;
320 /// use ironflow_engine::schedule::CronSchedule;
321 ///
322 /// let schedule = CronSchedule::new("0 0 * * *")?;
323 /// let info = WorkflowInfo::new("Nightly").with_schedule(schedule);
324 /// assert!(info.schedule.is_some());
325 /// # Ok::<(), String>(())
326 /// ```
327 pub fn with_schedule(mut self, schedule: CronSchedule) -> Self {
328 self.schedule = Some(schedule);
329 self
330 }
331
332 /// Set the default cumulative cost cap in USD.
333 ///
334 /// # Examples
335 ///
336 /// ```
337 /// use ironflow_engine::handler::WorkflowInfo;
338 /// use rust_decimal::Decimal;
339 ///
340 /// let info = WorkflowInfo::new("Analysis").with_default_max_cost_usd(Decimal::new(500, 2));
341 /// assert_eq!(info.default_max_cost_usd, Some(Decimal::new(500, 2)));
342 /// ```
343 pub fn with_default_max_cost_usd(mut self, cap: Decimal) -> Self {
344 self.default_max_cost_usd = Some(cap);
345 self
346 }
347
348 /// Set the default queue priority of the workflow runs.
349 ///
350 /// # Examples
351 ///
352 /// ```
353 /// use ironflow_engine::handler::WorkflowInfo;
354 ///
355 /// let info = WorkflowInfo::new("Hotfix").with_priority(50);
356 /// assert_eq!(info.priority, 50);
357 /// ```
358 pub fn with_priority(mut self, priority: i16) -> Self {
359 self.priority = priority;
360 self
361 }
362}
363
364/// A dynamic workflow handler with context-aware step chaining.
365///
366/// Implement this trait to define workflows where each step can use
367/// the output of previous steps. Register handlers with
368/// [`Engine::register`](crate::engine::Engine::register) and execute
369/// them by name.
370///
371/// # Why `Pin<Box<dyn Future>>` instead of `async fn`?
372///
373/// The handler must be object-safe (`dyn WorkflowHandler`) to allow
374/// registering different handler types in the engine's registry.
375pub trait WorkflowHandler: Send + Sync {
376 /// The workflow name used for registration and lookup.
377 fn name(&self) -> &str;
378
379 /// Handler version string, used to trace which code version produced a run.
380 ///
381 /// Override this to return a meaningful version (semver, git SHA, build
382 /// hash, etc.). The default is `"1"`.
383 ///
384 /// The engine records this value on every run it creates so that retries
385 /// can detect when the handler has changed since the original execution.
386 fn version(&self) -> Option<&str> {
387 Some("1")
388 }
389
390 /// Versions of this handler that can replay payloads produced by an
391 /// older run without requiring `force`.
392 ///
393 /// When a retry targets a run whose `handler_version` differs from
394 /// [`version`](Self::version), the engine checks this list. If the
395 /// run's version appears here, the retry proceeds normally; otherwise
396 /// it is refused with `409 HANDLER_VERSION_MISMATCH` unless the caller
397 /// passes `force=true`.
398 ///
399 /// The default is an empty slice (only the current version is accepted).
400 ///
401 /// # Examples
402 ///
403 /// ```
404 /// # use ironflow_engine::handler::{WorkflowHandler, HandlerFuture};
405 /// # use ironflow_engine::context::WorkflowContext;
406 /// struct MigratedHandler;
407 ///
408 /// impl WorkflowHandler for MigratedHandler {
409 /// fn name(&self) -> &str { "migrated" }
410 /// fn version(&self) -> Option<&str> { Some("2.0.0") }
411 /// fn compatible_versions(&self) -> &[&str] { &["1.0.0", "1.5.0"] }
412 /// fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
413 /// Box::pin(async move { Ok(()) })
414 /// }
415 /// }
416 ///
417 /// assert_eq!(MigratedHandler.compatible_versions(), &["1.0.0", "1.5.0"]);
418 /// ```
419 fn compatible_versions(&self) -> &[&str] {
420 &[]
421 }
422
423 /// Human-readable description shown in the dashboard and the CLI.
424 ///
425 /// The default is an empty string. Override this rather than
426 /// [`describe`](Self::describe): the default `describe` picks it up.
427 ///
428 /// # Examples
429 ///
430 /// ```
431 /// # use ironflow_engine::handler::{WorkflowHandler, HandlerFuture};
432 /// # use ironflow_engine::context::WorkflowContext;
433 /// struct Backup;
434 ///
435 /// impl WorkflowHandler for Backup {
436 /// fn name(&self) -> &str { "backup" }
437 /// fn description(&self) -> &str { "Nightly database backup" }
438 /// fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
439 /// Box::pin(async move { Ok(()) })
440 /// }
441 /// }
442 ///
443 /// assert_eq!(Backup.describe().description, "Nightly database backup");
444 /// ```
445 fn description(&self) -> &str {
446 ""
447 }
448
449 /// Rust source of the handler, displayed in the dashboard.
450 ///
451 /// Return `Some(include_str!("this_file.rs"))` to show the code next to
452 /// the run. The default is `None`.
453 ///
454 /// # Examples
455 ///
456 /// ```
457 /// # use ironflow_engine::handler::{WorkflowHandler, HandlerFuture};
458 /// # use ironflow_engine::context::WorkflowContext;
459 /// struct Backup;
460 ///
461 /// impl WorkflowHandler for Backup {
462 /// fn name(&self) -> &str { "backup" }
463 /// fn source_code(&self) -> Option<&str> { Some("struct Backup;") }
464 /// fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
465 /// Box::pin(async move { Ok(()) })
466 /// }
467 /// }
468 ///
469 /// assert_eq!(Backup.describe().source_code.as_deref(), Some("struct Backup;"));
470 /// ```
471 fn source_code(&self) -> Option<&str> {
472 None
473 }
474
475 /// Names of the sub-workflows this handler invokes through
476 /// [`WorkflowContext::workflow`](crate::context::WorkflowContext::workflow).
477 ///
478 /// Purely informational: the dashboard uses it to draw the call graph.
479 /// The default is empty. Build it with [`sub_workflow_names`] from the
480 /// child handlers themselves, never from hand-written names.
481 ///
482 /// # Examples
483 ///
484 /// ```
485 /// # use ironflow_engine::handler::{WorkflowHandler, HandlerFuture, sub_workflow_names};
486 /// # use ironflow_engine::context::WorkflowContext;
487 /// struct Collect;
488 ///
489 /// impl WorkflowHandler for Collect {
490 /// fn name(&self) -> &str { "collect" }
491 /// fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
492 /// Box::pin(async move { Ok(()) })
493 /// }
494 /// }
495 ///
496 /// struct Report;
497 ///
498 /// impl WorkflowHandler for Report {
499 /// fn name(&self) -> &str { "report" }
500 /// fn sub_workflows(&self) -> Vec<String> { sub_workflow_names(&[&Collect]) }
501 /// fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
502 /// Box::pin(async move { Ok(()) })
503 /// }
504 /// }
505 ///
506 /// assert_eq!(Report.describe().sub_workflows, vec!["collect".to_string()]);
507 /// ```
508 fn sub_workflows(&self) -> Vec<String> {
509 Vec::new()
510 }
511
512 /// Optional `/`-separated category path used to group workflows in the UI tree.
513 ///
514 /// Return a value like `"data/etl"` to place the workflow under `data` → `etl`.
515 /// The default is `None` (uncategorized).
516 ///
517 /// Validation (empty segments, leading or trailing `/`, `//`, whitespace
518 /// segments) is enforced at registration time by
519 /// [`Engine::register`](crate::engine::Engine::register).
520 fn category(&self) -> Option<&str> {
521 None
522 }
523
524 /// Return a JSON Schema describing the expected input payload.
525 ///
526 /// When present, the dashboard renders a dynamic form from this schema
527 /// and the engine validates the payload before creating a run.
528 /// The default is `None` (no schema, free-form payload). A handler that
529 /// implements [`TypedWorkflow`] returns
530 /// [`TypedWorkflow::typed_input_schema`].
531 fn input_schema(&self) -> Option<Value> {
532 None
533 }
534
535 /// Labels automatically applied to every run of this workflow.
536 ///
537 /// These are merged with any labels provided at run creation time.
538 /// User-provided labels take precedence over defaults.
539 fn default_labels(&self) -> HashMap<String, String> {
540 HashMap::new()
541 }
542
543 /// Optional cron schedule for automatic execution.
544 ///
545 /// Return a [`CronSchedule`] built from a cron expression
546 /// (5 or 6 fields, as supported by [`croner`]).
547 ///
548 /// When set, the engine exposes this handler via
549 /// [`Engine::scheduled_handlers`](crate::engine::Engine::scheduled_handlers)
550 /// so the runtime can wire it into a cron scheduler automatically.
551 ///
552 /// The default is `None` (no automatic scheduling).
553 ///
554 /// The schedule's policy decides what happens to the occurrences missed
555 /// while no server fired them ([`CronSchedule::with_catchup`]), whether a
556 /// run starts while the previous one is still active
557 /// ([`CronSchedule::with_overlap`]), and the timezone the expression is
558 /// evaluated in ([`CronSchedule::with_timezone`]). It is synced to the
559 /// stored schedule at startup, like the expression.
560 ///
561 /// # Examples
562 ///
563 /// ```
564 /// # use ironflow_engine::handler::{WorkflowHandler, HandlerFuture};
565 /// # use ironflow_engine::context::WorkflowContext;
566 /// # use ironflow_engine::schedule::CronSchedule;
567 /// struct HourlySync;
568 ///
569 /// impl WorkflowHandler for HourlySync {
570 /// fn name(&self) -> &str { "hourly-sync" }
571 /// fn schedule(&self) -> Option<&CronSchedule> {
572 /// // In practice, store as a field or use `std::sync::LazyLock`.
573 /// None
574 /// }
575 /// fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
576 /// Box::pin(async move { Ok(()) })
577 /// }
578 /// }
579 /// ```
580 ///
581 /// Every morning at 9:00 Paris time, catching up every missed morning,
582 /// and never two reports at once:
583 ///
584 /// ```
585 /// use std::sync::LazyLock;
586 ///
587 /// use ironflow_engine::context::WorkflowContext;
588 /// use ironflow_engine::handler::{HandlerFuture, WorkflowHandler};
589 /// use ironflow_engine::schedule::{CatchupPolicy, CronSchedule, OverlapPolicy};
590 ///
591 /// static MORNING: LazyLock<CronSchedule> = LazyLock::new(|| {
592 /// CronSchedule::new("0 9 * * *")
593 /// .and_then(|s| s.with_timezone("Europe/Paris"))
594 /// .map(|s| s.with_catchup(CatchupPolicy::All).with_overlap(OverlapPolicy::Skip))
595 /// .expect("valid schedule")
596 /// });
597 ///
598 /// struct MorningReport;
599 ///
600 /// impl WorkflowHandler for MorningReport {
601 /// fn name(&self) -> &str { "morning-report" }
602 /// fn schedule(&self) -> Option<&CronSchedule> {
603 /// Some(&MORNING)
604 /// }
605 /// fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
606 /// Box::pin(async move { Ok(()) })
607 /// }
608 /// }
609 ///
610 /// assert_eq!(MorningReport.schedule().map(|s| s.policy().timezone.name()), Some("Europe/Paris"));
611 /// ```
612 fn schedule(&self) -> Option<&CronSchedule> {
613 None
614 }
615
616 /// Default cumulative cost cap for runs of this workflow, in USD.
617 ///
618 /// Applied when the run creation request does not supply one. Takes
619 /// precedence over the server-wide
620 /// [`IRONFLOW_DEFAULT_RUN_MAX_COST_USD`](crate::budget::DEFAULT_RUN_MAX_COST_ENV).
621 /// The default is `None` (fall back to the server default, or no cap).
622 ///
623 /// # Examples
624 ///
625 /// ```
626 /// # use ironflow_engine::handler::{WorkflowHandler, HandlerFuture};
627 /// # use ironflow_engine::context::WorkflowContext;
628 /// use rust_decimal::Decimal;
629 ///
630 /// struct ExpensiveAnalysis;
631 ///
632 /// impl WorkflowHandler for ExpensiveAnalysis {
633 /// fn name(&self) -> &str { "expensive-analysis" }
634 /// fn default_max_cost_usd(&self) -> Option<Decimal> {
635 /// Some(Decimal::new(500, 2)) // $5.00
636 /// }
637 /// fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
638 /// Box::pin(async move { Ok(()) })
639 /// }
640 /// }
641 ///
642 /// assert_eq!(ExpensiveAnalysis.default_max_cost_usd(), Some(Decimal::new(500, 2)));
643 /// ```
644 fn default_max_cost_usd(&self) -> Option<Decimal> {
645 None
646 }
647
648 /// Default queue priority for runs of this workflow.
649 ///
650 /// Workers pick the pending run with the highest priority first, then the
651 /// oldest among equal priorities. Applied when the run creation request
652 /// does not supply a priority. The default is `0`. A value outside
653 /// `-100..=100` is clamped to the nearest bound.
654 ///
655 /// Priority only orders the queue: a running run is never preempted, and
656 /// nothing ages a low-priority run, so a steady flow of higher-priority
657 /// runs can delay it indefinitely.
658 ///
659 /// # Examples
660 ///
661 /// ```
662 /// # use ironflow_engine::handler::{WorkflowHandler, HandlerFuture};
663 /// # use ironflow_engine::context::WorkflowContext;
664 /// struct Hotfix;
665 ///
666 /// impl WorkflowHandler for Hotfix {
667 /// fn name(&self) -> &str { "hotfix" }
668 /// fn priority(&self) -> i16 { 50 }
669 /// fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
670 /// Box::pin(async move { Ok(()) })
671 /// }
672 /// }
673 ///
674 /// assert_eq!(Hotfix.priority(), 50);
675 /// assert_eq!(Hotfix.describe().priority, 50);
676 /// ```
677 fn priority(&self) -> i16 {
678 0
679 }
680
681 /// Optional guard configuration for this workflow.
682 ///
683 /// When present, overrides the engine's global guard configuration
684 /// for runs of this handler. The default is `None` (use the engine's
685 /// global configuration).
686 ///
687 /// # Examples
688 ///
689 /// ```
690 /// # use ironflow_engine::handler::{WorkflowHandler, HandlerFuture};
691 /// # use ironflow_engine::context::WorkflowContext;
692 /// use ironflow_engine::guard::WorkflowGuardConfig;
693 ///
694 /// struct StrictWorkflow;
695 ///
696 /// impl WorkflowHandler for StrictWorkflow {
697 /// fn name(&self) -> &str { "strict" }
698 /// fn guard_config(&self) -> Option<WorkflowGuardConfig> {
699 /// Some(WorkflowGuardConfig::new().with_max_depth(2))
700 /// }
701 /// fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
702 /// Box::pin(async move { Ok(()) })
703 /// }
704 /// }
705 ///
706 /// assert_eq!(StrictWorkflow.guard_config().unwrap().max_depth, 2);
707 /// ```
708 fn guard_config(&self) -> Option<WorkflowGuardConfig> {
709 None
710 }
711
712 /// Worker tags every run of this workflow requires.
713 ///
714 /// A run is only handed to a worker that carries all of these tags (see
715 /// `WorkerBuilder::tags` in `ironflow-worker`). Tags given at run creation
716 /// are merged with this list. The default is empty: any worker that
717 /// registered the workflow may take the run.
718 ///
719 /// # Examples
720 ///
721 /// ```
722 /// # use ironflow_engine::handler::{WorkflowHandler, HandlerFuture};
723 /// # use ironflow_engine::context::WorkflowContext;
724 /// struct Transcode;
725 ///
726 /// impl WorkflowHandler for Transcode {
727 /// fn name(&self) -> &str { "transcode" }
728 /// fn required_worker_tags(&self) -> Vec<String> {
729 /// vec!["gpu".into()]
730 /// }
731 /// fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
732 /// Box::pin(async move { Ok(()) })
733 /// }
734 /// }
735 ///
736 /// assert_eq!(Transcode.required_worker_tags(), vec!["gpu".to_string()]);
737 /// ```
738 fn required_worker_tags(&self) -> Vec<String> {
739 Vec::new()
740 }
741
742 /// Check whether a run carrying `run_version` can be replayed by this
743 /// handler without `force`.
744 ///
745 /// Compatibility rules:
746 /// - `run_version` is `None` (old run predating version tracking): always
747 /// compatible.
748 /// - `run_version` equals [`version`](Self::version): compatible.
749 /// - `run_version` appears in [`compatible_versions`](Self::compatible_versions):
750 /// compatible.
751 /// - Otherwise: incompatible.
752 ///
753 /// # Examples
754 ///
755 /// ```
756 /// # use ironflow_engine::handler::{WorkflowHandler, HandlerFuture};
757 /// # use ironflow_engine::context::WorkflowContext;
758 /// struct MyHandler;
759 ///
760 /// impl WorkflowHandler for MyHandler {
761 /// fn name(&self) -> &str { "my-handler" }
762 /// fn version(&self) -> Option<&str> { Some("2.0.0") }
763 /// fn compatible_versions(&self) -> &[&str] { &["1.0.0"] }
764 /// fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
765 /// Box::pin(async move { Ok(()) })
766 /// }
767 /// }
768 ///
769 /// assert!(MyHandler.is_version_compatible(None));
770 /// assert!(MyHandler.is_version_compatible(Some("2.0.0")));
771 /// assert!(MyHandler.is_version_compatible(Some("1.0.0")));
772 /// assert!(!MyHandler.is_version_compatible(Some("0.5.0")));
773 /// ```
774 fn is_version_compatible(&self, run_version: Option<&str>) -> bool {
775 let Some(rv) = run_version else {
776 return true;
777 };
778 if self.version() == Some(rv) {
779 return true;
780 }
781 self.compatible_versions().contains(&rv)
782 }
783
784 /// Return metadata about this workflow.
785 ///
786 /// The default assembles a [`WorkflowInfo`] from every other trait
787 /// method: [`description`](Self::description),
788 /// [`source_code`](Self::source_code),
789 /// [`sub_workflows`](Self::sub_workflows), [`category`](Self::category),
790 /// [`version`](Self::version),
791 /// [`compatible_versions`](Self::compatible_versions),
792 /// [`input_schema`](Self::input_schema),
793 /// [`default_labels`](Self::default_labels), [`schedule`](Self::schedule)
794 /// [`default_max_cost_usd`](Self::default_max_cost_usd)
795 /// and [`priority`](Self::priority). Override those instead of this
796 /// method; override `describe` only when the metadata cannot be
797 /// expressed through them.
798 fn describe(&self) -> WorkflowInfo {
799 WorkflowInfo {
800 description: self.description().to_string(),
801 source_code: self.source_code().map(str::to_string),
802 sub_workflows: self.sub_workflows(),
803 category: self.category().map(str::to_string),
804 version: self.version().map(str::to_string),
805 compatible_versions: self
806 .compatible_versions()
807 .iter()
808 .map(|s| s.to_string())
809 .collect(),
810 input_schema: self.input_schema(),
811 default_labels: self.default_labels(),
812 schedule: self.schedule().cloned(),
813 default_max_cost_usd: self.default_max_cost_usd(),
814 priority: self.priority(),
815 }
816 }
817
818 /// Create a run for this workflow, using handler metadata automatically.
819 ///
820 /// Assembles a [`NewRun`](ironflow_store::entities::NewRun) from [`name`](Self::name),
821 /// [`version`](Self::version), and [`default_max_cost_usd`](Self::default_max_cost_usd),
822 /// then delegates to the given [`RunCreator`].
823 ///
824 /// # Errors
825 ///
826 /// Returns [`EngineError`] if the underlying store rejects the run.
827 ///
828 /// # Examples
829 ///
830 /// ```no_run
831 /// # use ironflow_engine::handler::{WorkflowHandler, HandlerFuture};
832 /// # use ironflow_engine::context::WorkflowContext;
833 /// # use ironflow_engine::run_creator::{CreateRunOpts, RunCreator};
834 /// # use ironflow_store::entities::TriggerKind;
835 /// struct DeployWorkflow;
836 ///
837 /// impl WorkflowHandler for DeployWorkflow {
838 /// fn name(&self) -> &str { "deploy" }
839 /// fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
840 /// Box::pin(async move { Ok(()) })
841 /// }
842 /// }
843 ///
844 /// # async fn example(store: &dyn RunCreator) -> Result<(), ironflow_engine::error::EngineError> {
845 /// let opts = CreateRunOpts::new().trigger(TriggerKind::Api);
846 /// let run = DeployWorkflow.create_run(store, opts).await?.into_run();
847 /// assert_eq!(run.workflow_name, "deploy");
848 /// # Ok(())
849 /// # }
850 /// ```
851 fn create_run<'a>(
852 &self,
853 creator: &'a dyn RunCreator,
854 opts: CreateRunOpts,
855 ) -> RunCreatorFuture<'a> {
856 use tracing::{Instrument, info_span};
857
858 let new_run = opts
859 .default_priority(clamp_priority(self.priority()))
860 .worker_tags(self.required_worker_tags())
861 .build(self.name(), self.version(), self.default_max_cost_usd());
862 let span = info_span!("handler.create_run", workflow = %self.name());
863 Box::pin(creator.create_run(new_run).instrument(span))
864 }
865
866 /// Execute the workflow with the given context.
867 ///
868 /// The context provides [`shell`](WorkflowContext::shell),
869 /// [`http`](WorkflowContext::http), and [`agent`](WorkflowContext::agent)
870 /// methods that automatically persist each step.
871 ///
872 /// # Errors
873 ///
874 /// Return [`EngineError`] if any step fails. The engine will mark
875 /// the run as `Failed` and record the error.
876 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a>;
877}
878
879/// A boxed handler is a handler.
880///
881/// Lets a single `Vec<Box<dyn WorkflowHandler>>` feed both
882/// [`Engine::register`](crate::engine::Engine::register) and a worker
883/// builder, so the API server and the workers cannot drift apart in the
884/// list of workflows they know.
885///
886/// # Examples
887///
888/// ```
889/// use ironflow_engine::handler::{HandlerFuture, WorkflowHandler};
890/// use ironflow_engine::context::WorkflowContext;
891///
892/// struct Hello;
893///
894/// impl WorkflowHandler for Hello {
895/// fn name(&self) -> &str { "hello" }
896/// fn description(&self) -> &str { "Says hello" }
897/// fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
898/// Box::pin(async move { Ok(()) })
899/// }
900/// }
901///
902/// fn handlers() -> Vec<Box<dyn WorkflowHandler>> {
903/// vec![Box::new(Hello)]
904/// }
905///
906/// for handler in handlers() {
907/// assert_eq!(handler.name(), "hello");
908/// assert_eq!(handler.describe().description, "Says hello");
909/// }
910/// ```
911impl<T: WorkflowHandler + ?Sized> WorkflowHandler for Box<T> {
912 fn name(&self) -> &str {
913 (**self).name()
914 }
915
916 fn version(&self) -> Option<&str> {
917 (**self).version()
918 }
919
920 fn compatible_versions(&self) -> &[&str] {
921 (**self).compatible_versions()
922 }
923
924 fn description(&self) -> &str {
925 (**self).description()
926 }
927
928 fn source_code(&self) -> Option<&str> {
929 (**self).source_code()
930 }
931
932 fn sub_workflows(&self) -> Vec<String> {
933 (**self).sub_workflows()
934 }
935
936 fn category(&self) -> Option<&str> {
937 (**self).category()
938 }
939
940 fn input_schema(&self) -> Option<Value> {
941 (**self).input_schema()
942 }
943
944 fn default_labels(&self) -> HashMap<String, String> {
945 (**self).default_labels()
946 }
947
948 fn schedule(&self) -> Option<&CronSchedule> {
949 (**self).schedule()
950 }
951
952 fn default_max_cost_usd(&self) -> Option<Decimal> {
953 (**self).default_max_cost_usd()
954 }
955
956 fn priority(&self) -> i16 {
957 (**self).priority()
958 }
959
960 fn guard_config(&self) -> Option<WorkflowGuardConfig> {
961 (**self).guard_config()
962 }
963
964 fn required_worker_tags(&self) -> Vec<String> {
965 (**self).required_worker_tags()
966 }
967
968 fn is_version_compatible(&self, run_version: Option<&str>) -> bool {
969 (**self).is_version_compatible(run_version)
970 }
971
972 fn describe(&self) -> WorkflowInfo {
973 (**self).describe()
974 }
975
976 fn create_run<'a>(
977 &self,
978 creator: &'a dyn RunCreator,
979 opts: CreateRunOpts,
980 ) -> RunCreatorFuture<'a> {
981 (**self).create_run(creator, opts)
982 }
983
984 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
985 (**self).execute(ctx)
986 }
987}
988
989#[cfg(test)]
990mod tests {
991 use super::*;
992 use serde::{Deserialize, Serialize};
993
994 #[derive(Debug, Serialize, Deserialize, JsonSchema)]
995 struct TestInput {
996 environment: String,
997 #[serde(default)]
998 dry_run: bool,
999 }
1000
1001 struct MinimalHandler;
1002
1003 impl WorkflowHandler for MinimalHandler {
1004 fn name(&self) -> &str {
1005 "minimal"
1006 }
1007
1008 fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
1009 Box::pin(async { Ok(()) })
1010 }
1011 }
1012
1013 struct FullFeaturedHandler;
1014
1015 impl WorkflowHandler for FullFeaturedHandler {
1016 fn name(&self) -> &str {
1017 "full"
1018 }
1019
1020 fn version(&self) -> Option<&str> {
1021 Some("1.2.0")
1022 }
1023
1024 fn category(&self) -> Option<&str> {
1025 Some("data/etl")
1026 }
1027
1028 fn input_schema(&self) -> Option<Value> {
1029 Some(input_schema_for::<TestInput>())
1030 }
1031
1032 fn default_labels(&self) -> HashMap<String, String> {
1033 HashMap::from([
1034 ("team".to_string(), "platform".to_string()),
1035 ("env".to_string(), "prod".to_string()),
1036 ])
1037 }
1038
1039 fn default_max_cost_usd(&self) -> Option<Decimal> {
1040 Some(Decimal::new(750, 2))
1041 }
1042
1043 fn priority(&self) -> i16 {
1044 30
1045 }
1046
1047 fn describe(&self) -> WorkflowInfo {
1048 WorkflowInfo {
1049 description: "Full-featured test handler".to_string(),
1050 source_code: Some("fn test() {}".to_string()),
1051 sub_workflows: vec!["helper".to_string()],
1052 category: self.category().map(str::to_string),
1053 version: self.version().map(str::to_string),
1054 compatible_versions: self
1055 .compatible_versions()
1056 .iter()
1057 .map(|s| s.to_string())
1058 .collect(),
1059 input_schema: self.input_schema(),
1060 default_labels: self.default_labels(),
1061 schedule: self.schedule().cloned(),
1062 default_max_cost_usd: self.default_max_cost_usd(),
1063 priority: self.priority(),
1064 }
1065 }
1066
1067 fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
1068 Box::pin(async { Ok(()) })
1069 }
1070 }
1071
1072 #[test]
1073 fn minimal_handler_has_required_name() {
1074 let handler = MinimalHandler;
1075 assert_eq!(handler.name(), "minimal");
1076 }
1077
1078 #[test]
1079 fn minimal_handler_defaults_to_version_1() {
1080 let handler = MinimalHandler;
1081 assert_eq!(handler.version(), Some("1"));
1082 }
1083
1084 #[test]
1085 fn minimal_handler_defaults_to_no_compatible_versions() {
1086 let handler = MinimalHandler;
1087 assert!(handler.compatible_versions().is_empty());
1088 }
1089
1090 #[test]
1091 fn minimal_handler_defaults_to_no_category() {
1092 let handler = MinimalHandler;
1093 assert_eq!(handler.category(), None);
1094 }
1095
1096 #[test]
1097 fn minimal_handler_defaults_to_no_schema() {
1098 let handler = MinimalHandler;
1099 assert_eq!(handler.input_schema(), None);
1100 }
1101
1102 #[test]
1103 fn minimal_handler_defaults_to_empty_labels() {
1104 let handler = MinimalHandler;
1105 let labels = handler.default_labels();
1106 assert!(labels.is_empty());
1107 }
1108
1109 #[test]
1110 fn minimal_handler_defaults_to_no_schedule() {
1111 let handler = MinimalHandler;
1112 assert_eq!(handler.schedule(), None);
1113 }
1114
1115 #[test]
1116 fn minimal_handler_describe_reflects_defaults() {
1117 let handler = MinimalHandler;
1118 let info = handler.describe();
1119 assert_eq!(info.description, "");
1120 assert_eq!(info.source_code, None);
1121 assert_eq!(info.sub_workflows, Vec::<String>::new());
1122 assert_eq!(info.category, None);
1123 assert_eq!(info.version, Some("1".to_string()));
1124 assert!(info.compatible_versions.is_empty());
1125 assert_eq!(info.input_schema, None);
1126 assert!(info.default_labels.is_empty());
1127 assert_eq!(info.schedule, None);
1128 }
1129
1130 #[test]
1131 fn full_handler_returns_all_metadata() {
1132 let handler = FullFeaturedHandler;
1133 assert_eq!(handler.name(), "full");
1134 assert_eq!(handler.version(), Some("1.2.0"));
1135 assert_eq!(handler.category(), Some("data/etl"));
1136 assert!(handler.input_schema().is_some());
1137 }
1138
1139 #[test]
1140 fn full_handler_default_labels_are_set() {
1141 let handler = FullFeaturedHandler;
1142 let labels = handler.default_labels();
1143 assert_eq!(labels.get("team"), Some(&"platform".to_string()));
1144 assert_eq!(labels.get("env"), Some(&"prod".to_string()));
1145 }
1146
1147 #[test]
1148 fn full_handler_describe_includes_all_fields() {
1149 let handler = FullFeaturedHandler;
1150 let info = handler.describe();
1151 assert_eq!(info.description, "Full-featured test handler");
1152 assert_eq!(info.source_code, Some("fn test() {}".to_string()));
1153 assert_eq!(info.sub_workflows, vec!["helper".to_string()]);
1154 assert_eq!(info.category, Some("data/etl".to_string()));
1155 assert_eq!(info.version, Some("1.2.0".to_string()));
1156 assert!(info.input_schema.is_some());
1157 assert_eq!(info.default_labels.len(), 2);
1158 }
1159
1160 #[test]
1161 fn input_schema_for_generates_json_schema() {
1162 let schema = input_schema_for::<TestInput>();
1163 assert_eq!(schema["type"], "object");
1164 assert!(schema["properties"]["environment"].is_object());
1165 assert!(schema["properties"]["dry_run"].is_object());
1166 }
1167
1168 #[test]
1169 fn input_schema_for_preserves_serde_attributes() {
1170 let schema = input_schema_for::<TestInput>();
1171 let properties = &schema["properties"];
1172 assert!(properties.is_object());
1173 assert!(properties.get("environment").is_some());
1174 assert!(properties.get("dry_run").is_some());
1175 }
1176
1177 #[test]
1178 fn minimal_handler_defaults_to_no_max_cost() {
1179 assert!(MinimalHandler.default_max_cost_usd().is_none());
1180 assert!(MinimalHandler.describe().default_max_cost_usd.is_none());
1181 }
1182
1183 #[test]
1184 fn describe_propagates_handler_max_cost() {
1185 assert_eq!(
1186 FullFeaturedHandler.describe().default_max_cost_usd,
1187 Some(Decimal::new(750, 2))
1188 );
1189 }
1190
1191 #[test]
1192 fn workflow_info_omits_absent_max_cost_from_json() {
1193 let json = serde_json::to_value(MinimalHandler.describe()).expect("serialize");
1194 assert!(json.get("default_max_cost_usd").is_none());
1195 }
1196
1197 #[test]
1198 fn workflow_info_serializes_with_skip_empty() {
1199 let info = WorkflowInfo {
1200 description: "test".to_string(),
1201 source_code: None,
1202 sub_workflows: Vec::new(),
1203 category: None,
1204 version: None,
1205 compatible_versions: Vec::new(),
1206 input_schema: None,
1207 default_labels: HashMap::new(),
1208 schedule: None,
1209 default_max_cost_usd: None,
1210 priority: 0,
1211 };
1212
1213 let json = serde_json::to_value(&info).expect("serialize");
1214 assert_eq!(json["description"], "test");
1215 // Optional fields with skip_serializing_if may still be present or absent
1216 // depending on the serde configuration. Just verify the description is there.
1217 assert!(json.is_object());
1218 }
1219
1220 #[test]
1221 fn workflow_info_serializes_with_values() {
1222 let info = WorkflowInfo {
1223 description: "test".to_string(),
1224 source_code: Some("code".to_string()),
1225 sub_workflows: vec!["sub".to_string()],
1226 category: Some("cat".to_string()),
1227 version: Some("1.0.0".to_string()),
1228 compatible_versions: vec!["0.9.0".to_string()],
1229 input_schema: Some(serde_json::json!({"type": "object"})),
1230 default_labels: HashMap::from([("key".to_string(), "value".to_string())]),
1231 schedule: Some(CronSchedule::new("0 0 * * * *").unwrap()),
1232 default_max_cost_usd: Some(Decimal::new(750, 2)),
1233 priority: 0,
1234 };
1235
1236 let json = serde_json::to_value(&info).expect("serialize");
1237 assert_eq!(json["description"], "test");
1238 assert_eq!(json["source_code"], "code");
1239 assert_eq!(json["sub_workflows"][0], "sub");
1240 assert_eq!(json["category"], "cat");
1241 assert_eq!(json["version"], "1.0.0");
1242 assert_eq!(json["default_labels"]["key"], "value");
1243 assert_eq!(json["schedule"], "0 0 * * * *");
1244 assert_eq!(json["compatible_versions"][0], "0.9.0");
1245 }
1246
1247 // ---- is_version_compatible ----
1248
1249 struct VersionedHandler;
1250
1251 impl WorkflowHandler for VersionedHandler {
1252 fn name(&self) -> &str {
1253 "versioned"
1254 }
1255 fn version(&self) -> Option<&str> {
1256 Some("2.0.0")
1257 }
1258 fn compatible_versions(&self) -> &[&str] {
1259 &["1.5.0", "1.9.0"]
1260 }
1261 fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
1262 Box::pin(async { Ok(()) })
1263 }
1264 }
1265
1266 #[test]
1267 fn version_compatible_with_same_version() {
1268 assert!(VersionedHandler.is_version_compatible(Some("2.0.0")));
1269 }
1270
1271 #[test]
1272 fn version_compatible_with_none_run_version() {
1273 assert!(VersionedHandler.is_version_compatible(None));
1274 }
1275
1276 #[test]
1277 fn version_compatible_with_listed_version() {
1278 assert!(VersionedHandler.is_version_compatible(Some("1.5.0")));
1279 assert!(VersionedHandler.is_version_compatible(Some("1.9.0")));
1280 }
1281
1282 #[test]
1283 fn version_incompatible_with_unlisted_version() {
1284 assert!(!VersionedHandler.is_version_compatible(Some("1.0.0")));
1285 assert!(!VersionedHandler.is_version_compatible(Some("3.0.0")));
1286 }
1287
1288 #[test]
1289 fn minimal_handler_compatible_with_same_default() {
1290 assert!(MinimalHandler.is_version_compatible(Some("1")));
1291 }
1292
1293 #[test]
1294 fn minimal_handler_incompatible_with_different_version() {
1295 assert!(!MinimalHandler.is_version_compatible(Some("2")));
1296 }
1297
1298 // ---- WorkflowHandler::create_run ----
1299
1300 #[tokio::test]
1301 async fn handler_create_run_uses_handler_metadata() {
1302 use ironflow_store::entities::TriggerKind;
1303 use ironflow_store::memory::InMemoryStore;
1304
1305 let store = InMemoryStore::new();
1306
1307 let opts = CreateRunOpts::new().trigger(TriggerKind::Api);
1308 let creation = FullFeaturedHandler
1309 .create_run(&store, opts)
1310 .await
1311 .expect("create_run");
1312 let run = creation.into_run();
1313
1314 assert_eq!(run.workflow_name, "full");
1315 assert_eq!(run.handler_version, Some("1.2.0".to_string()));
1316 assert_eq!(run.max_cost_usd, Some(Decimal::new(750, 2)));
1317 assert_eq!(run.priority, 30);
1318 }
1319
1320 #[tokio::test]
1321 async fn handler_create_run_opts_override_handler_defaults() {
1322 use ironflow_store::memory::InMemoryStore;
1323
1324 let store = InMemoryStore::new();
1325
1326 let opts = CreateRunOpts::new()
1327 .max_cost_usd(Decimal::new(100, 2))
1328 .priority(-5);
1329 let creation = FullFeaturedHandler
1330 .create_run(&store, opts)
1331 .await
1332 .expect("create_run");
1333 let run = creation.into_run();
1334
1335 assert_eq!(run.max_cost_usd, Some(Decimal::new(100, 2)));
1336 assert_eq!(run.priority, -5);
1337 }
1338
1339 #[tokio::test]
1340 async fn handler_create_run_minimal_handler_defaults() {
1341 use ironflow_store::memory::InMemoryStore;
1342
1343 let store = InMemoryStore::new();
1344
1345 let opts = CreateRunOpts::new();
1346 let creation = MinimalHandler
1347 .create_run(&store, opts)
1348 .await
1349 .expect("create_run");
1350 let run = creation.into_run();
1351
1352 assert_eq!(run.workflow_name, "minimal");
1353 assert_eq!(run.handler_version, Some("1".to_string()));
1354 assert_eq!(run.max_cost_usd, None);
1355 assert_eq!(run.priority, 0);
1356 }
1357
1358 struct OutOfRangePriority(i16);
1359
1360 impl WorkflowHandler for OutOfRangePriority {
1361 fn name(&self) -> &str {
1362 "out-of-range-priority"
1363 }
1364
1365 fn priority(&self) -> i16 {
1366 self.0
1367 }
1368
1369 fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
1370 Box::pin(async { Ok(()) })
1371 }
1372 }
1373
1374 #[tokio::test]
1375 async fn handler_create_run_clamps_out_of_range_priority() {
1376 use ironflow_store::memory::InMemoryStore;
1377
1378 let store = InMemoryStore::new();
1379
1380 let high = OutOfRangePriority(i16::MAX)
1381 .create_run(&store, CreateRunOpts::new())
1382 .await
1383 .expect("create_run")
1384 .into_run();
1385 assert_eq!(high.priority, MAX_PRIORITY);
1386
1387 let low = OutOfRangePriority(i16::MIN)
1388 .create_run(&store, CreateRunOpts::new())
1389 .await
1390 .expect("create_run")
1391 .into_run();
1392 assert_eq!(low.priority, MIN_PRIORITY);
1393 }
1394
1395 #[test]
1396 fn workflow_info_priority_is_omitted_from_json_only_when_zero() {
1397 let json = serde_json::to_value(MinimalHandler.describe()).expect("serialize");
1398 assert!(json.get("priority").is_none());
1399
1400 let json = serde_json::to_value(FullFeaturedHandler.describe()).expect("serialize");
1401 assert_eq!(json["priority"], 30);
1402 }
1403
1404 #[test]
1405 fn clamp_priority_keeps_in_range_values() {
1406 assert_eq!(clamp_priority(0), 0);
1407 assert_eq!(clamp_priority(MAX_PRIORITY), MAX_PRIORITY);
1408 assert_eq!(clamp_priority(MIN_PRIORITY), MIN_PRIORITY);
1409 assert_eq!(clamp_priority(MAX_PRIORITY + 1), MAX_PRIORITY);
1410 assert_eq!(clamp_priority(MIN_PRIORITY - 1), MIN_PRIORITY);
1411 }
1412
1413 struct Documented;
1414
1415 impl WorkflowHandler for Documented {
1416 fn name(&self) -> &str {
1417 "documented"
1418 }
1419
1420 fn description(&self) -> &str {
1421 "A documented handler"
1422 }
1423
1424 fn source_code(&self) -> Option<&str> {
1425 Some("struct Documented;")
1426 }
1427
1428 fn sub_workflows(&self) -> Vec<String> {
1429 vec!["child".to_string()]
1430 }
1431
1432 fn category(&self) -> Option<&str> {
1433 Some("tests/handlers")
1434 }
1435
1436 fn version(&self) -> Option<&str> {
1437 Some("3.1.0")
1438 }
1439
1440 fn compatible_versions(&self) -> &[&str] {
1441 &["3.0.0"]
1442 }
1443
1444 fn input_schema(&self) -> Option<Value> {
1445 Some(input_schema_for::<TestInput>())
1446 }
1447
1448 fn default_labels(&self) -> HashMap<String, String> {
1449 HashMap::from([("team".to_string(), "core".to_string())])
1450 }
1451
1452 fn default_max_cost_usd(&self) -> Option<Decimal> {
1453 Some(Decimal::new(250, 2))
1454 }
1455
1456 fn priority(&self) -> i16 {
1457 -15
1458 }
1459
1460 fn guard_config(&self) -> Option<WorkflowGuardConfig> {
1461 Some(WorkflowGuardConfig::new().with_max_depth(4))
1462 }
1463
1464 fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
1465 Box::pin(async { Ok(()) })
1466 }
1467 }
1468
1469 #[test]
1470 fn default_describe_propagates_every_trait_method() {
1471 let info = Documented.describe();
1472 assert_eq!(info.description, "A documented handler");
1473 assert_eq!(info.source_code.as_deref(), Some("struct Documented;"));
1474 assert_eq!(info.sub_workflows, vec!["child".to_string()]);
1475 assert_eq!(info.category.as_deref(), Some("tests/handlers"));
1476 assert_eq!(info.version.as_deref(), Some("3.1.0"));
1477 assert_eq!(info.compatible_versions, vec!["3.0.0".to_string()]);
1478 assert!(info.input_schema.is_some());
1479 assert_eq!(info.default_labels["team"], "core");
1480 assert_eq!(info.default_max_cost_usd, Some(Decimal::new(250, 2)));
1481 assert_eq!(info.priority, -15);
1482 }
1483
1484 #[test]
1485 fn minimal_handler_describe_uses_defaults() {
1486 let info = MinimalHandler.describe();
1487 assert_eq!(info.description, "");
1488 assert!(info.source_code.is_none());
1489 assert!(info.sub_workflows.is_empty());
1490 assert!(info.category.is_none());
1491 assert_eq!(info.version.as_deref(), Some("1"));
1492 }
1493
1494 #[test]
1495 fn workflow_info_builder_sets_every_field() {
1496 let schedule = CronSchedule::new("0 0 * * *").expect("valid cron");
1497 let info = WorkflowInfo::new("desc")
1498 .with_source_code("code")
1499 .with_sub_workflows(["a", "b"])
1500 .with_category("cat/sub")
1501 .with_version("2")
1502 .with_compatible_versions(["1"])
1503 .with_input_schema(serde_json::json!({"type": "object"}))
1504 .with_default_labels(HashMap::from([("k".to_string(), "v".to_string())]))
1505 .with_schedule(schedule)
1506 .with_default_max_cost_usd(Decimal::ONE)
1507 .with_priority(-40);
1508
1509 assert_eq!(info.description, "desc");
1510 assert_eq!(info.source_code.as_deref(), Some("code"));
1511 assert_eq!(info.sub_workflows, vec!["a".to_string(), "b".to_string()]);
1512 assert_eq!(info.category.as_deref(), Some("cat/sub"));
1513 assert_eq!(info.version.as_deref(), Some("2"));
1514 assert_eq!(info.compatible_versions, vec!["1".to_string()]);
1515 assert_eq!(info.input_schema.unwrap()["type"], "object");
1516 assert_eq!(info.default_labels["k"], "v");
1517 assert!(info.schedule.is_some());
1518 assert_eq!(info.default_max_cost_usd, Some(Decimal::ONE));
1519 assert_eq!(info.priority, -40);
1520 }
1521
1522 #[test]
1523 fn workflow_info_new_matches_default_for_other_fields() {
1524 let info = WorkflowInfo::new("only description");
1525 let default = WorkflowInfo::default();
1526 assert_eq!(info.description, "only description");
1527 assert_eq!(default.description, "");
1528 assert_eq!(info.source_code, default.source_code);
1529 assert_eq!(info.sub_workflows, default.sub_workflows);
1530 assert_eq!(info.category, default.category);
1531 assert_eq!(info.version, default.version);
1532 assert_eq!(info.default_max_cost_usd, default.default_max_cost_usd);
1533 assert_eq!(info.priority, default.priority);
1534 }
1535
1536 #[test]
1537 fn boxed_handler_delegates_every_method() {
1538 let boxed: Box<dyn WorkflowHandler> = Box::new(Documented);
1539 assert_eq!(boxed.name(), "documented");
1540 assert_eq!(boxed.version(), Some("3.1.0"));
1541 assert_eq!(boxed.compatible_versions(), &["3.0.0"]);
1542 assert_eq!(boxed.description(), "A documented handler");
1543 assert_eq!(boxed.source_code(), Some("struct Documented;"));
1544 assert_eq!(boxed.sub_workflows(), vec!["child".to_string()]);
1545 assert_eq!(boxed.category(), Some("tests/handlers"));
1546 assert!(boxed.input_schema().is_some());
1547 assert_eq!(boxed.default_labels()["team"], "core");
1548 assert!(boxed.schedule().is_none());
1549 assert_eq!(boxed.default_max_cost_usd(), Some(Decimal::new(250, 2)));
1550 assert_eq!(boxed.priority(), -15);
1551 assert_eq!(boxed.guard_config().map(|g| g.max_depth), Some(4));
1552 assert!(boxed.is_version_compatible(Some("3.0.0")));
1553 assert!(!boxed.is_version_compatible(Some("0.1.0")));
1554 assert_eq!(boxed.describe().description, "A documented handler");
1555 }
1556
1557 #[test]
1558 fn boxed_handler_is_accepted_by_generic_register() {
1559 fn takes_handler(handler: impl WorkflowHandler + 'static) -> String {
1560 handler.name().to_string()
1561 }
1562 let boxed: Box<dyn WorkflowHandler> = Box::new(MinimalHandler);
1563 assert_eq!(takes_handler(boxed), "minimal");
1564 }
1565
1566 #[tokio::test]
1567 async fn boxed_handler_create_run_delegates_metadata() {
1568 use ironflow_store::memory::InMemoryStore;
1569 use ironflow_store::models::TriggerKind;
1570
1571 let store = InMemoryStore::new();
1572 let boxed: Box<dyn WorkflowHandler> = Box::new(Documented);
1573 let run = boxed
1574 .create_run(&store, CreateRunOpts::new().trigger(TriggerKind::Manual))
1575 .await
1576 .expect("run created")
1577 .into_run();
1578 assert_eq!(run.workflow_name, "documented");
1579 assert_eq!(run.handler_version.as_deref(), Some("3.1.0"));
1580 assert_eq!(run.max_cost_usd, Some(Decimal::new(250, 2)));
1581 }
1582}