diff --git a/backend/src/main.rs b/backend/src/main.rs index e90cb0114c..52001a4e0e 100644 --- a/backend/src/main.rs +++ b/backend/src/main.rs @@ -539,6 +539,11 @@ Windmill Community Edition {GIT_VERSION} loop { tokio::select! { + biased; + _ = monitor_killpill_rx.recv() => { + tracing::info!("received killpill for monitor job"); + break; + }, _ = tokio::time::sleep(Duration::from_secs(30)) => { monitor_db( &db, @@ -693,22 +698,28 @@ Windmill Community Edition {GIT_VERSION} }, Err(e) => { tracing::error!(error = %e, "Could not receive notification, attempting to reconnect listener"); - listener = retry_listen_pg(&db).await; - continue; + tokio::select! { + biased; + _ = monitor_killpill_rx.recv() => { + tracing::info!("received killpill for monitor job"); + break; + }, + new_listener = retry_listen_pg(&db) => { + listener = new_listener; + continue; + } + } } }; - }, - _ = monitor_killpill_rx.recv() => { - println!("received killpill for monitor job"); - break; } } } }); if let Err(e) = h.await { - tracing::error!("Error waiting for monitor handle:{e:#}") + tracing::error!("Error waiting for monitor handle: {e:#}") } + tracing::info!("Monitor exited"); Ok(()) as anyhow::Result<()> }; diff --git a/backend/src/monitor.rs b/backend/src/monitor.rs index cc695748aa..17a8133ff7 100644 --- a/backend/src/monitor.rs +++ b/backend/src/monitor.rs @@ -3,7 +3,10 @@ use std::{ fmt::Display, ops::Mul, str::FromStr, - sync::{atomic::Ordering, Arc}, + sync::{ + atomic::{AtomicU16, Ordering}, + Arc, + }, time::Duration, }; @@ -55,8 +58,8 @@ use windmill_common::{ }; use windmill_queue::cancel_job; use windmill_worker::{ - create_token_for_owner, handle_job_error, AuthedClient, SameWorkerPayload, SendResult, - BUNFIG_INSTALL_SCOPES, JOB_DEFAULT_TIMEOUT, KEEP_JOB_DIR, NPM_CONFIG_REGISTRY, + create_token_for_owner, handle_job_error, AuthedClient, SameWorkerPayload, SameWorkerSender, + SendResult, BUNFIG_INSTALL_SCOPES, JOB_DEFAULT_TIMEOUT, KEEP_JOB_DIR, NPM_CONFIG_REGISTRY, PIP_EXTRA_INDEX_URL, PIP_INDEX_URL, SCRIPT_TOKEN_EXPIRY, }; @@ -1309,6 +1312,8 @@ async fn handle_zombie_jobs // since the job is unrecoverable, the same worker queue should never be sent anything let (same_worker_tx_never_used, _same_worker_rx_never_used) = mpsc::channel::(1); + let same_worker_tx_never_used = + SameWorkerSender(same_worker_tx_never_used, Arc::new(AtomicU16::new(0))); let (send_result_never_used, _send_result_rx_never_used) = mpsc::channel::(1); let label = if job.permissioned_as != format!("u/{}", job.created_by) diff --git a/backend/windmill-api/src/lib.rs b/backend/windmill-api/src/lib.rs index 05202afd12..4a37ad8f3e 100644 --- a/backend/windmill-api/src/lib.rs +++ b/backend/windmill-api/src/lib.rs @@ -391,7 +391,7 @@ pub async fn run_server( let server = server.with_graceful_shutdown(async move { rx.recv().await.ok(); - println!("Graceful shutdown of server"); + tracing::info!("Graceful shutdown of server"); }); server.await?; diff --git a/backend/windmill-common/src/lib.rs b/backend/windmill-common/src/lib.rs index c50fbb6c53..ecdc4cd218 100644 --- a/backend/windmill-common/src/lib.rs +++ b/backend/windmill-common/src/lib.rs @@ -120,7 +120,7 @@ pub async fn shutdown_signal( }, } - println!("signal received, starting graceful shutdown"); + tracing::info!("signal received, starting graceful shutdown"); let _ = tx.send(()); Ok(()) } @@ -167,7 +167,7 @@ pub async fn serve_metrics( if let Err(e) = axum::serve(listener, router.into_make_service()) .with_graceful_shutdown(async move { rx.recv().await.ok(); - println!("Graceful shutdown of metrics"); + tracing::info!("Graceful shutdown of metrics"); }) .await { diff --git a/backend/windmill-worker/src/lib.rs b/backend/windmill-worker/src/lib.rs index f532f2f53f..92f6e0e74c 100644 --- a/backend/windmill-worker/src/lib.rs +++ b/backend/windmill-worker/src/lib.rs @@ -20,11 +20,14 @@ mod mysql_executor; mod pg_executor; mod php_executor; mod python_executor; +mod result_processor; mod rust_executor; mod worker; mod worker_flow; mod worker_lockfiles; pub use worker::*; +pub use result_processor::handle_job_error; + pub use bun_executor::{get_common_bun_proc_envs, install_bun_lockfile, prepare_job_dir}; pub use deno_executor::generate_deno_lock; diff --git a/backend/windmill-worker/src/result_processor.rs b/backend/windmill-worker/src/result_processor.rs new file mode 100644 index 0000000000..e749f28d3d --- /dev/null +++ b/backend/windmill-worker/src/result_processor.rs @@ -0,0 +1,435 @@ +use serde::Serialize; +use sqlx::{types::Json, Pool, Postgres}; +use std::sync::Arc; + +use uuid::Uuid; + +use windmill_common::{ + error::{self, Error}, + jobs::QueuedJob, + worker::to_raw_value, + DB, +}; + +use windmill_queue::{append_logs, get_queued_job, CanceledBy, WrappedError}; + +#[cfg(feature = "prometheus")] +use windmill_queue::register_metric; + +use serde_json::{json, value::RawValue}; + +use tokio::sync::mpsc::Sender; + +use windmill_queue::{add_completed_job, add_completed_job_error}; + +use crate::{ + bash_executor::ANSI_ESCAPE_RE, + common::{read_result, save_in_cache}, + worker_flow::update_flow_status_after_job_completion, + AuthedClient, Histo, JobCompleted, JobCompletedSender, SameWorkerSender, SendResult, +}; + +async fn send_job_completed( + job_completed_tx: JobCompletedSender, + job: Arc, + result: Arc>, + mem_peak: i32, + canceled_by: Option, + success: bool, + cached_res_path: Option, + token: String, +) { + let jc = JobCompleted { job, result, mem_peak, canceled_by, success, cached_res_path, token }; + job_completed_tx.send(jc).await.expect("send job completed") +} + +pub async fn process_result( + job: Arc, + result: error::Result>>, + job_dir: &str, + job_completed_tx: JobCompletedSender, + mem_peak: i32, + canceled_by: Option, + cached_res_path: Option, + token: String, + column_order: Option>, + db: &DB, +) -> error::Result<()> { + match result { + Ok(r) => { + let job = if let Some(column_order) = column_order { + let mut job_with_column_order = (*job).clone(); + match job_with_column_order.flow_status { + Some(_) => { + tracing::warn!("flow_status was expected to be none"); + } + None => { + job_with_column_order.flow_status = + Some(sqlx::types::Json(to_raw_value(&serde_json::json!({ + "_metadata": { + "column_order": column_order + } + })))); + } + } + Arc::new(job_with_column_order) + } else { + job + }; + + send_job_completed( + job_completed_tx, + job, + r, + mem_peak, + canceled_by, + true, + cached_res_path, + token, + ) + .await; + } + Err(e) => { + let error_value = match e { + Error::ExitStatus(i) => { + let res = read_result(job_dir).await.ok(); + + if res.as_ref().is_some_and(|x| !x.get().is_empty()) { + res.unwrap() + } else { + let last_10_log_lines = sqlx::query_scalar!( + "SELECT right(logs, 600) FROM job_logs WHERE job_id = $1 AND workspace_id = $2 ORDER BY created_at DESC LIMIT 1", + &job.id, + &job.workspace_id + ).fetch_one(db).await.ok().flatten().unwrap_or("".to_string()); + + let log_lines = last_10_log_lines + .split("CODE EXECUTION ---") + .last() + .unwrap_or(&last_10_log_lines); + + extract_error_value(log_lines, i, job.flow_step_id.clone()) + } + } + err @ _ => to_raw_value(&SerializedError { + message: format!("error during execution of the script:\n{}", err), + name: "ExecutionErr".to_string(), + step_id: job.flow_step_id.clone(), + }), + }; + + send_job_completed( + job_completed_tx, + job, + Arc::new(to_raw_value(&error_value)), + mem_peak, + canceled_by, + false, + cached_res_path, + token, + ) + .await; + } + }; + Ok(()) +} + +pub async fn handle_receive_completed_job< + R: rsmq_async::RsmqConnection + Send + Sync + Clone + 'static, +>( + jc: JobCompleted, + base_internal_url: &str, + db: &DB, + worker_dir: &str, + same_worker_tx: &SameWorkerSender, + rsmq: Option, + worker_name: &str, + worker_save_completed_job_duration: Option, + worker_flow_transition_duration: Option, + job_completed_tx: Sender, +) { + let token = jc.token.clone(); + let workspace = jc.job.workspace_id.clone(); + let client = AuthedClient { + base_internal_url: base_internal_url.to_string(), + workspace, + token, + force_client: None, + }; + let job = jc.job.clone(); + let mem_peak = jc.mem_peak.clone(); + let canceled_by = jc.canceled_by.clone(); + if let Err(err) = process_completed_job( + jc, + &client, + db, + &worker_dir, + same_worker_tx.clone(), + rsmq.clone(), + worker_name, + worker_save_completed_job_duration, + worker_flow_transition_duration, + job_completed_tx.clone(), + ) + .await + { + handle_job_error( + db, + &client, + job.as_ref(), + mem_peak, + canceled_by, + err, + false, + same_worker_tx.clone(), + &worker_dir, + rsmq.clone(), + worker_name, + job_completed_tx, + ) + .await; + } +} + +#[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, + client: &AuthedClient, + db: &DB, + worker_dir: &str, + same_worker_tx: SameWorkerSender, + rsmq: Option, + worker_name: &str, + _worker_save_completed_job_duration: Option, + _worker_flow_transition_duration: Option, + job_completed_tx: Sender, +) -> windmill_common::error::Result<()> { + if success { + // println!("bef completed job{:?}", SystemTime::now()); + if let Some(cached_path) = cached_res_path { + save_in_cache(db, client, &job, cached_path.to_string(), &result).await; + } + + let is_flow_step = job.is_flow_step; + let parent_job = job.parent_job.clone(); + let job_id = job.id.clone(); + let workspace_id = job.workspace_id.clone(); + #[cfg(feature = "prometheus")] + let timer = _worker_save_completed_job_duration + .as_ref() + .map(|x| x.start_timer()); + add_completed_job( + db, + &job, + true, + false, + Json(&result), + mem_peak.to_owned(), + canceled_by, + rsmq.clone(), + false, + ) + .await?; + drop(job); + + #[cfg(feature = "prometheus")] + timer.map(|x| x.stop_and_record()); + + if is_flow_step { + if let Some(parent_job) = parent_job { + #[cfg(feature = "prometheus")] + let timer = _worker_flow_transition_duration + .as_ref() + .map(|x| x.start_timer()); + tracing::info!(parent_flow = %parent_job, subflow = %job_id, "updating flow status (2)"); + update_flow_status_after_job_completion( + db, + client, + parent_job, + &job_id, + &workspace_id, + true, + result, + false, + same_worker_tx.clone(), + &worker_dir, + None, + rsmq.clone(), + worker_name, + job_completed_tx, + ) + .await?; + #[cfg(feature = "prometheus")] + timer.map(|x| x.stop_and_record()); + } + } + } else { + let result = add_completed_job_error( + db, + &job, + mem_peak.to_owned(), + canceled_by, + serde_json::from_str(result.get()).unwrap_or_else( + |_| json!({ "message": format!("Non serializable error: {}", result.get()) }), + ), + rsmq.clone(), + worker_name, + false, + ) + .await?; + if job.is_flow_step { + if let Some(parent_job) = job.parent_job { + tracing::error!(parent_flow = %parent_job, subflow = %job.id, "process completed job error, updating flow status"); + update_flow_status_after_job_completion( + db, + client, + parent_job, + &job.id, + &job.workspace_id, + false, + Arc::new(serde_json::value::to_raw_value(&result).unwrap()), + false, + same_worker_tx, + &worker_dir, + None, + rsmq, + worker_name, + job_completed_tx, + ) + .await?; + } + } + } + Ok(()) +} + +#[tracing::instrument(name = "job_error", level = "info", skip_all, fields(job_id = %job.id))] +pub async fn handle_job_error( + db: &Pool, + client: &AuthedClient, + job: &QueuedJob, + mem_peak: i32, + canceled_by: Option, + err: Error, + unrecoverable: bool, + same_worker_tx: SameWorkerSender, + worker_dir: &str, + rsmq: Option, + worker_name: &str, + job_completed_tx: Sender, +) { + let err = match err { + Error::JsonErr(err) => err, + _ => json!({"message": err.to_string(), "name": "InternalErr"}), + }; + + let rsmq_2 = rsmq.clone(); + let update_job_future = || async { + append_logs( + &job.id, + &job.workspace_id, + format!("Unexpected error during job execution:\n{err:#?}"), + db, + ) + .await; + add_completed_job_error( + db, + job, + mem_peak, + canceled_by.clone(), + err.clone(), + rsmq_2, + worker_name, + false, + ) + .await + }; + + let update_job_future = if job.is_flow_step || job.is_flow() { + let (flow, job_status_to_update) = if let Some(parent_job_id) = job.parent_job { + if let Err(e) = update_job_future().await { + tracing::error!( + "error updating job future for job {} for handle_job_error: {e:#}", + job.id + ); + } + (parent_job_id, job.id) + } else { + (job.id, Uuid::nil()) + }; + + let wrapped_error = WrappedError { error: err.clone() }; + tracing::error!(parent_flow = %flow, subflow = %job_status_to_update, "handle job error, updating flow status: {err:?}"); + let updated_flow = update_flow_status_after_job_completion( + db, + client, + flow, + &job_status_to_update, + &job.workspace_id, + false, + Arc::new(serde_json::value::to_raw_value(&wrapped_error).unwrap()), + unrecoverable, + same_worker_tx, + worker_dir, + None, + rsmq.clone(), + worker_name, + job_completed_tx.clone(), + ) + .await; + + if let Err(err) = updated_flow { + if let Some(parent_job_id) = job.parent_job { + if let Ok(Some(parent_job)) = + get_queued_job(&parent_job_id, &job.workspace_id, &db).await + { + let e = json!({"message": err.to_string(), "name": "InternalErr"}); + append_logs( + &parent_job.id, + &job.workspace_id, + format!("Unexpected error during flow job error handling:\n{err}"), + db, + ) + .await; + let _ = add_completed_job_error( + db, + &parent_job, + mem_peak, + canceled_by.clone(), + e, + rsmq, + worker_name, + false, + ) + .await; + } + } + } + + None + } else { + Some(update_job_future) + }; + if let Some(f) = update_job_future { + let _ = f().await; + } + tracing::error!(job_id = %job.id, "error handling job: {err:?} {} {} {}", job.id, job.workspace_id, job.created_by); +} + +#[derive(Debug, Serialize)] +pub struct SerializedError { + pub message: String, + pub name: String, + #[serde(skip_serializing_if = "Option::is_none")] + pub step_id: Option, +} +pub fn extract_error_value(log_lines: &str, i: i32, step_id: Option) -> Box { + return to_raw_value(&SerializedError { + message: format!( + "ExitCode: {i}, last log lines:\n{}", + ANSI_ESCAPE_RE.replace_all(log_lines.trim(), "").to_string() + ), + name: "ExecutionErr".to_string(), + step_id, + }); +} diff --git a/backend/windmill-worker/src/worker.rs b/backend/windmill-worker/src/worker.rs index 40aad31ebe..11f231d272 100644 --- a/backend/windmill-worker/src/worker.rs +++ b/backend/windmill-worker/src/worker.rs @@ -7,7 +7,12 @@ */ use windmill_common::{ - auth::{fetch_authed_from_permissioned_as, JWTAuthClaims, JobPerms, JWT_SECRET}, scripts::PREVIEW_IS_TAR_CODEBASE_HASH, worker::{get_memory, get_vcpus, get_windmill_memory_usage, get_worker_memory_usage, write_file, ROOT_CACHE_DIR, TMP_DIR} + auth::{fetch_authed_from_permissioned_as, JWTAuthClaims, JobPerms, JWT_SECRET}, + scripts::PREVIEW_IS_TAR_CODEBASE_HASH, + worker::{ + get_memory, get_vcpus, get_windmill_memory_usage, get_worker_memory_usage, write_file, + ROOT_CACHE_DIR, TMP_DIR, + }, }; use anyhow::{Context, Result}; @@ -30,7 +35,7 @@ use std::{ collections::{hash_map::DefaultHasher, HashMap}, hash::Hash, sync::{ - atomic::{AtomicBool, AtomicUsize, Ordering}, + atomic::{AtomicBool, AtomicU16, Ordering}, Arc, }, time::Duration, @@ -45,19 +50,19 @@ use windmill_common::{ scripts::{get_full_hub_script_by_path, ScriptHash, ScriptLang, PREVIEW_IS_CODEBASE_HASH}, users::SUPERADMIN_SECRET_EMAIL, utils::StripPath, - worker::{to_raw_value, update_ping, CLOUD_HOSTED, NO_LOGS, WORKER_CONFIG, WORKER_GROUP}, + worker::{update_ping, CLOUD_HOSTED, NO_LOGS, WORKER_CONFIG, WORKER_GROUP}, DB, IS_READY, }; use windmill_queue::{ - append_logs, canceled_job_to_result, empty_result, get_queued_job, pull, push, CanceledBy, - PushArgs, PushIsolationLevel, WrappedError, HTTP_CLIENT, + append_logs, canceled_job_to_result, empty_result, pull, push, CanceledBy, + PushArgs, PushIsolationLevel, HTTP_CLIENT, }; #[cfg(feature = "prometheus")] use windmill_queue::register_metric; -use serde_json::{json, value::RawValue}; +use serde_json::value::RawValue; #[cfg(any(target_os = "linux", target_os = "macos"))] use tokio::fs::symlink; @@ -77,23 +82,12 @@ use tokio::{ use rand::Rng; -use windmill_queue::{add_completed_job, add_completed_job_error}; use crate::{ - bash_executor::{handle_bash_job, handle_powershell_job, ANSI_ESCAPE_RE}, bun_executor::handle_bun_job, common::{ + bash_executor::{handle_bash_job, handle_powershell_job}, bun_executor::handle_bun_job, common::{ build_args_map, get_cached_resource_value_if_valid, get_reserved_variables, hash_args, - read_result, save_in_cache, NO_LOGS_AT_ALL, SLOW_LOGS, - }, - deno_executor::handle_deno_job, - go_executor::handle_go_job, - graphql_executor::do_graphql, - js_eval::{eval_fetch_timeout, transpile_ts}, - mysql_executor::do_mysql, - pg_executor::do_postgresql, - rust_executor::handle_rust_job, - php_executor::handle_php_job, - python_executor::handle_python_job, - worker_flow::{ + NO_LOGS_AT_ALL, SLOW_LOGS, + }, deno_executor::handle_deno_job, go_executor::handle_go_job, graphql_executor::do_graphql, js_eval::{eval_fetch_timeout, transpile_ts}, mysql_executor::do_mysql, pg_executor::do_postgresql, php_executor::handle_php_job, python_executor::handle_python_job, result_processor::{handle_job_error, handle_receive_completed_job, process_result}, rust_executor::handle_rust_job, worker_flow::{ handle_flow, update_flow_status_after_job_completion, update_flow_status_in_progress, }, worker_lockfiles::{ handle_app_dependency_job, handle_dependency_job, handle_flow_dependency_job, @@ -567,76 +561,24 @@ macro_rules! add_time { } #[cfg(feature = "prometheus")] -type Histo = Arc; +pub type Histo = Arc; #[cfg(feature = "prometheus")] type GGauge = Arc>; #[cfg(not(feature = "prometheus"))] -type Histo = (); +pub type Histo = (); #[cfg(not(feature = "prometheus"))] type GGauge = (); -async fn handle_receive_completed_job< - R: rsmq_async::RsmqConnection + Send + Sync + Clone + 'static, ->( - jc: JobCompleted, - base_internal_url: &str, - db: &Pool, - worker_dir: &str, - same_worker_tx: &Sender, - rsmq: Option, - worker_name: &str, - worker_save_completed_job_duration: Option, - worker_flow_transition_duration: Option, - job_completed_tx: Sender, -) { - let token = jc.token.clone(); - let workspace = jc.job.workspace_id.clone(); - let client = AuthedClient { - base_internal_url: base_internal_url.to_string(), - workspace, - token, - force_client: None, - }; - let job = jc.job.clone(); - let mem_peak = jc.mem_peak.clone(); - let canceled_by = jc.canceled_by.clone(); - if let Err(err) = process_completed_job( - jc, - &client, - db, - &worker_dir, - same_worker_tx.clone(), - rsmq.clone(), - worker_name, - worker_save_completed_job_duration, - worker_flow_transition_duration, - job_completed_tx.clone(), - ) - .await - { - handle_job_error( - db, - &client, - job.as_ref(), - mem_peak, - canceled_by, - err, - false, - same_worker_tx.clone(), - &worker_dir, - rsmq.clone(), - worker_name, - job_completed_tx, - ) - .await; - } -} - #[allow(dead_code)] #[derive(Clone)] pub struct JobCompletedSender(Sender, Option, Option); + +#[derive(Clone)] +pub struct SameWorkerSender(pub Sender, pub Arc); + + pub struct SameWorkerPayload { pub job_id: Uuid, pub recoverable: bool, @@ -660,6 +602,17 @@ impl JobCompletedSender { } } +impl SameWorkerSender { + pub async fn send( + &self, + payload: SameWorkerPayload, + ) -> Result<(), tokio::sync::mpsc::error::SendError> { + self.1.fetch_add(1, Ordering::Relaxed); + self.0.send(payload).await + } +} + + // on linux, we drop caches every DROP_CACHE_PERIOD to avoid OOM killer believing we are using too much memory just because we create lots of files when executing jobs #[cfg(any(target_os = "linux"))] pub async fn drop_cache() { @@ -774,8 +727,7 @@ pub async fn run_worker { - #[cfg(feature = "prometheus")] - if let Some(wj) = worker_job_completed_channel_queue2.as_ref() { - wj.dec(); - } - let rsmq2 = rsmq2.clone(); - let is_init_script_and_failure = - !jc.success && jc.job.tag.as_str() == INIT_SCRIPT_TAG; - let is_dependency_job = matches!( - jc.job.job_kind, - JobKind::Dependencies | JobKind::FlowDependencies); - handle_receive_completed_job( - jc, - &base_internal_url2, - &db2, - &worker_dir2, - &same_worker_tx2, - rsmq2, - &worker_name2, - worker_save_completed_job_duration2.clone(), - worker_flow_transition_duration2.clone(), - job_completed_sender.clone(), - ) - .await; - if is_init_script_and_failure { - tracing::error!("init script errored, exiting"); - killpill_tx2.send(()).unwrap_or_default(); - } - if is_dependency_job && is_dedicated_worker { - tracing::error!("Dedicated worker executed a dependency job, a new script has been deployed. Exiting expecting to be restarted."); - sqlx::query!("UPDATE config SET config = config WHERE name = $1", format!("worker__{}", *WORKER_GROUP)) - .execute(&db2) - .await - .expect("update config to trigger restart of all dedicated workers at that config"); - killpill_tx2.send(()).unwrap_or_default(); + let job_completed_processor_is_done = Arc::new(AtomicBool::new(false)); + let job_completed_processor_is_done2 = job_completed_processor_is_done.clone(); + let same_worker_queue_size2 = same_worker_queue_size.clone(); + let send_result = tokio::spawn( + (async move { + let mut has_been_killed = false; + + //if we have been killed, we want to drain the queue of jobs + while let Some(sr) = if has_been_killed && same_worker_queue_size2.load(Ordering::SeqCst) == 0 { job_completed_rx.try_recv().ok() } else { job_completed_rx.recv().await }{ + match sr { + SendResult::JobCompleted(jc) => { + #[cfg(feature = "prometheus")] + if let Some(wj) = worker_job_completed_channel_queue2.as_ref() { + wj.dec(); + } + let rsmq2 = rsmq2.clone(); + + let is_init_script_and_failure = !jc.success && jc.job.tag.as_str() == INIT_SCRIPT_TAG; + let is_dependency_job = matches!( + jc.job.job_kind, + JobKind::Dependencies | JobKind::FlowDependencies + ); + handle_receive_completed_job( + jc, + &base_internal_url2, + &db2, + &worker_dir2, + &same_worker_tx2, + rsmq2, + &worker_name2, + worker_save_completed_job_duration2.clone(), + worker_flow_transition_duration2.clone(), + job_completed_sender.clone(), + ) + .await; + if is_init_script_and_failure { + tracing::error!("init script errored, exiting"); + killpill_tx2.send(()).unwrap_or_default(); + } + if is_dependency_job && is_dedicated_worker { + tracing::error!("Dedicated worker executed a dependency job, a new script has been deployed. Exiting expecting to be restarted."); + sqlx::query!( + "UPDATE config SET config = config WHERE name = $1", + format!("worker__{}", *WORKER_GROUP) + ) + .execute(&db2) + .await + .expect("update config to trigger restart of all dedicated workers at that config"); + killpill_tx2.send(()).unwrap_or_default(); + } } - - } - SendResult::UpdateFlow { - flow, - w_id, - success, - result, - worker_dir, - stop_early_override, - token, - } => { - // let r; - tracing::info!(parent_flow = %flow, "updating flow status"); - if let Err(e) = update_flow_status_after_job_completion( - &db2, - &AuthedClient { - base_internal_url: base_internal_url2.to_string(), - workspace: w_id.clone(), - token: token.clone(), - force_client: None, - }, + SendResult::UpdateFlow { flow, - &Uuid::nil(), - &w_id, + w_id, success, - Arc::new(result), - true, - same_worker_tx2.clone(), - &worker_dir, + result, + worker_dir, stop_early_override, - rsmq2.clone(), - &worker_name2, - job_completed_sender.clone(), - ) - .await - { - tracing::error!("Error updating flow status after job completion for {flow} on {worker_name2}: {e:#}"); + token, + } => { + // let r; + tracing::info!(parent_flow = %flow, "updating flow status"); + if let Err(e) = update_flow_status_after_job_completion( + &db2, + &AuthedClient { + base_internal_url: base_internal_url2.to_string(), + workspace: w_id.clone(), + token: token.clone(), + force_client: None, + }, + flow, + &Uuid::nil(), + &w_id, + success, + Arc::new(result), + true, + same_worker_tx2.clone(), + &worker_dir, + stop_early_override, + rsmq2.clone(), + &worker_name2, + job_completed_sender.clone(), + ) + .await + { + tracing::error!("Error updating flow status after job completion for {flow} on {worker_name2}: {e:#}"); + } + } + SendResult::Kill => { + has_been_killed = true; } - } - SendResult::Kill => { - break; } } - } - tracing::info!("stopped processing new completed jobs"); - while thread_count.load(Ordering::SeqCst) > 0 { - tokio::time::sleep(Duration::from_millis(50)).await; - } - tracing::info!("finished processing all completed jobs"); - }).instrument(tracing::Span::current())); + job_completed_processor_is_done2.store(true, Ordering::SeqCst); + tracing::info!("finished processing all completed jobs"); + }) + .instrument(tracing::Span::current()), + ); let mut last_executed_job: Option = None; let mut last_checked_suspended = Instant::now(); @@ -1358,6 +1321,7 @@ pub async fn run_worker NUM_SECS_READINGS { + let (vcpus, memory) = if *REFRESH_CGROUP_READINGS + && last_reading.elapsed().as_secs() > NUM_SECS_READINGS + { last_reading = Instant::now(); (get_vcpus(), get_memory()) } else { @@ -1448,33 +1414,53 @@ pub async fn run_worker("UPDATE queue SET last_ping = now() WHERE id = $1 RETURNING *") - .bind(same_worker_job.job_id) - .fetch_optional(db) - .await - .map_err(|_| Error::InternalErr("Impossible to fetch same_worker job".to_string())); - if r.is_err() && !same_worker_job.recoverable { - tracing::error!("failed to fetch same_worker job on a non recoverable job, exiting"); - job_completed_tx.0.send(SendResult::Kill).await.unwrap(); - break; - } else { - r - } + same_worker_queue_size.fetch_sub(1, Ordering::SeqCst); + tracing::debug!( + "received {} from same worker channel", + same_worker_job.job_id + ); + let r = sqlx::query_as::<_, QueuedJob>( + "UPDATE queue SET last_ping = now() WHERE id = $1 RETURNING *", + ) + .bind(same_worker_job.job_id) + .fetch_optional(db) + .await + .map_err(|_| Error::InternalErr("Impossible to fetch same_worker job".to_string())); + if r.is_err() && !same_worker_job.recoverable { + tracing::error!( + "failed to fetch same_worker job on a non recoverable job, exiting" + ); + job_completed_tx.0.send(SendResult::Kill).await.expect("send kill to job completed tx"); + break; + } else { + r + } + } else if let Ok(_) = killpill_rx.try_recv() { + if !killed_but_draining_same_worker_jobs { + tracing::info!("received killpill for worker {}, processing only same worker jobs", i_worker); + killed_but_draining_same_worker_jobs = true; + job_completed_tx.0.send(SendResult::Kill).await.expect("send kill to job completed tx"); + } + continue; + } else if killed_but_draining_same_worker_jobs { + if job_completed_processor_is_done.load(Ordering::SeqCst) { + tracing::info!("all running jobs have completed and all completed jobs have been fully processed, exiting"); + break; + } else { + tracing::info!("there may be same_worker jobs to process later, waiting for job_completed_processor to finish progressing all remaining flows before exiting"); + tokio::time::sleep(Duration::from_millis(200)).await; + continue; + } } else { - let pull_time = Instant::now(); - let suspend_first = if suspend_first_success || last_checked_suspended.elapsed().as_secs() > 3 { - last_checked_suspended = Instant::now(); - true - } else { false }; + let suspend_first = + if suspend_first_success || last_checked_suspended.elapsed().as_secs() > 3 { + last_checked_suspended = Instant::now(); + true + } else { + false + }; let job = pull(&db, rsmq.clone(), suspend_first).await; add_time!(timing, loop_start, "post pull"); let duration_pull_s = pull_time.elapsed().as_secs_f64(); @@ -1491,7 +1477,6 @@ pub async fn run_worker 0.1 { tracing::warn!("pull took more than 0.1s ({duration_pull_s}) this is a sign that the database is undersized for this load. empty: {empty}, err: {err_pull}"); #[cfg(feature = "prometheus")] @@ -1845,7 +1830,7 @@ pub async fn run_worker( db: &Pool, - same_worker_tx: Sender, + same_worker_tx: SameWorkerSender, worker_name: &str, rsmq: Option, ) -> error::Result<()> { @@ -1914,118 +1899,6 @@ async fn queue_init_bash_maybe<'c, R: rsmq_async::RsmqConnection + Send + 'c>( // logs: String, // ) -> error::Result<()> { -#[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, - client: &AuthedClient, - db: &DB, - worker_dir: &str, - same_worker_tx: Sender, - rsmq: Option, - worker_name: &str, - _worker_save_completed_job_duration: Option, - _worker_flow_transition_duration: Option, - job_completed_tx: Sender, -) -> windmill_common::error::Result<()> { - if success { - // println!("bef completed job{:?}", SystemTime::now()); - if let Some(cached_path) = cached_res_path { - save_in_cache(db, client, &job, cached_path.to_string(), &result).await; - } - - let is_flow_step = job.is_flow_step; - let parent_job = job.parent_job.clone(); - let job_id = job.id.clone(); - let workspace_id = job.workspace_id.clone(); - #[cfg(feature = "prometheus")] - let timer = _worker_save_completed_job_duration - .as_ref() - .map(|x| x.start_timer()); - add_completed_job( - db, - &job, - true, - false, - Json(&result), - mem_peak.to_owned(), - canceled_by, - rsmq.clone(), - false, - ) - .await?; - drop(job); - - #[cfg(feature = "prometheus")] - timer.map(|x| x.stop_and_record()); - - if is_flow_step { - if let Some(parent_job) = parent_job { - #[cfg(feature = "prometheus")] - let timer = _worker_flow_transition_duration - .as_ref() - .map(|x| x.start_timer()); - tracing::info!(parent_flow = %parent_job, subflow = %job_id, "updating flow status (2)"); - update_flow_status_after_job_completion( - db, - client, - parent_job, - &job_id, - &workspace_id, - true, - result, - false, - same_worker_tx.clone(), - &worker_dir, - None, - rsmq.clone(), - worker_name, - job_completed_tx, - ) - .await?; - #[cfg(feature = "prometheus")] - timer.map(|x| x.stop_and_record()); - } - } - } else { - let result = add_completed_job_error( - db, - &job, - mem_peak.to_owned(), - canceled_by, - serde_json::from_str(result.get()).unwrap_or_else( - |_| json!({ "message": format!("Non serializable error: {}", result.get()) }), - ), - rsmq.clone(), - worker_name, - false, - ) - .await?; - if job.is_flow_step { - if let Some(parent_job) = job.parent_job { - tracing::error!(parent_flow = %parent_job, subflow = %job.id, "process completed job error, updating flow status"); - update_flow_status_after_job_completion( - db, - client, - parent_job, - &job.id, - &job.workspace_id, - false, - Arc::new(serde_json::value::to_raw_value(&result).unwrap()), - false, - same_worker_tx, - &worker_dir, - None, - rsmq, - worker_name, - job_completed_tx, - ) - .await?; - } - } - } - Ok(()) -} - // fn build_language_metrics( // worker_execution_failed: &HashMap< // Option, @@ -2062,136 +1935,6 @@ pub async fn process_completed_job( - db: &Pool, - client: &AuthedClient, - job: &QueuedJob, - mem_peak: i32, - canceled_by: Option, - err: Error, - unrecoverable: bool, - same_worker_tx: Sender, - worker_dir: &str, - rsmq: Option, - worker_name: &str, - job_completed_tx: Sender, -) { - let err = match err { - Error::JsonErr(err) => err, - _ => json!({"message": err.to_string(), "name": "InternalErr"}), - }; - - let rsmq_2 = rsmq.clone(); - let update_job_future = || async { - append_logs( - &job.id, - &job.workspace_id, - format!("Unexpected error during job execution:\n{err:#?}"), - db, - ) - .await; - add_completed_job_error( - db, - job, - mem_peak, - canceled_by.clone(), - err.clone(), - rsmq_2, - worker_name, - false, - ) - .await - }; - - let update_job_future = if job.is_flow_step || job.is_flow() { - let (flow, job_status_to_update) = if let Some(parent_job_id) = job.parent_job { - if let Err(e) = update_job_future().await { - tracing::error!( - "error updating job future for job {} for handle_job_error: {e:#}", - job.id - ); - } - (parent_job_id, job.id) - } else { - (job.id, Uuid::nil()) - }; - - let wrapped_error = WrappedError { error: err.clone() }; - tracing::error!(parent_flow = %flow, subflow = %job_status_to_update, "handle job error, updating flow status: {err:?}"); - let updated_flow = update_flow_status_after_job_completion( - db, - client, - flow, - &job_status_to_update, - &job.workspace_id, - false, - Arc::new(serde_json::value::to_raw_value(&wrapped_error).unwrap()), - unrecoverable, - same_worker_tx, - worker_dir, - None, - rsmq.clone(), - worker_name, - job_completed_tx.clone(), - ) - .await; - - if let Err(err) = updated_flow { - if let Some(parent_job_id) = job.parent_job { - if let Ok(Some(parent_job)) = - get_queued_job(&parent_job_id, &job.workspace_id, &db).await - { - let e = json!({"message": err.to_string(), "name": "InternalErr"}); - append_logs( - &parent_job.id, - &job.workspace_id, - format!("Unexpected error during flow job error handling:\n{err}"), - db, - ) - .await; - let _ = add_completed_job_error( - db, - &parent_job, - mem_peak, - canceled_by.clone(), - e, - rsmq, - worker_name, - false, - ) - .await; - } - } - } - - None - } else { - Some(update_job_future) - }; - if let Some(f) = update_job_future { - let _ = f().await; - } - tracing::error!(job_id = %job.id, "error handling job: {err:?} {} {} {}", job.id, job.workspace_id, job.created_by); -} - -#[derive(Debug, Serialize)] -struct SerializedError { - message: String, - name: String, - #[serde(skip_serializing_if = "Option::is_none")] - step_id: Option, -} -fn extract_error_value(log_lines: &str, i: i32, step_id: Option) -> Box { - return to_raw_value( - &SerializedError { - message: format!("ExitCode: {i}, last log lines:\n{}", ANSI_ESCAPE_RE.replace_all(log_lines.trim(), "").to_string()), - name: "ExecutionErr".to_string(), - step_id, - }, - ); -} - pub enum SendResult { JobCompleted(JobCompleted), UpdateFlow { @@ -2281,7 +2024,7 @@ async fn handle_queued_job( worker_name: &str, worker_dir: &str, job_dir: &str, - same_worker_tx: Sender, + same_worker_tx: SameWorkerSender, base_internal_url: &str, rsmq: Option, job_completed_tx: JobCompletedSender, @@ -2591,100 +2334,6 @@ async fn handle_queued_job( Ok(()) } -async fn process_result( - job: Arc, - result: error::Result>>, - job_dir: &str, - job_completed_tx: JobCompletedSender, - mem_peak: i32, - canceled_by: Option, - cached_res_path: Option, - token: String, - column_order: Option>, - db: &DB, -) -> error::Result<()> { - match result { - Ok(r) => { - let job = if let Some(column_order) = column_order { - let mut job_with_column_order = (*job).clone(); - match job_with_column_order.flow_status { - Some(_) => { - tracing::warn!("flow_status was expected to be none"); - } - None => { - job_with_column_order.flow_status = - Some(sqlx::types::Json(to_raw_value(&serde_json::json!({ - "_metadata": { - "column_order": column_order - } - })))); - } - } - Arc::new(job_with_column_order) - } else { - job - }; - job_completed_tx - .send(JobCompleted { - job, - result: r, - mem_peak, - canceled_by, - success: true, - cached_res_path, - token: token, - }) - .await - .expect("send job completed"); - } - Err(e) => { - let error_value = match e { - Error::ExitStatus(i) => { - let res = read_result(job_dir).await.ok(); - - if res.as_ref().is_some_and(|x| !x.get().is_empty()) { - res.unwrap() - } else { - let last_10_log_lines = sqlx::query_scalar!( - "SELECT right(logs, 600) FROM job_logs WHERE job_id = $1 AND workspace_id = $2 ORDER BY created_at DESC LIMIT 1", - &job.id, - &job.workspace_id - ).fetch_one(db).await.ok().flatten().unwrap_or("".to_string()); - - let log_lines = last_10_log_lines - .split("CODE EXECUTION ---") - .last() - .unwrap_or(&last_10_log_lines); - - extract_error_value(log_lines, i, job.flow_step_id.clone()) - } - } - err @ _ => to_raw_value( - &SerializedError { - message: format!("error during execution of the script:\n{}", err), - name: "ExecutionErr".to_string(), - step_id: job.flow_step_id.clone(), - }, - ), - }; - - // in the happy path and if job not a flow step, we can delegate updating the completed job in the background - job_completed_tx - .send(JobCompleted { - job: job, - result: Arc::new(to_raw_value(&error_value)), - mem_peak, - canceled_by, - success: false, - cached_res_path, - token: token, - }) - .await - .expect("send job completed"); - } - }; - Ok(()) -} pub fn build_envs( envs: Option>, @@ -2733,8 +2382,8 @@ pub async fn get_hub_script_content_and_requirements( .clone() .ok_or_else(|| Error::InternalErr(format!("expected script path for hub script")))?; - let script = get_full_hub_script_by_path(StripPath(script_path.to_string()), &HTTP_CLIENT, db) - .await?; + let script = + get_full_hub_script_by_path(StripPath(script_path.to_string()), &HTTP_CLIENT, db).await?; Ok(ContentReqLangEnvs { content: script.content, lockfile: script.lockfile, @@ -2828,7 +2477,7 @@ async fn handle_code_execution_job( Some(PREVIEW_IS_TAR_CODEBASE_HASH) => Some(format!("{}.tar", job.id)), _ => None, }; - + ContentReqLangEnvs { content: job .raw_code @@ -2837,8 +2486,9 @@ async fn handle_code_execution_job( lockfile: job.raw_lock.clone(), language: job.language.to_owned(), envs: None, - codebase - }}, + codebase, + } + } JobKind::Script_Hub => { get_hub_script_content_and_requirements(job.script_path.clone(), Some(db)).await? } @@ -3146,7 +2796,7 @@ mount {{ &shared_mount, ) .await - }, + } Some(ScriptLang::Rust) => { handle_rust_job( mem_peak, @@ -3162,7 +2812,7 @@ mount {{ worker_name, envs, ) - .await + .await } _ => panic!("unreachable, language is not supported: {language:#?}"), }; diff --git a/backend/windmill-worker/src/worker_flow.rs b/backend/windmill-worker/src/worker_flow.rs index 10edf6f86a..35732bda5c 100644 --- a/backend/windmill-worker/src/worker_flow.rs +++ b/backend/windmill-worker/src/worker_flow.rs @@ -15,7 +15,10 @@ use std::time::Duration; use crate::common::{hash_args, save_in_cache}; use crate::js_eval::{eval_timeout, IdContext}; -use crate::{AuthedClient, PreviousResult, SameWorkerPayload, SendResult, JOB_TOKEN, KEEP_JOB_DIR}; +use crate::{ + AuthedClient, PreviousResult, SameWorkerPayload, SameWorkerSender, SendResult, JOB_TOKEN, + KEEP_JOB_DIR, +}; use anyhow::Context; use mappable_rc::Marc; use serde::{Deserialize, Serialize}; @@ -67,7 +70,7 @@ pub async fn update_flow_status_after_job_completion< success: bool, result: Arc>, unrecoverable: bool, - same_worker_tx: Sender, + same_worker_tx: SameWorkerSender, worker_dir: &str, stop_early_override: Option, rsmq: Option, @@ -183,7 +186,7 @@ pub async fn update_flow_status_after_job_completion_internal< mut success: bool, result: Arc>, unrecoverable: bool, - same_worker_tx: Sender, + same_worker_tx: SameWorkerSender, worker_dir: &str, stop_early_override: Option, skip_error_handler: bool, @@ -1404,7 +1407,7 @@ pub async fn handle_flow( db: &sqlx::Pool, client: &AuthedClient, last_result: Option>>, - same_worker_tx: Sender, + same_worker_tx: SameWorkerSender, worker_dir: &str, rsmq: Option, job_completed_tx: Sender, @@ -1525,7 +1528,7 @@ async fn push_next_flow_job db: &sqlx::Pool, client: &AuthedClient, last_job_result: Option>>, - same_worker_tx: Sender, + same_worker_tx: SameWorkerSender, worker_dir: &str, rsmq: Option, job_completed_tx: Sender, diff --git a/frontend/src/lib/components/graph/renderers/edges/BaseEdge.svelte b/frontend/src/lib/components/graph/renderers/edges/BaseEdge.svelte index b29e481fb9..f27e89b8e0 100644 --- a/frontend/src/lib/components/graph/renderers/edges/BaseEdge.svelte +++ b/frontend/src/lib/components/graph/renderers/edges/BaseEdge.svelte @@ -55,7 +55,7 @@ {#if data?.insertable && !$useDataflow && !data?.moving}