Skip to main content

ironflow_api/routes/
retry_run.rs

1//! `POST /api/v1/runs/:id/retry` — Retry a failed run.
2
3use 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/// Query parameters for `POST /api/v1/runs/:id/retry`.
21#[derive(Debug, Deserialize, Default)]
22pub struct RetryQuery {
23    /// Force the retry even when the handler version has changed.
24    #[serde(default)]
25    pub force: bool,
26}
27
28/// Retry a failed run.
29///
30/// Creates a new `Pending` run with `TriggerKind::Retry` pointing to the
31/// original. Returns 400 if the run is not in a retryable state, 409 if an
32/// automatic retry is already armed, and 409 if the handler version has
33/// changed since the original run (pass `?force=true` to override).
34#[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    // Retrying means an automatic retry is already armed; creating a manual
68    // retry on top would produce a duplicate run once the timer fires.
69    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    // When the handler is still registered, use its current version.
105    // When unregistered (e.g. removed between deploys), preserve the
106    // original version so version-tracking information is not lost.
107    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            // The retry is attributed to the user who triggered it, not the
120            // original author, so the audit trail shows who actually acted.
121            created_by: Some(run_actor_of(&auth)),
122            // A retry must not inherit the parent's idempotency key: it is a
123            // new logical operation and must be eligible for its own dedup.
124            idempotency_key: None,
125            // Inherit the original cost cap so budget constraints survive retries.
126            max_cost_usd: original.max_cost_usd,
127        })
128        .await?
129        .into_run();
130
131    // Only emit RetryForced when both the handler is registered (so
132    // current_version is meaningful) and the versions actually differ.
133    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        // A run cancelled for reaching its cap is retryable.
284        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        // No duplicate run was created.
493        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    // ---- created_by ----
565
566    #[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    // ---- Idempotency-Key ----
645
646    #[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        // The original key still resolves to the original run.
700        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    // ---- Handler version compatibility ----
709
710    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}