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