fix: graceful worker exits for same worker jobs (#4371)

* all

* fix compile
This commit is contained in:
Ruben Fiszel
2024-09-11 18:32:08 +02:00
committed by GitHub
parent dd6643aab7
commit 7d29651d3e
9 changed files with 662 additions and 555 deletions
+18 -7
View File
@@ -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<()>
};
+8 -3
View File
@@ -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<R: rsmq_async::RsmqConnection + Send + Sync + Clone>
// 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::<SameWorkerPayload>(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::<SendResult>(1);
let label = if job.permissioned_as != format!("u/{}", job.created_by)
+1 -1
View File
@@ -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?;
+2 -2
View File
@@ -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
{
+3
View File
@@ -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;
@@ -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<QueuedJob>,
result: Arc<Box<RawValue>>,
mem_peak: i32,
canceled_by: Option<CanceledBy>,
success: bool,
cached_res_path: Option<String>,
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<QueuedJob>,
result: error::Result<Arc<Box<RawValue>>>,
job_dir: &str,
job_completed_tx: JobCompletedSender,
mem_peak: i32,
canceled_by: Option<CanceledBy>,
cached_res_path: Option<String>,
token: String,
column_order: Option<Vec<String>>,
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<R>,
worker_name: &str,
worker_save_completed_job_duration: Option<Histo>,
worker_flow_transition_duration: Option<Histo>,
job_completed_tx: Sender<SendResult>,
) {
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<R: rsmq_async::RsmqConnection + Send + Sync + Clone>(
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<R>,
worker_name: &str,
_worker_save_completed_job_duration: Option<Histo>,
_worker_flow_transition_duration: Option<Histo>,
job_completed_tx: Sender<SendResult>,
) -> 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<R: rsmq_async::RsmqConnection + Send + Sync + Clone>(
db: &Pool<Postgres>,
client: &AuthedClient,
job: &QueuedJob,
mem_peak: i32,
canceled_by: Option<CanceledBy>,
err: Error,
unrecoverable: bool,
same_worker_tx: SameWorkerSender,
worker_dir: &str,
rsmq: Option<R>,
worker_name: &str,
job_completed_tx: Sender<SendResult>,
) {
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<String>,
}
pub fn extract_error_value(log_lines: &str, i: i32, step_id: Option<String>) -> Box<RawValue> {
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,
});
}
+185 -535
View File
@@ -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<prometheus::Histogram>;
pub type Histo = Arc<prometheus::Histogram>;
#[cfg(feature = "prometheus")]
type GGauge = Arc<GenericGauge<AtomicI64>>;
#[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<Postgres>,
worker_dir: &str,
same_worker_tx: &Sender<SameWorkerPayload>,
rsmq: Option<R>,
worker_name: &str,
worker_save_completed_job_duration: Option<Histo>,
worker_flow_transition_duration: Option<Histo>,
job_completed_tx: Sender<SendResult>,
) {
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<SendResult>, Option<GGauge>, Option<Histo>);
#[derive(Clone)]
pub struct SameWorkerSender(pub Sender<SameWorkerPayload>, pub Arc<AtomicU16>);
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<SameWorkerPayload>> {
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<R: rsmq_async::RsmqConnection + Send + Sync + Clone + 's
if let Some(ref netrc) = *NETRC {
tracing::info!("Writing netrc at {}/.netrc", HOME_ENV.as_str());
write_file(&HOME_ENV, ".netrc", netrc)
.expect("could not write netrc");
write_file(&HOME_ENV, ".netrc", netrc).expect("could not write netrc");
}
DirBuilder::new()
@@ -1097,12 +1049,14 @@ pub async fn run_worker<R: rsmq_async::RsmqConnection + Send + Sync + Clone + 's
worker_completed_channel_queue_send_duration,
);
let same_worker_queue_size = Arc::new(AtomicU16::new(0));
let same_worker_tx = SameWorkerSender(same_worker_tx, same_worker_queue_size.clone());
let db2 = db.clone();
let base_internal_url2 = base_internal_url.to_string();
let same_worker_tx2 = same_worker_tx.clone();
let rsmq2 = rsmq.clone();
let worker_dir2 = worker_dir.clone();
let thread_count = Arc::new(AtomicUsize::new(0));
let is_dedicated_worker = WORKER_CONFIG.read().await.dedicated_worker.is_some();
@@ -1186,98 +1140,107 @@ pub async fn run_worker<R: rsmq_async::RsmqConnection + Send + Sync + Clone + 's
let worker_name2 = worker_name.clone();
let killpill_tx2 = killpill_tx.clone();
let job_completed_sender = job_completed_tx.0.clone();
let send_result = tokio::spawn((async move {
while let Some(sr) = 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();
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<Instant> = None;
let mut last_checked_suspended = Instant::now();
@@ -1358,6 +1321,7 @@ pub async fn run_worker<R: rsmq_async::RsmqConnection + Send + Sync + Clone + 's
let mut suspend_first_success = false;
let mut last_reading = Instant::now() - Duration::from_secs(NUM_SECS_READINGS + 1);
let mut killed_but_draining_same_worker_jobs = false;
loop {
#[cfg(feature = "benchmark")]
let loop_start = Instant::now();
@@ -1386,7 +1350,9 @@ pub async fn run_worker<R: rsmq_async::RsmqConnection + Send + Sync + Clone + 's
let memory_usage = get_worker_memory_usage();
let wm_memory_usage = get_windmill_memory_usage();
let (vcpus, memory) = if *REFRESH_CGROUP_READINGS && last_reading.elapsed().as_secs() > 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<R: rsmq_async::RsmqConnection + Send + Sync + Clone + 's
started = true
}
if let Ok(_) = killpill_rx.try_recv() {
println!("received killpill for worker {}", i_worker);
job_completed_tx.0.send(SendResult::Kill).await.unwrap();
break
}
if let Ok(same_worker_job) = same_worker_rx.try_recv() {
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.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<R: rsmq_async::RsmqConnection + Send + Sync + Clone + 's
} else if let Some(wp) = worker_pull_over_500_counter.as_ref() {
wp.inc();
}
} else if !agent_mode && duration_pull_s > 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<R: rsmq_async::RsmqConnection + Send + Sync + Clone + 's
async fn queue_init_bash_maybe<'c, R: rsmq_async::RsmqConnection + Send + 'c>(
db: &Pool<Postgres>,
same_worker_tx: Sender<SameWorkerPayload>,
same_worker_tx: SameWorkerSender,
worker_name: &str,
rsmq: Option<R>,
) -> 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<R: rsmq_async::RsmqConnection + Send + Sync + Clone>(
JobCompleted { job, result, mem_peak, success, cached_res_path, canceled_by, .. }: JobCompleted,
client: &AuthedClient,
db: &DB,
worker_dir: &str,
same_worker_tx: Sender<SameWorkerPayload>,
rsmq: Option<R>,
worker_name: &str,
_worker_save_completed_job_duration: Option<Histo>,
_worker_flow_transition_duration: Option<Histo>,
job_completed_tx: Sender<SendResult>,
) -> 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<ScriptLang>,
@@ -2062,136 +1935,6 @@ pub async fn process_completed_job<R: rsmq_async::RsmqConnection + Send + Sync +
// tracing::debug!("leader worker done waiting for");
// }
#[tracing::instrument(name = "job_error", level = "info", skip_all, fields(job_id = %job.id))]
pub async fn handle_job_error<R: rsmq_async::RsmqConnection + Send + Sync + Clone>(
db: &Pool<Postgres>,
client: &AuthedClient,
job: &QueuedJob,
mem_peak: i32,
canceled_by: Option<CanceledBy>,
err: Error,
unrecoverable: bool,
same_worker_tx: Sender<SameWorkerPayload>,
worker_dir: &str,
rsmq: Option<R>,
worker_name: &str,
job_completed_tx: Sender<SendResult>,
) {
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<String>,
}
fn extract_error_value(log_lines: &str, i: i32, step_id: Option<String>) -> Box<RawValue> {
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<R: rsmq_async::RsmqConnection + Send + Sync + Clone>(
worker_name: &str,
worker_dir: &str,
job_dir: &str,
same_worker_tx: Sender<SameWorkerPayload>,
same_worker_tx: SameWorkerSender,
base_internal_url: &str,
rsmq: Option<R>,
job_completed_tx: JobCompletedSender,
@@ -2591,100 +2334,6 @@ async fn handle_queued_job<R: rsmq_async::RsmqConnection + Send + Sync + Clone>(
Ok(())
}
async fn process_result(
job: Arc<QueuedJob>,
result: error::Result<Arc<Box<RawValue>>>,
job_dir: &str,
job_completed_tx: JobCompletedSender,
mem_peak: i32,
canceled_by: Option<CanceledBy>,
cached_res_path: Option<String>,
token: String,
column_order: Option<Vec<String>>,
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<Vec<String>>,
@@ -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:#?}"),
};
+8 -5
View File
@@ -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<Box<RawValue>>,
unrecoverable: bool,
same_worker_tx: Sender<SameWorkerPayload>,
same_worker_tx: SameWorkerSender,
worker_dir: &str,
stop_early_override: Option<bool>,
rsmq: Option<R>,
@@ -183,7 +186,7 @@ pub async fn update_flow_status_after_job_completion_internal<
mut success: bool,
result: Arc<Box<RawValue>>,
unrecoverable: bool,
same_worker_tx: Sender<SameWorkerPayload>,
same_worker_tx: SameWorkerSender,
worker_dir: &str,
stop_early_override: Option<bool>,
skip_error_handler: bool,
@@ -1404,7 +1407,7 @@ pub async fn handle_flow<R: rsmq_async::RsmqConnection + Send + Sync + Clone>(
db: &sqlx::Pool<sqlx::Postgres>,
client: &AuthedClient,
last_result: Option<Arc<Box<RawValue>>>,
same_worker_tx: Sender<SameWorkerPayload>,
same_worker_tx: SameWorkerSender,
worker_dir: &str,
rsmq: Option<R>,
job_completed_tx: Sender<SendResult>,
@@ -1525,7 +1528,7 @@ async fn push_next_flow_job<R: rsmq_async::RsmqConnection + Send + Sync + Clone>
db: &sqlx::Pool<sqlx::Postgres>,
client: &AuthedClient,
last_job_result: Option<Arc<Box<RawValue>>>,
same_worker_tx: Sender<SameWorkerPayload>,
same_worker_tx: SameWorkerSender,
worker_dir: &str,
rsmq: Option<R>,
job_completed_tx: Sender<SendResult>,
@@ -55,7 +55,7 @@
{#if data?.insertable && !$useDataflow && !data?.moving}
<div
class={twMerge('edgeButtonContainer nodrag nopan top-0', menuOpen ? 'z-50' : '')}
style:transform="translate(-50%, 50%) translate({sourceX}px,{sourceY}px)"
style:transform="translate(-50%, 50%) translate({sourceX}px,{sourceY + 2}px)"
>
<InsertModuleButton
disableAi={data.disableAi}
@@ -79,7 +79,7 @@
{#if data.enableTrigger}
<div
class="edgeButtonContainer nodrag nopan"
style:transform="translate(100%, 50%) translate({sourceX}px,{sourceY}px)"
style:transform="translate(100%, 50%) translate({sourceX}px,{sourceY + 2}px)"
>
<InsertTriggerButton
disableAi={data.disableAi}