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_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#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
30#[derive(Debug, Default, Deserialize, Serialize)]
31pub struct RejectHumanInputRequest {
32 #[serde(default)]
34 pub reason: Option<String>,
35}
36
37#[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 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#[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
195async 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
239fn 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
258fn resume_status(state: &AppState) -> RunStatus {
261 match state.engine.execution_mode() {
262 ExecutionMode::Local => RunStatus::Running,
263 ExecutionMode::Workers => RunStatus::Pending,
264 }
265}
266
267fn 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 #[allow(dead_code)]
313 #[derive(Deserialize, JsonSchema)]
314 struct Answers {
315 answers: Vec<String>,
316 }
317
318 #[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 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 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 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 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 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}