Skip to main content

ironflow_engine/
context.rs

1//! [`WorkflowContext`] — execution context for dynamic workflows.
2//!
3//! Provides step execution methods that automatically persist results to the
4//! store. Each call to [`shell`](WorkflowContext::shell),
5//! [`http`](WorkflowContext::http), [`agent`](WorkflowContext::agent), or
6//! [`workflow`](WorkflowContext::workflow) creates a step record, executes the
7//! operation, captures the output, and returns a
8//! [`StepOutput`](crate::executor::StepOutput) that the next step can
9//! reference.
10//!
11//! # Examples
12//!
13//! ```no_run
14//! use ironflow_engine::context::WorkflowContext;
15//! use ironflow_engine::config::{ShellConfig, AgentStepConfig};
16//! use ironflow_engine::error::EngineError;
17//!
18//! # async fn example(ctx: &mut WorkflowContext) -> Result<(), EngineError> {
19//! let build = ctx.shell("build", ShellConfig::new("cargo build")).await?;
20//! let review = ctx.agent("review", AgentStepConfig::new(
21//!     &format!("Build output:\n{}", build.output["stdout"])
22//! )).await?;
23//! # Ok(())
24//! # }
25//! ```
26
27mod accessors;
28mod artifacts;
29mod error_handlers;
30mod failure;
31mod guard;
32mod lifecycle;
33mod steps;
34
35#[cfg(test)]
36mod tests;
37
38use std::collections::HashMap;
39use std::fmt;
40use std::sync::Arc;
41
42use rust_decimal::Decimal;
43use uuid::Uuid;
44
45use ironflow_core::decision::DecisionProvider;
46use ironflow_core::provider::AgentProvider;
47use ironflow_core::trace_context::WorkflowTraceContext;
48use ironflow_store::models::Step;
49use ironflow_store::store::Store;
50
51use crate::artifact::ArtifactSink;
52use crate::config::StepConfig;
53use crate::executor::StepResult;
54use crate::guard::{SharedGuardState, WorkflowGuardConfig};
55use crate::handler::WorkflowHandler;
56use crate::log_sender::LogSender;
57use crate::notify::WorkflowEventBus;
58use crate::operation::OperationContext;
59
60/// Callback type for resolving workflow handlers by name.
61pub(crate) type HandlerResolver =
62    Arc<dyn Fn(&str) -> Option<Arc<dyn WorkflowHandler>> + Send + Sync>;
63
64/// Execution context for a single workflow run.
65///
66/// Tracks the current step position and provides convenience methods
67/// for executing operations with automatic persistence.
68///
69/// # Examples
70///
71/// ```no_run
72/// use ironflow_engine::context::WorkflowContext;
73/// use ironflow_engine::config::ShellConfig;
74/// use ironflow_engine::error::EngineError;
75///
76/// # async fn example(ctx: &mut WorkflowContext) -> Result<(), EngineError> {
77/// let result = ctx.shell("greet", ShellConfig::new("echo hello")).await?;
78/// assert!(result.output["stdout"].as_str().unwrap().contains("hello"));
79/// # Ok(())
80/// # }
81/// ```
82pub struct WorkflowContext {
83    run_id: Uuid,
84    workflow_name: String,
85    store: Arc<dyn Store>,
86    provider: Arc<dyn AgentProvider>,
87    /// Optional decision backend (System One / Jev) for `ctx.decision(...)`.
88    /// `None` when no decision provider was wired: a decision step then fails
89    /// explicitly instead of silently doing nothing.
90    decision_provider: Option<Arc<dyn DecisionProvider>>,
91    handler_resolver: Option<HandlerResolver>,
92    position: u32,
93    /// IDs of the last executed step(s) -- used to record DAG dependencies.
94    last_step_ids: Vec<Uuid>,
95    /// Accumulated cost across all steps in this run.
96    total_cost_usd: Decimal,
97    /// Accumulated duration across all steps.
98    total_duration_ms: u64,
99    /// Cumulative cost cap for this run, resolved at creation. `None` = no cap.
100    max_cost_usd: Option<Decimal>,
101    /// Cost already spent by ancestor runs when this context belongs to a
102    /// sub-workflow. Zero for a top-level run.
103    inherited_cost_usd: Decimal,
104    /// Steps from a previous execution of the *same* attempt, keyed by position.
105    /// Used when resuming after approval to replay completed steps.
106    replay_steps: HashMap<u32, Step>,
107    /// Approvals granted in an *earlier* attempt, keyed by position, holding the
108    /// attempt that granted them. An approval is carried by the run, not by the
109    /// attempt, so a retry never asks a human to approve the same gate twice.
110    granted_approvals: HashMap<u32, u32>,
111    /// Which run attempt this context is executing (1-based).
112    attempt: u32,
113    /// Wall-clock duration already recorded on the run by previous attempts.
114    /// Added to this attempt's duration when the run is finalized.
115    carried_duration_ms: u64,
116    /// Optional sender for real-time log streaming.
117    log_sender: Option<LogSender>,
118    /// Where artifact bytes are read and written. `None` when no artifact
119    /// storage is configured: steps that declare artifacts then fail explicitly
120    /// instead of silently dropping their files.
121    artifact_sink: Option<Arc<dyn ArtifactSink>>,
122    /// Set to `true` when at least one `allow_failure` step failed.
123    has_allowed_failure: bool,
124    /// Error handlers registered via [`on_error`](Self::on_error).
125    error_handlers: Vec<OnErrorHandler>,
126    /// Shared guard state for workflow execution limits.
127    guard_state: Option<SharedGuardState>,
128    /// Guard configuration for this workflow run.
129    guard_config: Option<WorkflowGuardConfig>,
130    /// Accumulated step results for post-execution inspection.
131    step_results: Vec<StepResult>,
132    /// Optional event bus for per-run real-time monitoring.
133    event_bus: Option<WorkflowEventBus>,
134    /// W3C trace context for distributed tracing propagation.
135    trace_context: WorkflowTraceContext,
136    /// Shared operation context for custom operations.
137    operation_ctx: Option<OperationContext>,
138}
139
140/// A registered error handler that fires when a subsequent step fails.
141struct OnErrorHandler {
142    name: String,
143    config: StepConfig,
144}
145
146impl fmt::Debug for WorkflowContext {
147    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
148        f.debug_struct("WorkflowContext")
149            .field("run_id", &self.run_id)
150            .field("position", &self.position)
151            .field("total_cost_usd", &self.total_cost_usd)
152            .field("inherited_cost_usd", &self.inherited_cost_usd)
153            .field("max_cost_usd", &self.max_cost_usd)
154            .finish_non_exhaustive()
155    }
156}