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) <noreply@anthropic.com>

* docs: state what advance_flow_status returning None means

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

* 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) <noreply@anthropic.com>
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) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_014FbA4shbpCnRsUxir1GEfm

---------

Co-authored-by: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
This commit is contained in:
Ruben Fiszel
2026-09-26 00:18:59 +02:00
committed by GitHub
co-authored by Claude Opus 5.5
parent bebd762194
commit 65cba2dbb7
8 changed files with 543 additions and 121 deletions
@@ -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"
}
@@ -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"
}
@@ -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"
}
@@ -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"
}
@@ -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"
}
@@ -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<HashMap<String, Box<RawValue>>>\",\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<Box<RawValue>>\",\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<HashMap<String, Box<RawValue>>>",
"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<Box<RawValue>>",
"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"
}
@@ -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"
}
+187 -74
View File
@@ -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<i32>,
module_status: Option<(String, Value)>,
leaf_job: Option<&(String, Value)>,
remove_retry: bool,
) -> error::Result<Option<MiniPulledJob>> {
// 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<HashMap<String, Box<RawValue>>>\",
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<Box<RawValue>>\",
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<String>,
concurrency_settings: ConcurrencySettingsWithCustom,
module: &FlowModule,