feat(backend): flow duration is now computed as the sum of every child

This commit is contained in:
Ruben Fiszel
2022-11-05 22:31:54 +01:00
parent cddec6469e
commit badc60193c
2 changed files with 115 additions and 63 deletions
+82 -61
View File
@@ -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": [
+33 -2
View File
@@ -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<Uuid, Error> {
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::<Vec<_>>()
});
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<i64>
)
.execute(&mut tx)
.await