From 2fd32aba35c71e456c07db0b8f5ffed3c121621b Mon Sep 17 00:00:00 2001 From: hugocasa Date: Thu, 23 Oct 2025 10:36:21 +0200 Subject: [PATCH] feat(backend): stream early return (#6896) * feat(backend): stream early return * revert early return change * nit * sqlx * fix build --- ...8d6fc79c331431326250c1a59674054ceaabd.json | 22 -- ...4628b45a3c2ed6a6eb65ac8815992d5114e70.json | 22 ++ ...84aad096911efd453dd50f6865e53c1e77881.json | 46 ++++ ...acde9f57cfaba047c5b7680e0dccd2c1507df.json | 40 ---- backend/windmill-worker/src/common.rs | 212 ++++++++++++------ 5 files changed, 209 insertions(+), 133 deletions(-) delete mode 100644 backend/.sqlx/query-0997b46bae6e2374b568e8367898d6fc79c331431326250c1a59674054ceaabd.json create mode 100644 backend/.sqlx/query-50261db5d492cb9ba56e6ebc3eb4628b45a3c2ed6a6eb65ac8815992d5114e70.json create mode 100644 backend/.sqlx/query-5c13c681df0c57fcdbd1364f0f084aad096911efd453dd50f6865e53c1e77881.json delete mode 100644 backend/.sqlx/query-d8ef35b4990eb9b2a306494d5b9acde9f57cfaba047c5b7680e0dccd2c1507df.json diff --git a/backend/.sqlx/query-0997b46bae6e2374b568e8367898d6fc79c331431326250c1a59674054ceaabd.json b/backend/.sqlx/query-0997b46bae6e2374b568e8367898d6fc79c331431326250c1a59674054ceaabd.json deleted file mode 100644 index ef5348cf50..0000000000 --- a/backend/.sqlx/query-0997b46bae6e2374b568e8367898d6fc79c331431326250c1a59674054ceaabd.json +++ /dev/null @@ -1,22 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "SELECT parent_job FROM v2_job WHERE id = $1", - "describe": { - "columns": [ - { - "ordinal": 0, - "name": "parent_job", - "type_info": "Uuid" - } - ], - "parameters": { - "Left": [ - "Uuid" - ] - }, - "nullable": [ - true - ] - }, - "hash": "0997b46bae6e2374b568e8367898d6fc79c331431326250c1a59674054ceaabd" -} diff --git a/backend/.sqlx/query-50261db5d492cb9ba56e6ebc3eb4628b45a3c2ed6a6eb65ac8815992d5114e70.json b/backend/.sqlx/query-50261db5d492cb9ba56e6ebc3eb4628b45a3c2ed6a6eb65ac8815992d5114e70.json new file mode 100644 index 0000000000..31f9825cda --- /dev/null +++ b/backend/.sqlx/query-50261db5d492cb9ba56e6ebc3eb4628b45a3c2ed6a6eb65ac8815992d5114e70.json @@ -0,0 +1,22 @@ +{ + "db_name": "PostgreSQL", + "query": "\n SELECT fv.value->>'early_return' as \"early_return\"\n FROM v2_job j\n INNER JOIN flow_version fv ON fv.id = j.runnable_id\n WHERE j.id = $1\n ", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "early_return", + "type_info": "Text" + } + ], + "parameters": { + "Left": [ + "Uuid" + ] + }, + "nullable": [ + null + ] + }, + "hash": "50261db5d492cb9ba56e6ebc3eb4628b45a3c2ed6a6eb65ac8815992d5114e70" +} diff --git a/backend/.sqlx/query-5c13c681df0c57fcdbd1364f0f084aad096911efd453dd50f6865e53c1e77881.json b/backend/.sqlx/query-5c13c681df0c57fcdbd1364f0f084aad096911efd453dd50f6865e53c1e77881.json new file mode 100644 index 0000000000..4d756dbaa7 --- /dev/null +++ b/backend/.sqlx/query-5c13c681df0c57fcdbd1364f0f084aad096911efd453dd50f6865e53c1e77881.json @@ -0,0 +1,46 @@ +{ + "db_name": "PostgreSQL", + "query": "SELECT\n flow_step_id,\n (flow_status->'step')::integer as step,\n jsonb_array_length(flow_status->'modules') as len,\n runnable_path ~ '/branchone-\\d+$' as is_branch_one,\n parent_job as next_parent\n FROM v2_job\n LEFT JOIN v2_job_status USING (id)\n WHERE v2_job.id = $1", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "flow_step_id", + "type_info": "Varchar" + }, + { + "ordinal": 1, + "name": "step", + "type_info": "Int4" + }, + { + "ordinal": 2, + "name": "len", + "type_info": "Int4" + }, + { + "ordinal": 3, + "name": "is_branch_one", + "type_info": "Bool" + }, + { + "ordinal": 4, + "name": "next_parent", + "type_info": "Uuid" + } + ], + "parameters": { + "Left": [ + "Uuid" + ] + }, + "nullable": [ + true, + null, + null, + null, + true + ] + }, + "hash": "5c13c681df0c57fcdbd1364f0f084aad096911efd453dd50f6865e53c1e77881" +} diff --git a/backend/.sqlx/query-d8ef35b4990eb9b2a306494d5b9acde9f57cfaba047c5b7680e0dccd2c1507df.json b/backend/.sqlx/query-d8ef35b4990eb9b2a306494d5b9acde9f57cfaba047c5b7680e0dccd2c1507df.json deleted file mode 100644 index 43ce3fada4..0000000000 --- a/backend/.sqlx/query-d8ef35b4990eb9b2a306494d5b9acde9f57cfaba047c5b7680e0dccd2c1507df.json +++ /dev/null @@ -1,40 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "SELECT \n (flow_status->'step')::integer as step,\n jsonb_array_length(flow_status->'modules') as len,\n flow_status->'modules'->-1->>'branch_chosen' IS NOT NULL as is_branch_one,\n parent_job as ppp_job\n FROM v2_job \n LEFT JOIN v2_job_status USING (id)\n WHERE v2_job.id = $1", - "describe": { - "columns": [ - { - "ordinal": 0, - "name": "step", - "type_info": "Int4" - }, - { - "ordinal": 1, - "name": "len", - "type_info": "Int4" - }, - { - "ordinal": 2, - "name": "is_branch_one", - "type_info": "Bool" - }, - { - "ordinal": 3, - "name": "ppp_job", - "type_info": "Uuid" - } - ], - "parameters": { - "Left": [ - "Uuid" - ] - }, - "nullable": [ - null, - null, - null, - true - ] - }, - "hash": "d8ef35b4990eb9b2a306494d5b9acde9f57cfaba047c5b7680e0dccd2c1507df" -} diff --git a/backend/windmill-worker/src/common.rs b/backend/windmill-worker/src/common.rs index e7592584fd..b64b9558bf 100644 --- a/backend/windmill-worker/src/common.rs +++ b/backend/windmill-worker/src/common.rs @@ -34,7 +34,6 @@ use windmill_common::{ use anyhow::{anyhow, Result}; use windmill_parser_sql::{s3_mode_extension, S3ModeArgs, S3ModeFormat}; -use windmill_queue::flow_status::get_step_of_flow_status; use windmill_queue::MiniPulledJob; use std::collections::HashSet; @@ -1055,12 +1054,14 @@ pub fn use_flow_root_path(flow_path: &str) -> String { } pub fn build_http_client(timeout_duration: std::time::Duration) -> error::Result { - configure_client(reqwest::ClientBuilder::new() - .user_agent("windmill/beta") - .timeout(timeout_duration) - .connect_timeout(std::time::Duration::from_secs(10))) - .build() - .map_err(|e| Error::internal_err(format!("Error building http client: {e:#}"))) + configure_client( + reqwest::ClientBuilder::new() + .user_agent("windmill/beta") + .timeout(timeout_duration) + .connect_timeout(std::time::Duration::from_secs(10)), + ) + .build() + .map_err(|e| Error::internal_err(format!("Error building http client: {e:#}"))) } pub fn get_root_job_id(job: &MiniPulledJob) -> uuid::Uuid { @@ -1077,76 +1078,115 @@ pub struct StreamNotifier { job_id: uuid::Uuid, parent_job: uuid::Uuid, root_job: uuid::Uuid, + flow_step_id: Option, } -#[async_recursion] -async fn check_if_nested_step_is_last( - db: &DB, - parent_job: Uuid, - parent_of_parent_job: Option, - root_job: Uuid, - visited: Option>, -) -> error::Result { - // Initialize or use the provided visited set for cycle detection - let mut visited = visited.unwrap_or_else(HashSet::new); +// Helper struct to hold parent job information +#[derive(Debug)] +struct JobInfo { + flow_step_id: Option, + step: Option, + len: Option, + is_branch_one: Option, + next_parent: Option, +} - // Check for cycles - if we've already visited this job, return false to break the recursion - if !visited.insert(parent_job) { - return Ok(false); - } - - // get parent of parent job to get step of parent job - let parent_of_parent_job = parent_of_parent_job.or(sqlx::query_scalar!( - "SELECT parent_job FROM v2_job WHERE id = $1", - parent_job +// Shared helper function to fetch parent job info with flow status +async fn get_job_info(db: &DB, job_id: Uuid) -> error::Result { + sqlx::query_as!( + JobInfo, + r#"SELECT + flow_step_id, + (flow_status->'step')::integer as step, + jsonb_array_length(flow_status->'modules') as len, + runnable_path ~ '/branchone-\d+$' as is_branch_one, + parent_job as next_parent + FROM v2_job + LEFT JOIN v2_job_status USING (id) + WHERE v2_job.id = $1"#, + job_id ) .fetch_one(db) - .await?); - if let Some(parent_of_parent_job) = parent_of_parent_job { - // Check for cycles again with the parent_of_parent_job - if !visited.insert(parent_of_parent_job) { + .await + .map_err(|e| Error::internal_err(format!("fetching parent job info: {e:#}"))) +} + +async fn check_if_last_step( + db: &DB, + mut next_parent: Option, + root_job: Uuid, +) -> error::Result { + let mut visited = HashSet::new(); + loop { + // Get parent of current job + let Some(parent_job) = next_parent else { + return Ok(false); + }; + + // Check for cycles + if !visited.insert(parent_job) { return Ok(false); } - let r = sqlx::query!( - r#"SELECT - (flow_status->'step')::integer as step, - jsonb_array_length(flow_status->'modules') as len, - flow_status->'modules'->-1->>'branch_chosen' IS NOT NULL as is_branch_one, - parent_job as ppp_job - FROM v2_job - LEFT JOIN v2_job_status USING (id) - WHERE v2_job.id = $1"#, - parent_of_parent_job - ) - .fetch_one(db) - .await - .map_err(|e| Error::internal_err(format!("fetching step flow status: {e:#}")))?; + let parent_info = get_job_info(db, parent_job).await?; - if let Some(step) = r.step { - let step = Step::from_i32_and_len(step, r.len.unwrap_or(0) as usize); - - // if parent job is last and a branch one and - // - root_job is equal to parent of parent job, return true - // - root job is not equal to parent of parent job, recursively check if the parent of parent job is a branch one and last - if step.is_last_step() && r.is_branch_one.unwrap_or(false) { - if parent_of_parent_job == root_job { + if let Some(step) = parent_info.step { + let step = Step::from_i32_and_len(step, parent_info.len.unwrap_or(0) as usize); + if step.is_last_step() { + if parent_job == root_job { return Ok(true); - } else { - return check_if_nested_step_is_last( - db, - parent_of_parent_job, - r.ppp_job, - root_job, - Some(visited), - ) - .await; + } else if parent_info.is_branch_one.unwrap_or(false) { + next_parent = parent_info.next_parent; + continue; } } } - } - Ok(false) + return Ok(false); + } +} + +// Iterative implementation to avoid stack overflow from async_recursion +// Checks if we're at last step in nested branches AND if any parent's flow_step_id matches early_return_id +async fn check_if_early_return_or_last_in_early_return_parent( + db: &DB, + mut next_parent: Option, + mut step_id: Option, + early_return_id: &str, + root_job: Uuid, +) -> error::Result { + let mut visited = HashSet::new(); + loop { + if step_id.is_none() { + return Ok(false); + } + + let Some(parent_job) = next_parent else { + return Ok(false); + }; + + // Check for cycles + if !visited.insert(parent_job) { + return Ok(false); + } + + // If the parent's flow_step_id matches early_return_id, we found it! + if step_id.as_deref() == Some(early_return_id) && parent_job == root_job { + return Ok(true); + } else { + let parent_info = get_job_info(db, parent_job).await?; + if let Some(step) = parent_info.step { + let step = Step::from_i32_and_len(step, parent_info.len.unwrap_or(0) as usize); + // we only continue if we are at the last step and the parent is a branch one + if step.is_last_step() && parent_info.is_branch_one.unwrap_or(false) { + next_parent = parent_info.next_parent; + step_id = parent_info.flow_step_id; + continue; + } + } + return Ok(false); + } + } } impl StreamNotifier { @@ -1159,6 +1199,7 @@ impl StreamNotifier { parent_job: job.parent_job.unwrap(), job_id: job.id, root_job, + flow_step_id: job.flow_step_id.clone(), }), Connection::Http(_) => { tracing::warn!( @@ -1177,13 +1218,36 @@ impl StreamNotifier { parent_job: Uuid, job_id: Uuid, root_job: Uuid, + flow_step_id: Option, ) -> Result<(), Error> { - let step = get_step_of_flow_status(&db, parent_job).await?; + // Check if early_return is set at the flow level + let early_return_node_id = sqlx::query_scalar!( + r#" + SELECT fv.value->>'early_return' as "early_return" + FROM v2_job j + INNER JOIN flow_version fv ON fv.id = j.runnable_id + WHERE j.id = $1 + "#, + root_job + ) + .fetch_optional(&db) + .await? + .flatten(); - if step.is_last_step() - && (parent_job == root_job - || check_if_nested_step_is_last(&db, parent_job, None, root_job, None).await?) - { + let should_set_stream_job = if let Some(ref early_return_id) = early_return_node_id { + check_if_early_return_or_last_in_early_return_parent( + &db, + Some(parent_job), + flow_step_id, + early_return_id, + root_job, + ) + .await? + } else { + check_if_last_step(&db, Some(parent_job), root_job).await? + }; + + if should_set_stream_job { sqlx::query!(r#" UPDATE v2_job_status SET flow_status = jsonb_set(flow_status, array['stream_job'], to_jsonb($1::UUID::TEXT)) @@ -1203,10 +1267,16 @@ impl StreamNotifier { let parent_job = self.parent_job; let job_id = self.job_id; let root_job = self.root_job; + let flow_step_id = self.flow_step_id.clone(); tokio::spawn(async move { - if let Err(err) = - Self::update_flow_status_with_stream_job_inner(db, parent_job, job_id, root_job) - .await + if let Err(err) = Self::update_flow_status_with_stream_job_inner( + db, + parent_job, + job_id, + root_job, + flow_step_id, + ) + .await { tracing::error!("Could not notify about stream job {}: {err:#?}", parent_job); }