improve error handling after refactor

This commit is contained in:
Ruben Fiszel
2023-10-16 18:40:01 +02:00
parent b18ab1d154
commit 719cde60cc
6 changed files with 127 additions and 20 deletions
+8 -10
View File
@@ -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<T: Serialize> {
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<R: rsmq_async::RsmqConnection + Clone + Send>(
db: &Pool<Postgres>,
queued_job: &QueuedJob,
logs: String,
mem_peak: i32,
e: T,
e: serde_json::Value,
metrics: Option<Metrics>,
rsmq: Option<R>,
) -> Result<WrappedError<T>, Error> {
) -> Result<WrappedError, Error> {
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,
+8 -6
View File
@@ -1490,7 +1490,9 @@ pub async fn process_completed_job<R: rsmq_async::RsmqConnection + Send + Sync +
&job,
logs.to_string(),
mem_peak.to_owned(),
result,
serde_json::from_str(result.get()).unwrap_or_else(
|_| json!({"message": format!("Non serializable error: {}", result.get())}),
),
metrics.clone(),
rsmq.clone(),
)
@@ -1576,9 +1578,9 @@ pub async fn handle_job_error<R: rsmq_async::RsmqConnection + Send + Sync + Clon
add_completed_job_error(
db,
job,
format!("Unexpected error during job execution:\n{err}"),
format!("Unexpected error during job execution:\n{err:#?}"),
mem_peak,
&err,
err.clone(),
metrics.clone(),
rsmq_2,
)
@@ -1596,7 +1598,7 @@ pub async fn handle_job_error<R: rsmq_async::RsmqConnection + Send + Sync + Clon
(job.id, Uuid::nil(), Some(update_job_future))
};
let wrapped_error = WrappedError { error: json!(err) };
let wrapped_error = WrappedError { error: err.clone() };
let updated_flow = update_flow_status_after_job_completion(
db,
client,
@@ -1626,7 +1628,7 @@ pub async fn handle_job_error<R: rsmq_async::RsmqConnection + Send + Sync + Clon
&parent_job,
format!("Unexpected error during flow job error handling:\n{err}"),
mem_peak,
&e,
e,
metrics.clone(),
rsmq,
)
@@ -1953,7 +1955,7 @@ async fn process_result(
Error::ExitStatus(i) => {
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
+3 -3
View File
@@ -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(),
)
@@ -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}
<Tabs bind:selected>
<Tab value="graph"><span class="font-semibold text-md">Graph</span></Tab>
<!-- <Tab value="timeline"><span class="font-semibold">Timeline</span></Tab> -->
<Tab value="sequence"><span class="font-semibold">Details</span></Tab>
</Tabs>
{/if}
{/if}
<div class={selected == 'graph' ? 'hidden' : ''}>
{#if render && selected == 'timeline'}
<FlowTimeline {flowModuleStates} />
{/if}
<div class={selected != 'sequence' ? 'hidden' : ''}>
{#if isListJob}
<h3 class="text-md leading-6 font-bold text-tertiary border-b mb-4">
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
@@ -0,0 +1,92 @@
<script lang="ts">
import { FlowStatusModule } from '$lib/gen'
import { Skeleton } from './common'
import type { GraphModuleState } from './graph'
export let flowModuleStates: Record<string, GraphModuleState> = {}
let items = flowModuleStates ? computeItems(flowModuleStates) : undefined
$: computeItems(flowModuleStates)
let min: undefined | number = undefined
let max: undefined | number = undefined
function computeItems(flowModuleStates: Record<string, GraphModuleState>): any {
let nmin: undefined | number = undefined
let nmax: undefined | number = undefined
let isStillRunning = false
Object.values(flowModuleStates).forEach((v) => {
if (v.started_at) {
if (!nmin) {
nmin = v.started_at
} else {
nmin = Math.min(nmin, v.started_at)
}
}
if (v.type == FlowStatusModule.type.IN_PROGRESS) {
isStillRunning = true
}
if (!isStillRunning) {
if (v.started_at && v.duration_ms) {
let lmax = v.started_at + v.duration_ms
if (!nmax) {
nmax = lmax
} else {
nmax = Math.max(nmax, lmax)
}
}
}
})
let total = (isStillRunning || !nmax ? Date.now() : nmax) - (nmin ?? Date.now())
const nentries = Object.entries(flowModuleStates).map(([k, v]) => {
let started_at = v.started_at && nmin ? ((v.started_at - nmin) / total) * 100 : undefined
let len = 0
if (v.duration_ms) {
len = v.duration_ms
} else if (v.started_at) {
len = Date.now() - v.started_at
} else {
len = 0
}
if (total) {
len *= 100 / total
} else {
len = 0
}
return { name: k, started_at, len }
})
min = nmin
max = nmax
return nentries
}
</script>
{#if items}
<div class="border rounded-md divide-y">
<div class="px-2 py-2 grid grid-cols-12 w-full"
><div />
<div class="col-span-11 pt-1 px-2 flex justify-between"><div>{min}</div><div>{max}</div></div>
</div>
{#each items as item}
<div class="px-2 py-2 grid grid-cols-12 w-full"
><div>{item.name}</div>
<div class="col-span-11 pt-1 px-2 flex"
><div style="width: {item.started_at}%" class="h-4" /><div
style="width: {item.len}%"
class="h-4 bg-blue-600 border center-center text-white text-xs"
>{Math.ceil(item.len) ?? 0}%</div
></div
></div
>
{/each}
</div>
{:else}
<div class="mt-4" />
<Skeleton layout={[[2], 1]} />
{#each new Array(6) as _}
<Skeleton layout={[[4], 0.5]} />
{/each}
{/if}
@@ -36,6 +36,7 @@ export type GraphModuleState = {
iteration_total?: number
retries?: number
duration_ms?: number
started_at?: number
}
export type NestedNodes = GraphItem[]