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