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