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