From 1cac7bf2391500556a34f6defca44ff669a41dae Mon Sep 17 00:00:00 2001 From: Ruben Fiszel Date: Tue, 30 Apr 2024 17:18:54 +0200 Subject: [PATCH] more logs for dedicated workers --- backend/windmill-worker/src/bun_executor.rs | 1 + backend/windmill-worker/src/dedicated_worker.rs | 5 ++++- backend/windmill-worker/src/deno_executor.rs | 1 + backend/windmill-worker/src/python_executor.rs | 1 + backend/windmill-worker/src/worker.rs | 3 +++ 5 files changed, 10 insertions(+), 1 deletion(-) diff --git a/backend/windmill-worker/src/bun_executor.rs b/backend/windmill-worker/src/bun_executor.rs index 526f503abc..34b053dbb2 100644 --- a/backend/windmill-worker/src/bun_executor.rs +++ b/backend/windmill-worker/src/bun_executor.rs @@ -983,6 +983,7 @@ plugin(p) jobs_rx, worker_name, db, + script_path, ) .await } diff --git a/backend/windmill-worker/src/dedicated_worker.rs b/backend/windmill-worker/src/dedicated_worker.rs index 75e123d7d3..669c65f820 100644 --- a/backend/windmill-worker/src/dedicated_worker.rs +++ b/backend/windmill-worker/src/dedicated_worker.rs @@ -58,6 +58,7 @@ pub async fn handle_dedicated_process( mut jobs_rx: Receiver>, worker_name: &str, db: &DB, + script_path: &str, ) -> std::result::Result<(), error::Error> { //do not cache local dependencies let mut child = { @@ -147,6 +148,7 @@ pub async fn handle_dedicated_process( tracing::debug!("processed job: {line}"); if line.starts_with("wm_res[") { let job: Arc = jobs.pop_front().expect("pop"); + tracing::info!("job completed on dedicated worker {script_path}: {}", job.id); match serde_json::from_str::>(&line.replace("wm_res[success]:", "").replace("wm_res[error]:", "")) { Ok(result) => { append_logs(job.id, job.workspace_id.clone(), logs.clone(), db).await; @@ -174,8 +176,9 @@ pub async fn handle_dedicated_process( job = conditional_polling(jobs_rx.recv(), alive && jobs.len() < MAX_BUFFERED_DEDICATED_JOBS) => { // i += 1; if let Some(job) = job { - tracing::debug!("received job"); jobs.push_back(job.clone()); + tracing::info!("received job and adding to queue on dedicated worker for {script_path}: {} (queue_size: {})", job.id, jobs.len()); + // write_stdin(&mut stdin, &serde_json::to_string(&job.args.unwrap_or_else(|| serde_json::json!({"x": job.id}))).expect("serialize")).await?; write_stdin(&mut stdin, &serde_json::to_string(&job.args).expect("serialize")).await?; stdin.flush().await.context("stdin flush")?; diff --git a/backend/windmill-worker/src/deno_executor.rs b/backend/windmill-worker/src/deno_executor.rs index d83632bfb0..b6e8b958f2 100644 --- a/backend/windmill-worker/src/deno_executor.rs +++ b/backend/windmill-worker/src/deno_executor.rs @@ -527,6 +527,7 @@ for await (const chunk of Deno.stdin.readable) {{ jobs_rx, worker_name, db, + script_path, ) .await } diff --git a/backend/windmill-worker/src/python_executor.rs b/backend/windmill-worker/src/python_executor.rs index 3e97d9a007..3f8abaf890 100644 --- a/backend/windmill-worker/src/python_executor.rs +++ b/backend/windmill-worker/src/python_executor.rs @@ -1193,6 +1193,7 @@ for line in sys.stdin: jobs_rx, worker_name, db, + script_path, ) .await } diff --git a/backend/windmill-worker/src/worker.rs b/backend/windmill-worker/src/worker.rs index d4801cd617..95dcaf41fa 100644 --- a/backend/windmill-worker/src/worker.rs +++ b/backend/windmill-worker/src/worker.rs @@ -2155,6 +2155,9 @@ async fn spawn_dedicated_worker( } { tracing::error!("error in dedicated worker: {:?}", e); }; + if let Err(e) = killpill_tx.clone().send(()) { + tracing::error!("failed to send final killpill to dedicated worker: {:?}", e); + } }); return Some((node_id.unwrap_or(path2), dedicated_worker_tx, Some(handle))); // (Some(dedi_path), Some(dedicated_worker_tx), Some(handle))