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