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::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}