diff --git a/backend/windmill-queue/src/jobs.rs b/backend/windmill-queue/src/jobs.rs index f90a0097e0..36f0b5c872 100644 --- a/backend/windmill-queue/src/jobs.rs +++ b/backend/windmill-queue/src/jobs.rs @@ -132,7 +132,7 @@ pub async fn cancel_job<'c: 'async_recursion>( &job_running, format!("canceled by {username}: (force cancel: {force_cancel})"), job_running.mem_peak.unwrap_or(0), - &e, + e, None, rsmq.clone(), ) @@ -178,28 +178,26 @@ pub async fn cancel_job<'c: 'async_recursion>( Ok((tx, Some(id))) } -#[derive(Serialize)] -pub struct WrappedError { - pub error: T, +#[derive(Serialize, Debug)] +pub struct WrappedError { + pub error: serde_json::Value, } #[instrument(level = "trace", skip_all)] -pub async fn add_completed_job_error< - T: Serialize + Send + Sync, - R: rsmq_async::RsmqConnection + Clone + Send, ->( +pub async fn add_completed_job_error( db: &Pool, queued_job: &QueuedJob, logs: String, mem_peak: i32, - e: T, + e: serde_json::Value, metrics: Option, rsmq: Option, -) -> Result, Error> { +) -> Result { if *METRICS_ENABLED { metrics.map(|m| m.worker_execution_failed.inc()); } let result = WrappedError { error: e }; + tracing::error!("FOO {:?}", serde_json::to_string(&result)); let _ = add_completed_job( db, &queued_job, diff --git a/backend/windmill-worker/src/worker.rs b/backend/windmill-worker/src/worker.rs index 0d9ff22ca4..00d3c1403e 100644 --- a/backend/windmill-worker/src/worker.rs +++ b/backend/windmill-worker/src/worker.rs @@ -1490,7 +1490,9 @@ pub async fn process_completed_job { let res = read_result(job_dir).await.ok(); - if res.is_some() { + if res.as_ref().is_some_and(|x| !x.get().is_empty()) { res.unwrap() } else { let last_10_log_lines = logs diff --git a/backend/windmill-worker/src/worker_flow.rs b/backend/windmill-worker/src/worker_flow.rs index 18c9f2ad02..a6dfc7d355 100644 --- a/backend/windmill-worker/src/worker_flow.rs +++ b/backend/windmill-worker/src/worker_flow.rs @@ -110,7 +110,7 @@ pub async fn update_flow_status_after_job_completion< &nrec.job_id_for_status, w_id, false, - &to_raw_value(&Json(&WrappedError { error: e.to_string() })), + &to_raw_value(&Json(&WrappedError { error: json!(e.to_string()) })), metrics.clone(), true, same_worker_tx.clone(), @@ -635,7 +635,7 @@ pub async fn update_flow_status_after_job_completion_internal< &flow_job, logs, 0, - &canceled_job_to_result(&flow_job), + canceled_job_to_result(&flow_job), metrics.clone(), rsmq.clone(), ) @@ -682,7 +682,7 @@ pub async fn update_flow_status_after_job_completion_internal< &flow_job, "Unexpected error during flow chaining:\n".to_string(), 0, - &e, + e, metrics.clone(), rsmq.clone(), ) diff --git a/frontend/src/lib/components/FlowStatusViewer.svelte b/frontend/src/lib/components/FlowStatusViewer.svelte index a7663d2b2e..b431f665a6 100644 --- a/frontend/src/lib/components/FlowStatusViewer.svelte +++ b/frontend/src/lib/components/FlowStatusViewer.svelte @@ -18,6 +18,7 @@ import { Loader2 } from 'lucide-svelte' import FlowStatusWaitingForEvents from './FlowStatusWaitingForEvents.svelte' import { deepEqual } from 'fast-equals' + import FlowTimeline from './FlowTimeline.svelte' const dispatch = createEventDispatcher() @@ -159,6 +160,7 @@ async function updateJobId() { if (jobId !== job?.id) { retry_status = {} + flowModuleStates = {} localFlowModuleStates = {} await loadJobInProgress() job?.script_path && loadOwner(job.script_path) @@ -198,6 +200,7 @@ type: FlowStatusModule.type.IN_PROGRESS, logs: job.logs, args: job.args, + started_at: job.started_at ? new Date(job.started_at).getTime() : undefined, parent_module: mod['parent_module'] } } else { @@ -209,6 +212,7 @@ job_id: job.id, parent_module: mod['parent_module'], duration_ms: job['duration_ms'], + started_at: job.started_at ? new Date(job.started_at).getTime() : undefined, iteration_total: mod.iterator?.itered?.length // retries: flowState?.raw_flow } @@ -291,11 +295,15 @@ {#if innerModules.length > 0 && !isListJob} Graph + Details {/if} {/if} -
+ {#if render && selected == 'timeline'} + + {/if} +
{#if isListJob}

Embedded flows: ({flowJobIds?.flowJobs.length} items) @@ -369,6 +377,9 @@ if (e.detail.type == 'QueuedJob') { localFlowModuleStates[flowJobIds.moduleId] = { type: FlowStatusModule.type.IN_PROGRESS, + started_at: e.detail.started_at + ? new Date(e.detail.started_at).getTime() + : undefined, logs: e.detail.logs, job_id: e.detail.id, args: e.detail.args, @@ -376,6 +387,9 @@ } } else { localFlowModuleStates[flowJobIds.moduleId] = { + started_at: e.detail.started_at + ? new Date(e.detail.started_at).getTime() + : undefined, args: e.detail.args, type: e.detail.success ? FlowStatusModule.type.SUCCESS diff --git a/frontend/src/lib/components/FlowTimeline.svelte b/frontend/src/lib/components/FlowTimeline.svelte new file mode 100644 index 0000000000..0981760eb5 --- /dev/null +++ b/frontend/src/lib/components/FlowTimeline.svelte @@ -0,0 +1,92 @@ + + +{#if items} +
+
+
{min}
{max}
+
+ {#each items as item} +
{item.name}
+
{Math.ceil(item.len) ?? 0}%
+ {/each} +
+{:else} +
+ + {#each new Array(6) as _} + + {/each} +{/if} diff --git a/frontend/src/lib/components/graph/model.ts b/frontend/src/lib/components/graph/model.ts index d7f4ba9a78..244bc7c315 100644 --- a/frontend/src/lib/components/graph/model.ts +++ b/frontend/src/lib/components/graph/model.ts @@ -36,6 +36,7 @@ export type GraphModuleState = { iteration_total?: number retries?: number duration_ms?: number + started_at?: number } export type NestedNodes = GraphItem[]