Skip to main content

ironflow_api/routes/
human_input.rs

1//! `POST /api/v1/runs/:id/steps/:step_id/input` -- Answer a human input step.
2//!
3//! `POST /api/v1/runs/:id/steps/:step_id/reject` -- Reject a human input step.
4
5use axum::Json;
6use axum::extract::{Path, State};
7use axum::response::IntoResponse;
8use chrono::Utc;
9use ironflow_auth::extractor::Authenticated;
10use ironflow_engine::config::HUMAN_INPUT_SCHEMA_KEY;
11use ironflow_engine::engine::ExecutionMode;
12use ironflow_store::models::{
13    Run, RunStatus, Step, StepApproval, StepKind, StepStatus, StepUpdate,
14};
15use jsonschema::validator_for;
16use serde::{Deserialize, Serialize};
17use serde_json::Value;
18use tokio::spawn;
19use tracing::error;
20use uuid::Uuid;
21
22use crate::entities::RunResponse;
23use crate::error::ApiError;
24use crate::response::ok;
25use crate::routes::approve_run::authorize_gate;
26use crate::state::AppState;
27
28/// Request body for rejecting a human input step.
29#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
30#[derive(Debug, Default, Deserialize, Serialize)]
31pub struct RejectHumanInputRequest {
32    /// Why the input is refused. Passed to the handler.
33    #[serde(default)]
34    pub reason: Option<String>,
35}
36
37/// Answer a human input step with a value matching its JSON schema.
38///
39/// The answer is validated against the schema stored on the step. The caller
40/// must be allowed to resolve the step like an approval gate (admin, approver
41/// groups, assignee or delegation). The first valid answer wins: the step
42/// completes, the run moves from `AwaitingApproval` back to `Running` and
43/// resumes, and the handler receives the typed answer.
44///
45/// # Errors
46///
47/// - 400 if the step is not a human input or is not awaiting input
48/// - 403 if the caller may not answer the step
49/// - 404 if the run or the step does not exist
50/// - 409 if the input was already answered or rejected
51/// - 422 if the answer does not match the schema
52#[cfg_attr(
53    feature = "openapi",
54    utoipa::path(
55        post,
56        path = "/api/v1/runs/{id}/steps/{step_id}/input",
57        tags = ["runs"],
58        params(
59            ("id" = Uuid, Path, description = "Run ID"),
60            ("step_id" = Uuid, Path, description = "Step ID")
61        ),
62        request_body(content = Value, content_type = "application/json", description = "Answer matching the JSON schema stored on the step"),
63        responses(
64            (status = 200, description = "Answer recorded, the run resumes", body = RunResponse),
65            (status = 400, description = "Step is not a human input awaiting input"),
66            (status = 401, description = "Unauthorized"),
67            (status = 403, description = "Forbidden"),
68            (status = 404, description = "Run or step not found"),
69            (status = 409, description = "Input already answered or rejected"),
70            (status = 422, description = "Answer does not match the schema")
71        ),
72        security(("Bearer" = []))
73    )
74)]
75pub async fn submit_human_input(
76    auth: Authenticated,
77    State(state): State<AppState>,
78    Path((id, step_id)): Path<(Uuid, Uuid)>,
79    Json(answer): Json<Value>,
80) -> Result<impl IntoResponse, ApiError> {
81    let (run, step) = open_input_step(&state, id, step_id).await?;
82    let actor = authorize_gate(&auth, &state, &run, &step).await?;
83    validate_answer(step.input.as_ref(), &answer)?;
84
85    // Like approve, two concurrent answers are not serialized: both may pass
86    // the status check above before either is written. This is an accepted
87    // limitation of the gate routes.
88    let now = Utc::now();
89    state
90        .store
91        .record_step_approval(
92            step.id,
93            StepApproval {
94                user_id: auth.user_id,
95                approved_by: actor,
96                at: now,
97            },
98        )
99        .await?;
100    state
101        .store
102        .update_step(
103            step.id,
104            StepUpdate {
105                status: Some(StepStatus::Completed),
106                output: Some(answer),
107                completed_at: Some(now),
108                clear_approval_deadline: true,
109                ..StepUpdate::default()
110            },
111        )
112        .await?;
113    state
114        .store
115        .update_run_status(id, resume_status(&state))
116        .await?;
117
118    resume_in_background(&state, id);
119
120    Ok(ok(RunResponse::from(state.get_run_or_404(id).await?)))
121}
122
123/// Reject a human input step.
124///
125/// The step is marked `Rejected` with the reason, the run moves from
126/// `AwaitingApproval` back to `Running` and resumes: the handler receives a
127/// `HumanInputRejected` error and decides what happens next. The body is
128/// optional; without a reason, one naming the caller is recorded.
129///
130/// # Errors
131///
132/// - 400 if the step is not a human input or is not awaiting input
133/// - 403 if the caller may not resolve the step
134/// - 404 if the run or the step does not exist
135/// - 409 if the input was already answered or rejected
136#[cfg_attr(
137    feature = "openapi",
138    utoipa::path(
139        post,
140        path = "/api/v1/runs/{id}/steps/{step_id}/reject",
141        tags = ["runs"],
142        params(
143            ("id" = Uuid, Path, description = "Run ID"),
144            ("step_id" = Uuid, Path, description = "Step ID")
145        ),
146        request_body(content = RejectHumanInputRequest, description = "Optional rejection reason"),
147        responses(
148            (status = 200, description = "Input rejected, the run resumes", body = RunResponse),
149            (status = 400, description = "Step is not a human input awaiting input"),
150            (status = 401, description = "Unauthorized"),
151            (status = 403, description = "Forbidden"),
152            (status = 404, description = "Run or step not found"),
153            (status = 409, description = "Input already answered or rejected")
154        ),
155        security(("Bearer" = []))
156    )
157)]
158pub async fn reject_human_input(
159    auth: Authenticated,
160    State(state): State<AppState>,
161    Path((id, step_id)): Path<(Uuid, Uuid)>,
162    body: Option<Json<RejectHumanInputRequest>>,
163) -> Result<impl IntoResponse, ApiError> {
164    let (run, step) = open_input_step(&state, id, step_id).await?;
165    let actor = authorize_gate(&auth, &state, &run, &step).await?;
166
167    let reason = body
168        .and_then(|Json(b)| b.reason)
169        .filter(|r| !r.trim().is_empty())
170        .unwrap_or_else(|| format!("input rejected by {actor}"));
171
172    state
173        .store
174        .update_step(
175            step.id,
176            StepUpdate {
177                status: Some(StepStatus::Rejected),
178                error: Some(reason),
179                completed_at: Some(Utc::now()),
180                clear_approval_deadline: true,
181                ..StepUpdate::default()
182            },
183        )
184        .await?;
185    state
186        .store
187        .update_run_status(id, resume_status(&state))
188        .await?;
189
190    resume_in_background(&state, id);
191
192    Ok(ok(RunResponse::from(state.get_run_or_404(id).await?)))
193}
194
195/// Load the run and the human input step, and check the step can be resolved.
196///
197/// A resolved step returns 409 whatever the run status, so replaying an
198/// accepted request is reported as a conflict rather than a bad request.
199async fn open_input_step(
200    state: &AppState,
201    run_id: Uuid,
202    step_id: Uuid,
203) -> Result<(Run, Step), ApiError> {
204    let run = state.get_run_or_404(run_id).await?;
205    let step = state
206        .store
207        .get_step(step_id)
208        .await?
209        .filter(|s| s.run_id == run_id)
210        .ok_or(ApiError::StepNotFound(step_id))?;
211
212    if step.kind != StepKind::HumanInput {
213        return Err(ApiError::BadRequest(
214            "step does not wait for input".to_string(),
215        ));
216    }
217
218    match step.status.state {
219        StepStatus::Completed => {
220            return Err(ApiError::Conflict("input already provided".to_string()));
221        }
222        StepStatus::Rejected => {
223            return Err(ApiError::Conflict("input already rejected".to_string()));
224        }
225        _ => {}
226    }
227
228    if step.status.state != StepStatus::AwaitingApproval
229        || run.status.state != RunStatus::AwaitingApproval
230    {
231        return Err(ApiError::BadRequest(
232            "step is not awaiting input".to_string(),
233        ));
234    }
235
236    Ok((run, step))
237}
238
239/// Validate `answer` against the JSON schema stored in the step input.
240fn validate_answer(input: Option<&Value>, answer: &Value) -> Result<(), ApiError> {
241    let schema = input
242        .and_then(|i| i.get(HUMAN_INPUT_SCHEMA_KEY))
243        .ok_or_else(|| ApiError::Internal("input step has no stored schema".to_string()))?;
244    let validator = validator_for(schema)
245        .map_err(|e| ApiError::Internal(format!("input step has an invalid schema: {e}")))?;
246
247    let errors: Vec<String> = validator
248        .iter_errors(answer)
249        .map(|e| e.to_string())
250        .collect();
251    if errors.is_empty() {
252        Ok(())
253    } else {
254        Err(ApiError::InvalidInput(errors))
255    }
256}
257
258/// The status a run moves to once its input is answered or rejected:
259/// `Running` to resume it here, `Pending` to requeue it for a worker.
260fn resume_status(state: &AppState) -> RunStatus {
261    match state.engine.execution_mode() {
262        ExecutionMode::Local => RunStatus::Running,
263        ExecutionMode::Workers => RunStatus::Pending,
264    }
265}
266
267/// Under [`ExecutionMode::Local`], resume the run in the background; the
268/// handler replays up to the input. Under [`ExecutionMode::Workers`], do
269/// nothing: the run is already `Pending` and a worker picks it up.
270fn resume_in_background(state: &AppState, id: Uuid) {
271    if !matches!(state.engine.execution_mode(), ExecutionMode::Local) {
272        return;
273    }
274    let engine = state.engine.clone();
275    spawn(async move {
276        if let Err(err) = engine.resume_run(id).await {
277            error!(run_id = %id, error = %err, "failed to resume run after human input");
278        }
279    });
280}
281
282#[cfg(test)]
283mod tests {
284    use std::collections::HashMap;
285    use std::sync::Arc;
286
287    use axum::Router;
288    use axum::body::Body;
289    use axum::http::{Request, StatusCode};
290    use axum::routing::post;
291    use http_body_util::BodyExt;
292    use ironflow_auth::jwt::{AccessToken, JwtConfig};
293    use ironflow_auth::password;
294    use ironflow_core::providers::claude::ClaudeCodeProvider;
295    use ironflow_engine::engine::Engine;
296    use ironflow_engine::handler::input_schema_for;
297    use ironflow_engine::notify::Event;
298    use ironflow_store::memory::InMemoryStore;
299    use ironflow_store::models::{
300        Assignee, NewRun, NewStep, NewUser, TriggerKind, User, step_trace_id,
301    };
302    use ironflow_store::store::RunStore;
303    use ironflow_store::user_store::UserStore;
304    use schemars::JsonSchema;
305    use serde_json::{from_slice, json};
306    use tokio::sync::broadcast;
307    use tower::ServiceExt;
308
309    use super::*;
310
311    /// The answer type the test steps ask for.
312    #[allow(dead_code)]
313    #[derive(Deserialize, JsonSchema)]
314    struct Answers {
315        answers: Vec<String>,
316    }
317
318    /// An answer type whose schema carries a `$defs` reference.
319    #[allow(dead_code)]
320    #[derive(Deserialize, JsonSchema)]
321    struct Nested {
322        target: Target,
323    }
324
325    #[allow(dead_code)]
326    #[derive(Deserialize, JsonSchema)]
327    struct Target {
328        env: String,
329    }
330
331    fn test_state(store: Arc<InMemoryStore>) -> AppState {
332        let provider = Arc::new(ClaudeCodeProvider::new());
333        let engine = Arc::new(Engine::new(store.clone(), provider));
334        let jwt_config = Arc::new(JwtConfig {
335            secret: "test-secret".to_string(),
336            access_token_ttl_secs: 900,
337            refresh_token_ttl_secs: 604800,
338            cookie_domain: None,
339            cookie_secure: false,
340        });
341        let (event_sender, _) = broadcast::channel::<Event>(1);
342        AppState::new(
343            store,
344            engine,
345            jwt_config,
346            "test-worker-token".to_string(),
347            event_sender,
348        )
349    }
350
351    fn admin_header(state: &AppState) -> String {
352        let token =
353            AccessToken::for_user(Uuid::now_v7(), "admin", true, &state.jwt_config).unwrap();
354        format!("Bearer {}", token.0)
355    }
356
357    async fn member(store: &Arc<InMemoryStore>, username: &str) -> User {
358        let password_hash = password::hash("password123").expect("hash");
359        store
360            .create_user(NewUser {
361                email: format!("{username}@example.com"),
362                username: username.to_string(),
363                password_hash,
364                // The first user would otherwise become an implicit admin.
365                is_admin: Some(false),
366            })
367            .await
368            .expect("create user")
369    }
370
371    fn member_header(user: &User, state: &AppState) -> String {
372        let token = AccessToken::for_user(user.id, &user.username, false, &state.jwt_config)
373            .expect("token");
374        format!("Bearer {}", token.0)
375    }
376
377    async fn awaiting_run(store: &Arc<InMemoryStore>) -> Uuid {
378        let run = store
379            .create_run(NewRun {
380                created_by: None,
381                workflow_name: "test".to_string(),
382                trigger: TriggerKind::Manual,
383                payload: json!({}),
384                max_retries: 0,
385                handler_version: None,
386                labels: HashMap::new(),
387                scheduled_at: None,
388                idempotency_key: None,
389                max_cost_usd: None,
390            })
391            .await
392            .unwrap()
393            .into_run();
394        store
395            .update_run_status(run.id, RunStatus::Running)
396            .await
397            .unwrap();
398        store
399            .update_run_status(run.id, RunStatus::AwaitingApproval)
400            .await
401            .unwrap();
402        run.id
403    }
404
405    /// A step of `kind` awaiting approval on `run_id`.
406    async fn awaiting_step(
407        store: &Arc<InMemoryStore>,
408        run_id: Uuid,
409        kind: StepKind,
410        assignee: Option<Assignee>,
411    ) -> Uuid {
412        let step = store
413            .create_step(NewStep {
414                run_id,
415                trace_id: step_trace_id(run_id, "clarify", 0),
416                name: "clarify".to_string(),
417                kind,
418                position: 0,
419                input: Some(json!({
420                    "message": "Answer the questions",
421                    HUMAN_INPUT_SCHEMA_KEY: input_schema_for::<Answers>(),
422                })),
423                is_error_handler: false,
424            })
425            .await
426            .unwrap();
427        store
428            .update_step(
429                step.id,
430                StepUpdate {
431                    status: Some(StepStatus::Running),
432                    ..StepUpdate::default()
433                },
434            )
435            .await
436            .unwrap();
437        store
438            .update_step(
439                step.id,
440                StepUpdate {
441                    status: Some(StepStatus::AwaitingApproval),
442                    approval_assignee: assignee,
443                    ..StepUpdate::default()
444                },
445            )
446            .await
447            .unwrap();
448        step.id
449    }
450
451    /// A run awaiting a human input step, with an admin header.
452    async fn setup() -> (Arc<InMemoryStore>, AppState, Uuid, Uuid) {
453        let store = Arc::new(InMemoryStore::new());
454        let run_id = awaiting_run(&store).await;
455        let step_id = awaiting_step(&store, run_id, StepKind::HumanInput, None).await;
456        let state = test_state(store.clone());
457        (store, state, run_id, step_id)
458    }
459
460    fn app(state: AppState) -> Router {
461        Router::new()
462            .route("/{id}/steps/{step_id}/input", post(submit_human_input))
463            .route("/{id}/steps/{step_id}/reject", post(reject_human_input))
464            .with_state(state)
465    }
466
467    /// POST `body` (or nothing) to `/{run_id}/steps/{step_id}/{verb}`.
468    async fn call(
469        state: &AppState,
470        auth: &str,
471        run_id: Uuid,
472        step_id: Uuid,
473        verb: &str,
474        body: Option<Value>,
475    ) -> (StatusCode, Value) {
476        let builder = Request::builder()
477            .method("POST")
478            .uri(format!("/{run_id}/steps/{step_id}/{verb}"))
479            .header("authorization", auth);
480        let req = match body {
481            Some(body) => builder
482                .header("content-type", "application/json")
483                .body(Body::from(body.to_string()))
484                .unwrap(),
485            None => builder.body(Body::empty()).unwrap(),
486        };
487
488        let resp = app(state.clone()).oneshot(req).await.unwrap();
489        let status = resp.status();
490        let bytes = resp.into_body().collect().await.unwrap().to_bytes();
491        let json_val = if bytes.is_empty() {
492            Value::Null
493        } else {
494            from_slice(&bytes).unwrap()
495        };
496        (status, json_val)
497    }
498
499    #[tokio::test]
500    async fn human_input_submit_valid_answer_resumes_the_run() {
501        let (store, state, run_id, step_id) = setup().await;
502        let auth = admin_header(&state);
503        let answer = json!({"answers": ["staging"]});
504
505        let (status, body) = call(
506            &state,
507            &auth,
508            run_id,
509            step_id,
510            "input",
511            Some(answer.clone()),
512        )
513        .await;
514
515        assert_eq!(status, StatusCode::OK);
516        assert_eq!(body["data"]["status"], "running");
517        let step = store.get_step(step_id).await.unwrap().unwrap();
518        assert_eq!(step.status.state, StepStatus::Completed);
519        assert_eq!(step.output, Some(answer));
520        assert_eq!(step.approvals.len(), 1);
521        assert_eq!(step.approvals[0].approved_by, "admin");
522    }
523
524    #[tokio::test]
525    async fn human_input_submit_invalid_answer_returns_422() {
526        let (store, state, run_id, step_id) = setup().await;
527        let auth = admin_header(&state);
528
529        for answer in [json!({"answers": 3}), json!({})] {
530            let (status, body) = call(&state, &auth, run_id, step_id, "input", Some(answer)).await;
531
532            assert_eq!(status, StatusCode::UNPROCESSABLE_ENTITY);
533            assert_eq!(body["error"]["code"], "INVALID_INPUT");
534            let errors = body["error"]["details"]["errors"]
535                .as_array()
536                .expect("errors array");
537            assert!(!errors.is_empty());
538        }
539
540        let step = store.get_step(step_id).await.unwrap().unwrap();
541        assert_eq!(step.status.state, StepStatus::AwaitingApproval);
542        let run = store.get_run(run_id).await.unwrap().unwrap();
543        assert_eq!(run.status.state, RunStatus::AwaitingApproval);
544    }
545
546    #[tokio::test]
547    async fn human_input_submit_twice_returns_409() {
548        let (_store, state, run_id, step_id) = setup().await;
549        let auth = admin_header(&state);
550        let answer = json!({"answers": ["yes"]});
551
552        let (first, _) = call(
553            &state,
554            &auth,
555            run_id,
556            step_id,
557            "input",
558            Some(answer.clone()),
559        )
560        .await;
561        assert_eq!(first, StatusCode::OK);
562
563        let (second, body) = call(&state, &auth, run_id, step_id, "input", Some(answer)).await;
564        assert_eq!(second, StatusCode::CONFLICT);
565        assert_eq!(body["error"]["message"], "input already provided");
566    }
567
568    #[tokio::test]
569    async fn human_input_submit_on_an_approval_step_returns_400() {
570        let store = Arc::new(InMemoryStore::new());
571        let run_id = awaiting_run(&store).await;
572        let step_id = awaiting_step(&store, run_id, StepKind::Approval, None).await;
573        let state = test_state(store);
574        let auth = admin_header(&state);
575
576        let answer = json!({"answers": []});
577        let (status, _) = call(&state, &auth, run_id, step_id, "input", Some(answer)).await;
578        assert_eq!(status, StatusCode::BAD_REQUEST);
579    }
580
581    #[tokio::test]
582    async fn human_input_submit_on_a_run_not_awaiting_returns_400() {
583        let (store, state, run_id, step_id) = setup().await;
584        store
585            .update_run_status(run_id, RunStatus::Running)
586            .await
587            .unwrap();
588        let auth = admin_header(&state);
589
590        let answer = json!({"answers": []});
591        let (status, _) = call(&state, &auth, run_id, step_id, "input", Some(answer)).await;
592        assert_eq!(status, StatusCode::BAD_REQUEST);
593    }
594
595    #[tokio::test]
596    async fn human_input_unknown_step_or_step_of_another_run_returns_404() {
597        let (store, state, run_id, step_id) = setup().await;
598        let auth = admin_header(&state);
599        let answer = json!({"answers": []});
600
601        let unknown = Uuid::now_v7();
602        let (status, body) = call(
603            &state,
604            &auth,
605            run_id,
606            unknown,
607            "input",
608            Some(answer.clone()),
609        )
610        .await;
611        assert_eq!(status, StatusCode::NOT_FOUND);
612        assert_eq!(body["error"]["code"], "STEP_NOT_FOUND");
613
614        let other_run = awaiting_run(&store).await;
615        let (status, _) = call(&state, &auth, other_run, step_id, "input", Some(answer)).await;
616        assert_eq!(status, StatusCode::NOT_FOUND);
617    }
618
619    #[tokio::test]
620    async fn human_input_unknown_run_returns_404() {
621        let (_store, state, _run_id, step_id) = setup().await;
622        let auth = admin_header(&state);
623
624        let answer = json!({"answers": []});
625        let (status, body) = call(
626            &state,
627            &auth,
628            Uuid::now_v7(),
629            step_id,
630            "input",
631            Some(answer),
632        )
633        .await;
634        assert_eq!(status, StatusCode::NOT_FOUND);
635        assert_eq!(body["error"]["code"], "RUN_NOT_FOUND");
636    }
637
638    #[tokio::test]
639    async fn human_input_non_admin_non_assignee_is_forbidden() {
640        let store = Arc::new(InMemoryStore::new());
641        let alice = member(&store, "alice").await;
642        let bob = member(&store, "bob").await;
643        let run_id = awaiting_run(&store).await;
644        let assignee = Some(Assignee::user(&alice.username));
645        let step_id = awaiting_step(&store, run_id, StepKind::HumanInput, assignee).await;
646        let state = test_state(store.clone());
647
648        let answer = json!({"answers": ["yes"]});
649        let bob_auth = member_header(&bob, &state);
650        let (status, _) = call(
651            &state,
652            &bob_auth,
653            run_id,
654            step_id,
655            "input",
656            Some(answer.clone()),
657        )
658        .await;
659        assert_eq!(status, StatusCode::FORBIDDEN);
660
661        let alice_auth = member_header(&alice, &state);
662        let (status, _) = call(&state, &alice_auth, run_id, step_id, "input", Some(answer)).await;
663        assert_eq!(status, StatusCode::OK);
664        let step = store.get_step(step_id).await.unwrap().unwrap();
665        assert_eq!(step.approvals[0].user_id, alice.id);
666    }
667
668    #[tokio::test]
669    async fn human_input_reject_marks_the_step_and_resumes_the_run() {
670        let (store, state, run_id, step_id) = setup().await;
671        let auth = admin_header(&state);
672
673        let body = json!({"reason": "out of scope"});
674        let (status, resp) = call(&state, &auth, run_id, step_id, "reject", Some(body)).await;
675
676        assert_eq!(status, StatusCode::OK);
677        assert_eq!(resp["data"]["status"], "running");
678        let step = store.get_step(step_id).await.unwrap().unwrap();
679        assert_eq!(step.status.state, StepStatus::Rejected);
680        assert_eq!(step.error.as_deref(), Some("out of scope"));
681        assert!(step.approval_deadline_at.is_none());
682    }
683
684    /// A run awaiting a human input assigned to a member, on an engine that
685    /// requeues resumed runs for a worker.
686    async fn workers_setup() -> (Arc<InMemoryStore>, AppState, Uuid, Uuid, String) {
687        let store = Arc::new(InMemoryStore::new());
688        let alice = member(&store, "alice").await;
689        let run_id = awaiting_run(&store).await;
690        let assignee = Some(Assignee::user(&alice.username));
691        let step_id = awaiting_step(&store, run_id, StepKind::HumanInput, assignee).await;
692        let provider = Arc::new(ClaudeCodeProvider::new());
693        let engine =
694            Engine::new(store.clone(), provider).with_execution_mode(ExecutionMode::Workers);
695        let jwt_config = Arc::new(JwtConfig {
696            secret: "test-secret".to_string(),
697            access_token_ttl_secs: 900,
698            refresh_token_ttl_secs: 604800,
699            cookie_domain: None,
700            cookie_secure: false,
701        });
702        let (event_sender, _) = broadcast::channel::<Event>(1);
703        let state = AppState::new(
704            store.clone(),
705            Arc::new(engine),
706            jwt_config,
707            "test-worker-token".to_string(),
708            event_sender,
709        );
710        let auth = member_header(&alice, &state);
711        (store, state, run_id, step_id, auth)
712    }
713
714    #[tokio::test]
715    async fn human_input_submit_in_workers_mode_requeues_the_run() {
716        let (store, state, run_id, step_id, auth) = workers_setup().await;
717
718        let answer = json!({"answers": ["ok"]});
719        let (status, body) = call(&state, &auth, run_id, step_id, "input", Some(answer)).await;
720
721        assert_eq!(status, StatusCode::OK);
722        assert_eq!(body["data"]["status"], "pending");
723        let run = store.get_run(run_id).await.unwrap().unwrap();
724        assert_eq!(run.status.state, RunStatus::Pending);
725        let step = store.get_step(step_id).await.unwrap().unwrap();
726        assert_eq!(step.status.state, StepStatus::Completed);
727    }
728
729    #[tokio::test]
730    async fn human_input_reject_in_workers_mode_requeues_the_run() {
731        let (store, state, run_id, step_id, auth) = workers_setup().await;
732
733        let body = json!({"reason": "out of scope"});
734        let (status, resp) = call(&state, &auth, run_id, step_id, "reject", Some(body)).await;
735
736        assert_eq!(status, StatusCode::OK);
737        assert_eq!(resp["data"]["status"], "pending");
738        let run = store.get_run(run_id).await.unwrap().unwrap();
739        assert_eq!(run.status.state, RunStatus::Pending);
740        let step = store.get_step(step_id).await.unwrap().unwrap();
741        assert_eq!(step.status.state, StepStatus::Rejected);
742    }
743
744    #[tokio::test]
745    async fn human_input_reject_without_body_uses_the_default_reason() {
746        let (store, state, run_id, step_id) = setup().await;
747        let auth = admin_header(&state);
748
749        let (status, _) = call(&state, &auth, run_id, step_id, "reject", None).await;
750
751        assert_eq!(status, StatusCode::OK);
752        let step = store.get_step(step_id).await.unwrap().unwrap();
753        assert_eq!(step.error.as_deref(), Some("input rejected by admin"));
754    }
755
756    #[tokio::test]
757    async fn human_input_reject_with_a_blank_reason_uses_the_default_reason() {
758        let (store, state, run_id, step_id) = setup().await;
759        let auth = admin_header(&state);
760
761        let body = json!({"reason": "   "});
762        let (status, _) = call(&state, &auth, run_id, step_id, "reject", Some(body)).await;
763
764        assert_eq!(status, StatusCode::OK);
765        let step = store.get_step(step_id).await.unwrap().unwrap();
766        assert_eq!(step.error.as_deref(), Some("input rejected by admin"));
767    }
768
769    #[tokio::test]
770    async fn human_input_reject_after_submit_returns_409() {
771        let (_store, state, run_id, step_id) = setup().await;
772        let auth = admin_header(&state);
773
774        let answer = json!({"answers": ["yes"]});
775        let (first, _) = call(&state, &auth, run_id, step_id, "input", Some(answer)).await;
776        assert_eq!(first, StatusCode::OK);
777
778        let (second, _) = call(&state, &auth, run_id, step_id, "reject", None).await;
779        assert_eq!(second, StatusCode::CONFLICT);
780    }
781
782    #[tokio::test]
783    async fn human_input_submit_after_reject_returns_409() {
784        let (_store, state, run_id, step_id) = setup().await;
785        let auth = admin_header(&state);
786
787        let (first, _) = call(&state, &auth, run_id, step_id, "reject", None).await;
788        assert_eq!(first, StatusCode::OK);
789
790        let answer = json!({"answers": ["yes"]});
791        let (second, body) = call(&state, &auth, run_id, step_id, "input", Some(answer)).await;
792        assert_eq!(second, StatusCode::CONFLICT);
793        assert_eq!(body["error"]["message"], "input already rejected");
794    }
795
796    #[test]
797    fn human_input_validate_answer_without_a_schema_is_internal() {
798        let err = validate_answer(Some(&json!({"message": "Answer?"})), &json!({}))
799            .expect_err("no schema");
800        assert!(matches!(err, ApiError::Internal(_)), "got {err:?}");
801
802        let err = validate_answer(None, &json!({})).expect_err("no input");
803        assert!(matches!(err, ApiError::Internal(_)), "got {err:?}");
804    }
805
806    #[test]
807    fn human_input_validate_answer_resolves_defs_references() {
808        let schema = input_schema_for::<Nested>();
809        assert!(schema.get("$defs").is_some(), "schema: {schema}");
810        let input = json!({ HUMAN_INPUT_SCHEMA_KEY: schema });
811
812        let valid = validate_answer(Some(&input), &json!({"target": {"env": "prod"}}));
813        assert!(valid.is_ok(), "got {valid:?}");
814
815        let err = validate_answer(Some(&input), &json!({"target": {"env": 1}}))
816            .expect_err("wrong nested type");
817        match err {
818            ApiError::InvalidInput(errors) => assert_eq!(errors.len(), 1, "got {errors:?}"),
819            other => panic!("expected InvalidInput, got {other:?}"),
820        }
821    }
822}