fix: show start to finish time for flows instead of cumulative (#3486)

* fix: show start to finish time for flows instead of cumulative

* fix: build
This commit is contained in:
HugoCasa
2024-03-28 15:50:18 +01:00
committed by GitHub
parent 0fc22938c9
commit dd4f48d244
7 changed files with 9 additions and 97 deletions
@@ -1,6 +1,6 @@
{
"db_name": "PostgreSQL",
"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 , raw_code\n , raw_lock\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 , email\n , visible_to_owner\n , mem_peak\n , tag\n , priority\n )\n VALUES ($1, $2, $3, $4, $5, COALESCE($6, now()), COALESCE($25, (EXTRACT('epoch' FROM (now())) - EXTRACT('epoch' FROM (COALESCE($6, now()))))*1000), $7, $8, $9,$10, $11, $12, $13, $14, $15, $16, $17, $18, $19, $20, $21, $22, $23, $24, $26, $27, $28, $29, $30)\n ON CONFLICT (id) DO UPDATE SET success = $7, result = $11 RETURNING duration_ms",
"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 , raw_code\n , raw_lock\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 , email\n , visible_to_owner\n , mem_peak\n , tag\n , priority\n )\n VALUES ($1, $2, $3, $4, $5, COALESCE($6, now()), (EXTRACT('epoch' FROM (now())) - EXTRACT('epoch' FROM (COALESCE($6, now()))))*1000, $7, $8, $9,$10, $11, $12, $13, $14, $15, $16, $17, $18, $19, $20, $21, $22, $23, $24, $25, $26, $27, $28, $29)\n ON CONFLICT (id) DO UPDATE SET success = $7, result = $11 RETURNING duration_ms",
"describe": {
"columns": [
{
@@ -79,7 +79,6 @@
}
}
},
"Numeric",
"Varchar",
"Bool",
"Int4",
@@ -91,5 +90,5 @@
false
]
},
"hash": "2ade671449393541fa565088b21268dad137314d250f7ded502defb9a6de0b2f"
"hash": "d5a8614286c170e0d175903cd1b53ff66b37ed8110a0b67aedb9f25e6a7383e1"
}
@@ -1,22 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "SELECT SUM(duration_ms) as duration FROM completed_job WHERE id = ANY($1)",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "duration",
"type_info": "Numeric"
}
],
"parameters": {
"Left": [
"UuidArray"
]
},
"nullable": [
null
]
},
"hash": "f516ca558816c2cbab3c8ae865ef8f764aa4686e7df9c53d30d49a3dbbf36af4"
}
+2 -46
View File
@@ -21,7 +21,6 @@ use axum::{
http::{request::Parts, Request, Uri},
response::{IntoResponse, Response},
};
use bigdecimal::ToPrimitive;
use chrono::{DateTime, Duration, Utc};
use itertools::Itertools;
#[cfg(feature = "prometheus")]
@@ -440,24 +439,6 @@ pub async fn add_completed_job_error<R: rsmq_async::RsmqConnection + Clone + Sen
Ok(result)
}
fn flatten_jobs(modules: Vec<FlowStatusModule>) -> Vec<Uuid> {
modules
.into_iter()
.filter_map(|m| match m {
FlowStatusModule::Success { job, flow_jobs, .. }
| FlowStatusModule::Failure { job, flow_jobs, .. } => {
if let Some(flow_jobs) = flow_jobs {
Some(flow_jobs)
} else {
Some(vec![job])
}
}
_ => None,
})
.flatten()
.collect::<Vec<_>>()
}
lazy_static::lazy_static! {
pub static ref GLOBAL_ERROR_HANDLER_PATH_IN_ADMINS_WORKSPACE: Option<String> = std::env::var("GLOBAL_ERROR_HANDLER_PATH_IN_ADMINS_WORKSPACE").ok();
}
@@ -487,30 +468,6 @@ pub async fn add_completed_job<
}
let is_flow = queued_job.is_flow();
let duration = if is_flow {
let jobs = queued_job.parse_flow_status().map(|s| {
let mut modules = s.modules;
modules.extend([s.failure_module.module_status]);
flatten_jobs(modules)
});
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()
.map(|x| x.to_i64())
.flatten()
} else {
tracing::warn!("Could not parse flow status");
None
}
} else {
None
};
let mut tx: QueueTransaction<'_, R> = (rsmq.clone(), db.begin().await?).into();
let job_id = queued_job.id;
@@ -556,8 +513,8 @@ pub async fn add_completed_job<
, tag
, priority
)
VALUES ($1, $2, $3, $4, $5, COALESCE($6, now()), COALESCE($25, (EXTRACT('epoch' FROM (now())) - EXTRACT('epoch' FROM (COALESCE($6, now()))))*1000), $7, $8, $9,\
$10, $11, $12, $13, $14, $15, $16, $17, $18, $19, $20, $21, $22, $23, $24, $26, $27, $28, $29, $30)
VALUES ($1, $2, $3, $4, $5, COALESCE($6, now()), (EXTRACT('epoch' FROM (now())) - EXTRACT('epoch' FROM (COALESCE($6, now()))))*1000, $7, $8, $9,\
$10, $11, $12, $13, $14, $15, $16, $17, $18, $19, $20, $21, $22, $23, $24, $25, $26, $27, $28, $29)
ON CONFLICT (id) DO UPDATE SET success = $7, result = $11 RETURNING duration_ms",
queued_job.workspace_id,
queued_job.id,
@@ -583,7 +540,6 @@ pub async fn add_completed_job<
queued_job.is_flow_step,
skipped,
queued_job.language.clone() as Option<ScriptLang>,
duration as Option<i64>,
queued_job.email,
queued_job.visible_to_owner,
if mem_peak > 0 { Some(mem_peak) } else { None },
+1 -10
View File
@@ -1,20 +1,11 @@
<script lang="ts">
import { msToSec } from '$lib/utils'
import { Badge } from './common'
import Tooltip from './Tooltip.svelte'
import { Hourglass } from 'lucide-svelte'
export let duration_ms: number
export let flow: boolean
</script>
<Badge large icon={{ icon: Hourglass, position: 'left' }}>
Ran in {msToSec(duration_ms)}s {#if flow}(sum){/if}
{#if flow}
<Tooltip>
Cumulative time of the execution of each steps. Suspend/sleep/transition times are not
accounted for, and each step's duration in a parallel branch would be added. Hence the time
here can differ from the time it took for the flow to have a result from start.
</Tooltip>
{/if}
Ran in {msToSec(duration_ms)}s
</Badge>
+2 -8
View File
@@ -17,18 +17,12 @@
<Badge large color="green" icon={{ icon: CheckCircle2, position: 'left' }}>
Success {job.is_skipped ? '(Skipped)' : ''}
</Badge>
<DurationMs
flow={job.job_kind == 'flow' || job?.job_kind == 'flowpreview'}
duration_ms={job.duration_ms}
/>
<DurationMs duration_ms={job.duration_ms} />
</div>
{:else if job && 'success' in job}
<div class="flex flex-row flex-wrap gap-y-1 mb-1 gap-x-2">
<Badge large color="red" icon={{ icon: XCircle, position: 'left' }}>Failed</Badge>
<DurationMs
flow={job.job_kind == 'flow' || job?.job_kind == 'flowpreview'}
duration_ms={job.duration_ms}
/>
<DurationMs duration_ms={job.duration_ms} />
</div>
{:else if job && 'running' in job && job.running}
<div>
@@ -108,10 +108,7 @@
Mem: {job?.['mem_peak'] ? `${(job['mem_peak'] / 1024).toPrecision(4)}MB` : 'N/A'}
</Badge>
{#if job?.['duration_ms']}
<DurationMs
flow={job.job_kind == 'flow' || job?.job_kind == 'flowpreview'}
duration_ms={job?.['duration_ms']}
/>
<DurationMs duration_ms={job?.['duration_ms']} />
{/if}
</div>
<div class="w-1/2 h-full overflow-auto">
@@ -62,10 +62,7 @@
</Badge>
{/if}
{#if job && 'duration_ms' in job && job.duration_ms != undefined}
<DurationMs
flow={job.job_kind == 'flow' || job?.job_kind == 'flowpreview'}
duration_ms={job.duration_ms}
/>
<DurationMs duration_ms={job.duration_ms} />
{/if}
{#if job?.['mem_peak']}
<Badge large>