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.stdout())
22//! )).await?;
23//! # Ok(())
24//! # }
25//! ```
26
27mod accessors;
28mod artifacts;
29mod error_handlers;
30mod failure;
31mod guard;
32mod lifecycle;
33mod session;
34mod steps;
35
36#[cfg(test)]
37mod tests;
38
39pub(crate) use lifecycle::interrupt_running_steps;
40
41use std::collections::HashMap;
42use std::fmt;
43use std::sync::Arc;
44
45use chrono::{DateTime, Utc};
46use rust_decimal::Decimal;
47use serde_json::Value;
48use uuid::Uuid;
49
50use ironflow_core::decision::DecisionProvider;
51use ironflow_core::provider::AgentProvider;
52use ironflow_core::trace_context::WorkflowTraceContext;
53use ironflow_store::models::Step;
54use ironflow_store::store::Store;
55
56use crate::artifact::ArtifactSink;
57use crate::config::StepConfig;
58use crate::executor::{StepInterceptor, StepResult};
59use crate::guard::{SharedGuardState, WorkflowGuardConfig};
60use crate::handler::WorkflowHandler;
61use crate::log_sender::LogSender;
62use crate::notify::WorkflowEventBus;
63use crate::operation::OperationContext;
64use crate::plan::SharedPlanRecorder;
65
66/// Label set on every child run of a sub-workflow step, holding the id of the
67/// run that started it. Defined by the store, which follows it to list the
68/// descendants of a run; re-exported here where workflow code looks for it.
69///
70/// The root of the chain is recorded under
71/// [`LABEL_ROOT_RUN_ID`](ironflow_core::provider::LABEL_ROOT_RUN_ID).
72///
73/// # Examples
74///
75/// ```no_run
76/// use std::collections::HashMap;
77/// use ironflow_engine::context::PARENT_RUN_ID_LABEL;
78/// use ironflow_store::models::RunFilter;
79/// use uuid::Uuid;
80///
81/// # fn example(parent: Uuid) {
82/// let children = RunFilter {
83///     labels: Some(HashMap::from([(PARENT_RUN_ID_LABEL.to_string(), parent.to_string())])),
84///     ..RunFilter::default()
85/// };
86/// # }
87/// ```
88pub use ironflow_store::entities::PARENT_RUN_ID_LABEL;
89
90pub use steps::agent::AgentReply;
91
92/// Callback type for resolving workflow handlers by name.
93pub(crate) type HandlerResolver =
94    Arc<dyn Fn(&str) -> Option<Arc<dyn WorkflowHandler>> + Send + Sync>;
95
96/// Execution context for a single workflow run.
97///
98/// Tracks the current step position and provides convenience methods
99/// for executing operations with automatic persistence.
100///
101/// # Examples
102///
103/// ```no_run
104/// use ironflow_engine::context::WorkflowContext;
105/// use ironflow_engine::config::ShellConfig;
106/// use ironflow_engine::error::EngineError;
107///
108/// # async fn example(ctx: &mut WorkflowContext) -> Result<(), EngineError> {
109/// let result = ctx.shell("greet", ShellConfig::new("echo hello")).await?;
110/// assert!(result.stdout().contains("hello"));
111/// # Ok(())
112/// # }
113/// ```
114pub struct WorkflowContext {
115    run_id: Uuid,
116    /// The top-level run: `run_id` itself, or the parent's root for a
117    /// sub-workflow. Stamped on agent pods so a retry releases the children.
118    root_run_id: Uuid,
119    workflow_name: String,
120    store: Arc<dyn Store>,
121    provider: Arc<dyn AgentProvider>,
122    /// Optional decision backend (System One / Jev) for `ctx.decision(...)`.
123    /// `None` when no decision provider was wired: a decision step then fails
124    /// explicitly instead of silently doing nothing.
125    decision_provider: Option<Arc<dyn DecisionProvider>>,
126    handler_resolver: Option<HandlerResolver>,
127    position: u32,
128    /// IDs of the last executed step(s) -- used to record DAG dependencies.
129    last_step_ids: Vec<Uuid>,
130    /// Accumulated cost across all steps in this run.
131    total_cost_usd: Decimal,
132    /// Accumulated duration across all steps.
133    total_duration_ms: u64,
134    /// Cumulative cost cap for this run, resolved at creation. `None` = no cap.
135    max_cost_usd: Option<Decimal>,
136    /// Cost already spent by ancestor runs when this context belongs to a
137    /// sub-workflow. Zero for a top-level run.
138    inherited_cost_usd: Decimal,
139    /// Steps from a previous execution of the *same* attempt, keyed by position.
140    /// Used when resuming after approval to replay completed steps.
141    replay_steps: HashMap<u32, Step>,
142    /// All steps of a previous execution of the *same* attempt, keyed by
143    /// `(position, step name)`. A `parallel` wave shares one position across
144    /// several steps, which `replay_steps` cannot represent -- this index lets
145    /// `parallel()` check that every step of a wave already completed before
146    /// replaying the whole wave from the store, without re-running any item.
147    replay_wave_steps: HashMap<(u32, String), Step>,
148    /// Approvals granted in an *earlier* attempt, keyed by position, holding the
149    /// attempt that granted them. An approval is carried by the run, not by the
150    /// attempt, so a retry never asks a human to approve the same gate twice.
151    granted_approvals: HashMap<u32, u32>,
152    /// Human inputs answered in an *earlier* attempt, keyed by position:
153    /// (attempt, answer). Like an approval, an answer is carried by the run, so
154    /// a retry never asks a human to answer the same input twice.
155    answered_inputs: HashMap<u32, (u32, Value)>,
156    /// `Workflow` steps of the current attempt interrupted by a lost worker
157    /// lease, keyed by position. The step executed again at that position
158    /// re-enters the child run they recorded instead of starting a new one.
159    interrupted_children: HashMap<u32, Step>,
160    /// Which run attempt this context is executing (1-based).
161    attempt: u32,
162    /// Wall-clock duration already recorded on the run by previous attempts.
163    /// Added to this attempt's duration when the run is finalized.
164    carried_duration_ms: u64,
165    /// Optional sender for real-time log streaming.
166    log_sender: Option<LogSender>,
167    /// Where artifact bytes are read and written. `None` when no artifact
168    /// storage is configured: steps that declare artifacts then fail explicitly
169    /// instead of silently dropping their files.
170    artifact_sink: Option<Arc<dyn ArtifactSink>>,
171    /// Set to `true` when at least one `allow_failure` step failed.
172    has_allowed_failure: bool,
173    /// Error handlers registered via [`on_error`](Self::on_error).
174    error_handlers: Vec<OnErrorHandler>,
175    /// Shared guard state for workflow execution limits.
176    guard_state: Option<SharedGuardState>,
177    /// Guard configuration for this workflow run.
178    guard_config: Option<WorkflowGuardConfig>,
179    /// Accumulated step results for post-execution inspection.
180    step_results: Vec<StepResult>,
181    /// Optional event bus for per-run real-time monitoring.
182    event_bus: Option<WorkflowEventBus>,
183    /// Optional hook that resolves steps without executing them. `None` in
184    /// production; set by [`crate::testing::TestEngine`].
185    interceptor: Option<Arc<dyn StepInterceptor>>,
186    /// W3C trace context for distributed tracing propagation.
187    trace_context: WorkflowTraceContext,
188    /// Shared operation context for custom operations.
189    operation_ctx: Option<OperationContext>,
190    /// When the run was created, set by the engine. Bounds the signals a wait
191    /// step accepts: a signal received before the run existed is not for it.
192    /// `None` falls back to reading the run from the store.
193    run_created_at: Option<DateTime<Utc>>,
194    /// Set when the context is recording an execution plan instead of running.
195    /// Every step method checks this first and records intent without executing.
196    plan: Option<SharedPlanRecorder>,
197    /// Output set by the handler with [`set_output`](Self::set_output),
198    /// persisted on the run when the execution ends.
199    output: Option<Value>,
200    /// Worker tags carried by the process running this context, set by the
201    /// engine of a tagged worker. A sub-workflow requiring a tag missing from
202    /// this list is refused. `None` outside a worker: no check.
203    worker_tags: Option<Arc<Vec<String>>>,
204}
205
206/// A registered error handler that fires when a subsequent step fails.
207struct OnErrorHandler {
208    name: String,
209    config: StepConfig,
210}
211
212impl fmt::Debug for WorkflowContext {
213    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
214        f.debug_struct("WorkflowContext")
215            .field("run_id", &self.run_id)
216            .field("position", &self.position)
217            .field("total_cost_usd", &self.total_cost_usd)
218            .field("inherited_cost_usd", &self.inherited_cost_usd)
219            .field("max_cost_usd", &self.max_cost_usd)
220            .field("planning", &self.plan.is_some())
221            .finish_non_exhaustive()
222    }
223}