ironflow_engine/context/steps/human_input.rs
1//! Typed human input step for [`WorkflowContext`].
2
3use chrono::{TimeDelta, Utc};
4use schemars::{JsonSchema, schema_for};
5use serde::de::DeserializeOwned;
6use serde_json::{Value, from_value, json, to_value};
7use tracing::info;
8
9use ironflow_store::models::{NewStep, Step, StepKind, StepStatus, StepUpdate, step_trace_id};
10
11use crate::config::{Approvers, HUMAN_INPUT_SCHEMA_KEY, HumanInputConfig};
12use crate::context::WorkflowContext;
13use crate::context::lifecycle::check_replay_identity;
14use crate::error::EngineError;
15use crate::executor::HumanInputOutcome;
16use crate::notify::{WorkflowEvent, WorkflowInputRequiredEvent};
17use crate::plan::lock_plan;
18
19impl WorkflowContext {
20 /// Ask a human for a typed answer and suspend the run until it is given.
21 ///
22 /// On first execution, records a human input step carrying the JSON schema
23 /// of `T` and returns [`EngineError::HumanInputRequired`] to suspend the
24 /// run. The engine transitions the run to `AwaitingApproval`. The answer is
25 /// posted to `POST /api/v1/runs/{id}/steps/{step_id}/input`, validated
26 /// against the schema, and the run resumes.
27 ///
28 /// On resume, the step is replayed and the stored answer is deserialized
29 /// into `T`. A rejected input (`POST .../steps/{step_id}/reject`) returns
30 /// [`EngineError::HumanInputRejected`] so the handler decides what happens
31 /// next. An answer given in an earlier attempt is carried over to a retry.
32 ///
33 /// The config reuses the approval gate machinery: deadline, escalation
34 /// policy, assignee and the [`Approvers`] allowed to answer.
35 ///
36 /// While planning, the step is recorded and never suspends: `T` is
37 /// deserialized from `{}` when it accepts that (for example with
38 /// `#[serde(default)]`).
39 ///
40 /// # Errors
41 ///
42 /// Returns [`EngineError::HumanInputRequired`] to pause the run until an
43 /// answer is given. Returns [`EngineError::HumanInputRejected`] when the
44 /// input was rejected. Returns [`EngineError::StepConfig`] when the stored
45 /// answer does not match `T`, or while planning when `T` cannot be built
46 /// from `{}`. Returns [`EngineError::ReplayDivergence`] when the step
47 /// recorded at this position has a different name or kind. Returns other
48 /// [`EngineError`] variants on store failures.
49 ///
50 /// # Examples
51 ///
52 /// ```no_run
53 /// use ironflow_engine::config::HumanInputConfig;
54 /// use ironflow_engine::context::WorkflowContext;
55 /// use ironflow_engine::error::EngineError;
56 /// use schemars::JsonSchema;
57 /// use serde::Deserialize;
58 ///
59 /// #[derive(Deserialize, JsonSchema)]
60 /// struct Answers {
61 /// answers: Vec<String>,
62 /// }
63 ///
64 /// # async fn example(ctx: &mut WorkflowContext) -> Result<(), EngineError> {
65 /// let answers: Answers = ctx
66 /// .human_input("clarify", HumanInputConfig::new("Answer the clarification questions"))
67 /// .await?;
68 /// assert!(answers.answers.len() < 100);
69 /// # Ok(())
70 /// # }
71 /// ```
72 pub async fn human_input<T: DeserializeOwned + JsonSchema>(
73 &mut self,
74 name: &str,
75 config: HumanInputConfig,
76 ) -> Result<T, EngineError> {
77 let schema = to_value(schema_for!(T))?;
78
79 // Plan mode: record the step and continue. Planning must never suspend.
80 if let Some(plan) = self.plan().cloned() {
81 self.position += 1;
82 {
83 let mut recorder = lock_plan(&plan);
84 if recorder.record(name, StepKind::HumanInput, &self.workflow_name, None) {
85 recorder.set_last(vec![name.to_string()]);
86 }
87 }
88 return from_value(json!({})).map_err(|e| {
89 EngineError::StepConfig(format!(
90 "human input '{name}' has no answer while planning: {e}"
91 ))
92 });
93 }
94
95 let value = self.human_input_value(name, &config, schema).await?;
96 from_value::<T>(value).map_err(|e| {
97 EngineError::StepConfig(format!(
98 "human input '{name}' answer does not match the expected type: {e}"
99 ))
100 })
101 }
102
103 /// Replay, carry over, intercept or open a human input step and return the
104 /// raw answer.
105 async fn human_input_value(
106 &mut self,
107 name: &str,
108 config: &HumanInputConfig,
109 schema: Value,
110 ) -> Result<Value, EngineError> {
111 let position = self.position;
112 self.position += 1;
113
114 // Replay: the step exists from a prior execution of this attempt.
115 if let Some(existing) = self.replay_steps.get(&position).cloned() {
116 check_replay_identity(&existing, position, name, &StepKind::HumanInput)?;
117 return self
118 .human_input_replay(name, config, position, existing)
119 .await;
120 }
121
122 // Carried over: a human already answered this input in an earlier
123 // attempt. Record a fresh completed step so each attempt keeps a
124 // complete DAG, and continue with the same answer.
125 if let Some((answered_in, value)) = self.answered_inputs.get(&position).cloned() {
126 let step = self
127 .create_human_input_step(name, position, config, &schema)
128 .await?;
129 let now = Utc::now();
130 self.start_step(step.id, now).await?;
131 self.store
132 .update_step(
133 step.id,
134 StepUpdate {
135 status: Some(StepStatus::Completed),
136 output: Some(value.clone()),
137 completed_at: Some(now),
138 ..StepUpdate::default()
139 },
140 )
141 .await?;
142
143 self.last_step_ids = vec![step.id];
144 info!(
145 run_id = %self.run_id,
146 step = %name,
147 position,
148 answered_in_attempt = answered_in,
149 attempt = self.attempt,
150 "human input carried over from a previous attempt"
151 );
152 return Ok(value);
153 }
154
155 // Recorded only when the input opens: the stored requirement is the
156 // source of truth from here.
157 let requirement = config.approvers().map(Approvers::to_requirement);
158
159 // An interceptor answers inline: the run neither suspends nor waits.
160 if let Some(interceptor) = self.interceptor.clone()
161 && let Some(outcome) = interceptor.intercept_human_input(name, config, &schema)
162 {
163 let step = self
164 .create_human_input_step(name, position, config, &schema)
165 .await?;
166 let now = Utc::now();
167 self.start_step(step.id, now).await?;
168 self.last_step_ids = vec![step.id];
169
170 return match outcome {
171 HumanInputOutcome::Provided(value) => {
172 self.store
173 .update_step(
174 step.id,
175 StepUpdate {
176 status: Some(StepStatus::Completed),
177 output: Some(value.clone()),
178 approval_requirement: requirement,
179 completed_at: Some(now),
180 ..StepUpdate::default()
181 },
182 )
183 .await?;
184 info!(
185 run_id = %self.run_id,
186 step = %name,
187 position,
188 "human input provided by the step interceptor"
189 );
190 Ok(value)
191 }
192 HumanInputOutcome::Rejected { reason } => {
193 // The step FSM only reaches Rejected from AwaitingApproval.
194 self.store
195 .update_step(
196 step.id,
197 StepUpdate {
198 status: Some(StepStatus::AwaitingApproval),
199 approval_requirement: requirement,
200 ..StepUpdate::default()
201 },
202 )
203 .await?;
204 self.store
205 .update_step(
206 step.id,
207 StepUpdate {
208 status: Some(StepStatus::Rejected),
209 error: Some(reason.clone()),
210 completed_at: Some(Utc::now()),
211 ..StepUpdate::default()
212 },
213 )
214 .await?;
215 info!(
216 run_id = %self.run_id,
217 step = %name,
218 position,
219 %reason,
220 "human input rejected by the step interceptor"
221 );
222 Err(EngineError::HumanInputRejected {
223 run_id: self.run_id,
224 step_id: step.id,
225 reason,
226 })
227 }
228 };
229 }
230
231 // First execution: create the step, arm the gate and suspend.
232 let step = self
233 .create_human_input_step(name, position, config, &schema)
234 .await?;
235 self.start_step(step.id, Utc::now()).await?;
236
237 let deadline_at = config
238 .effective_deadline_secs()
239 .map(|secs| Utc::now() + TimeDelta::seconds(secs as i64));
240
241 self.store
242 .update_step(
243 step.id,
244 StepUpdate {
245 status: Some(StepStatus::AwaitingApproval),
246 approval_deadline_at: deadline_at,
247 approval_stage: Some(0),
248 approval_assignee: config.assignee().cloned(),
249 approval_requirement: requirement,
250 ..StepUpdate::default()
251 },
252 )
253 .await?;
254
255 self.last_step_ids = vec![step.id];
256
257 if let Some(ref bus) = self.event_bus {
258 bus.publish(
259 self.run_id,
260 WorkflowEvent::InputRequired(WorkflowInputRequiredEvent {
261 run_id: self.run_id,
262 step_id: step.id,
263 step_name: name.to_string(),
264 step_index: position,
265 message: config.message().to_string(),
266 schema,
267 }),
268 );
269 }
270
271 info!(
272 run_id = %self.run_id,
273 step = %name,
274 position,
275 "human input requested"
276 );
277 Err(EngineError::HumanInputRequired {
278 run_id: self.run_id,
279 step_id: step.id,
280 message: config.message().to_string(),
281 })
282 }
283
284 /// Replay a human input step recorded in this attempt.
285 async fn human_input_replay(
286 &mut self,
287 name: &str,
288 config: &HumanInputConfig,
289 position: u32,
290 existing: Step,
291 ) -> Result<Value, EngineError> {
292 self.last_step_ids = vec![existing.id];
293
294 match existing.status.state {
295 StepStatus::Completed => {
296 info!(
297 run_id = %self.run_id,
298 step = %name,
299 position,
300 "human input replayed (answered)"
301 );
302 existing.output.ok_or_else(|| {
303 EngineError::StepConfig(format!("human input '{name}' has no stored answer"))
304 })
305 }
306 StepStatus::Rejected => {
307 info!(
308 run_id = %self.run_id,
309 step = %name,
310 position,
311 "human input replayed (rejected)"
312 );
313 Err(EngineError::HumanInputRejected {
314 run_id: self.run_id,
315 step_id: existing.id,
316 reason: existing
317 .error
318 .unwrap_or_else(|| "input rejected".to_string()),
319 })
320 }
321 state => {
322 // Resumed without an answer (or after a crash mid-open): the
323 // step keeps waiting, no new step is created.
324 if state == StepStatus::Running {
325 self.store
326 .update_step(
327 existing.id,
328 StepUpdate {
329 status: Some(StepStatus::AwaitingApproval),
330 ..StepUpdate::default()
331 },
332 )
333 .await?;
334 }
335 info!(
336 run_id = %self.run_id,
337 step = %name,
338 position,
339 "human input still unanswered, suspending again"
340 );
341 Err(EngineError::HumanInputRequired {
342 run_id: self.run_id,
343 step_id: existing.id,
344 message: config.message().to_string(),
345 })
346 }
347 }
348 }
349
350 /// Create the step record of a human input.
351 async fn create_human_input_step(
352 &self,
353 name: &str,
354 position: u32,
355 config: &HumanInputConfig,
356 schema: &Value,
357 ) -> Result<Step, EngineError> {
358 let trace_id = step_trace_id(self.run_id, name, position);
359 Ok(self
360 .store
361 .create_step(NewStep {
362 run_id: self.run_id,
363 trace_id,
364 name: name.to_string(),
365 kind: StepKind::HumanInput,
366 position,
367 input: Some(stored_input(config, schema)?),
368 is_error_handler: false,
369 })
370 .await?)
371 }
372}
373
374/// The stored step input: the flattened config plus the answer schema under
375/// [`HUMAN_INPUT_SCHEMA_KEY`].
376fn stored_input(config: &HumanInputConfig, schema: &Value) -> Result<Value, EngineError> {
377 let mut input = to_value(config)?;
378 if let Some(object) = input.as_object_mut() {
379 object.insert(HUMAN_INPUT_SCHEMA_KEY.to_string(), schema.clone());
380 }
381 Ok(input)
382}