diff --git a/backend/windmill-api/src/jobs.rs b/backend/windmill-api/src/jobs.rs index c29e4e6f02..5332bcb6c4 100644 --- a/backend/windmill-api/src/jobs.rs +++ b/backend/windmill-api/src/jobs.rs @@ -3741,7 +3741,7 @@ async fn batch_rerun_handle_job( PushArgsOwned { extra: None, args }, ) .await; - if let Ok(uuid) = result { + if let Ok((uuid, _)) = result { return Ok(uuid.to_string()); } } @@ -3798,7 +3798,7 @@ pub async fn run_flow_by_path( ) .await?; - let uuid = + let (uuid, _) = run_flow_by_path_inner(authed, db, user_db, w_id, flow_path, run_query, args).await?; Ok((StatusCode::CREATED, uuid.to_string())) @@ -3812,7 +3812,7 @@ pub async fn run_flow_by_path_inner( flow_path: StripPath, run_query: RunJobQuery, args: PushArgsOwned, -) -> error::Result { +) -> error::Result<(Uuid, Option)> { #[cfg(feature = "enterprise")] check_license_key_valid().await?; @@ -3828,6 +3828,7 @@ pub async fn run_flow_by_path_inner( has_preprocessor, on_behalf_of_email, edited_by, + early_return, .. } = get_latest_flow_version_info_for_path(&mut *tx, &w_id, &flow_path, true).await?; @@ -3889,7 +3890,7 @@ pub async fn run_flow_by_path_inner( ) .await?; tx.commit().await?; - Ok(uuid) + Ok((uuid, early_return)) } #[cfg(not(feature = "enterprise"))] diff --git a/backend/windmill-api/src/trigger_helpers.rs b/backend/windmill-api/src/trigger_helpers.rs index 723a966bca..1c2bdc029c 100644 --- a/backend/windmill-api/src/trigger_helpers.rs +++ b/backend/windmill-api/src/trigger_helpers.rs @@ -472,12 +472,12 @@ async fn trigger_runnable_inner( error_handler_path: Option<&str>, error_handler_args: Option<&sqlx::types::Json>>>, trigger_path: String, -) -> Result<(Uuid, Option)> { +) -> Result<(Uuid, Option, Option)> { let user_db = user_db.unwrap_or_else(|| UserDB::new(db.clone())); - let (uuid, delete_after_use) = if is_flow { + let (uuid, delete_after_use, early_return) = if is_flow { let run_query = RunJobQuery::default(); let path = StripPath(runnable_path.to_string()); - let uuid = run_flow_by_path_inner( + let (uuid, early_return) = run_flow_by_path_inner( authed, db.clone(), user_db, @@ -487,9 +487,9 @@ async fn trigger_runnable_inner( args, ) .await?; - (uuid, None) + (uuid, None, early_return) } else { - trigger_script_internal( + let (uuid, delete_after_use) = trigger_script_internal( db, user_db, authed, @@ -501,10 +501,11 @@ async fn trigger_runnable_inner( error_handler_args, trigger_path, ) - .await? + .await?; + (uuid, delete_after_use, None) }; - Ok((uuid, delete_after_use)) + Ok((uuid, delete_after_use, early_return)) } #[allow(dead_code)] @@ -521,7 +522,7 @@ pub async fn trigger_runnable( error_handler_args: Option<&sqlx::types::Json>>>, trigger_path: String, ) -> Result { - let (uuid, _) = trigger_runnable_inner( + let (uuid, _, _) = trigger_runnable_inner( db, user_db, authed, @@ -553,7 +554,7 @@ pub async fn trigger_runnable_and_wait_for_result( trigger_path: String, ) -> Result { let username = authed.username.clone(); - let (uuid, delete_after_use) = trigger_runnable_inner( + let (uuid, delete_after_use, early_return) = trigger_runnable_inner( db, user_db, authed, @@ -568,7 +569,8 @@ pub async fn trigger_runnable_and_wait_for_result( ) .await?; let (result, success) = - run_wait_result_internal(db, uuid, workspace_id.to_string(), None, &username).await?; + run_wait_result_internal(db, uuid, workspace_id.to_string(), early_return, &username) + .await?; if delete_after_use.unwrap_or(false) { delete_job_metadata_after_use(&db, uuid).await?; @@ -592,7 +594,7 @@ pub async fn trigger_runnable_and_wait_for_raw_result( trigger_path: String, ) -> Result> { let username = authed.username.clone(); - let (uuid, delete_after_use) = trigger_runnable_inner( + let (uuid, delete_after_use, early_return) = trigger_runnable_inner( db, user_db, authed, @@ -607,23 +609,6 @@ pub async fn trigger_runnable_and_wait_for_raw_result( ) .await?; - let early_return = if is_flow { - sqlx::query_scalar!( - r#"SELECT flow_version.value->>'early_return' as early_return - FROM flow - LEFT JOIN flow_version - ON flow_version.id = flow.versions[array_upper(flow.versions, 1)] - WHERE flow.path = $1 and flow.workspace_id = $2"#, - runnable_path, - workspace_id, - ) - .fetch_optional(db) - .await? - .flatten() - } else { - None - }; - let (result, success) = run_wait_result_internal(db, uuid, workspace_id.to_string(), early_return, &username) .await