diff --git a/backend/sqlx-data.json b/backend/sqlx-data.json index 993c1f4fe5..752bd3286e 100644 --- a/backend/sqlx-data.json +++ b/backend/sqlx-data.json @@ -2306,6 +2306,28 @@ }, "query": "INSERT INTO workspace_settings\n (workspace_id, slack_team_id, slack_name)\n VALUES ($1, $2, $3) ON CONFLICT (workspace_id) DO UPDATE SET slack_team_id = $2, slack_name = $3" }, + "a241c56415759105ccbcbf7fff77287fa4ec2cc096c0060d14db421115d63e2d": { + "describe": { + "columns": [ + { + "name": "exists", + "ordinal": 0, + "type_info": "Bool" + } + ], + "nullable": [ + null + ], + "parameters": { + "Left": [ + "Text", + "Text", + "Timestamptz" + ] + } + }, + "query": "SELECT EXISTS (SELECT 1 FROM queue WHERE workspace_id = $1 AND schedule_path = $2 AND scheduled_for = $3)" + }, "a34066d4a1578a13b2e322e6936ae80a0239a79148f3edce65f51a93910a1a4b": { "describe": { "columns": [ diff --git a/backend/src/schedule.rs b/backend/src/schedule.rs index 22944d9d76..ab7833512e 100644 --- a/backend/src/schedule.rs +++ b/backend/src/schedule.rs @@ -25,7 +25,7 @@ use axum::{ use chrono::{DateTime, Duration, FixedOffset}; use serde::{Deserialize, Serialize}; use serde_json::{Map, Value}; -use sqlx::{FromRow, Postgres, Transaction}; +use sqlx::{query_scalar, FromRow, Postgres, Transaction}; pub fn workspaced_service() -> Router { Router::new() @@ -83,6 +83,20 @@ pub async fn push_scheduled_job<'c>( .expect("a schedule should have a next event") + offset; + let already_exists: bool = query_scalar!( + "SELECT EXISTS (SELECT 1 FROM queue WHERE workspace_id = $1 AND schedule_path = $2 AND scheduled_for = $3)", + &schedule.workspace_id, + &schedule.path, + next + ) + .fetch_one(&mut tx) + .await? + .unwrap_or(false); + + if already_exists { + return Ok(tx); + } + let mut args: Option> = None; if let Some(args_v) = schedule.args { @@ -108,6 +122,7 @@ pub async fn push_scheduled_job<'c>( path: schedule.script_path, } }; + let (_, tx) = push( tx, &schedule.workspace_id, diff --git a/backend/src/worker_flow.rs b/backend/src/worker_flow.rs index 378611eabe..5070c602ab 100644 --- a/backend/src/worker_flow.rs +++ b/backend/src/worker_flow.rs @@ -6,8 +6,8 @@ use crate::{ error::{self, Error}, flows::{FlowModule, FlowModuleValue, FlowValue, InputTransform, Retry, Suspend}, jobs::{ - add_completed_job, add_completed_job_error, get_queued_job, postprocess_queued_job, push, - script_path_to_payload, JobPayload, QueuedJob, RawCode, + add_completed_job, add_completed_job_error, delete_job, get_queued_job, push, + schedule_again_if_scheduled, script_path_to_payload, JobPayload, QueuedJob, RawCode, }, js_eval::{eval_timeout, EvalCreds, IdContext}, more_serde::is_default, @@ -374,6 +374,16 @@ pub async fn update_flow_status_after_job_completion( tx.commit().await?; + if old_status.step == 0 && !flow_job.is_flow_step { + schedule_again_if_scheduled( + flow_job.schedule_path.clone(), + flow_job.script_path.clone(), + &w_id, + db, + ) + .await?; + } + let done = if !should_continue_flow { let logs = if flow_job.canceled { "Flow job canceled".to_string() @@ -421,15 +431,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, - flow, - db, - ) - .await?; + let _ = delete_job(db, w_id, flow).await?; if flow_job.same_worker && !keep_job_dir { let _ = tokio::fs::remove_dir_all(format!("{worker_dir}/{}", flow_job.id)).await; @@ -883,15 +885,7 @@ async fn push_next_flow_job( let result = is_cancelled.unwrap_or(json!({ "error": logs })); let _uuid = add_completed_job(db, &flow_job, success, skipped, result, logs).await?; - postprocess_queued_job( - false, - flow_job.schedule_path.clone(), - flow_job.script_path.clone(), - &flow_job.workspace_id, - flow_job.id, - db, - ) - .await?; + let _ = delete_job(db, &flow_job.workspace_id, flow_job.id).await?; return Ok(()); } @@ -1236,15 +1230,7 @@ async fn jump_to_next_step( let skipped = false; let logs = "Forloop completed without iteration".to_string(); let _uuid = add_completed_job(db, &new_job, success, skipped, json!([]), logs).await?; - postprocess_queued_job( - false, - new_job.schedule_path, - new_job.script_path, - &new_job.workspace_id, - new_job.id, - db, - ) - .await?; + let _ = delete_job(db, &new_job.workspace_id, new_job.id).await?; return Ok(()); } }