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