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, 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/// 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
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        // A run cancelled for reaching its cap is retryable.
290        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        // No duplicate run was created.
499        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    // ---- created_by ----
571
572    #[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    // ---- Idempotency-Key ----
651
652    #[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        // The original key still resolves to the original run.
706        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    // ---- Handler version compatibility ----
715
716    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}