From badc60193c2480f93056eee5be6548bcf49fc1fc Mon Sep 17 00:00:00 2001 From: Ruben Fiszel Date: Sat, 5 Nov 2022 22:31:54 +0100 Subject: [PATCH] feat(backend): flow duration is now computed as the sum of every child --- backend/sqlx-data.json | 143 ++++++++++++++++------------ backend/windmill-worker/src/jobs.rs | 35 ++++++- 2 files changed, 115 insertions(+), 63 deletions(-) diff --git a/backend/sqlx-data.json b/backend/sqlx-data.json index 11a42c3f27..df8e6e447d 100644 --- a/backend/sqlx-data.json +++ b/backend/sqlx-data.json @@ -1916,67 +1916,6 @@ }, "query": "UPDATE group_ SET summary = $1 WHERE name = $2 AND workspace_id = $3" }, - "850e97afdf07f7d2bbfe33b0aaecceef6d7c736b91aaaaf60d764bd070250cba": { - "describe": { - "columns": [], - "nullable": [], - "parameters": { - "Left": [ - "Varchar", - "Uuid", - "Uuid", - "Varchar", - "Timestamptz", - "Timestamptz", - "Bool", - "Int8", - "Varchar", - "Jsonb", - "Jsonb", - "Text", - "Text", - "Bool", - "Varchar", - "Text", - { - "Custom": { - "kind": { - "Enum": [ - "script", - "preview", - "flow", - "dependencies", - "flowpreview", - "script_hub", - "identity" - ] - }, - "name": "job_kind" - } - }, - "Varchar", - "Varchar", - "Jsonb", - "Jsonb", - "Bool", - "Bool", - { - "Custom": { - "kind": { - "Enum": [ - "python3", - "deno", - "go" - ] - }, - "name": "script_lang" - } - } - ] - } - }, - "query": "INSERT INTO completed_job AS cj\n ( workspace_id\n , id\n , parent_job\n , created_by\n , created_at\n , started_at\n , duration_ms\n , success\n , script_hash\n , script_path\n , args\n , result\n , logs\n , raw_code\n , canceled\n , canceled_by\n , canceled_reason\n , job_kind\n , schedule_path\n , permissioned_as\n , flow_status\n , raw_flow\n , is_flow_step\n , is_skipped\n , language )\n VALUES ($1, $2, $3, $4, $5, $6, EXTRACT(milliseconds FROM (now() - $6)), $7, $8, $9,$10, $11, $12, $13, $14, $15, $16, $17, $18, $19, $20, $21, $22, $23, $24)\n ON CONFLICT (id) DO UPDATE SET success = $7, result = $11, logs = concat(cj.logs, $12)" - }, "853788436dbe987853433e8dc83665f68bd127de31d4c807abafeead896f6ac4": { "describe": { "columns": [ @@ -2552,6 +2491,68 @@ }, "query": "\n UPDATE queue\n SET flow_status = JSONB_SET(flow_status, ARRAY['failure_module'], $1)\n WHERE id = $2\n " }, + "a643c7ace2e4231c099e5a4d278ec31e327a71e6125ad59e38691b2b3b80a3bf": { + "describe": { + "columns": [], + "nullable": [], + "parameters": { + "Left": [ + "Varchar", + "Uuid", + "Uuid", + "Varchar", + "Timestamptz", + "Timestamptz", + "Bool", + "Int8", + "Varchar", + "Jsonb", + "Jsonb", + "Text", + "Text", + "Bool", + "Varchar", + "Text", + { + "Custom": { + "kind": { + "Enum": [ + "script", + "preview", + "flow", + "dependencies", + "flowpreview", + "script_hub", + "identity" + ] + }, + "name": "job_kind" + } + }, + "Varchar", + "Varchar", + "Jsonb", + "Jsonb", + "Bool", + "Bool", + { + "Custom": { + "kind": { + "Enum": [ + "python3", + "deno", + "go" + ] + }, + "name": "script_lang" + } + }, + "Numeric" + ] + } + }, + "query": "INSERT INTO completed_job AS cj\n ( workspace_id\n , id\n , parent_job\n , created_by\n , created_at\n , started_at\n , duration_ms\n , success\n , script_hash\n , script_path\n , args\n , result\n , logs\n , raw_code\n , canceled\n , canceled_by\n , canceled_reason\n , job_kind\n , schedule_path\n , permissioned_as\n , flow_status\n , raw_flow\n , is_flow_step\n , is_skipped\n , language )\n VALUES ($1, $2, $3, $4, $5, $6, COALESCE($25, EXTRACT(milliseconds FROM (now() - $6))), $7, $8, $9,$10, $11, $12, $13, $14, $15, $16, $17, $18, $19, $20, $21, $22, $23, $24)\n ON CONFLICT (id) DO UPDATE SET success = $7, result = $11, logs = concat(cj.logs, $12)" + }, "a98b2d68f023f46ab91167d3147416df672c2aed2ba5ab70e98a9da5fa47255a": { "describe": { "columns": [], @@ -3732,6 +3733,26 @@ }, "query": "INSERT INTO schedule (workspace_id, path, schedule, offset_, edited_by, script_path, is_flow, args, enabled) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9) RETURNING *" }, + "f516ca558816c2cbab3c8ae865ef8f764aa4686e7df9c53d30d49a3dbbf36af4": { + "describe": { + "columns": [ + { + "name": "duration", + "ordinal": 0, + "type_info": "Int8" + } + ], + "nullable": [ + null + ], + "parameters": { + "Left": [ + "UuidArray" + ] + } + }, + "query": "SELECT SUM(duration_ms) as duration FROM completed_job WHERE id = ANY($1)" + }, "f7906298e4204ad55ec84021bb2461f369386493519637279f1188227230c580": { "describe": { "columns": [ diff --git a/backend/windmill-worker/src/jobs.rs b/backend/windmill-worker/src/jobs.rs index 77dea01606..dc38c385ee 100644 --- a/backend/windmill-worker/src/jobs.rs +++ b/backend/windmill-worker/src/jobs.rs @@ -10,7 +10,7 @@ use serde_json::{Map, Value}; use sqlx::{Pool, Postgres, Transaction}; use tracing::instrument; use uuid::Uuid; -use windmill_common::error::Error; +use windmill_common::{error::Error, flow_status::FlowStatusModule}; use windmill_queue::{delete_job, JobKind, QueuedJob}; #[instrument(level = "trace", skip_all)] @@ -58,6 +58,36 @@ pub async fn add_completed_job( result: serde_json::Value, logs: String, ) -> Result { + let duration = + if queued_job.job_kind == JobKind::Flow || queued_job.job_kind == JobKind::FlowPreview { + let jobs = queued_job.parse_flow_status().map(|s| { + let mut modules = s.modules; + modules.extend([s.failure_module]); + modules + .into_iter() + .filter_map(|m| match m { + FlowStatusModule::Success { job, .. } + | FlowStatusModule::Failure { job, .. } => Some(job), + _ => None, + }) + .collect::>() + }); + if let Some(jobs) = jobs { + sqlx::query_scalar!( + "SELECT SUM(duration_ms) as duration FROM completed_job WHERE id = ANY($1)", + jobs.as_slice() + ) + .fetch_one(db) + .await + .ok() + .flatten() + } else { + tracing::warn!("Could not parse flow status"); + None + } + } else { + None + }; let mut tx = db.begin().await?; let job_id = queued_job.id.clone(); sqlx::query!( @@ -87,7 +117,7 @@ pub async fn add_completed_job( , is_flow_step , is_skipped , language ) - VALUES ($1, $2, $3, $4, $5, $6, EXTRACT(milliseconds FROM (now() - $6)), $7, $8, $9,\ + VALUES ($1, $2, $3, $4, $5, $6, COALESCE($25, EXTRACT(milliseconds FROM (now() - $6))), $7, $8, $9,\ $10, $11, $12, $13, $14, $15, $16, $17, $18, $19, $20, $21, $22, $23, $24) ON CONFLICT (id) DO UPDATE SET success = $7, result = $11, logs = concat(cj.logs, $12)", queued_job.workspace_id, @@ -114,6 +144,7 @@ pub async fn add_completed_job( queued_job.is_flow_step, skipped, queued_job.language: ScriptLang, + duration: Option ) .execute(&mut tx) .await