From 884aba4811eeb57c83aae25908fd545fcb5ff6d3 Mon Sep 17 00:00:00 2001 From: HugoCasa Date: Wed, 19 Mar 2025 10:51:30 +0100 Subject: [PATCH] perf: improve perf of get completed flow node (#5418) * perf: improve perf of get completed flow node * better * nit * nits --------- Co-authored-by: Ruben Fiszel --- ...79e9bd7e1b8348e91b33030df7181a55798a6.json | 30 ++++++++++++++++ ...fab3b6f26ccbd822fc05b7a0fb3d6a17c5435.json | 35 +++++++++++++++++++ ...0c36c3726b2da9ff3592589ac7e83df1c537c.json | 29 --------------- backend/windmill-queue/src/jobs.rs | 33 +++++++++-------- 4 files changed, 84 insertions(+), 43 deletions(-) create mode 100644 backend/.sqlx/query-433da02be85347333323480f7f279e9bd7e1b8348e91b33030df7181a55798a6.json create mode 100644 backend/.sqlx/query-77f12f22d7c9da1799e5720c9a8fab3b6f26ccbd822fc05b7a0fb3d6a17c5435.json delete mode 100644 backend/.sqlx/query-8123ba05f6e7b9bd395175ee4ec0c36c3726b2da9ff3592589ac7e83df1c537c.json diff --git a/backend/.sqlx/query-433da02be85347333323480f7f279e9bd7e1b8348e91b33030df7181a55798a6.json b/backend/.sqlx/query-433da02be85347333323480f7f279e9bd7e1b8348e91b33030df7181a55798a6.json new file mode 100644 index 0000000000..bee9c556b2 --- /dev/null +++ b/backend/.sqlx/query-433da02be85347333323480f7f279e9bd7e1b8348e91b33030df7181a55798a6.json @@ -0,0 +1,30 @@ +{ + "db_name": "PostgreSQL", + "query": "SELECT j.id, jc.flow_status AS \"flow_status!: Json\"\n FROM v2_job j\n JOIN v2_job_completed jc ON j.id = jc.id\n WHERE j.parent_job = $1 AND j.workspace_id = $2 AND j.created_at >= $3 AND jc.flow_status IS NOT NULL", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "id", + "type_info": "Uuid" + }, + { + "ordinal": 1, + "name": "flow_status!: Json", + "type_info": "Jsonb" + } + ], + "parameters": { + "Left": [ + "Uuid", + "Text", + "Timestamptz" + ] + }, + "nullable": [ + false, + true + ] + }, + "hash": "433da02be85347333323480f7f279e9bd7e1b8348e91b33030df7181a55798a6" +} diff --git a/backend/.sqlx/query-77f12f22d7c9da1799e5720c9a8fab3b6f26ccbd822fc05b7a0fb3d6a17c5435.json b/backend/.sqlx/query-77f12f22d7c9da1799e5720c9a8fab3b6f26ccbd822fc05b7a0fb3d6a17c5435.json new file mode 100644 index 0000000000..42289a7791 --- /dev/null +++ b/backend/.sqlx/query-77f12f22d7c9da1799e5720c9a8fab3b6f26ccbd822fc05b7a0fb3d6a17c5435.json @@ -0,0 +1,35 @@ +{ + "db_name": "PostgreSQL", + "query": "SELECT jc.id, jc.flow_status AS \"flow_status!: Json\", j.created_at\n FROM v2_job_completed jc\n JOIN v2_job j ON j.id = jc.id\n WHERE jc.id = $1 AND jc.workspace_id = $2 AND jc.flow_status IS NOT NULL", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "id", + "type_info": "Uuid" + }, + { + "ordinal": 1, + "name": "flow_status!: Json", + "type_info": "Jsonb" + }, + { + "ordinal": 2, + "name": "created_at", + "type_info": "Timestamptz" + } + ], + "parameters": { + "Left": [ + "Uuid", + "Text" + ] + }, + "nullable": [ + false, + true, + false + ] + }, + "hash": "77f12f22d7c9da1799e5720c9a8fab3b6f26ccbd822fc05b7a0fb3d6a17c5435" +} diff --git a/backend/.sqlx/query-8123ba05f6e7b9bd395175ee4ec0c36c3726b2da9ff3592589ac7e83df1c537c.json b/backend/.sqlx/query-8123ba05f6e7b9bd395175ee4ec0c36c3726b2da9ff3592589ac7e83df1c537c.json deleted file mode 100644 index 1126431b46..0000000000 --- a/backend/.sqlx/query-8123ba05f6e7b9bd395175ee4ec0c36c3726b2da9ff3592589ac7e83df1c537c.json +++ /dev/null @@ -1,29 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "SELECT id, flow_status AS \"flow_status!: Json\"\n FROM v2_job_completed WHERE id = $1 AND workspace_id = $2", - "describe": { - "columns": [ - { - "ordinal": 0, - "name": "id", - "type_info": "Uuid" - }, - { - "ordinal": 1, - "name": "flow_status!: Json", - "type_info": "Jsonb" - } - ], - "parameters": { - "Left": [ - "Uuid", - "Text" - ] - }, - "nullable": [ - false, - true - ] - }, - "hash": "8123ba05f6e7b9bd395175ee4ec0c36c3726b2da9ff3592589ac7e83df1c537c" -} diff --git a/backend/windmill-queue/src/jobs.rs b/backend/windmill-queue/src/jobs.rs index a8f883367e..7f1473055a 100644 --- a/backend/windmill-queue/src/jobs.rs +++ b/backend/windmill-queue/src/jobs.rs @@ -2591,7 +2591,8 @@ pub async fn get_result_by_id_from_running_flow_inner( async fn get_completed_flow_node_result_rec( db: &Pool, w_id: &str, - subflows: impl std::iter::Iterator, + created_at: DateTime, + subflows: Vec<(Uuid, FlowStatus)>, node_id: &str, ) -> error::Result> { for (id, flow_status) in subflows { @@ -2611,19 +2612,19 @@ async fn get_completed_flow_node_result_rec( }; } else { let subflows = sqlx::query!( - "SELECT v2_job_completed.id AS \"id!\", flow_status AS \"flow_status!: Json\" - FROM v2_job_completed - INNER JOIN v2_job ON (v2_job_completed.id = v2_job.id) - WHERE parent_job = $1 AND v2_job_completed.workspace_id = $2 AND flow_status IS NOT NULL", + "SELECT j.id, jc.flow_status AS \"flow_status!: Json\" + FROM v2_job j + JOIN v2_job_completed jc ON j.id = jc.id + WHERE j.parent_job = $1 AND j.workspace_id = $2 AND j.created_at >= $3 AND jc.flow_status IS NOT NULL", id, - w_id + w_id, + created_at ) .map(|record| (record.id, record.flow_status.0)) .fetch_all(db) - .await? - .into_iter(); + .await?; match Box::pin(get_completed_flow_node_result_rec( - db, w_id, subflows, node_id, + db, w_id, created_at, subflows, node_id, )) .await? { @@ -2643,22 +2644,26 @@ async fn get_result_by_id_from_original_flow_inner( node_id: &str, ) -> error::Result { let flow_job = sqlx::query!( - "SELECT id, flow_status AS \"flow_status!: Json\" - FROM v2_job_completed WHERE id = $1 AND workspace_id = $2", + "SELECT jc.id, jc.flow_status AS \"flow_status!: Json\", j.created_at + FROM v2_job_completed jc + JOIN v2_job j ON j.id = jc.id + WHERE jc.id = $1 AND jc.workspace_id = $2 AND jc.flow_status IS NOT NULL", completed_flow_id, w_id ) - .map(|record| (record.id, record.flow_status.0)) + .map(|record| (record.id, record.flow_status.0, record.created_at)) .fetch_optional(db) .await?; - let flow_job = not_found_if_none( + let (id, flow_status, created_at) = not_found_if_none( flow_job, "Root completed job", format!("root: {}, id: {}", completed_flow_id, node_id), )?; - match get_completed_flow_node_result_rec(db, w_id, [flow_job].into_iter(), node_id).await? { + match get_completed_flow_node_result_rec(db, w_id, created_at, vec![(id, flow_status)], node_id) + .await? + { Some(res) => Ok(res), None => Err(Error::NotFound(format!( "Flow result by id not found going top-down from {}, (id: {})",