Skip to main content

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}