From f6d9cd062d3b1e17cfd64cd4938122ac76ae19f7 Mon Sep 17 00:00:00 2001 From: yacine Bouraroui Date: Wed, 2 Oct 2024 10:28:01 +0200 Subject: [PATCH] write to completed_jobs_result on job completion --- backend/windmill-api/src/jobs.rs | 40 ++++++++++++++-------------- backend/windmill-queue/src/jobs.rs | 42 +++++++++++++++++++++--------- 2 files changed, 51 insertions(+), 31 deletions(-) diff --git a/backend/windmill-api/src/jobs.rs b/backend/windmill-api/src/jobs.rs index 6113bdeaf7..c8fc83efdd 100644 --- a/backend/windmill-api/src/jobs.rs +++ b/backend/windmill-api/src/jobs.rs @@ -1349,6 +1349,8 @@ async fn cancel_jobs( ) -> error::JsonResult> { let mut uuids = vec![]; let mut tx = db.begin().await?; + + let result = serde_json::json!({"error": { "message": format!("Job canceled: cancel all by {username}"), "name": "Canceled", "reason": "cancel all", "canceler": username}}); let trivial_jobs = sqlx::query!("INSERT INTO completed_job AS cj ( workspace_id , id @@ -1412,10 +1414,19 @@ async fn cancel_jobs( , tag , priority FROM queue WHERE id = any($2) AND running = false AND parent_job IS NULL AND workspace_id = $3 AND schedule_path IS NULL FOR UPDATE SKIP LOCKED - ON CONFLICT (id) DO NOTHING RETURNING id", username, &jobs, w_id, serde_json::json!({"error": { "message": format!("Job canceled: cancel all by {username}"), "name": "Canceled", "reason": "cancel all", "canceler": username}})) + ON CONFLICT (id) DO NOTHING RETURNING id", username, &jobs, w_id, &result) .fetch_all(&mut *tx) .await?.into_iter().map(|x| x.id).collect::>(); + sqlx::query!( + "INSERT INTO completed_jobs_result(id, result, tag, workspace_id) SELECT id, $1, tag, $2 FROM completed_job WHERE id = any($3) ON CONFLICT (id) DO NOTHING", + result, + w_id, + &trivial_jobs, + ) + .execute(&mut *tx) + .await?; + sqlx::query!( "DELETE FROM queue WHERE id = any($1) AND workspace_id = $2", &trivial_jobs, @@ -1423,6 +1434,7 @@ async fn cancel_jobs( ) .execute(&mut *tx) .await?; + tx.commit().await?; // sqlx::query!( @@ -1436,18 +1448,10 @@ async fn cancel_jobs( continue; } let rsmq = rsmq.clone(); - match tokio::time::timeout(tokio::time::Duration::from_secs(5), async move { + if let Ok(result) = tokio::time::timeout(tokio::time::Duration::from_secs(5), async move { let tx = db.begin().await?; let (tx, _) = windmill_queue::cancel_job( - username, - None, - job_id.clone(), - w_id, - tx, - db, - rsmq, - false, - false, + username, None, job_id, w_id, tx, db, rsmq, false, false, ) .await?; tx.commit().await?; @@ -1455,20 +1459,19 @@ async fn cancel_jobs( }) .await { - Ok(result) => match result { + match result { Ok(_) => { uuids.push(job_id); } Err(e) => { tracing::error!("Failed to cancel job {:?}: {:?}", job_id, e); } - }, - Err(_) => { - tracing::error!( - "Timeout while trying to cancel job {:?} after 5 seconds", - job_id - ); } + } else { + tracing::error!( + "Timeout while trying to cancel job {:?} after 5 seconds", + job_id + ); } } @@ -4686,7 +4689,6 @@ async fn get_job_update( &w_id, job_id, "progress_perc" - ) .fetch_optional(&db) .await?.and_then(|inner| inner) diff --git a/backend/windmill-queue/src/jobs.rs b/backend/windmill-queue/src/jobs.rs index 98b0a004b4..eaa421fb3c 100644 --- a/backend/windmill-queue/src/jobs.rs +++ b/backend/windmill-queue/src/jobs.rs @@ -588,6 +588,8 @@ pub async fn add_completed_job< let mem_peak = mem_peak.max(queued_job.mem_peak.unwrap_or(0)); add_time!(bench, "add_completed_job query START"); + + // On conflict (when id already exists), update the success and result fields. let _duration: i64 = sqlx::query_scalar!( "INSERT INTO completed_job AS cj ( workspace_id @@ -624,8 +626,8 @@ pub async fn add_completed_job< 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, + &queued_job.workspace_id, + &queued_job.id, queued_job.parent_job, queued_job.created_by, queued_job.created_at, @@ -634,7 +636,7 @@ pub async fn add_completed_job< queued_job.script_hash.map(|x| x.0), queued_job.script_path, &queued_job.args as &Option>>>, - result as Json<&T>, + &result as &Json<&T>, queued_job.raw_code, queued_job.raw_lock, canceled_by.is_some(), @@ -651,16 +653,32 @@ pub async fn add_completed_job< queued_job.email, queued_job.visible_to_owner, if mem_peak > 0 { Some(mem_peak) } else { None }, - queued_job.tag, + &queued_job.tag, queued_job.priority, ) .fetch_one(&mut tx) .await .map_err(|e| Error::InternalErr(format!("Could not add completed job {job_id}: {e:#}")))?; - // tracing::error!("2 {:?}", start.elapsed()); - add_time!(bench, "add_completed_job query END"); + add_time!(bench, "completed_jobs_result query START"); + sqlx::query!( + "INSERT INTO completed_jobs_result(id, result, tag, workspace_id) VALUES($1, $2, $3, $4) ON CONFLICT (id) DO UPDATE SET result = $2", + queued_job.id, + &result as &Json<&T>, + queued_job.tag, + queued_job.workspace_id, + ) + .execute(&mut tx) + .await + .map_err(|e| { + Error::InternalErr(format!( + "Could not add completed job result {job_id}: {e:#}" + )) + })?; + + add_time!(bench, "completed_jobs_result query END"); + if !queued_job.is_flow_step { if _duration > 500 && (queued_job.job_kind == JobKind::Script || queued_job.job_kind == JobKind::Preview) @@ -2486,12 +2504,14 @@ async fn extract_result_from_job_result( Some(json_path) => { let mut parts = json_path.split("."); - let Some(idx) = parts.next().map(|x| x.parse::().ok()).flatten() else { + let Some(idx) = parts.next().and_then(|x| x.parse::().ok()) else { return Ok(to_raw_value(&serde_json::Value::Null)); }; let Some(job_id) = job_ids.get(idx).cloned() else { return Ok(to_raw_value(&serde_json::Value::Null)); }; + + // ici Ok(sqlx::query_as::<_, ResultR>( "SELECT result #> $3 as result FROM completed_job WHERE id = $1 AND workspace_id = $2", ) @@ -2502,8 +2522,7 @@ async fn extract_result_from_job_result( ) .fetch_optional(db) .await? - .map(|r| r.result.map(|x| x.0)) - .flatten() + .and_then(|r| r.result.map(|x| x.0)) .unwrap_or_else(|| to_raw_value(&serde_json::Value::Null))) } None => { @@ -2528,7 +2547,7 @@ async fn extract_result_from_job_result( Ok(to_raw_value(&result)) } }, - + // ici JobResult::SingleJob(x) => Ok(sqlx::query_as::<_, ResultR>( "SELECT result #> $3 as result FROM completed_job WHERE id = $1 AND workspace_id = $2", ) @@ -2541,8 +2560,7 @@ async fn extract_result_from_job_result( ) .fetch_optional(db) .await? - .map(|r| r.result.map(|x| x.0)) - .flatten() + .and_then(|r| r.result.map(|x| x.0)) .unwrap_or_else(|| to_raw_value(&serde_json::Value::Null))), } }