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