diff --git a/backend/src/jobs.rs b/backend/src/jobs.rs index 6ceca3cd00..3c4a5894d2 100644 --- a/backend/src/jobs.rs +++ b/backend/src/jobs.rs @@ -1331,6 +1331,7 @@ pub async fn add_completed_job( #[instrument(level = "trace", skip_all)] pub async fn postprocess_queued_job( + is_flow_step: bool, schedule_path: Option, script_path: Option, w_id: &str, @@ -1338,7 +1339,9 @@ pub async fn postprocess_queued_job( db: &DB, ) -> crate::error::Result<()> { let _ = delete_job(db, w_id, job_id).await?; - schedule_again_if_scheduled(schedule_path, script_path, &w_id, db).await?; + if !is_flow_step { + schedule_again_if_scheduled(schedule_path, script_path, &w_id, db).await?; + } Ok(()) } diff --git a/backend/src/worker.rs b/backend/src/worker.rs index 9c9efa8708..930f6e21fb 100644 --- a/backend/src/worker.rs +++ b/backend/src/worker.rs @@ -139,6 +139,7 @@ pub async fn run_worker( { let job = job2.clone(); let _ = postprocess_queued_job( + job.is_flow_step, job.schedule_path, job.script_path, &job2.workspace_id, @@ -249,8 +250,15 @@ async fn handle_queued_job( } }; - let _ = - postprocess_queued_job(job.schedule_path, job.script_path, &w_id, job_id, db).await; + let _ = postprocess_queued_job( + job.is_flow_step, + job.schedule_path, + job.script_path, + &w_id, + job_id, + db, + ) + .await; } } Ok(()) diff --git a/backend/src/worker_flow.rs b/backend/src/worker_flow.rs index 62d373f9ff..48fec1b331 100644 --- a/backend/src/worker_flow.rs +++ b/backend/src/worker_flow.rs @@ -206,6 +206,7 @@ pub async fn update_flow_status_after_job_completion( if done { postprocess_queued_job( + flow_job.is_flow_step, flow_job.schedule_path.clone(), flow_job.script_path.clone(), &w_id,