diff --git a/backend/.sqlx/query-02424907504848e983bfa89eec343061932dc5b4b17cf13d5cf8d833aedbe6d5.json b/backend/.sqlx/query-02424907504848e983bfa89eec343061932dc5b4b17cf13d5cf8d833aedbe6d5.json index 043d236a6b..c8b5e3086f 100644 --- a/backend/.sqlx/query-02424907504848e983bfa89eec343061932dc5b4b17cf13d5cf8d833aedbe6d5.json +++ b/backend/.sqlx/query-02424907504848e983bfa89eec343061932dc5b4b17cf13d5cf8d833aedbe6d5.json @@ -5,7 +5,7 @@ "columns": [ { "ordinal": 0, - "name": "bool", + "name": "?column?", "type_info": "Bool" } ], diff --git a/backend/.sqlx/query-661f472ff3860983322162420457f5033b9c9afc344d9c3e385ba20a3ad2197a.json b/backend/.sqlx/query-661f472ff3860983322162420457f5033b9c9afc344d9c3e385ba20a3ad2197a.json index 1fa370e682..75b8108281 100644 --- a/backend/.sqlx/query-661f472ff3860983322162420457f5033b9c9afc344d9c3e385ba20a3ad2197a.json +++ b/backend/.sqlx/query-661f472ff3860983322162420457f5033b9c9afc344d9c3e385ba20a3ad2197a.json @@ -5,7 +5,7 @@ "columns": [ { "ordinal": 0, - "name": "bool", + "name": "?column?", "type_info": "Bool" } ], diff --git a/backend/.sqlx/query-d6615719bf8db4b333ed55c9a3c160ad1962668fc3305d8e80ae91ef73614a80.json b/backend/.sqlx/query-d6615719bf8db4b333ed55c9a3c160ad1962668fc3305d8e80ae91ef73614a80.json index 103219fe05..271c395f9a 100644 --- a/backend/.sqlx/query-d6615719bf8db4b333ed55c9a3c160ad1962668fc3305d8e80ae91ef73614a80.json +++ b/backend/.sqlx/query-d6615719bf8db4b333ed55c9a3c160ad1962668fc3305d8e80ae91ef73614a80.json @@ -5,7 +5,7 @@ "columns": [ { "ordinal": 0, - "name": "bool", + "name": "?column?", "type_info": "Bool" } ], diff --git a/backend/.sqlx/query-d5a8614286c170e0d175903cd1b53ff66b37ed8110a0b67aedb9f25e6a7383e1.json b/backend/.sqlx/query-f9fc0084fe086ef80005bb64a8bb6b493e53583017c18e2ab44f44125c52d548.json similarity index 58% rename from backend/.sqlx/query-d5a8614286c170e0d175903cd1b53ff66b37ed8110a0b67aedb9f25e6a7383e1.json rename to backend/.sqlx/query-f9fc0084fe086ef80005bb64a8bb6b493e53583017c18e2ab44f44125c52d548.json index 164157b02f..e5a7c10282 100644 --- a/backend/.sqlx/query-d5a8614286c170e0d175903cd1b53ff66b37ed8110a0b67aedb9f25e6a7383e1.json +++ b/backend/.sqlx/query-f9fc0084fe086ef80005bb64a8bb6b493e53583017c18e2ab44f44125c52d548.json @@ -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()), (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", + "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($30::bigint, (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": [ { @@ -87,12 +87,13 @@ "Bool", "Int4", "Varchar", - "Int2" + "Int2", + "Int8" ] }, "nullable": [ false ] }, - "hash": "d5a8614286c170e0d175903cd1b53ff66b37ed8110a0b67aedb9f25e6a7383e1" + "hash": "f9fc0084fe086ef80005bb64a8bb6b493e53583017c18e2ab44f44125c52d548" } diff --git a/backend/Cargo.lock b/backend/Cargo.lock index 99989b1f52..0c6bfe70a1 100644 --- a/backend/Cargo.lock +++ b/backend/Cargo.lock @@ -10916,6 +10916,7 @@ version = "1.430.1" dependencies = [ "anyhow", "async-recursion", + "backon", "base64 0.22.1", "bit-vec", "bytes", diff --git a/backend/windmill-queue/src/jobs.rs b/backend/windmill-queue/src/jobs.rs index 9da04e0265..7ee95c373a 100644 --- a/backend/windmill-queue/src/jobs.rs +++ b/backend/windmill-queue/src/jobs.rs @@ -39,7 +39,6 @@ use windmill_audit::audit_ee::{audit_log, AuditAuthor}; use windmill_audit::ActionKind; use windmill_common::{ - add_time, auth::{fetch_authed_from_permissioned_as, permissioned_as_to_username}, db::{Authed, UserDB}, error::{self, to_anyhow, Error}, @@ -65,6 +64,9 @@ use windmill_common::{ DB, METRICS_ENABLED, }; +use backon::ConstantBuilder; +use backon::{BackoffBuilder, Retryable}; + #[cfg(feature = "enterprise")] use windmill_common::BASE_URL; @@ -173,8 +175,7 @@ pub async fn cancel_single_job<'c>( e, "server", false, - #[cfg(feature = "benchmark")] - &mut windmill_common::bench::BenchmarkIter::new(), + None, ) .await; @@ -462,7 +463,7 @@ pub async fn add_completed_job_error( e: serde_json::Value, _worker_name: &str, flow_is_done: bool, - #[cfg(feature = "benchmark")] bench: &mut windmill_common::bench::BenchmarkIter, + duration: Option, ) -> Result { #[cfg(feature = "prometheus")] register_metric( @@ -499,8 +500,7 @@ pub async fn add_completed_job_error( mem_peak, canceled_by, flow_is_done, - #[cfg(feature = "benchmark")] - bench, + duration, ) .await?; Ok(result) @@ -520,134 +520,139 @@ pub async fn add_completed_job( mem_peak: i32, canceled_by: Option, flow_is_done: bool, - #[cfg(feature = "benchmark")] bench: &mut windmill_common::bench::BenchmarkIter, + duration: Option, ) -> Result { // tracing::error!("Start"); // let start = tokio::time::Instant::now(); - add_time!(bench, "add_completed_job start"); + // add_time!(bench, "add_completed_job start"); if !result.is_valid_json() { return Err(Error::InternalErr( "Result of job is invalid json (empty)".to_string(), )); } - let mut tx = db.begin().await?; + let _job_id = queued_job.id; + let (opt_uuid, _duration, _skip_downstream_error_handlers) = (|| async { + let mut tx = db.begin().await?; - let job_id = queued_job.id; - // tracing::error!("1 {:?}", start.elapsed()); + let job_id = queued_job.id; + // tracing::error!("1 {:?}", start.elapsed()); - tracing::debug!( - "completed job {} {}", - queued_job.id, - serde_json::to_string(&result).unwrap_or_else(|_| "".to_string()) - ); + tracing::debug!( + "completed job {} {}", + queued_job.id, + serde_json::to_string(&result).unwrap_or_else(|_| "".to_string()) + ); - let (raw_code, raw_lock, raw_flow) = if !*MIN_VERSION_IS_AT_LEAST_1_427.read().await { - sqlx::query!( - "SELECT raw_code, raw_lock, raw_flow AS \"raw_flow: Json>\" - FROM job WHERE id = $1 AND workspace_id = $2 LIMIT 1", - &job_id, - &queued_job.workspace_id - ) - .fetch_one(db) - .map_ok(|record| (record.raw_code, record.raw_lock, record.raw_flow)) - .or_else(|_| { + let (raw_code, raw_lock, raw_flow) = if !*MIN_VERSION_IS_AT_LEAST_1_427.read().await { sqlx::query!( "SELECT raw_code, raw_lock, raw_flow AS \"raw_flow: Json>\" - FROM queue WHERE id = $1 AND workspace_id = $2 LIMIT 1", + FROM job WHERE id = $1 AND workspace_id = $2 LIMIT 1", &job_id, &queued_job.workspace_id ) .fetch_one(db) .map_ok(|record| (record.raw_code, record.raw_lock, record.raw_flow)) - }) - .await - .unwrap_or_default() - } else { - (None, None, None) - }; - - let mem_peak = mem_peak.max(queued_job.mem_peak.unwrap_or(0)); - add_time!(bench, "add_completed_job query START"); - let _duration: i64 = sqlx::query_scalar!( - "INSERT INTO completed_job AS cj - ( workspace_id - , id - , parent_job - , created_by - , created_at - , started_at - , duration_ms - , success - , script_hash - , script_path - , args - , result - , raw_code - , raw_lock - , canceled - , canceled_by - , canceled_reason - , job_kind - , schedule_path - , permissioned_as - , flow_status - , raw_flow - , is_flow_step - , is_skipped - , language - , email - , visible_to_owner - , mem_peak - , tag - , priority + .or_else(|_| { + sqlx::query!( + "SELECT raw_code, raw_lock, raw_flow AS \"raw_flow: Json>\" + FROM queue WHERE id = $1 AND workspace_id = $2 LIMIT 1", + &job_id, + &queued_job.workspace_id ) - 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.parent_job, - queued_job.created_by, - queued_job.created_at, - queued_job.started_at, - success, - queued_job.script_hash.map(|x| x.0), - queued_job.script_path, - &queued_job.args as &Option>>>, - result as Json<&T>, - raw_code, - raw_lock, - canceled_by.is_some(), - canceled_by.clone().map(|cb| cb.username).flatten(), - canceled_by.clone().map(|cb| cb.reason).flatten(), - queued_job.job_kind.clone() as JobKind, - queued_job.schedule_path, - queued_job.permissioned_as, - &queued_job.flow_status as &Option>>, - &raw_flow as &Option>>, - queued_job.is_flow_step, - skipped, - queued_job.language.clone() as Option, - queued_job.email, - queued_job.visible_to_owner, - if mem_peak > 0 { Some(mem_peak) } else { None }, - 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()); + .fetch_one(db) + .map_ok(|record| (record.raw_code, record.raw_lock, record.raw_flow)) + }) + .await + .unwrap_or_default() + } else { + (None, None, None) + }; - add_time!(bench, "add_completed_job query END"); + let mem_peak = mem_peak.max(queued_job.mem_peak.unwrap_or(0)); + // add_time!(bench, "add_completed_job query START"); - if !queued_job.is_flow_step { - if _duration > 500 - && (queued_job.job_kind == JobKind::Script || queued_job.job_kind == JobKind::Preview) - { - if let Err(e) = sqlx::query!( + let _duration = sqlx::query_scalar!( + "INSERT INTO completed_job AS cj + ( workspace_id + , id + , parent_job + , created_by + , created_at + , started_at + , duration_ms + , success + , script_hash + , script_path + , args + , result + , raw_code + , raw_lock + , canceled + , canceled_by + , canceled_reason + , job_kind + , schedule_path + , permissioned_as + , flow_status + , raw_flow + , is_flow_step + , is_skipped + , language + , email + , visible_to_owner + , mem_peak + , tag + , priority + ) + VALUES ($1, $2, $3, $4, $5, COALESCE($6, now()), COALESCE($30::bigint, (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.parent_job, + queued_job.created_by, + queued_job.created_at, + queued_job.started_at, + success, + queued_job.script_hash.map(|x| x.0), + queued_job.script_path, + &queued_job.args as &Option>>>, + result as Json<&T>, + raw_code, + raw_lock, + canceled_by.is_some(), + canceled_by.clone().map(|cb| cb.username).flatten(), + canceled_by.clone().map(|cb| cb.reason).flatten(), + queued_job.job_kind.clone() as JobKind, + queued_job.schedule_path, + queued_job.permissioned_as, + &queued_job.flow_status as &Option>>, + &raw_flow as &Option>>, + queued_job.is_flow_step, + skipped, + queued_job.language.clone() as Option, + queued_job.email, + queued_job.visible_to_owner, + if mem_peak > 0 { Some(mem_peak) } else { None }, + queued_job.tag, + queued_job.priority, + duration, + ) + .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"); + + if !queued_job.is_flow_step { + if _duration > 500 + && (queued_job.job_kind == JobKind::Script + || queued_job.job_kind == JobKind::Preview) + { + if let Err(e) = sqlx::query!( "UPDATE completed_job SET flow_status = q.flow_status FROM queue q WHERE completed_job.id = $1 AND q.id = $1 AND q.workspace_id = $2 AND completed_job.workspace_id = $2 AND q.flow_status IS NOT NULL", &queued_job.id, &queued_job.workspace_id @@ -656,9 +661,9 @@ pub async fn add_completed_job( .await { tracing::error!("Could not update job duration: {}", e); } - } - if let Some(parent_job) = queued_job.parent_job { - if let Err(e) = sqlx::query_scalar!( + } + if let Some(parent_job) = queued_job.parent_job { + if let Err(e) = sqlx::query_scalar!( "UPDATE queue SET flow_status = jsonb_set(jsonb_set(COALESCE(flow_status, '{}'::jsonb), array[$1], COALESCE(flow_status->$1, '{}'::jsonb)), array[$1, 'duration_ms'], to_jsonb($2::bigint)) WHERE id = $3 AND workspace_id = $4", &queued_job.id.to_string(), _duration, @@ -669,106 +674,106 @@ pub async fn add_completed_job( .await { tracing::error!("Could not update parent job flow_status: {}", e); } + } } - } - // tracing::error!("Added completed job {:#?}", queued_job); - #[cfg(feature = "enterprise")] - let mut skip_downstream_error_handlers = false; - tx = delete_job(tx, &queued_job.workspace_id, job_id).await?; - // tracing::error!("3 {:?}", start.elapsed()); + // tracing::error!("Added completed job {:#?}", queued_job); - if queued_job.is_flow_step { - if let Some(parent_job) = queued_job.parent_job { - // persist the flow last progress timestamp to avoid zombie flow jobs - tracing::debug!( - "Persisting flow last progress timestamp to flow job: {:?}", - parent_job - ); - sqlx::query!( + let mut _skip_downstream_error_handlers = false; + tx = delete_job(tx, &queued_job.workspace_id, job_id).await?; + // tracing::error!("3 {:?}", start.elapsed()); + + if queued_job.is_flow_step { + if let Some(parent_job) = queued_job.parent_job { + // persist the flow last progress timestamp to avoid zombie flow jobs + tracing::debug!( + "Persisting flow last progress timestamp to flow job: {:?}", + parent_job + ); + sqlx::query!( "UPDATE queue SET last_ping = now() WHERE id = $1 AND workspace_id = $2 AND canceled = false", parent_job, &queued_job.workspace_id ) .execute(&mut *tx) .await?; - if flow_is_done { - let r = sqlx::query_scalar!( + if flow_is_done { + let r = sqlx::query_scalar!( "UPDATE parallel_monitor_lock SET last_ping = now() WHERE parent_flow_id = $1 and job_id = $2 RETURNING 1", parent_job, &queued_job.id ).fetch_optional(&mut *tx).await?; - if r.is_some() { - tracing::info!( + if r.is_some() { + tracing::info!( "parallel flow iteration is done, setting parallel monitor last ping lock for job {}", &queued_job.id ); + } } } - } - } else { - if queued_job.schedule_path.is_some() && queued_job.script_path.is_some() { - let schedule_path = queued_job.schedule_path.as_ref().unwrap(); - let script_path = queued_job.script_path.as_ref().unwrap(); + } else { + if queued_job.schedule_path.is_some() && queued_job.script_path.is_some() { + let schedule_path = queued_job.schedule_path.as_ref().unwrap(); + let script_path = queued_job.script_path.as_ref().unwrap(); - let schedule = - get_schedule_opt(&mut tx, &queued_job.workspace_id, schedule_path).await?; + let schedule = + get_schedule_opt(&mut tx, &queued_job.workspace_id, schedule_path).await?; - if let Some(schedule) = schedule { - #[cfg(feature = "enterprise")] - { - skip_downstream_error_handlers = schedule.ws_error_handler_muted; - } + if let Some(schedule) = schedule { + #[cfg(feature = "enterprise")] + { + _skip_downstream_error_handlers = schedule.ws_error_handler_muted; + } - // script or flow that failed on start and might not have been rescheduled - let schedule_next_tick = !queued_job.is_flow() - || { - let flow_status = queued_job.parse_flow_status(); - flow_status.is_some_and(|fs| { + // script or flow that failed on start and might not have been rescheduled + let schedule_next_tick = !queued_job.is_flow() + || { + let flow_status = queued_job.parse_flow_status(); + flow_status.is_some_and(|fs| { fs.step == 0 && fs.modules.get(0).is_some_and(|m| { matches!(m, FlowStatusModule::WaitingForPriorSteps { .. }) || matches!(m, FlowStatusModule::Failure { job, ..} if job == &Uuid::nil()) }) }) - }; + }; - if schedule_next_tick { - if let Err(err) = handle_maybe_scheduled_job( + if schedule_next_tick { + if let Err(err) = handle_maybe_scheduled_job( + db, + queued_job, + &schedule, + script_path, + &queued_job.workspace_id, + ) + .await + { + match err { + Error::QuotaExceeded(_) => (), + // scheduling next job failed and could not disable schedule => make zombie job to retry + _ => return Ok((Some(job_id), 0, true)), + } + }; + } + + #[cfg(feature = "enterprise")] + if let Err(err) = apply_schedule_handlers( db, - queued_job, &schedule, script_path, &queued_job.workspace_id, + success, + result, + job_id, + queued_job.started_at.unwrap_or(chrono::Utc::now()), + queued_job.priority, ) .await { - match err { - Error::QuotaExceeded(_) => return Err(err.into()), - // scheduling next job failed and could not disable schedule => make zombie job to retry - _ => return Ok(job_id), - } - }; - } - - #[cfg(feature = "enterprise")] - if let Err(err) = apply_schedule_handlers( - db, - &schedule, - script_path, - &queued_job.workspace_id, - success, - result, - job_id, - queued_job.started_at.unwrap_or(chrono::Utc::now()), - queued_job.priority, - ) - .await - { - if !success { - tracing::error!("Could not apply schedule error handler: {}", err); - let base_url = BASE_URL.read().await; - let w_id: &String = &queued_job.workspace_id; - if !matches!(err, Error::QuotaExceeded(_)) { - report_error_to_workspace_handler_or_critical_side_channel( + if !success { + tracing::error!("Could not apply schedule error handler: {}", err); + let base_url = BASE_URL.read().await; + let w_id: &String = &queued_job.workspace_id; + if !matches!(err, Error::QuotaExceeded(_)) { + report_error_to_workspace_handler_or_critical_side_channel( &queued_job, db, format!( @@ -778,82 +783,104 @@ pub async fn add_completed_job( ), ) .await; + } + } else { + tracing::error!("Could not apply schedule recovery handler: {}", err); } - } else { - tracing::error!("Could not apply schedule recovery handler: {}", err); - } - }; - } else { - tracing::error!( + }; + } else { + tracing::error!( "Schedule {schedule_path} in {} not found. Impossible to schedule again and apply schedule handlers", &queued_job.workspace_id ); + } } } - } - if queued_job.concurrent_limit.is_some() { - let concurrency_key = match concurrency_key(db, queued_job).await { - Ok(c) => c, - Err(e) => { - tracing::error!( - "Could not get concurrency key for job {} defaulting to default key: {e:?}", - queued_job.id - ); - legacy_concurrency_key(db, queued_job) - .await - .unwrap_or_else(|| queued_job.full_path_with_workspace()) - } - }; - if let Err(e) = sqlx::query_scalar!( + if queued_job.concurrent_limit.is_some() { + let concurrency_key = match concurrency_key(db, queued_job).await { + Ok(c) => c, + Err(e) => { + tracing::error!( + "Could not get concurrency key for job {} defaulting to default key: {e:?}", + queued_job.id + ); + legacy_concurrency_key(db, queued_job) + .await + .unwrap_or_else(|| queued_job.full_path_with_workspace()) + } + }; + if let Err(e) = sqlx::query_scalar!( "UPDATE concurrency_counter SET job_uuids = job_uuids - $2 WHERE concurrency_id = $1", concurrency_key, queued_job.id.hyphenated().to_string(), ) - .execute(&mut *tx) - .await - { - tracing::error!("Could not decrement concurrency counter: {}", e); - } + .execute(&mut *tx) + .await + { + tracing::error!("Could not decrement concurrency counter: {}", e); + } - if let Err(e) = sqlx::query_scalar!( - "UPDATE concurrency_key SET ended_at = now() WHERE job_id = $1", - queued_job.id, - ) - .execute(&mut *tx) - .await - .map_err(|e| { - Error::InternalErr(format!( + if let Err(e) = sqlx::query_scalar!( + "UPDATE concurrency_key SET ended_at = now() WHERE job_id = $1", + queued_job.id, + ) + .execute(&mut *tx) + .await + .map_err(|e| { + Error::InternalErr(format!( "Error updating to add ended_at timestamp concurrency_key={concurrency_key}: {e:#}" )) - }) { - tracing::error!("Could not update concurrency_key: {}", e); + }) { + tracing::error!("Could not update concurrency_key: {}", e); + } + tracing::debug!("decremented concurrency counter"); } - tracing::debug!("decremented concurrency counter"); - } - if JOB_TOKEN.is_none() { - sqlx::query!("DELETE FROM job_perms WHERE job_id = $1", job_id) - .execute(&mut *tx) - .await?; - } + if JOB_TOKEN.is_none() { + sqlx::query!("DELETE FROM job_perms WHERE job_id = $1", job_id) + .execute(&mut *tx) + .await?; + } - tx.commit().await?; - tracing::info!( - %job_id, - root_job = ?queued_job.root_job.map(|x| x.to_string()).unwrap_or_else(|| String::new()), - path = &queued_job.script_path(), - job_kind = ?queued_job.job_kind, - started_at = ?queued_job.started_at.map(|x| x.to_string()).unwrap_or_else(|| String::new()), - duration = ?_duration, - permissioned_as = ?queued_job.permissioned_as, - email = ?queued_job.email, - created_by = queued_job.created_by, - is_flow_step = queued_job.is_flow_step, - language = ?queued_job.language, - success, - "inserted completed job: {} (success: {success})", - queued_job.id - ); + tx.commit().await?; + + tracing::info!( + %job_id, + root_job = ?queued_job.root_job.map(|x| x.to_string()).unwrap_or_else(|| String::new()), + path = &queued_job.script_path(), + job_kind = ?queued_job.job_kind, + started_at = ?queued_job.started_at.map(|x| x.to_string()).unwrap_or_else(|| String::new()), + duration = ?_duration, + permissioned_as = ?queued_job.permissioned_as, + email = ?queued_job.email, + created_by = queued_job.created_by, + is_flow_step = queued_job.is_flow_step, + language = ?queued_job.language, + success, + "inserted completed job: {} (success: {success})", + queued_job.id + ); + Ok((None, _duration, _skip_downstream_error_handlers)) as windmill_common::error::Result<(Option, i64, bool)> + }) + .retry( + ConstantBuilder::default() + .with_delay(std::time::Duration::from_secs(3)) + .with_max_times(5) + .build(), + ) + .when(|err| !matches!(err, Error::QuotaExceeded(_))) + .notify(|err, dur| { + tracing::error!( + "Could not insert completed job, retrying in {dur:#?}, err: {err:#?}" + ); + }) + .sleep(tokio::time::sleep) + .await?; + + // if scheduling next job failed, return the job_id early to ensure the job get retried after a timeout + if let Some(job_id) = opt_uuid { + return Ok(job_id); + } #[cfg(feature = "cloud")] if *CLOUD_HOSTED && !queued_job.is_flow() && _duration > 1000 { @@ -938,10 +965,10 @@ pub async fn add_completed_job( ), ) .await; - } else if !skip_downstream_error_handlers + } else if !_skip_downstream_error_handlers && (matches!(queued_job.job_kind, JobKind::Script) || matches!(queued_job.job_kind, JobKind::Flow) - && !has_failure_module(db, job_id).await) + && !has_failure_module(db, _job_id).await) && queued_job.parent_job.is_none() { let result = serde_json::from_str( @@ -1232,9 +1259,6 @@ pub async fn send_error_to_workspace_handler<'a, 'c, T: Serialize + Send + Sync> Ok(()) } -use backon::ConstantBuilder; -use backon::{BackoffBuilder, Retryable}; - #[instrument(level = "trace", skip_all)] pub async fn handle_maybe_scheduled_job<'c>( db: &Pool, diff --git a/backend/windmill-worker/Cargo.toml b/backend/windmill-worker/Cargo.toml index f0265accdf..cb64757a7f 100644 --- a/backend/windmill-worker/Cargo.toml +++ b/backend/windmill-worker/Cargo.toml @@ -90,7 +90,7 @@ object_store = { workspace = true, optional = true} convert_case.workspace = true yaml-rust.workspace = true swc_ecma_parser.workspace = true - +backon.workspace = true [build-dependencies] deno_fetch = { workspace = true, optional = true } diff --git a/backend/windmill-worker/src/dedicated_worker.rs b/backend/windmill-worker/src/dedicated_worker.rs index 1a332e76ef..0a39a33610 100644 --- a/backend/windmill-worker/src/dedicated_worker.rs +++ b/backend/windmill-worker/src/dedicated_worker.rs @@ -185,14 +185,14 @@ pub async fn handle_dedicated_process( let result = Arc::new(result); append_logs(&job.id, &job.workspace_id, logs.clone(), db).await; if line.starts_with("wm_res[success]:") { - job_completed_tx.send(JobCompleted { job , result, mem_peak: 0, canceled_by: None, success: true, cached_res_path: None, token: token.to_string() }).await.unwrap() + job_completed_tx.send(JobCompleted { job , result, mem_peak: 0, canceled_by: None, success: true, cached_res_path: None, token: token.to_string(), duration: None }).await.unwrap() } else { - job_completed_tx.send(JobCompleted { job , result, mem_peak: 0, canceled_by: None, success: false, cached_res_path: None, token: token.to_string() }).await.unwrap() + job_completed_tx.send(JobCompleted { job , result, mem_peak: 0, canceled_by: None, success: false, cached_res_path: None, token: token.to_string(), duration: None }).await.unwrap() } }, Err(e) => { tracing::error!("Could not deserialize job result `{line}`: {e:?}"); - job_completed_tx.send(JobCompleted { job , result: Arc::new(to_raw_value(&serde_json::json!({"error": format!("Could not deserialize job result `{line}`: {e:?}")}))), mem_peak: 0, canceled_by: None, success: false, cached_res_path: None, token: token.to_string() }).await.unwrap(); + job_completed_tx.send(JobCompleted { job , result: Arc::new(to_raw_value(&serde_json::json!({"error": format!("Could not deserialize job result `{line}`: {e:?}")}))), mem_peak: 0, canceled_by: None, success: false, cached_res_path: None, token: token.to_string(), duration: None }).await.unwrap(); }, }; logs = init_log.clone(); diff --git a/backend/windmill-worker/src/result_processor.rs b/backend/windmill-worker/src/result_processor.rs index 883953ecec..ea8a883407 100644 --- a/backend/windmill-worker/src/result_processor.rs +++ b/backend/windmill-worker/src/result_processor.rs @@ -186,8 +186,18 @@ async fn send_job_completed( success: bool, cached_res_path: Option, token: String, + duration: Option, ) { - let jc = JobCompleted { job, result, mem_peak, canceled_by, success, cached_res_path, token }; + let jc = JobCompleted { + job, + result, + mem_peak, + canceled_by, + success, + cached_res_path, + token, + duration, + }; job_completed_tx.send(jc).await.expect("send job completed") } @@ -203,6 +213,7 @@ pub async fn process_result( column_order: Option>, new_args: Option>>, db: &DB, + duration: Option, ) -> error::Result { match result { Ok(r) => { @@ -257,6 +268,7 @@ pub async fn process_result( true, cached_res_path, token, + duration, ) .await; Ok(true) @@ -300,6 +312,7 @@ pub async fn process_result( false, cached_res_path, token, + duration, ) .await; Ok(false) @@ -362,7 +375,7 @@ pub async fn handle_receive_completed_job( #[tracing::instrument(name = "completed_job", level = "info", skip_all, fields(job_id = %job.id))] pub async fn process_completed_job( - JobCompleted { job, result, mem_peak, success, cached_res_path, canceled_by, .. }: JobCompleted, + JobCompleted { job, result, mem_peak, success, cached_res_path, canceled_by, duration, .. }: JobCompleted, client: &AuthedClient, db: &DB, worker_dir: &str, @@ -391,8 +404,7 @@ pub async fn process_completed_job( mem_peak.to_owned(), canceled_by, false, - #[cfg(feature = "benchmark")] - bench, + duration, ) .await?; drop(job); @@ -434,8 +446,7 @@ pub async fn process_completed_job( ), worker_name, false, - #[cfg(feature = "benchmark")] - bench, + None, ) .await?; if job.is_flow_step { @@ -501,8 +512,7 @@ pub async fn handle_job_error( err.clone(), worker_name, false, - #[cfg(feature = "benchmark")] - bench, + None, ) .await }; @@ -562,8 +572,7 @@ pub async fn handle_job_error( e, worker_name, false, - #[cfg(feature = "benchmark")] - bench, + None, ) .await; } diff --git a/backend/windmill-worker/src/worker.rs b/backend/windmill-worker/src/worker.rs index a865a6df5a..959d749cf8 100644 --- a/backend/windmill-worker/src/worker.rs +++ b/backend/windmill-worker/src/worker.rs @@ -107,6 +107,9 @@ use crate::{ }, }; +use backon::ConstantBuilder; +use backon::{BackoffBuilder, Retryable}; + #[cfg(feature = "enterprise")] use crate::dedicated_worker::create_dedicated_worker_map; @@ -1107,7 +1110,7 @@ pub async fn run_worker( let (occupancy_rate, occupancy_rate_15s, occupancy_rate_5m, occupancy_rate_30m) = occupancy_metrics.update_occupancy_metrics(); - if let Err(e) = sqlx::query!( + if let Err(e) = (|| sqlx::query!( "UPDATE worker_ping SET ping_at = now(), jobs_executed = $1, custom_tags = $2, occupancy_rate = $3, memory_usage = $4, wm_memory_usage = $5, vcpus = COALESCE($7, vcpus), memory = COALESCE($8, memory), occupancy_rate_15s = $9, occupancy_rate_5m = $10, occupancy_rate_30m = $11 WHERE worker = $6", @@ -1122,7 +1125,19 @@ pub async fn run_worker( occupancy_rate_15s, occupancy_rate_5m, occupancy_rate_30m - ).execute(db).await { + ).execute(db)).retry( + ConstantBuilder::default() + .with_delay(std::time::Duration::from_secs(2)) + .with_max_times(10) + .build(), + ) + .notify(|err, dur| { + tracing::error!( + "retrying updating worker ping in {dur:#?}, err: {err:#?}" + ); + }) + .sleep(tokio::time::sleep) + .await { tracing::error!("failed to update worker ping, exiting: {}", e); killpill_tx.send(()).unwrap_or_default(); } @@ -1349,6 +1364,7 @@ pub async fn run_worker( cached_res_path: None, token: "".to_string(), canceled_by: None, + duration: None, }) .await .expect("send job completed END"); @@ -1704,6 +1720,7 @@ pub struct JobCompleted { pub cached_res_path: Option, pub token: String, pub canceled_by: Option, + pub duration: Option, } async fn do_nativets( @@ -1828,6 +1845,7 @@ async fn handle_queued_job( None }; + let started = Instant::now(); let (raw_code, raw_lock, raw_flow) = match (raw_code, raw_lock, raw_flow) { (None, None, None) => sqlx::query!( "SELECT raw_code, raw_lock, raw_flow AS \"raw_flow: Json>\" @@ -1922,6 +1940,7 @@ async fn handle_queued_job( success: true, cached_res_path: None, token: authed_client.token, + duration: None, }) .await .expect("send job completed"); @@ -2087,6 +2106,7 @@ async fn handle_queued_job( column_order, new_args, db, + Some(started.elapsed().as_millis() as i64), ) .await } diff --git a/backend/windmill-worker/src/worker_flow.rs b/backend/windmill-worker/src/worker_flow.rs index bb9c82f39d..9a5e734dd2 100644 --- a/backend/windmill-worker/src/worker_flow.rs +++ b/backend/windmill-worker/src/worker_flow.rs @@ -1033,8 +1033,7 @@ pub async fn update_flow_status_after_job_completion_internal( canceled_job_to_result(&flow_job), worker_name, true, - #[cfg(feature = "benchmark")] - bench, + None, ) .await?; } else { @@ -1074,9 +1073,8 @@ pub async fn update_flow_status_after_job_completion_internal( Json(&nresult), 0, None, - true, - #[cfg(feature = "benchmark")] - bench, + true, + None, ) .await?; } else { @@ -1092,9 +1090,8 @@ pub async fn update_flow_status_after_job_completion_internal( ), 0, None, - true, - #[cfg(feature = "benchmark")] - bench, + true, + None, ) .await?; } @@ -1131,8 +1128,7 @@ pub async fn update_flow_status_after_job_completion_internal( e, worker_name, true, - #[cfg(feature = "benchmark")] - bench, + None, ) .await; true diff --git a/frontend/src/lib/components/CronInput.svelte b/frontend/src/lib/components/CronInput.svelte index 1feedad79e..3ce3caf385 100644 --- a/frontend/src/lib/components/CronInput.svelte +++ b/frontend/src/lib/components/CronInput.svelte @@ -210,10 +210,11 @@
Invalid cron syntax
{/if} -
+
CronerCroner