Skip to main content

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