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    /// Check whether a run carrying `run_version` can be replayed by this
610    /// handler without `force`.
611    ///
612    /// Compatibility rules:
613    /// - `run_version` is `None` (old run predating version tracking): always
614    ///   compatible.
615    /// - `run_version` equals [`version`](Self::version): compatible.
616    /// - `run_version` appears in [`compatible_versions`](Self::compatible_versions):
617    ///   compatible.
618    /// - Otherwise: incompatible.
619    ///
620    /// # Examples
621    ///
622    /// ```
623    /// # use ironflow_engine::handler::{WorkflowHandler, HandlerFuture};
624    /// # use ironflow_engine::context::WorkflowContext;
625    /// struct MyHandler;
626    ///
627    /// impl WorkflowHandler for MyHandler {
628    ///     fn name(&self) -> &str { "my-handler" }
629    ///     fn version(&self) -> Option<&str> { Some("2.0.0") }
630    ///     fn compatible_versions(&self) -> &[&str] { &["1.0.0"] }
631    ///     fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
632    ///         Box::pin(async move { Ok(()) })
633    ///     }
634    /// }
635    ///
636    /// assert!(MyHandler.is_version_compatible(None));
637    /// assert!(MyHandler.is_version_compatible(Some("2.0.0")));
638    /// assert!(MyHandler.is_version_compatible(Some("1.0.0")));
639    /// assert!(!MyHandler.is_version_compatible(Some("0.5.0")));
640    /// ```
641    fn is_version_compatible(&self, run_version: Option<&str>) -> bool {
642        let Some(rv) = run_version else {
643            return true;
644        };
645        if self.version() == Some(rv) {
646            return true;
647        }
648        self.compatible_versions().contains(&rv)
649    }
650
651    /// Return metadata about this workflow.
652    ///
653    /// The default assembles a [`WorkflowInfo`] from every other trait
654    /// method: [`description`](Self::description),
655    /// [`source_code`](Self::source_code),
656    /// [`sub_workflows`](Self::sub_workflows), [`category`](Self::category),
657    /// [`version`](Self::version),
658    /// [`compatible_versions`](Self::compatible_versions),
659    /// [`input_schema`](Self::input_schema),
660    /// [`default_labels`](Self::default_labels), [`schedule`](Self::schedule)
661    /// and [`default_max_cost_usd`](Self::default_max_cost_usd). Override
662    /// those instead of this method; override `describe` only when the
663    /// metadata cannot be expressed through them.
664    fn describe(&self) -> WorkflowInfo {
665        WorkflowInfo {
666            description: self.description().to_string(),
667            source_code: self.source_code().map(str::to_string),
668            sub_workflows: self.sub_workflows(),
669            category: self.category().map(str::to_string),
670            version: self.version().map(str::to_string),
671            compatible_versions: self
672                .compatible_versions()
673                .iter()
674                .map(|s| s.to_string())
675                .collect(),
676            input_schema: self.input_schema(),
677            default_labels: self.default_labels(),
678            schedule: self.schedule().cloned(),
679            default_max_cost_usd: self.default_max_cost_usd(),
680        }
681    }
682
683    /// Create a run for this workflow, using handler metadata automatically.
684    ///
685    /// Assembles a [`NewRun`](ironflow_store::entities::NewRun) from [`name`](Self::name),
686    /// [`version`](Self::version), and [`default_max_cost_usd`](Self::default_max_cost_usd),
687    /// then delegates to the given [`RunCreator`].
688    ///
689    /// # Errors
690    ///
691    /// Returns [`EngineError`] if the underlying store rejects the run.
692    ///
693    /// # Examples
694    ///
695    /// ```no_run
696    /// # use ironflow_engine::handler::{WorkflowHandler, HandlerFuture};
697    /// # use ironflow_engine::context::WorkflowContext;
698    /// # use ironflow_engine::run_creator::{CreateRunOpts, RunCreator};
699    /// # use ironflow_store::entities::TriggerKind;
700    /// struct DeployWorkflow;
701    ///
702    /// impl WorkflowHandler for DeployWorkflow {
703    ///     fn name(&self) -> &str { "deploy" }
704    ///     fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
705    ///         Box::pin(async move { Ok(()) })
706    ///     }
707    /// }
708    ///
709    /// # async fn example(store: &dyn RunCreator) -> Result<(), ironflow_engine::error::EngineError> {
710    /// let opts = CreateRunOpts::new().trigger(TriggerKind::Api);
711    /// let run = DeployWorkflow.create_run(store, opts).await?.into_run();
712    /// assert_eq!(run.workflow_name, "deploy");
713    /// # Ok(())
714    /// # }
715    /// ```
716    fn create_run<'a>(
717        &self,
718        creator: &'a dyn RunCreator,
719        opts: CreateRunOpts,
720    ) -> RunCreatorFuture<'a> {
721        use tracing::{Instrument, info_span};
722
723        let new_run = opts.build(self.name(), self.version(), self.default_max_cost_usd());
724        let span = info_span!("handler.create_run", workflow = %self.name());
725        Box::pin(creator.create_run(new_run).instrument(span))
726    }
727
728    /// Execute the workflow with the given context.
729    ///
730    /// The context provides [`shell`](WorkflowContext::shell),
731    /// [`http`](WorkflowContext::http), and [`agent`](WorkflowContext::agent)
732    /// methods that automatically persist each step.
733    ///
734    /// # Errors
735    ///
736    /// Return [`EngineError`] if any step fails. The engine will mark
737    /// the run as `Failed` and record the error.
738    fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a>;
739}
740
741/// A boxed handler is a handler.
742///
743/// Lets a single `Vec<Box<dyn WorkflowHandler>>` feed both
744/// [`Engine::register`](crate::engine::Engine::register) and a worker
745/// builder, so the API server and the workers cannot drift apart in the
746/// list of workflows they know.
747///
748/// # Examples
749///
750/// ```
751/// use ironflow_engine::handler::{HandlerFuture, WorkflowHandler};
752/// use ironflow_engine::context::WorkflowContext;
753///
754/// struct Hello;
755///
756/// impl WorkflowHandler for Hello {
757///     fn name(&self) -> &str { "hello" }
758///     fn description(&self) -> &str { "Says hello" }
759///     fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
760///         Box::pin(async move { Ok(()) })
761///     }
762/// }
763///
764/// fn handlers() -> Vec<Box<dyn WorkflowHandler>> {
765///     vec![Box::new(Hello)]
766/// }
767///
768/// for handler in handlers() {
769///     assert_eq!(handler.name(), "hello");
770///     assert_eq!(handler.describe().description, "Says hello");
771/// }
772/// ```
773impl<T: WorkflowHandler + ?Sized> WorkflowHandler for Box<T> {
774    fn name(&self) -> &str {
775        (**self).name()
776    }
777
778    fn version(&self) -> Option<&str> {
779        (**self).version()
780    }
781
782    fn compatible_versions(&self) -> &[&str] {
783        (**self).compatible_versions()
784    }
785
786    fn description(&self) -> &str {
787        (**self).description()
788    }
789
790    fn source_code(&self) -> Option<&str> {
791        (**self).source_code()
792    }
793
794    fn sub_workflows(&self) -> Vec<String> {
795        (**self).sub_workflows()
796    }
797
798    fn category(&self) -> Option<&str> {
799        (**self).category()
800    }
801
802    fn input_schema(&self) -> Option<Value> {
803        (**self).input_schema()
804    }
805
806    fn default_labels(&self) -> HashMap<String, String> {
807        (**self).default_labels()
808    }
809
810    fn schedule(&self) -> Option<&CronSchedule> {
811        (**self).schedule()
812    }
813
814    fn default_max_cost_usd(&self) -> Option<Decimal> {
815        (**self).default_max_cost_usd()
816    }
817
818    fn guard_config(&self) -> Option<WorkflowGuardConfig> {
819        (**self).guard_config()
820    }
821
822    fn is_version_compatible(&self, run_version: Option<&str>) -> bool {
823        (**self).is_version_compatible(run_version)
824    }
825
826    fn describe(&self) -> WorkflowInfo {
827        (**self).describe()
828    }
829
830    fn create_run<'a>(
831        &self,
832        creator: &'a dyn RunCreator,
833        opts: CreateRunOpts,
834    ) -> RunCreatorFuture<'a> {
835        (**self).create_run(creator, opts)
836    }
837
838    fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
839        (**self).execute(ctx)
840    }
841}
842
843#[cfg(test)]
844mod tests {
845    use super::*;
846    use serde::{Deserialize, Serialize};
847
848    #[derive(Debug, Serialize, Deserialize, JsonSchema)]
849    struct TestInput {
850        environment: String,
851        #[serde(default)]
852        dry_run: bool,
853    }
854
855    struct MinimalHandler;
856
857    impl WorkflowHandler for MinimalHandler {
858        fn name(&self) -> &str {
859            "minimal"
860        }
861
862        fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
863            Box::pin(async { Ok(()) })
864        }
865    }
866
867    struct FullFeaturedHandler;
868
869    impl WorkflowHandler for FullFeaturedHandler {
870        fn name(&self) -> &str {
871            "full"
872        }
873
874        fn version(&self) -> Option<&str> {
875            Some("1.2.0")
876        }
877
878        fn category(&self) -> Option<&str> {
879            Some("data/etl")
880        }
881
882        fn input_schema(&self) -> Option<Value> {
883            Some(input_schema_for::<TestInput>())
884        }
885
886        fn default_labels(&self) -> HashMap<String, String> {
887            HashMap::from([
888                ("team".to_string(), "platform".to_string()),
889                ("env".to_string(), "prod".to_string()),
890            ])
891        }
892
893        fn default_max_cost_usd(&self) -> Option<Decimal> {
894            Some(Decimal::new(750, 2))
895        }
896
897        fn describe(&self) -> WorkflowInfo {
898            WorkflowInfo {
899                description: "Full-featured test handler".to_string(),
900                source_code: Some("fn test() {}".to_string()),
901                sub_workflows: vec!["helper".to_string()],
902                category: self.category().map(str::to_string),
903                version: self.version().map(str::to_string),
904                compatible_versions: self
905                    .compatible_versions()
906                    .iter()
907                    .map(|s| s.to_string())
908                    .collect(),
909                input_schema: self.input_schema(),
910                default_labels: self.default_labels(),
911                schedule: self.schedule().cloned(),
912                default_max_cost_usd: self.default_max_cost_usd(),
913            }
914        }
915
916        fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
917            Box::pin(async { Ok(()) })
918        }
919    }
920
921    #[test]
922    fn minimal_handler_has_required_name() {
923        let handler = MinimalHandler;
924        assert_eq!(handler.name(), "minimal");
925    }
926
927    #[test]
928    fn minimal_handler_defaults_to_version_1() {
929        let handler = MinimalHandler;
930        assert_eq!(handler.version(), Some("1"));
931    }
932
933    #[test]
934    fn minimal_handler_defaults_to_no_compatible_versions() {
935        let handler = MinimalHandler;
936        assert!(handler.compatible_versions().is_empty());
937    }
938
939    #[test]
940    fn minimal_handler_defaults_to_no_category() {
941        let handler = MinimalHandler;
942        assert_eq!(handler.category(), None);
943    }
944
945    #[test]
946    fn minimal_handler_defaults_to_no_schema() {
947        let handler = MinimalHandler;
948        assert_eq!(handler.input_schema(), None);
949    }
950
951    #[test]
952    fn minimal_handler_defaults_to_empty_labels() {
953        let handler = MinimalHandler;
954        let labels = handler.default_labels();
955        assert!(labels.is_empty());
956    }
957
958    #[test]
959    fn minimal_handler_defaults_to_no_schedule() {
960        let handler = MinimalHandler;
961        assert_eq!(handler.schedule(), None);
962    }
963
964    #[test]
965    fn minimal_handler_describe_reflects_defaults() {
966        let handler = MinimalHandler;
967        let info = handler.describe();
968        assert_eq!(info.description, "");
969        assert_eq!(info.source_code, None);
970        assert_eq!(info.sub_workflows, Vec::<String>::new());
971        assert_eq!(info.category, None);
972        assert_eq!(info.version, Some("1".to_string()));
973        assert!(info.compatible_versions.is_empty());
974        assert_eq!(info.input_schema, None);
975        assert!(info.default_labels.is_empty());
976        assert_eq!(info.schedule, None);
977    }
978
979    #[test]
980    fn full_handler_returns_all_metadata() {
981        let handler = FullFeaturedHandler;
982        assert_eq!(handler.name(), "full");
983        assert_eq!(handler.version(), Some("1.2.0"));
984        assert_eq!(handler.category(), Some("data/etl"));
985        assert!(handler.input_schema().is_some());
986    }
987
988    #[test]
989    fn full_handler_default_labels_are_set() {
990        let handler = FullFeaturedHandler;
991        let labels = handler.default_labels();
992        assert_eq!(labels.get("team"), Some(&"platform".to_string()));
993        assert_eq!(labels.get("env"), Some(&"prod".to_string()));
994    }
995
996    #[test]
997    fn full_handler_describe_includes_all_fields() {
998        let handler = FullFeaturedHandler;
999        let info = handler.describe();
1000        assert_eq!(info.description, "Full-featured test handler");
1001        assert_eq!(info.source_code, Some("fn test() {}".to_string()));
1002        assert_eq!(info.sub_workflows, vec!["helper".to_string()]);
1003        assert_eq!(info.category, Some("data/etl".to_string()));
1004        assert_eq!(info.version, Some("1.2.0".to_string()));
1005        assert!(info.input_schema.is_some());
1006        assert_eq!(info.default_labels.len(), 2);
1007    }
1008
1009    #[test]
1010    fn input_schema_for_generates_json_schema() {
1011        let schema = input_schema_for::<TestInput>();
1012        assert_eq!(schema["type"], "object");
1013        assert!(schema["properties"]["environment"].is_object());
1014        assert!(schema["properties"]["dry_run"].is_object());
1015    }
1016
1017    #[test]
1018    fn input_schema_for_preserves_serde_attributes() {
1019        let schema = input_schema_for::<TestInput>();
1020        let properties = &schema["properties"];
1021        assert!(properties.is_object());
1022        assert!(properties.get("environment").is_some());
1023        assert!(properties.get("dry_run").is_some());
1024    }
1025
1026    #[test]
1027    fn minimal_handler_defaults_to_no_max_cost() {
1028        assert!(MinimalHandler.default_max_cost_usd().is_none());
1029        assert!(MinimalHandler.describe().default_max_cost_usd.is_none());
1030    }
1031
1032    #[test]
1033    fn describe_propagates_handler_max_cost() {
1034        assert_eq!(
1035            FullFeaturedHandler.describe().default_max_cost_usd,
1036            Some(Decimal::new(750, 2))
1037        );
1038    }
1039
1040    #[test]
1041    fn workflow_info_omits_absent_max_cost_from_json() {
1042        let json = serde_json::to_value(MinimalHandler.describe()).expect("serialize");
1043        assert!(json.get("default_max_cost_usd").is_none());
1044    }
1045
1046    #[test]
1047    fn workflow_info_serializes_with_skip_empty() {
1048        let info = WorkflowInfo {
1049            description: "test".to_string(),
1050            source_code: None,
1051            sub_workflows: Vec::new(),
1052            category: None,
1053            version: None,
1054            compatible_versions: Vec::new(),
1055            input_schema: None,
1056            default_labels: HashMap::new(),
1057            schedule: None,
1058            default_max_cost_usd: None,
1059        };
1060
1061        let json = serde_json::to_value(&info).expect("serialize");
1062        assert_eq!(json["description"], "test");
1063        // Optional fields with skip_serializing_if may still be present or absent
1064        // depending on the serde configuration. Just verify the description is there.
1065        assert!(json.is_object());
1066    }
1067
1068    #[test]
1069    fn workflow_info_serializes_with_values() {
1070        let info = WorkflowInfo {
1071            description: "test".to_string(),
1072            source_code: Some("code".to_string()),
1073            sub_workflows: vec!["sub".to_string()],
1074            category: Some("cat".to_string()),
1075            version: Some("1.0.0".to_string()),
1076            compatible_versions: vec!["0.9.0".to_string()],
1077            input_schema: Some(serde_json::json!({"type": "object"})),
1078            default_labels: HashMap::from([("key".to_string(), "value".to_string())]),
1079            schedule: Some(CronSchedule::new("0 0 * * * *").unwrap()),
1080            default_max_cost_usd: Some(Decimal::new(750, 2)),
1081        };
1082
1083        let json = serde_json::to_value(&info).expect("serialize");
1084        assert_eq!(json["description"], "test");
1085        assert_eq!(json["source_code"], "code");
1086        assert_eq!(json["sub_workflows"][0], "sub");
1087        assert_eq!(json["category"], "cat");
1088        assert_eq!(json["version"], "1.0.0");
1089        assert_eq!(json["default_labels"]["key"], "value");
1090        assert_eq!(json["schedule"], "0 0 * * * *");
1091        assert_eq!(json["compatible_versions"][0], "0.9.0");
1092    }
1093
1094    // ---- is_version_compatible ----
1095
1096    struct VersionedHandler;
1097
1098    impl WorkflowHandler for VersionedHandler {
1099        fn name(&self) -> &str {
1100            "versioned"
1101        }
1102        fn version(&self) -> Option<&str> {
1103            Some("2.0.0")
1104        }
1105        fn compatible_versions(&self) -> &[&str] {
1106            &["1.5.0", "1.9.0"]
1107        }
1108        fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
1109            Box::pin(async { Ok(()) })
1110        }
1111    }
1112
1113    #[test]
1114    fn version_compatible_with_same_version() {
1115        assert!(VersionedHandler.is_version_compatible(Some("2.0.0")));
1116    }
1117
1118    #[test]
1119    fn version_compatible_with_none_run_version() {
1120        assert!(VersionedHandler.is_version_compatible(None));
1121    }
1122
1123    #[test]
1124    fn version_compatible_with_listed_version() {
1125        assert!(VersionedHandler.is_version_compatible(Some("1.5.0")));
1126        assert!(VersionedHandler.is_version_compatible(Some("1.9.0")));
1127    }
1128
1129    #[test]
1130    fn version_incompatible_with_unlisted_version() {
1131        assert!(!VersionedHandler.is_version_compatible(Some("1.0.0")));
1132        assert!(!VersionedHandler.is_version_compatible(Some("3.0.0")));
1133    }
1134
1135    #[test]
1136    fn minimal_handler_compatible_with_same_default() {
1137        assert!(MinimalHandler.is_version_compatible(Some("1")));
1138    }
1139
1140    #[test]
1141    fn minimal_handler_incompatible_with_different_version() {
1142        assert!(!MinimalHandler.is_version_compatible(Some("2")));
1143    }
1144
1145    // ---- WorkflowHandler::create_run ----
1146
1147    #[tokio::test]
1148    async fn handler_create_run_uses_handler_metadata() {
1149        use ironflow_store::entities::TriggerKind;
1150        use ironflow_store::memory::InMemoryStore;
1151
1152        let store = InMemoryStore::new();
1153
1154        let opts = CreateRunOpts::new().trigger(TriggerKind::Api);
1155        let creation = FullFeaturedHandler
1156            .create_run(&store, opts)
1157            .await
1158            .expect("create_run");
1159        let run = creation.into_run();
1160
1161        assert_eq!(run.workflow_name, "full");
1162        assert_eq!(run.handler_version, Some("1.2.0".to_string()));
1163        assert_eq!(run.max_cost_usd, Some(Decimal::new(750, 2)));
1164    }
1165
1166    #[tokio::test]
1167    async fn handler_create_run_opts_override_handler_defaults() {
1168        use ironflow_store::memory::InMemoryStore;
1169
1170        let store = InMemoryStore::new();
1171
1172        let opts = CreateRunOpts::new().max_cost_usd(Decimal::new(100, 2));
1173        let creation = FullFeaturedHandler
1174            .create_run(&store, opts)
1175            .await
1176            .expect("create_run");
1177        let run = creation.into_run();
1178
1179        assert_eq!(run.max_cost_usd, Some(Decimal::new(100, 2)));
1180    }
1181
1182    #[tokio::test]
1183    async fn handler_create_run_minimal_handler_defaults() {
1184        use ironflow_store::memory::InMemoryStore;
1185
1186        let store = InMemoryStore::new();
1187
1188        let opts = CreateRunOpts::new();
1189        let creation = MinimalHandler
1190            .create_run(&store, opts)
1191            .await
1192            .expect("create_run");
1193        let run = creation.into_run();
1194
1195        assert_eq!(run.workflow_name, "minimal");
1196        assert_eq!(run.handler_version, Some("1".to_string()));
1197        assert_eq!(run.max_cost_usd, None);
1198    }
1199
1200    struct Documented;
1201
1202    impl WorkflowHandler for Documented {
1203        fn name(&self) -> &str {
1204            "documented"
1205        }
1206
1207        fn description(&self) -> &str {
1208            "A documented handler"
1209        }
1210
1211        fn source_code(&self) -> Option<&str> {
1212            Some("struct Documented;")
1213        }
1214
1215        fn sub_workflows(&self) -> Vec<String> {
1216            vec!["child".to_string()]
1217        }
1218
1219        fn category(&self) -> Option<&str> {
1220            Some("tests/handlers")
1221        }
1222
1223        fn version(&self) -> Option<&str> {
1224            Some("3.1.0")
1225        }
1226
1227        fn compatible_versions(&self) -> &[&str] {
1228            &["3.0.0"]
1229        }
1230
1231        fn input_schema(&self) -> Option<Value> {
1232            Some(input_schema_for::<TestInput>())
1233        }
1234
1235        fn default_labels(&self) -> HashMap<String, String> {
1236            HashMap::from([("team".to_string(), "core".to_string())])
1237        }
1238
1239        fn default_max_cost_usd(&self) -> Option<Decimal> {
1240            Some(Decimal::new(250, 2))
1241        }
1242
1243        fn guard_config(&self) -> Option<WorkflowGuardConfig> {
1244            Some(WorkflowGuardConfig::new().with_max_depth(4))
1245        }
1246
1247        fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
1248            Box::pin(async { Ok(()) })
1249        }
1250    }
1251
1252    #[test]
1253    fn default_describe_propagates_every_trait_method() {
1254        let info = Documented.describe();
1255        assert_eq!(info.description, "A documented handler");
1256        assert_eq!(info.source_code.as_deref(), Some("struct Documented;"));
1257        assert_eq!(info.sub_workflows, vec!["child".to_string()]);
1258        assert_eq!(info.category.as_deref(), Some("tests/handlers"));
1259        assert_eq!(info.version.as_deref(), Some("3.1.0"));
1260        assert_eq!(info.compatible_versions, vec!["3.0.0".to_string()]);
1261        assert!(info.input_schema.is_some());
1262        assert_eq!(info.default_labels["team"], "core");
1263        assert_eq!(info.default_max_cost_usd, Some(Decimal::new(250, 2)));
1264    }
1265
1266    #[test]
1267    fn minimal_handler_describe_uses_defaults() {
1268        let info = MinimalHandler.describe();
1269        assert_eq!(info.description, "");
1270        assert!(info.source_code.is_none());
1271        assert!(info.sub_workflows.is_empty());
1272        assert!(info.category.is_none());
1273        assert_eq!(info.version.as_deref(), Some("1"));
1274    }
1275
1276    #[test]
1277    fn workflow_info_builder_sets_every_field() {
1278        let schedule = CronSchedule::new("0 0 * * *").expect("valid cron");
1279        let info = WorkflowInfo::new("desc")
1280            .with_source_code("code")
1281            .with_sub_workflows(["a", "b"])
1282            .with_category("cat/sub")
1283            .with_version("2")
1284            .with_compatible_versions(["1"])
1285            .with_input_schema(serde_json::json!({"type": "object"}))
1286            .with_default_labels(HashMap::from([("k".to_string(), "v".to_string())]))
1287            .with_schedule(schedule)
1288            .with_default_max_cost_usd(Decimal::ONE);
1289
1290        assert_eq!(info.description, "desc");
1291        assert_eq!(info.source_code.as_deref(), Some("code"));
1292        assert_eq!(info.sub_workflows, vec!["a".to_string(), "b".to_string()]);
1293        assert_eq!(info.category.as_deref(), Some("cat/sub"));
1294        assert_eq!(info.version.as_deref(), Some("2"));
1295        assert_eq!(info.compatible_versions, vec!["1".to_string()]);
1296        assert_eq!(info.input_schema.unwrap()["type"], "object");
1297        assert_eq!(info.default_labels["k"], "v");
1298        assert!(info.schedule.is_some());
1299        assert_eq!(info.default_max_cost_usd, Some(Decimal::ONE));
1300    }
1301
1302    #[test]
1303    fn workflow_info_new_matches_default_for_other_fields() {
1304        let info = WorkflowInfo::new("only description");
1305        let default = WorkflowInfo::default();
1306        assert_eq!(info.description, "only description");
1307        assert_eq!(default.description, "");
1308        assert_eq!(info.source_code, default.source_code);
1309        assert_eq!(info.sub_workflows, default.sub_workflows);
1310        assert_eq!(info.category, default.category);
1311        assert_eq!(info.version, default.version);
1312        assert_eq!(info.default_max_cost_usd, default.default_max_cost_usd);
1313    }
1314
1315    #[test]
1316    fn boxed_handler_delegates_every_method() {
1317        let boxed: Box<dyn WorkflowHandler> = Box::new(Documented);
1318        assert_eq!(boxed.name(), "documented");
1319        assert_eq!(boxed.version(), Some("3.1.0"));
1320        assert_eq!(boxed.compatible_versions(), &["3.0.0"]);
1321        assert_eq!(boxed.description(), "A documented handler");
1322        assert_eq!(boxed.source_code(), Some("struct Documented;"));
1323        assert_eq!(boxed.sub_workflows(), vec!["child".to_string()]);
1324        assert_eq!(boxed.category(), Some("tests/handlers"));
1325        assert!(boxed.input_schema().is_some());
1326        assert_eq!(boxed.default_labels()["team"], "core");
1327        assert!(boxed.schedule().is_none());
1328        assert_eq!(boxed.default_max_cost_usd(), Some(Decimal::new(250, 2)));
1329        assert_eq!(boxed.guard_config().map(|g| g.max_depth), Some(4));
1330        assert!(boxed.is_version_compatible(Some("3.0.0")));
1331        assert!(!boxed.is_version_compatible(Some("0.1.0")));
1332        assert_eq!(boxed.describe().description, "A documented handler");
1333    }
1334
1335    #[test]
1336    fn boxed_handler_is_accepted_by_generic_register() {
1337        fn takes_handler(handler: impl WorkflowHandler + 'static) -> String {
1338            handler.name().to_string()
1339        }
1340        let boxed: Box<dyn WorkflowHandler> = Box::new(MinimalHandler);
1341        assert_eq!(takes_handler(boxed), "minimal");
1342    }
1343
1344    #[tokio::test]
1345    async fn boxed_handler_create_run_delegates_metadata() {
1346        use ironflow_store::memory::InMemoryStore;
1347        use ironflow_store::models::TriggerKind;
1348
1349        let store = InMemoryStore::new();
1350        let boxed: Box<dyn WorkflowHandler> = Box::new(Documented);
1351        let run = boxed
1352            .create_run(&store, CreateRunOpts::new().trigger(TriggerKind::Manual))
1353            .await
1354            .expect("run created")
1355            .into_run();
1356        assert_eq!(run.workflow_name, "documented");
1357        assert_eq!(run.handler_version.as_deref(), Some("3.1.0"));
1358        assert_eq!(run.max_cost_usd, Some(Decimal::new(250, 2)));
1359    }
1360}