Skip to main content

ironflow_engine/executor/
interceptor.rs

1//! Hook that resolves a step without executing it.
2//!
3//! A [`StepInterceptor`] is consulted by
4//! [`execute_step_config_intercepted`](crate::executor::execute_step_config_intercepted)
5//! before the dispatcher picks an executor, and by
6//! [`WorkflowContext::approval`](crate::context::WorkflowContext::approval)
7//! before the gate suspends the run, and by
8//! [`WorkflowContext::human_input`](crate::context::WorkflowContext::human_input)
9//! before the input request suspends the run, and by
10//! [`WorkflowContext::wait_for_signal`](crate::context::WorkflowContext::wait_for_signal)
11//! before the signal step suspends the run. Returning `Some(..)` short-circuits the
12//! step: no process is spawned, no request is sent, no human is asked.
13//!
14//! Production wiring leaves the hook unset. In practice the only implementor is
15//! [`crate::testing`], which uses it to run a handler's real logic against
16//! canned step results.
17
18use serde_json::Value;
19
20use crate::config::{ApprovalConfig, HumanInputConfig, StepConfig};
21use crate::error::EngineError;
22use crate::executor::StepOutput;
23
24/// Decision applied to an approval gate by a [`StepInterceptor`].
25///
26/// # Examples
27///
28/// ```
29/// use ironflow_engine::executor::ApprovalOutcome;
30///
31/// let granted = ApprovalOutcome::Approved;
32/// assert_eq!(granted, ApprovalOutcome::Approved);
33///
34/// let refused = ApprovalOutcome::Rejected { reason: "not on a Friday".to_string() };
35/// assert_ne!(granted, refused);
36/// ```
37#[derive(Debug, Clone, PartialEq, Eq)]
38pub enum ApprovalOutcome {
39    /// The gate is granted; execution continues past it.
40    Approved,
41    /// The gate is refused; the run fails with [`EngineError::ApprovalRejected`].
42    Rejected {
43        /// Human-readable reason recorded on the step and on the run.
44        reason: String,
45    },
46}
47
48impl ApprovalOutcome {
49    /// Build an [`ApprovalOutcome::Rejected`] with the given reason.
50    ///
51    /// # Examples
52    ///
53    /// ```
54    /// use ironflow_engine::executor::ApprovalOutcome;
55    ///
56    /// let outcome = ApprovalOutcome::reject("budget freeze");
57    /// assert_eq!(
58    ///     outcome,
59    ///     ApprovalOutcome::Rejected { reason: "budget freeze".to_string() }
60    /// );
61    /// ```
62    pub fn reject(reason: &str) -> Self {
63        Self::Rejected {
64            reason: reason.to_string(),
65        }
66    }
67}
68
69/// Answer applied to a human input step by a [`StepInterceptor`].
70///
71/// # Examples
72///
73/// ```
74/// use ironflow_engine::executor::HumanInputOutcome;
75/// use serde_json::json;
76///
77/// let answered = HumanInputOutcome::Provided(json!({"answers": ["yes"]}));
78/// let refused = HumanInputOutcome::reject("out of scope");
79/// assert_ne!(answered, refused);
80/// ```
81#[derive(Debug, Clone, PartialEq)]
82pub enum HumanInputOutcome {
83    /// The input is answered with this value; it must match the expected type.
84    Provided(Value),
85    /// The input is refused; the handler receives
86    /// [`EngineError::HumanInputRejected`].
87    Rejected {
88        /// Human-readable reason recorded on the step.
89        reason: String,
90    },
91}
92
93impl HumanInputOutcome {
94    /// Build a [`HumanInputOutcome::Rejected`] with the given reason.
95    ///
96    /// # Examples
97    ///
98    /// ```
99    /// use ironflow_engine::executor::HumanInputOutcome;
100    ///
101    /// let outcome = HumanInputOutcome::reject("out of scope");
102    /// assert_eq!(
103    ///     outcome,
104    ///     HumanInputOutcome::Rejected { reason: "out of scope".to_string() }
105    /// );
106    /// ```
107    pub fn reject(reason: &str) -> Self {
108        Self::Rejected {
109            reason: reason.to_string(),
110        }
111    }
112}
113
114/// Outcome applied to a signal step by a [`StepInterceptor`].
115///
116/// # Examples
117///
118/// ```
119/// use ironflow_engine::executor::SignalOutcome;
120/// use serde_json::json;
121///
122/// let received = SignalOutcome::Received(json!({"status": "success"}));
123/// assert_ne!(received, SignalOutcome::TimedOut);
124/// ```
125#[derive(Debug, Clone, PartialEq)]
126pub enum SignalOutcome {
127    /// A signal arrived with this payload; it must match the expected type.
128    Received(Value),
129    /// No signal arrived before the deadline: the handler receives `None`.
130    TimedOut,
131}
132
133/// Resolves steps without executing them.
134///
135/// # Examples
136///
137/// ```
138/// use ironflow_engine::config::{ShellConfig, StepConfig};
139/// use ironflow_engine::error::EngineError;
140/// use ironflow_engine::executor::{StepArtifacts, StepInterceptor, StepOutput};
141/// use rust_decimal::Decimal;
142/// use serde_json::json;
143///
144/// struct AlwaysOk;
145///
146/// impl StepInterceptor for AlwaysOk {
147///     fn intercept(&self, config: &StepConfig) -> Option<Result<StepOutput, EngineError>> {
148///         match config {
149///             StepConfig::Shell(_) => Some(Ok(StepOutput {
150///                 output: json!({"stdout": "ok", "stderr": "", "exit_code": 0}),
151///                 duration_ms: 0,
152///                 cost_usd: Decimal::ZERO,
153///                 input_tokens: None,
154///                 cache_read_input_tokens: None,
155///                 cache_creation_input_tokens: None,
156///                 output_tokens: None,
157///                 model: None,
158///                 debug_messages: None,
159///                 artifacts: StepArtifacts::default(),
160///                 account_id: None,
161///             })),
162///             _ => None,
163///         }
164///     }
165/// }
166///
167/// let config = StepConfig::Shell(ShellConfig::new("./deploy.sh"));
168/// let intercepted = AlwaysOk.intercept(&config).expect("shell is intercepted");
169/// assert_eq!(intercepted.expect("canned output").stdout(), "ok");
170/// ```
171pub trait StepInterceptor: Send + Sync {
172    /// Return `Some(result)` to short-circuit this step, `None` to execute it
173    /// for real.
174    fn intercept(&self, config: &StepConfig) -> Option<Result<StepOutput, EngineError>>;
175
176    /// Resolve an approval gate instead of suspending the run.
177    ///
178    /// The default implementation returns `None`: the gate suspends the run.
179    fn intercept_approval(&self, name: &str, config: &ApprovalConfig) -> Option<ApprovalOutcome> {
180        let _ = (name, config);
181        None
182    }
183
184    /// Answer a human input step instead of suspending the run.
185    ///
186    /// `schema` is the JSON schema of the expected answer. The default
187    /// implementation returns `None`: the input suspends the run.
188    fn intercept_human_input(
189        &self,
190        name: &str,
191        config: &HumanInputConfig,
192        schema: &Value,
193    ) -> Option<HumanInputOutcome> {
194        let _ = (name, config, schema);
195        None
196    }
197
198    /// Resolve a signal step instead of suspending the run.
199    ///
200    /// `name` is the step name, `signal_name` and `key` identify the awaited
201    /// signal, and `schema` is the JSON schema of its payload. The default
202    /// implementation returns `None`: the step waits for a real signal.
203    fn intercept_signal(
204        &self,
205        name: &str,
206        signal_name: &str,
207        key: &str,
208        schema: &Value,
209    ) -> Option<SignalOutcome> {
210        let _ = (name, signal_name, key, schema);
211        None
212    }
213}