From 65cba2dbb7bf5de88ee474f3849eba3b1bcd1232 Mon Sep 17 00:00:00 2001 From: Ruben Fiszel Date: Sat, 26 Sep 2026 00:18:59 +0200 Subject: [PATCH] perf: advance a flow step with one v2_job_status update (#11357) * perf: advance a flow step with one v2_job_status update Co-Authored-By: Claude Opus 5.5 (1M context) * docs: state what advance_flow_status returning None means Co-Authored-By: Claude Opus 5.5 (1M context) * fix: keep the merged flow advance identical for rows without a status row or with a malformed status Co-Authored-By: Claude Opus 5.5 (1M context) Claude-Session: https://claude.ai/code/session_014FbA4shbpCnRsUxir1GEfm * docs: note the JSON null invariant behind the empty-path no-op Co-Authored-By: Claude Opus 5.5 (1M context) Claude-Session: https://claude.ai/code/session_014FbA4shbpCnRsUxir1GEfm --------- Co-authored-by: Claude Opus 5.5 (1M context) --- ...316237eb1a547678858c1a1e45985035b3468.json | 14 - ...49b5bacac241549614bd39e4b055db7fc1ea.json} | 4 +- ...efe0ca030638ded82c2ffd2155cacdf36ec11.json | 15 - ...614fe8c00a264667ce48fb1f0b45662cef43d.json | 14 + ...e6e6180a48ab792d305d04146e81b06b0d69b.json | 15 + ...a208814ce37fcde72e8ef61904a4f41fb4579.json | 325 ++++++++++++++++++ ...fb77dc03469355d2a0da0b2d6b4aeeea37d3e.json | 16 - backend/windmill-worker/src/worker_flow.rs | 261 ++++++++++---- 8 files changed, 543 insertions(+), 121 deletions(-) delete mode 100644 backend/.sqlx/query-06db0e720dd59a7c52c0a98ea7b316237eb1a547678858c1a1e45985035b3468.json rename backend/.sqlx/{query-0aaec91ab06753e46c595d82469924a98f28b0dead245df7248a9ccb8a5f20c3.json => query-3c74bcf357bb87e2ee960473d40949b5bacac241549614bd39e4b055db7fc1ea.json} (56%) delete mode 100644 backend/.sqlx/query-4622d28e2fa09bc60b9d0c79397efe0ca030638ded82c2ffd2155cacdf36ec11.json create mode 100644 backend/.sqlx/query-5ffba9baef525e1e46e9a258e39614fe8c00a264667ce48fb1f0b45662cef43d.json create mode 100644 backend/.sqlx/query-9fb1a1e4e39e0e8a0934fcda7b3e6e6180a48ab792d305d04146e81b06b0d69b.json create mode 100644 backend/.sqlx/query-ab05c1359039ae99a2aa6156eb2a208814ce37fcde72e8ef61904a4f41fb4579.json delete mode 100644 backend/.sqlx/query-aed8bd751c3e988f422216e74acfb77dc03469355d2a0da0b2d6b4aeeea37d3e.json diff --git a/backend/.sqlx/query-06db0e720dd59a7c52c0a98ea7b316237eb1a547678858c1a1e45985035b3468.json b/backend/.sqlx/query-06db0e720dd59a7c52c0a98ea7b316237eb1a547678858c1a1e45985035b3468.json deleted file mode 100644 index aea8b8302b..0000000000 --- a/backend/.sqlx/query-06db0e720dd59a7c52c0a98ea7b316237eb1a547678858c1a1e45985035b3468.json +++ /dev/null @@ -1,14 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "UPDATE v2_job_status\n SET flow_status = flow_status - 'retry'\n WHERE id = $1", - "describe": { - "columns": [], - "parameters": { - "Left": [ - "Uuid" - ] - }, - "nullable": [] - }, - "hash": "06db0e720dd59a7c52c0a98ea7b316237eb1a547678858c1a1e45985035b3468" -} diff --git a/backend/.sqlx/query-0aaec91ab06753e46c595d82469924a98f28b0dead245df7248a9ccb8a5f20c3.json b/backend/.sqlx/query-3c74bcf357bb87e2ee960473d40949b5bacac241549614bd39e4b055db7fc1ea.json similarity index 56% rename from backend/.sqlx/query-0aaec91ab06753e46c595d82469924a98f28b0dead245df7248a9ccb8a5f20c3.json rename to backend/.sqlx/query-3c74bcf357bb87e2ee960473d40949b5bacac241549614bd39e4b055db7fc1ea.json index 04f5ba5b84..aae7f79b15 100644 --- a/backend/.sqlx/query-0aaec91ab06753e46c595d82469924a98f28b0dead245df7248a9ccb8a5f20c3.json +++ b/backend/.sqlx/query-3c74bcf357bb87e2ee960473d40949b5bacac241549614bd39e4b055db7fc1ea.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "UPDATE v2_job_status\n SET flow_status = JSONB_SET(flow_status, ARRAY['modules', $1::TEXT], $2)\n WHERE id = $3", + "query": "UPDATE v2_job_status\n SET flow_leaf_jobs = JSONB_SET(coalesce(flow_leaf_jobs, '{}'::jsonb), ARRAY[$1::TEXT], $2)\n WHERE id = $3", "describe": { "columns": [], "parameters": { @@ -12,5 +12,5 @@ }, "nullable": [] }, - "hash": "0aaec91ab06753e46c595d82469924a98f28b0dead245df7248a9ccb8a5f20c3" + "hash": "3c74bcf357bb87e2ee960473d40949b5bacac241549614bd39e4b055db7fc1ea" } diff --git a/backend/.sqlx/query-4622d28e2fa09bc60b9d0c79397efe0ca030638ded82c2ffd2155cacdf36ec11.json b/backend/.sqlx/query-4622d28e2fa09bc60b9d0c79397efe0ca030638ded82c2ffd2155cacdf36ec11.json deleted file mode 100644 index f30ab8d49d..0000000000 --- a/backend/.sqlx/query-4622d28e2fa09bc60b9d0c79397efe0ca030638ded82c2ffd2155cacdf36ec11.json +++ /dev/null @@ -1,15 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "UPDATE v2_job_status\n SET flow_status = JSONB_SET(flow_status, ARRAY['step'], $1)\n WHERE id = $2", - "describe": { - "columns": [], - "parameters": { - "Left": [ - "Jsonb", - "Uuid" - ] - }, - "nullable": [] - }, - "hash": "4622d28e2fa09bc60b9d0c79397efe0ca030638ded82c2ffd2155cacdf36ec11" -} diff --git a/backend/.sqlx/query-5ffba9baef525e1e46e9a258e39614fe8c00a264667ce48fb1f0b45662cef43d.json b/backend/.sqlx/query-5ffba9baef525e1e46e9a258e39614fe8c00a264667ce48fb1f0b45662cef43d.json new file mode 100644 index 0000000000..8381e1ca55 --- /dev/null +++ b/backend/.sqlx/query-5ffba9baef525e1e46e9a258e39614fe8c00a264667ce48fb1f0b45662cef43d.json @@ -0,0 +1,14 @@ +{ + "db_name": "PostgreSQL", + "query": "UPDATE v2_job_status\n SET flow_status = flow_status - 'retry'\n WHERE id = $1", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Uuid" + ] + }, + "nullable": [] + }, + "hash": "5ffba9baef525e1e46e9a258e39614fe8c00a264667ce48fb1f0b45662cef43d" +} diff --git a/backend/.sqlx/query-9fb1a1e4e39e0e8a0934fcda7b3e6e6180a48ab792d305d04146e81b06b0d69b.json b/backend/.sqlx/query-9fb1a1e4e39e0e8a0934fcda7b3e6e6180a48ab792d305d04146e81b06b0d69b.json new file mode 100644 index 0000000000..6ae18315bf --- /dev/null +++ b/backend/.sqlx/query-9fb1a1e4e39e0e8a0934fcda7b3e6e6180a48ab792d305d04146e81b06b0d69b.json @@ -0,0 +1,15 @@ +{ + "db_name": "PostgreSQL", + "query": "UPDATE v2_job_status\n SET flow_status = JSONB_SET(flow_status, ARRAY['step'], $1)\n WHERE id = $2", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Jsonb", + "Uuid" + ] + }, + "nullable": [] + }, + "hash": "9fb1a1e4e39e0e8a0934fcda7b3e6e6180a48ab792d305d04146e81b06b0d69b" +} diff --git a/backend/.sqlx/query-ab05c1359039ae99a2aa6156eb2a208814ce37fcde72e8ef61904a4f41fb4579.json b/backend/.sqlx/query-ab05c1359039ae99a2aa6156eb2a208814ce37fcde72e8ef61904a4f41fb4579.json new file mode 100644 index 0000000000..bb358546ac --- /dev/null +++ b/backend/.sqlx/query-ab05c1359039ae99a2aa6156eb2a208814ce37fcde72e8ef61904a4f41fb4579.json @@ -0,0 +1,325 @@ +{ + "db_name": "PostgreSQL", + "query": "UPDATE v2_job_status SET\n flow_status = JSONB_SET(\n JSONB_SET(v2_job_status.flow_status, $2::TEXT[], $3),\n $4::TEXT[], $5\n ) - $6::TEXT[],\n flow_leaf_jobs = CASE\n WHEN $7::TEXT IS NULL\n OR COALESCE(v2_job.flow_innermost_root_job, v2_job_status.id) <> v2_job_status.id\n THEN v2_job_status.flow_leaf_jobs\n ELSE JSONB_SET(COALESCE(v2_job_status.flow_leaf_jobs, '{}'::JSONB), ARRAY[$7::TEXT], $8) END\n FROM v2_job_queue INNER JOIN v2_job ON v2_job.id = v2_job_queue.id\n WHERE v2_job_status.id = $1 AND v2_job_queue.id = $1\n RETURNING\n v2_job_queue.workspace_id,\n v2_job_queue.id,\n v2_job.args as \"args: sqlx::types::Json>>\",\n v2_job.parent_job,\n v2_job.created_by,\n v2_job_queue.started_at,\n v2_job_queue.runnable_settings_handle,\n v2_job_queue.scheduled_for,\n v2_job.runnable_path,\n v2_job.kind as \"kind: JobKind\",\n v2_job.runnable_id as \"runnable_id: ScriptHash\",\n v2_job_queue.canceled_reason,\n v2_job_queue.canceled_by,\n v2_job.permissioned_as,\n v2_job.permissioned_as_email,\n v2_job_status.flow_status as \"flow_status: sqlx::types::Json>\",\n v2_job.tag,\n v2_job.script_lang as \"script_lang: ScriptLang\",\n v2_job.same_worker,\n v2_job.pre_run_error,\n v2_job.concurrent_limit,\n v2_job.concurrency_time_window_s,\n v2_job.flow_innermost_root_job,\n v2_job.root_job,\n v2_job.timeout,\n v2_job.flow_step_id,\n v2_job.cache_ttl,\n v2_job_queue.cache_ignore_s3_path,\n v2_job_queue.priority,\n v2_job.preprocessed,\n v2_job.script_entrypoint_override,\n v2_job.trigger,\n v2_job.trigger_kind as \"trigger_kind: TriggerKindLabel\",\n v2_job.visible_to_owner,\n NULL as permissioned_as_end_user_email", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "workspace_id", + "type_info": "Varchar" + }, + { + "ordinal": 1, + "name": "id", + "type_info": "Uuid" + }, + { + "ordinal": 2, + "name": "args: sqlx::types::Json>>", + "type_info": "Jsonb" + }, + { + "ordinal": 3, + "name": "parent_job", + "type_info": "Uuid" + }, + { + "ordinal": 4, + "name": "created_by", + "type_info": "Varchar" + }, + { + "ordinal": 5, + "name": "started_at", + "type_info": "Timestamptz" + }, + { + "ordinal": 6, + "name": "runnable_settings_handle", + "type_info": "Int8" + }, + { + "ordinal": 7, + "name": "scheduled_for", + "type_info": "Timestamptz" + }, + { + "ordinal": 8, + "name": "runnable_path", + "type_info": "Varchar" + }, + { + "ordinal": 9, + "name": "kind: JobKind", + "type_info": { + "Custom": { + "name": "job_kind", + "kind": { + "Enum": [ + "script", + "preview", + "flow", + "dependencies", + "flowpreview", + "script_hub", + "identity", + "flowdependencies", + "http", + "graphql", + "postgresql", + "noop", + "appdependencies", + "deploymentcallback", + "singlestepflow", + "flowscript", + "flownode", + "appscript", + "aiagent", + "unassigned_script", + "unassigned_flow", + "unassigned_singlestepflow" + ] + } + } + } + }, + { + "ordinal": 10, + "name": "runnable_id: ScriptHash", + "type_info": "Int8" + }, + { + "ordinal": 11, + "name": "canceled_reason", + "type_info": "Text" + }, + { + "ordinal": 12, + "name": "canceled_by", + "type_info": "Varchar" + }, + { + "ordinal": 13, + "name": "permissioned_as", + "type_info": "Varchar" + }, + { + "ordinal": 14, + "name": "permissioned_as_email", + "type_info": "Varchar" + }, + { + "ordinal": 15, + "name": "flow_status: sqlx::types::Json>", + "type_info": "Jsonb" + }, + { + "ordinal": 16, + "name": "tag", + "type_info": "Varchar" + }, + { + "ordinal": 17, + "name": "script_lang: ScriptLang", + "type_info": { + "Custom": { + "name": "script_lang", + "kind": { + "Enum": [ + "python3", + "deno", + "go", + "bash", + "postgresql", + "nativets", + "bun", + "mysql", + "bigquery", + "snowflake", + "graphql", + "powershell", + "mssql", + "php", + "bunnative", + "rust", + "ansible", + "csharp", + "oracledb", + "nu", + "java", + "duckdb", + "ruby", + "rlang", + "dbt" + ] + } + } + } + }, + { + "ordinal": 18, + "name": "same_worker", + "type_info": "Bool" + }, + { + "ordinal": 19, + "name": "pre_run_error", + "type_info": "Text" + }, + { + "ordinal": 20, + "name": "concurrent_limit", + "type_info": "Int4" + }, + { + "ordinal": 21, + "name": "concurrency_time_window_s", + "type_info": "Int4" + }, + { + "ordinal": 22, + "name": "flow_innermost_root_job", + "type_info": "Uuid" + }, + { + "ordinal": 23, + "name": "root_job", + "type_info": "Uuid" + }, + { + "ordinal": 24, + "name": "timeout", + "type_info": "Int4" + }, + { + "ordinal": 25, + "name": "flow_step_id", + "type_info": "Varchar" + }, + { + "ordinal": 26, + "name": "cache_ttl", + "type_info": "Int4" + }, + { + "ordinal": 27, + "name": "cache_ignore_s3_path", + "type_info": "Bool" + }, + { + "ordinal": 28, + "name": "priority", + "type_info": "Int2" + }, + { + "ordinal": 29, + "name": "preprocessed", + "type_info": "Bool" + }, + { + "ordinal": 30, + "name": "script_entrypoint_override", + "type_info": "Varchar" + }, + { + "ordinal": 31, + "name": "trigger", + "type_info": "Varchar" + }, + { + "ordinal": 32, + "name": "trigger_kind: TriggerKindLabel", + "type_info": { + "Custom": { + "name": "job_trigger_kind", + "kind": { + "Enum": [ + "webhook", + "http", + "websocket", + "kafka", + "email", + "nats", + "schedule", + "app", + "ui", + "postgres", + "sqs", + "gcp", + "mqtt", + "nextcloud", + "google", + "ci_test", + "github", + "azure", + "asset", + "freshness", + "amqp" + ] + } + } + } + }, + { + "ordinal": 33, + "name": "visible_to_owner", + "type_info": "Bool" + }, + { + "ordinal": 34, + "name": "permissioned_as_end_user_email", + "type_info": "Text" + } + ], + "parameters": { + "Left": [ + "Uuid", + "TextArray", + "Jsonb", + "TextArray", + "Jsonb", + "TextArray", + "Text", + "Jsonb" + ] + }, + "nullable": [ + false, + false, + true, + true, + false, + true, + true, + false, + true, + false, + true, + true, + true, + false, + false, + true, + false, + true, + false, + true, + true, + true, + true, + true, + true, + true, + true, + true, + true, + true, + true, + true, + true, + false, + null + ] + }, + "hash": "ab05c1359039ae99a2aa6156eb2a208814ce37fcde72e8ef61904a4f41fb4579" +} diff --git a/backend/.sqlx/query-aed8bd751c3e988f422216e74acfb77dc03469355d2a0da0b2d6b4aeeea37d3e.json b/backend/.sqlx/query-aed8bd751c3e988f422216e74acfb77dc03469355d2a0da0b2d6b4aeeea37d3e.json deleted file mode 100644 index d8f93f74c4..0000000000 --- a/backend/.sqlx/query-aed8bd751c3e988f422216e74acfb77dc03469355d2a0da0b2d6b4aeeea37d3e.json +++ /dev/null @@ -1,16 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "UPDATE v2_job_status\n SET flow_leaf_jobs = JSONB_SET(coalesce(flow_leaf_jobs, '{}'::jsonb), ARRAY[$1::TEXT], $2)\n WHERE COALESCE((SELECT flow_innermost_root_job FROM v2_job WHERE id = $3), $3) = id", - "describe": { - "columns": [], - "parameters": { - "Left": [ - "Text", - "Jsonb", - "Uuid" - ] - }, - "nullable": [] - }, - "hash": "aed8bd751c3e988f422216e74acfb77dc03469355d2a0da0b2d6b4aeeea37d3e" -} diff --git a/backend/windmill-worker/src/worker_flow.rs b/backend/windmill-worker/src/worker_flow.rs index 6e01378448..a3257ace36 100644 --- a/backend/windmill-worker/src/worker_flow.rs +++ b/backend/windmill-worker/src/worker_flow.rs @@ -44,12 +44,13 @@ use windmill_common::flow_status::{ }; use windmill_common::flows::{add_virtual_items_if_necessary, Branch, FlowNodeId, StopAfterIf}; use windmill_common::jobs::{ - script_path_to_payload, JobKind, JobPayload, OnBehalfOf, RawCode, ENTRYPOINT_OVERRIDE, + script_path_to_payload, JobKind, JobPayload, OnBehalfOf, RawCode, TriggerKindLabel, + ENTRYPOINT_OVERRIDE, }; use windmill_common::runnable_settings::{ ConcurrencySettingsWithCustom, DebouncingSettings, RunnableSettingsTrait, }; -use windmill_common::scripts::{ScriptHash, ScriptRunnableSettingsInline}; +use windmill_common::scripts::{ScriptHash, ScriptLang, ScriptRunnableSettingsInline}; use windmill_common::utils::WarnAfterExt; use windmill_common::worker::{error_to_value, to_raw_value, Connection}; use windmill_common::{ @@ -1373,36 +1374,52 @@ pub async fn update_flow_status_after_job_completion_internal( } let step_counter = if inc_step_counter { - sqlx::query!( - "UPDATE v2_job_status - SET flow_status = JSONB_SET(flow_status, ARRAY['step'], $1) - WHERE id = $2", - json!(old_status.step + 1), - flow - ) - .execute(&mut *tx) - .await - .map_err(|e| { - Error::internal_err(format!("error while setting flow index for {flow}: {e:#}")) - })?; old_status.step + 1 } else { old_status.step }; - // tracing::error!( - // "step_counter: {:?} {} {inc_step_counter} {flow}", - // step_counter, - // old_status.step, - // ); - // panic!("stop"); - /* is_last_step is true when the step_counter (the next step index) is an invalid index */ let is_last_step = usize::try_from(step_counter) .map(|i| !(..old_status.modules.len()).contains(&i)) .unwrap_or(true); - if let Some(new_status) = new_status.as_ref() { + let nresult = if let Some(nresult) = nresult { + // can be some either with early stop error or with the flow jobs results (was fetched to evaluate stop_early_after_all_iters but evaluated to false) + nresult + } else { + match &new_status { + Some(FlowStatusModule::Success { flow_jobs: Some(jobs), .. }) + | Some(FlowStatusModule::Failure { flow_jobs: Some(jobs), .. }) => { + Arc::new(retrieve_flow_jobs_results(&mut *tx, w_id, jobs).await?) + } + _ => result.clone(), + } + }; + + let remove_retry = old_status.retry.fail_count > 0 + && matches!(&new_status, Some(FlowStatusModule::Success { .. })); + + let special_step_status = new_status + .as_ref() + .filter(|_| is_failure_step || module_step.is_preprocessor_step()); + + let flow_job = if let Some(new_status) = special_step_status { + if inc_step_counter { + sqlx::query!( + "UPDATE v2_job_status + SET flow_status = JSONB_SET(flow_status, ARRAY['step'], $1) + WHERE id = $2", + json!(step_counter), + flow + ) + .execute(&mut *tx) + .await + .map_err(|e| { + Error::internal_err(format!("error while setting flow index for {flow}: {e:#}")) + })?; + } + if is_failure_step { let parent_module = sqlx::query_scalar!( "SELECT flow_status->'failure_module'->>'parent_module' FROM v2_job_status WHERE id = $1", @@ -1432,7 +1449,7 @@ pub async fn update_flow_status_after_job_completion_internal( "error while setting flow status in failure step: {e:#}" )) })?; - } else if module_step.is_preprocessor_step() { + } else { sqlx::query!( "UPDATE v2_job_status SET flow_status = JSONB_SET(flow_status, ARRAY['preprocessor_module'], $1) @@ -1447,69 +1464,71 @@ pub async fn update_flow_status_after_job_completion_internal( "error while setting flow status in preprocessing step: {e:#}" )) })?; - } else { + } + + if remove_retry { sqlx::query!( "UPDATE v2_job_status - SET flow_status = JSONB_SET(flow_status, ARRAY['modules', $1::TEXT], $2) - WHERE id = $3", - old_status.step.to_string(), - json!(new_status), + SET flow_status = flow_status - 'retry' + WHERE id = $1", flow ) .execute(&mut *tx) .await - .map_err(|e| { - Error::internal_err(format!("error while setting new flow status: {e:#}")) - })?; - - if let Some(job_result) = new_status.job_result() { - sqlx::query!( - "UPDATE v2_job_status - SET flow_leaf_jobs = JSONB_SET(coalesce(flow_leaf_jobs, '{}'::jsonb), ARRAY[$1::TEXT], $2) - WHERE COALESCE((SELECT flow_innermost_root_job FROM v2_job WHERE id = $3), $3) = id", - new_status.id(), - json!(job_result), - flow - ) - .execute(&mut *tx) - .await.map_err(|e| { - Error::internal_err(format!( - "error while setting leaf jobs: {e:#}" - )) - })?; - } + .context("remove flow status retry")?; } - } - let nresult = if let Some(nresult) = nresult { - // can be some either with early stop error or with the flow jobs results (was fetched to evaluate stop_early_after_all_iters but evaluated to false) - nresult + get_mini_pulled_job(&mut *tx, &flow).await? } else { - match &new_status { - Some(FlowStatusModule::Success { flow_jobs: Some(jobs), .. }) - | Some(FlowStatusModule::Failure { flow_jobs: Some(jobs), .. }) => { - Arc::new(retrieve_flow_jobs_results(&mut *tx, w_id, jobs).await?) + let module_status = new_status + .as_ref() + .map(|s| (old_status.step.to_string(), json!(s))); + let leaf_job = new_status + .as_ref() + .and_then(|s| s.job_result().map(|r| (s.id(), json!(r)))); + + let flow_job = if inc_step_counter || module_status.is_some() || remove_retry { + match advance_flow_status( + &mut tx, + flow, + inc_step_counter.then_some(step_counter), + module_status, + leaf_job.as_ref(), + remove_retry, + ) + .await? + { + Some(flow_job) => Some(flow_job), + None => get_mini_pulled_job(&mut *tx, &flow).await?, } - _ => result.clone(), + } else { + get_mini_pulled_job(&mut *tx, &flow).await? + }; + + // The leaf jobs of a subflow are kept on its innermost root's row, which the + // statement above does not touch. + let innermost_root = flow_job + .as_ref() + .and_then(|j| j.flow_innermost_root_job) + .filter(|root| *root != flow); + if let (Some((leaf_id, leaf_result)), Some(root)) = (leaf_job, innermost_root) { + sqlx::query!( + "UPDATE v2_job_status + SET flow_leaf_jobs = JSONB_SET(coalesce(flow_leaf_jobs, '{}'::jsonb), ARRAY[$1::TEXT], $2) + WHERE id = $3", + leaf_id, + leaf_result, + root + ) + .execute(&mut *tx) + .await + .map_err(|e| Error::internal_err(format!("error while setting leaf jobs: {e:#}")))?; } + + flow_job }; - if old_status.retry.fail_count > 0 - && matches!(&new_status, Some(FlowStatusModule::Success { .. })) - { - sqlx::query!( - "UPDATE v2_job_status - SET flow_status = flow_status - 'retry' - WHERE id = $1", - flow - ) - .execute(&mut *tx) - .await - .context("remove flow status retry")?; - } - - let flow_job = get_mini_pulled_job(&mut *tx, &flow) - .await? + let flow_job = flow_job .ok_or_else(|| Error::internal_err(format!("requiring flow to be in the queue")))?; tx.commit().await?; @@ -2304,6 +2323,100 @@ async fn set_success_and_duration_in_flow_job_success<'c>( Ok(()) } +/// Applies a step's edits to the flow's `v2_job_status` row and reads the flow back, in one +/// statement: every UPDATE of that row writes a new row version carrying the whole +/// `flow_status`, so the edits must not be split. `leaf_job` is only written here when the flow +/// is its own innermost root; otherwise it belongs on the root's row and the caller writes it +/// there. `Ok(None)` means nothing was written because the flow has no queued `v2_job_status` +/// row; the caller then reads the flow with `get_mini_pulled_job`, whose LEFT JOIN still returns +/// a queued flow that lacks a status row. +async fn advance_flow_status( + tx: &mut Transaction<'_, Postgres>, + flow: Uuid, + step: Option, + module_status: Option<(String, Value)>, + leaf_job: Option<&(String, Value)>, + remove_retry: bool, +) -> error::Result> { + // An edit that does not apply gets an empty path, which JSONB_SET treats as a no-op. Its + // value must stay JSON `null`, never SQL NULL: JSONB_SET is strict and would null the whole + // `flow_status`. Each edit stays a JSONB_SET on its own path so that a malformed + // `flow_status` fails exactly as the equivalent separate UPDATEs would. + let (step_path, step) = match step { + Some(step) => (vec!["step"], json!(step)), + None => (vec![], Value::Null), + }; + let (module_path, module_status) = match &module_status { + Some((index, status)) => (vec!["modules", index.as_str()], status), + None => (vec![], &Value::Null), + }; + let (leaf_id, leaf_result) = leaf_job.map(|(id, r)| (id.as_str(), r)).unzip(); + let removed_keys: &[&str] = if remove_retry { &["retry"] } else { &[] }; + + sqlx::query_as!( + MiniPulledJob, + "UPDATE v2_job_status SET + flow_status = JSONB_SET( + JSONB_SET(v2_job_status.flow_status, $2::TEXT[], $3), + $4::TEXT[], $5 + ) - $6::TEXT[], + flow_leaf_jobs = CASE + WHEN $7::TEXT IS NULL + OR COALESCE(v2_job.flow_innermost_root_job, v2_job_status.id) <> v2_job_status.id + THEN v2_job_status.flow_leaf_jobs + ELSE JSONB_SET(COALESCE(v2_job_status.flow_leaf_jobs, '{}'::JSONB), ARRAY[$7::TEXT], $8) END + FROM v2_job_queue INNER JOIN v2_job ON v2_job.id = v2_job_queue.id + WHERE v2_job_status.id = $1 AND v2_job_queue.id = $1 + RETURNING + v2_job_queue.workspace_id, + v2_job_queue.id, + v2_job.args as \"args: sqlx::types::Json>>\", + v2_job.parent_job, + v2_job.created_by, + v2_job_queue.started_at, + v2_job_queue.runnable_settings_handle, + v2_job_queue.scheduled_for, + v2_job.runnable_path, + v2_job.kind as \"kind: JobKind\", + v2_job.runnable_id as \"runnable_id: ScriptHash\", + v2_job_queue.canceled_reason, + v2_job_queue.canceled_by, + v2_job.permissioned_as, + v2_job.permissioned_as_email, + v2_job_status.flow_status as \"flow_status: sqlx::types::Json>\", + v2_job.tag, + v2_job.script_lang as \"script_lang: ScriptLang\", + v2_job.same_worker, + v2_job.pre_run_error, + v2_job.concurrent_limit, + v2_job.concurrency_time_window_s, + v2_job.flow_innermost_root_job, + v2_job.root_job, + v2_job.timeout, + v2_job.flow_step_id, + v2_job.cache_ttl, + v2_job_queue.cache_ignore_s3_path, + v2_job_queue.priority, + v2_job.preprocessed, + v2_job.script_entrypoint_override, + v2_job.trigger, + v2_job.trigger_kind as \"trigger_kind: TriggerKindLabel\", + v2_job.visible_to_owner, + NULL as permissioned_as_end_user_email", + flow, + &step_path as &[&str], + step, + &module_path as &[&str], + module_status, + removed_keys as &[&str], + leaf_id, + leaf_result, + ) + .fetch_optional(&mut **tx) + .await + .map_err(|e| Error::internal_err(format!("error while setting new flow status: {e:#}"))) +} + async fn retrieve_flow_jobs_results<'c>( e: impl sqlx::PgExecutor<'c>, w_id: &str, @@ -6185,7 +6298,7 @@ async fn payload_from_simple_module( pub fn raw_script_to_payload( path: String, content: String, - language: windmill_common::scripts::ScriptLang, + language: ScriptLang, lock: Option, concurrency_settings: ConcurrencySettingsWithCustom, module: &FlowModule,