1use axum::extract::{Path, Query, State};
4use axum::http::StatusCode;
5use axum::response::IntoResponse;
6use chrono::Utc;
7use ironflow_auth::extractor::Authenticated;
8use ironflow_engine::error::HANDLER_VERSION_MISMATCH_CODE;
9use ironflow_engine::notify::{Event, RetryForcedEvent, RunCreatedEvent};
10use ironflow_store::models::{NewRun, RunStatus, TriggerKind};
11use serde::Deserialize;
12use uuid::Uuid;
13
14use crate::actor::run_actor_of;
15use crate::entities::RunResponse;
16use crate::error::ApiError;
17use crate::response::ok;
18use crate::state::AppState;
19
20#[derive(Debug, Deserialize, Default)]
22pub struct RetryQuery {
23 #[serde(default)]
25 pub force: bool,
26}
27
28#[cfg_attr(
35 feature = "openapi",
36 utoipa::path(
37 post,
38 path = "/api/v1/runs/{id}/retry",
39 tags = ["runs"],
40 params(
41 ("id" = Uuid, Path, description = "Run ID"),
42 ("force" = Option<bool>, Query, description = "Force retry despite handler version mismatch"),
43 ),
44 responses(
45 (status = 201, description = "Run retry created successfully", body = RunResponse),
46 (status = 400, description = "Run cannot be retried"),
47 (status = 401, description = "Unauthorized"),
48 (status = 403, description = "Forbidden"),
49 (status = 404, description = "Run not found"),
50 (status = 409, description = "Version mismatch or automatic retry already armed")
51 ),
52 security(("Bearer" = []))
53 )
54)]
55pub async fn retry_run(
56 auth: Authenticated,
57 State(state): State<AppState>,
58 Path(id): Path<Uuid>,
59 Query(query): Query<RetryQuery>,
60) -> Result<impl IntoResponse, ApiError> {
61 if !auth.is_admin() {
62 return Err(ApiError::Forbidden);
63 }
64
65 let original = state.get_run_or_404(id).await?;
66
67 if original.status.state == RunStatus::Retrying {
70 return Err(ApiError::Conflict(
71 "run is already waiting for an automatic retry; cancel it first to retry manually"
72 .to_string(),
73 ));
74 }
75
76 if !matches!(
77 original.status.state,
78 RunStatus::Failed | RunStatus::Cancelled | RunStatus::Warning
79 ) {
80 return Err(ApiError::BadRequest(format!(
81 "cannot retry run in {} state",
82 original.status.state
83 )));
84 }
85
86 let force = query.force;
87 let handler = state.engine.get_handler(&original.workflow_name);
88 let current_version = handler.and_then(|h| h.version().map(str::to_string));
89
90 if let Some(handler) = &handler
91 && !handler.is_version_compatible(original.handler_version.as_deref())
92 && !force
93 {
94 return Err(ApiError::Conflict(format!(
95 "{}: handler '{}' is now at version {}, but the run was created \
96 with version {}. Pass ?force=true to override.",
97 HANDLER_VERSION_MISMATCH_CODE,
98 original.workflow_name,
99 current_version.as_deref().unwrap_or("unknown"),
100 original.handler_version.as_deref().unwrap_or("unknown"),
101 )));
102 }
103
104 let effective_version = current_version.clone().or(original.handler_version.clone());
108
109 let new_run = state
110 .store
111 .create_run(NewRun {
112 workflow_name: original.workflow_name.clone(),
113 trigger: TriggerKind::Retry { parent_run_id: id },
114 payload: original.payload,
115 max_retries: original.max_retries,
116 handler_version: effective_version,
117 labels: original.labels,
118 scheduled_at: None,
119 created_by: Some(run_actor_of(&auth)),
122 idempotency_key: None,
125 max_cost_usd: original.max_cost_usd,
127 })
128 .await?
129 .into_run();
130
131 if force && handler.is_some() && original.handler_version != current_version {
134 state
135 .engine
136 .event_publisher()
137 .publish(Event::RetryForced(RetryForcedEvent {
138 run_id: new_run.id,
139 workflow_name: original.workflow_name.clone(),
140 original_version: original.handler_version.unwrap_or_default(),
141 current_version: current_version.unwrap_or_default(),
142 at: Utc::now(),
143 }));
144 }
145
146 state
147 .engine
148 .event_publisher()
149 .publish(Event::RunCreated(RunCreatedEvent {
150 run_id: new_run.id,
151 workflow_name: new_run.workflow_name.clone(),
152 at: Utc::now(),
153 }));
154
155 Ok((StatusCode::CREATED, ok(RunResponse::from(new_run))))
156}
157
158#[cfg(test)]
159mod tests {
160 use std::collections::HashMap;
161
162 use axum::Router;
163 use axum::body::Body;
164 use axum::http::{Request, StatusCode as HttpStatusCode};
165 use axum::routing::post;
166 use http_body_util::BodyExt;
167 use ironflow_auth::jwt::AccessToken;
168 use ironflow_core::providers::claude::ClaudeCodeProvider;
169 use ironflow_engine::context::WorkflowContext;
170 use ironflow_engine::engine::Engine;
171 use ironflow_engine::handler::{HandlerFuture, WorkflowHandler};
172 use ironflow_engine::notify::Event;
173 use ironflow_store::memory::InMemoryStore;
174 use ironflow_store::models::{NewRun, NewUser, RunActor, RunStatus, TriggerKind};
175 use ironflow_store::store::RunStore;
176 use ironflow_store::user_store::UserStore;
177 use rust_decimal::Decimal;
178 use serde_json::{Value as JsonValue, from_slice, from_value, json};
179 use std::sync::Arc;
180 use tokio::sync::broadcast;
181 use tower::ServiceExt;
182 use uuid::Uuid;
183
184 use super::*;
185
186 fn make_auth_header(state: &AppState) -> String {
187 let user_id = Uuid::now_v7();
188 let token = AccessToken::for_user(user_id, "testuser", true, &state.jwt_config).unwrap();
189 format!("Bearer {}", token.0)
190 }
191
192 fn test_state(store: Arc<InMemoryStore>) -> AppState {
193 Arc::new(InMemoryStore::new());
194 let provider = Arc::new(ClaudeCodeProvider::new());
195 let engine = Arc::new(Engine::new(store.clone(), provider));
196 let jwt_config = Arc::new(ironflow_auth::jwt::JwtConfig {
197 secret: "test-secret".to_string(),
198 access_token_ttl_secs: 900,
199 refresh_token_ttl_secs: 604800,
200 cookie_domain: None,
201 cookie_secure: false,
202 });
203 let (event_sender, _) = broadcast::channel::<Event>(1);
204 AppState::new(
205 store,
206 engine,
207 jwt_config,
208 "test-worker-token".to_string(),
209 event_sender,
210 )
211 }
212
213 #[tokio::test]
214 async fn retry_failed_run() {
215 let store = Arc::new(InMemoryStore::new());
216 let run = store
217 .create_run(NewRun {
218 created_by: None,
219 workflow_name: "test".to_string(),
220 trigger: TriggerKind::Manual,
221 payload: json!({"key": "value"}),
222 max_retries: 3,
223 handler_version: None,
224 labels: HashMap::new(),
225 scheduled_at: None,
226 idempotency_key: None,
227 max_cost_usd: None,
228 })
229 .await
230 .unwrap()
231 .into_run();
232
233 store
234 .update_run_status(run.id, RunStatus::Running)
235 .await
236 .unwrap();
237 store
238 .update_run_status(run.id, RunStatus::Failed)
239 .await
240 .unwrap();
241
242 let state = test_state(store.clone());
243 let auth_header = make_auth_header(&state);
244 let app = Router::new()
245 .route("/{id}/retry", post(retry_run))
246 .with_state(state);
247
248 let req = Request::builder()
249 .method("POST")
250 .uri(format!("/{}/retry", run.id))
251 .header("content-type", "application/json")
252 .header("authorization", auth_header)
253 .body(Body::from("{}"))
254 .unwrap();
255
256 let resp = app.oneshot(req).await.unwrap();
257 assert_eq!(resp.status(), HttpStatusCode::CREATED);
258
259 let body = resp.into_body().collect().await.unwrap().to_bytes();
260 let json_val: JsonValue = from_slice(&body).unwrap();
261 let new_id: Uuid = from_value(json_val["data"]["id"].clone()).unwrap();
262
263 let new_run = store.get_run(new_id).await.unwrap().unwrap();
264 assert_eq!(new_run.status.state, RunStatus::Pending);
265 assert!(matches!(new_run.trigger, TriggerKind::Retry { .. }));
266 }
267
268 #[tokio::test]
269 async fn retry_inherits_the_original_cost_cap() {
270 let store = Arc::new(InMemoryStore::new());
271 let cap = Decimal::new(250, 2);
272 let run = store
273 .create_run(NewRun {
274 workflow_name: "test".to_string(),
275 trigger: TriggerKind::Manual,
276 payload: json!({}),
277 max_retries: 1,
278 handler_version: None,
279 labels: HashMap::new(),
280 scheduled_at: None,
281 created_by: None,
282 idempotency_key: None,
283 max_cost_usd: Some(cap),
284 })
285 .await
286 .unwrap()
287 .into_run();
288
289 store
291 .update_run_status(run.id, RunStatus::Running)
292 .await
293 .unwrap();
294 store
295 .update_run_status(run.id, RunStatus::Cancelled)
296 .await
297 .unwrap();
298
299 let state = test_state(store.clone());
300 let auth_header = make_auth_header(&state);
301 let app = Router::new()
302 .route("/{id}/retry", post(retry_run))
303 .with_state(state);
304
305 let req = Request::builder()
306 .method("POST")
307 .uri(format!("/{}/retry", run.id))
308 .header("content-type", "application/json")
309 .header("authorization", auth_header)
310 .body(Body::from("{}"))
311 .unwrap();
312
313 let resp = app.oneshot(req).await.unwrap();
314 assert_eq!(resp.status(), HttpStatusCode::CREATED);
315
316 let body = resp.into_body().collect().await.unwrap().to_bytes();
317 let json_val: JsonValue = from_slice(&body).unwrap();
318 let new_id: Uuid = from_value(json_val["data"]["id"].clone()).unwrap();
319
320 let new_run = store.get_run(new_id).await.unwrap().unwrap();
321 assert_eq!(new_run.max_cost_usd, Some(cap));
322 }
323
324 #[tokio::test]
325 async fn retry_pending_run_returns_400() {
326 let store = Arc::new(InMemoryStore::new());
327 let run = store
328 .create_run(NewRun {
329 created_by: None,
330 workflow_name: "test".to_string(),
331 trigger: TriggerKind::Manual,
332 payload: json!({}),
333 max_retries: 0,
334 handler_version: None,
335 labels: HashMap::new(),
336 scheduled_at: None,
337 idempotency_key: None,
338 max_cost_usd: None,
339 })
340 .await
341 .unwrap()
342 .into_run();
343
344 let state = test_state(store);
345 let auth_header = make_auth_header(&state);
346 let app = Router::new()
347 .route("/{id}/retry", post(retry_run))
348 .with_state(state);
349
350 let req = Request::builder()
351 .method("POST")
352 .uri(format!("/{}/retry", run.id))
353 .header("content-type", "application/json")
354 .header("authorization", auth_header)
355 .body(Body::from("{}"))
356 .unwrap();
357
358 let resp = app.oneshot(req).await.unwrap();
359 assert_eq!(resp.status(), HttpStatusCode::BAD_REQUEST);
360 }
361
362 #[tokio::test]
363 async fn retry_completed_run_returns_400() {
364 let store = Arc::new(InMemoryStore::new());
365 let run = store
366 .create_run(NewRun {
367 created_by: None,
368 workflow_name: "test".to_string(),
369 trigger: TriggerKind::Manual,
370 payload: json!({}),
371 max_retries: 0,
372 handler_version: None,
373 labels: HashMap::new(),
374 scheduled_at: None,
375 idempotency_key: None,
376 max_cost_usd: None,
377 })
378 .await
379 .unwrap()
380 .into_run();
381
382 store
383 .update_run_status(run.id, RunStatus::Running)
384 .await
385 .unwrap();
386 store
387 .update_run_status(run.id, RunStatus::Completed)
388 .await
389 .unwrap();
390
391 let state = test_state(store);
392 let auth_header = make_auth_header(&state);
393 let app = Router::new()
394 .route("/{id}/retry", post(retry_run))
395 .with_state(state);
396
397 let req = Request::builder()
398 .method("POST")
399 .uri(format!("/{}/retry", run.id))
400 .header("content-type", "application/json")
401 .header("authorization", auth_header)
402 .body(Body::from("{}"))
403 .unwrap();
404
405 let resp = app.oneshot(req).await.unwrap();
406 assert_eq!(resp.status(), HttpStatusCode::BAD_REQUEST);
407 }
408
409 #[tokio::test]
410 async fn retry_running_run_returns_400() {
411 let store = Arc::new(InMemoryStore::new());
412 let run = store
413 .create_run(NewRun {
414 created_by: None,
415 workflow_name: "test".to_string(),
416 trigger: TriggerKind::Manual,
417 payload: json!({}),
418 max_retries: 0,
419 handler_version: None,
420 labels: HashMap::new(),
421 scheduled_at: None,
422 idempotency_key: None,
423 max_cost_usd: None,
424 })
425 .await
426 .unwrap()
427 .into_run();
428
429 store
430 .update_run_status(run.id, RunStatus::Running)
431 .await
432 .unwrap();
433
434 let state = test_state(store);
435 let auth_header = make_auth_header(&state);
436 let app = Router::new()
437 .route("/{id}/retry", post(retry_run))
438 .with_state(state);
439
440 let req = Request::builder()
441 .method("POST")
442 .uri(format!("/{}/retry", run.id))
443 .header("content-type", "application/json")
444 .header("authorization", auth_header)
445 .body(Body::from("{}"))
446 .unwrap();
447
448 let resp = app.oneshot(req).await.unwrap();
449 assert_eq!(resp.status(), HttpStatusCode::BAD_REQUEST);
450 }
451
452 #[tokio::test]
453 async fn retry_run_awaiting_automatic_retry_returns_409() {
454 let store = Arc::new(InMemoryStore::new());
455 let run = store
456 .create_run(NewRun {
457 created_by: None,
458 workflow_name: "test".to_string(),
459 trigger: TriggerKind::Manual,
460 payload: json!({}),
461 max_retries: 3,
462 handler_version: None,
463 labels: HashMap::new(),
464 scheduled_at: None,
465 idempotency_key: None,
466 max_cost_usd: None,
467 })
468 .await
469 .unwrap()
470 .into_run();
471
472 store
473 .update_run_status(run.id, RunStatus::Running)
474 .await
475 .unwrap();
476 store
477 .update_run_status(run.id, RunStatus::Retrying)
478 .await
479 .unwrap();
480
481 let state = test_state(store.clone());
482 let auth_header = make_auth_header(&state);
483 let app = Router::new()
484 .route("/{id}/retry", post(retry_run))
485 .with_state(state);
486
487 let req = Request::builder()
488 .method("POST")
489 .uri(format!("/{}/retry", run.id))
490 .header("content-type", "application/json")
491 .header("authorization", auth_header)
492 .body(Body::from("{}"))
493 .unwrap();
494
495 let resp = app.oneshot(req).await.unwrap();
496 assert_eq!(resp.status(), HttpStatusCode::CONFLICT);
497
498 let runs = store
500 .list_runs(ironflow_store::models::RunFilter::default(), 1, 50)
501 .await
502 .unwrap();
503 assert_eq!(runs.total, 1);
504 }
505
506 #[tokio::test]
507 async fn retry_cancelled_run_is_allowed() {
508 let store = Arc::new(InMemoryStore::new());
509 let run = store
510 .create_run(NewRun {
511 created_by: None,
512 workflow_name: "test".to_string(),
513 trigger: TriggerKind::Manual,
514 payload: json!({}),
515 max_retries: 0,
516 handler_version: None,
517 labels: HashMap::new(),
518 scheduled_at: None,
519 idempotency_key: None,
520 max_cost_usd: None,
521 })
522 .await
523 .unwrap()
524 .into_run();
525
526 store
527 .update_run_status(run.id, RunStatus::Cancelled)
528 .await
529 .unwrap();
530
531 let state = test_state(store);
532 let auth_header = make_auth_header(&state);
533 let app = Router::new()
534 .route("/{id}/retry", post(retry_run))
535 .with_state(state);
536
537 let req = Request::builder()
538 .method("POST")
539 .uri(format!("/{}/retry", run.id))
540 .header("content-type", "application/json")
541 .header("authorization", auth_header)
542 .body(Body::from("{}"))
543 .unwrap();
544
545 let resp = app.oneshot(req).await.unwrap();
546 assert_eq!(resp.status(), HttpStatusCode::CREATED);
547 }
548
549 #[tokio::test]
550 async fn retry_nonexistent_run_returns_404() {
551 let store = Arc::new(InMemoryStore::new());
552 let state = test_state(store);
553 let auth_header = make_auth_header(&state);
554 let app = Router::new()
555 .route("/{id}/retry", post(retry_run))
556 .with_state(state);
557
558 let req = Request::builder()
559 .method("POST")
560 .uri(format!("/{}/retry", Uuid::now_v7()))
561 .header("content-type", "application/json")
562 .header("authorization", auth_header)
563 .body(Body::from("{}"))
564 .unwrap();
565
566 let resp = app.oneshot(req).await.unwrap();
567 assert_eq!(resp.status(), HttpStatusCode::NOT_FOUND);
568 }
569
570 #[tokio::test]
573 async fn retry_attributes_the_new_run_to_the_caller_not_the_original_author() {
574 let store = Arc::new(InMemoryStore::new());
575 let original_author = store
576 .create_user(NewUser {
577 email: "alice@example.com".to_string(),
578 username: "alice".to_string(),
579 password_hash: "hash".to_string(),
580 is_admin: Some(true),
581 })
582 .await
583 .unwrap();
584 let retrying_user = store
585 .create_user(NewUser {
586 email: "bob@example.com".to_string(),
587 username: "bob".to_string(),
588 password_hash: "hash".to_string(),
589 is_admin: Some(true),
590 })
591 .await
592 .unwrap();
593
594 let run = store
595 .create_run(NewRun {
596 workflow_name: "test".to_string(),
597 trigger: TriggerKind::Api,
598 payload: json!({}),
599 max_retries: 3,
600 handler_version: None,
601 labels: HashMap::new(),
602 scheduled_at: None,
603 created_by: Some(RunActor::User {
604 user_id: original_author.id,
605 }),
606 idempotency_key: None,
607 max_cost_usd: None,
608 })
609 .await
610 .unwrap()
611 .into_run();
612
613 store
614 .update_run_status(run.id, RunStatus::Running)
615 .await
616 .unwrap();
617 store
618 .update_run_status(run.id, RunStatus::Failed)
619 .await
620 .unwrap();
621
622 let state = test_state(store.clone());
623 let token =
624 AccessToken::for_user(retrying_user.id, "bob", true, &state.jwt_config).unwrap();
625 let app = Router::new()
626 .route("/{id}/retry", post(retry_run))
627 .with_state(state);
628
629 let req = Request::builder()
630 .method("POST")
631 .uri(format!("/{}/retry", run.id))
632 .header("content-type", "application/json")
633 .header("authorization", format!("Bearer {}", token.0))
634 .body(Body::from("{}"))
635 .unwrap();
636
637 let resp = app.oneshot(req).await.unwrap();
638 assert_eq!(resp.status(), HttpStatusCode::CREATED);
639
640 let body = resp.into_body().collect().await.unwrap().to_bytes();
641 let json_val: JsonValue = from_slice(&body).unwrap();
642 assert_eq!(json_val["data"]["created_by"]["kind"], "user");
643 assert_eq!(
644 json_val["data"]["created_by"]["id"],
645 retrying_user.id.to_string()
646 );
647 assert_eq!(json_val["data"]["created_by"]["label"], "bob");
648 }
649
650 #[tokio::test]
653 async fn retry_does_not_inherit_the_idempotency_key() {
654 let store = Arc::new(InMemoryStore::new());
655 let run = store
656 .create_run(NewRun {
657 workflow_name: "test".to_string(),
658 trigger: TriggerKind::Api,
659 payload: json!({"key": "value"}),
660 max_retries: 3,
661 handler_version: None,
662 labels: HashMap::new(),
663 scheduled_at: None,
664 created_by: None,
665 idempotency_key: Some("github:abc-123".to_string()),
666 max_cost_usd: None,
667 })
668 .await
669 .unwrap()
670 .into_run();
671
672 store
673 .update_run_status(run.id, RunStatus::Running)
674 .await
675 .unwrap();
676 store
677 .update_run_status(run.id, RunStatus::Failed)
678 .await
679 .unwrap();
680
681 let state = test_state(store.clone());
682 let auth_header = make_auth_header(&state);
683 let app = Router::new()
684 .route("/{id}/retry", post(retry_run))
685 .with_state(state);
686
687 let req = Request::builder()
688 .method("POST")
689 .uri(format!("/{}/retry", run.id))
690 .header("content-type", "application/json")
691 .header("authorization", auth_header)
692 .body(Body::from("{}"))
693 .unwrap();
694
695 let resp = app.oneshot(req).await.unwrap();
696 assert_eq!(resp.status(), HttpStatusCode::CREATED);
697
698 let body = resp.into_body().collect().await.unwrap().to_bytes();
699 let json_val: JsonValue = from_slice(&body).unwrap();
700 assert!(
701 json_val["data"].get("idempotency_key").is_none(),
702 "the retry run must not carry the original key"
703 );
704
705 let bound = store
707 .find_run_by_idempotency_key("github:abc-123")
708 .await
709 .unwrap()
710 .expect("key still bound");
711 assert_eq!(bound.id, run.id);
712 }
713
714 struct V2Handler;
717 impl WorkflowHandler for V2Handler {
718 fn name(&self) -> &str {
719 "versioned-wf"
720 }
721 fn version(&self) -> Option<&str> {
722 Some("2.0.0")
723 }
724 fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
725 Box::pin(async { Ok(()) })
726 }
727 }
728
729 struct V2CompatHandler;
730 impl WorkflowHandler for V2CompatHandler {
731 fn name(&self) -> &str {
732 "compat-wf"
733 }
734 fn version(&self) -> Option<&str> {
735 Some("2.0.0")
736 }
737 fn compatible_versions(&self) -> &[&str] {
738 &["1.0.0"]
739 }
740 fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
741 Box::pin(async { Ok(()) })
742 }
743 }
744
745 fn test_state_with_handlers(store: Arc<InMemoryStore>) -> AppState {
746 let provider = Arc::new(ClaudeCodeProvider::new());
747 let mut engine = Engine::new(store.clone(), provider);
748 engine.register(V2Handler).unwrap();
749 engine.register(V2CompatHandler).unwrap();
750 let jwt_config = Arc::new(ironflow_auth::jwt::JwtConfig {
751 secret: "test-secret".to_string(),
752 access_token_ttl_secs: 900,
753 refresh_token_ttl_secs: 604800,
754 cookie_domain: None,
755 cookie_secure: false,
756 });
757 let (event_sender, _) = broadcast::channel::<Event>(1);
758 AppState::new(
759 store,
760 Arc::new(engine),
761 jwt_config,
762 "test-worker-token".to_string(),
763 event_sender,
764 )
765 }
766
767 async fn create_failed_run(
768 store: &Arc<InMemoryStore>,
769 workflow_name: &str,
770 handler_version: Option<&str>,
771 ) -> ironflow_store::models::Run {
772 let run = store
773 .create_run(NewRun {
774 created_by: None,
775 workflow_name: workflow_name.to_string(),
776 trigger: TriggerKind::Manual,
777 payload: json!({}),
778 max_retries: 0,
779 handler_version: handler_version.map(str::to_string),
780 labels: HashMap::new(),
781 scheduled_at: None,
782 idempotency_key: None,
783 max_cost_usd: None,
784 })
785 .await
786 .unwrap()
787 .into_run();
788
789 store
790 .update_run_status(run.id, RunStatus::Running)
791 .await
792 .unwrap();
793 store
794 .update_run_status(run.id, RunStatus::Failed)
795 .await
796 .unwrap();
797 run
798 }
799
800 #[tokio::test]
801 async fn retry_same_version_succeeds() {
802 let store = Arc::new(InMemoryStore::new());
803 let run = create_failed_run(&store, "versioned-wf", Some("2.0.0")).await;
804
805 let state = test_state_with_handlers(store);
806 let auth_header = make_auth_header(&state);
807 let app = Router::new()
808 .route("/{id}/retry", post(retry_run))
809 .with_state(state);
810
811 let req = Request::builder()
812 .method("POST")
813 .uri(format!("/{}/retry", run.id))
814 .header("content-type", "application/json")
815 .header("authorization", auth_header)
816 .body(Body::from("{}"))
817 .unwrap();
818
819 let resp = app.oneshot(req).await.unwrap();
820 assert_eq!(resp.status(), HttpStatusCode::CREATED);
821 }
822
823 #[tokio::test]
824 async fn retry_version_mismatch_returns_409() {
825 let store = Arc::new(InMemoryStore::new());
826 let run = create_failed_run(&store, "versioned-wf", Some("1.0.0")).await;
827
828 let state = test_state_with_handlers(store);
829 let auth_header = make_auth_header(&state);
830 let app = Router::new()
831 .route("/{id}/retry", post(retry_run))
832 .with_state(state);
833
834 let req = Request::builder()
835 .method("POST")
836 .uri(format!("/{}/retry", run.id))
837 .header("content-type", "application/json")
838 .header("authorization", auth_header)
839 .body(Body::from("{}"))
840 .unwrap();
841
842 let resp = app.oneshot(req).await.unwrap();
843 assert_eq!(resp.status(), HttpStatusCode::CONFLICT);
844
845 let body = resp.into_body().collect().await.unwrap().to_bytes();
846 let json_val: JsonValue = from_slice(&body).unwrap();
847 let msg = json_val["error"]["message"].as_str().unwrap();
848 assert!(msg.contains("HANDLER_VERSION_MISMATCH"));
849 assert!(msg.contains("2.0.0"));
850 assert!(msg.contains("1.0.0"));
851 }
852
853 #[tokio::test]
854 async fn retry_version_mismatch_with_force_succeeds() {
855 let store = Arc::new(InMemoryStore::new());
856 let run = create_failed_run(&store, "versioned-wf", Some("1.0.0")).await;
857
858 let state = test_state_with_handlers(store.clone());
859 let auth_header = make_auth_header(&state);
860 let app = Router::new()
861 .route("/{id}/retry", post(retry_run))
862 .with_state(state);
863
864 let req = Request::builder()
865 .method("POST")
866 .uri(format!("/{}/retry?force=true", run.id))
867 .header("content-type", "application/json")
868 .header("authorization", auth_header)
869 .body(Body::from("{}"))
870 .unwrap();
871
872 let resp = app.oneshot(req).await.unwrap();
873 assert_eq!(resp.status(), HttpStatusCode::CREATED);
874
875 let body = resp.into_body().collect().await.unwrap().to_bytes();
876 let json_val: JsonValue = from_slice(&body).unwrap();
877 let new_id: Uuid = from_value(json_val["data"]["id"].clone()).unwrap();
878 let new_run = store.get_run(new_id).await.unwrap().unwrap();
879 assert_eq!(new_run.handler_version, Some("2.0.0".to_string()));
880 }
881
882 #[tokio::test]
883 async fn retry_compatible_version_succeeds() {
884 let store = Arc::new(InMemoryStore::new());
885 let run = create_failed_run(&store, "compat-wf", Some("1.0.0")).await;
886
887 let state = test_state_with_handlers(store);
888 let auth_header = make_auth_header(&state);
889 let app = Router::new()
890 .route("/{id}/retry", post(retry_run))
891 .with_state(state);
892
893 let req = Request::builder()
894 .method("POST")
895 .uri(format!("/{}/retry", run.id))
896 .header("content-type", "application/json")
897 .header("authorization", auth_header)
898 .body(Body::from("{}"))
899 .unwrap();
900
901 let resp = app.oneshot(req).await.unwrap();
902 assert_eq!(resp.status(), HttpStatusCode::CREATED);
903 }
904
905 #[tokio::test]
906 async fn retry_null_version_is_compatible() {
907 let store = Arc::new(InMemoryStore::new());
908 let run = create_failed_run(&store, "versioned-wf", None).await;
909
910 let state = test_state_with_handlers(store);
911 let auth_header = make_auth_header(&state);
912 let app = Router::new()
913 .route("/{id}/retry", post(retry_run))
914 .with_state(state);
915
916 let req = Request::builder()
917 .method("POST")
918 .uri(format!("/{}/retry", run.id))
919 .header("content-type", "application/json")
920 .header("authorization", auth_header)
921 .body(Body::from("{}"))
922 .unwrap();
923
924 let resp = app.oneshot(req).await.unwrap();
925 assert_eq!(resp.status(), HttpStatusCode::CREATED);
926 }
927
928 #[tokio::test]
929 async fn retry_writes_current_handler_version() {
930 let store = Arc::new(InMemoryStore::new());
931 let run = create_failed_run(&store, "versioned-wf", Some("2.0.0")).await;
932
933 let state = test_state_with_handlers(store.clone());
934 let auth_header = make_auth_header(&state);
935 let app = Router::new()
936 .route("/{id}/retry", post(retry_run))
937 .with_state(state);
938
939 let req = Request::builder()
940 .method("POST")
941 .uri(format!("/{}/retry", run.id))
942 .header("content-type", "application/json")
943 .header("authorization", auth_header)
944 .body(Body::from("{}"))
945 .unwrap();
946
947 let resp = app.oneshot(req).await.unwrap();
948 assert_eq!(resp.status(), HttpStatusCode::CREATED);
949
950 let body = resp.into_body().collect().await.unwrap().to_bytes();
951 let json_val: JsonValue = from_slice(&body).unwrap();
952 let new_id: Uuid = from_value(json_val["data"]["id"].clone()).unwrap();
953 let new_run = store.get_run(new_id).await.unwrap().unwrap();
954 assert_eq!(new_run.handler_version, Some("2.0.0".to_string()));
955 }
956}