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