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}