From 5f85b67dfcf063fd8a3c3f69f0e7605fc40e473d Mon Sep 17 00:00:00 2001 From: Guillaume Bouvignies Date: Thu, 21 Dec 2023 10:29:26 +0100 Subject: [PATCH] fix: Failing jobs in dedicated worker mode are now marked as failing (#2894) * fix: Failing jobs in dedicated worker mode are now marked as failing * remove log line --------- Co-authored-by: Ruben Fiszel --- backend/windmill-worker/src/bun_executor.rs | 4 ++-- backend/windmill-worker/src/dedicated_worker.rs | 15 ++++++++++----- backend/windmill-worker/src/deno_executor.rs | 4 ++-- backend/windmill-worker/src/python_executor.rs | 6 +++--- 4 files changed, 17 insertions(+), 12 deletions(-) diff --git a/backend/windmill-worker/src/bun_executor.rs b/backend/windmill-worker/src/bun_executor.rs index 22e762459f..61e3776731 100644 --- a/backend/windmill-worker/src/bun_executor.rs +++ b/backend/windmill-worker/src/bun_executor.rs @@ -680,9 +680,9 @@ for await (const chunk of Bun.stdin.stream()) {{ try {{ let {{ {spread} }} = JSON.parse(line) let res: any = await main(...[ {spread} ]); - stdout.write("wm_res:" + JSON.stringify(res ?? null, (key, value) => typeof value === 'undefined' ? null : value) + '\n'); + stdout.write("wm_res[success]:" + JSON.stringify(res ?? null, (key, value) => typeof value === 'undefined' ? null : value) + '\n'); }} catch (e) {{ - stdout.write("wm_res:" + JSON.stringify({{ error: {{ message: e.message, name: e.name, stack: e.stack, line: line }}}}) + '\n'); + stdout.write("wm_res[error]:" + JSON.stringify({{ message: e.message, name: e.name, stack: e.stack, line: line }}) + '\n'); }} stdout.flush(); }} diff --git a/backend/windmill-worker/src/dedicated_worker.rs b/backend/windmill-worker/src/dedicated_worker.rs index 738d21a230..62e6a6a24f 100644 --- a/backend/windmill-worker/src/dedicated_worker.rs +++ b/backend/windmill-worker/src/dedicated_worker.rs @@ -104,8 +104,7 @@ pub async fn handle_dedicated_process( .wait() .await .expect("child process encountered an error"); - - println!("child status was: {}", status); + tracing::info!("child status was: {}", status); }); let mut jobs = VecDeque::with_capacity(MAX_BUFFERED_DEDICATED_JOBS); @@ -144,10 +143,16 @@ pub async fn handle_dedicated_process( continue; } tracing::debug!("processed job: {line}"); - if line.starts_with("wm_res:") { + if line.starts_with("wm_res[") { let job: Arc = jobs.pop_front().expect("pop"); - match serde_json::from_str::>(&line.replace("wm_res:", "")) { - Ok(result) => job_completed_tx.send(JobCompleted { job , result, logs: logs, mem_peak: 0, canceled_by: None, success: true, cached_res_path: None, token: token.to_string() }).await.unwrap(), + match serde_json::from_str::>(&line.replace("wm_res[success]:", "").replace("wm_res[error]:", "")) { + Ok(result) => { + if line.starts_with("wm_res[success]:") { + job_completed_tx.send(JobCompleted { job , result, logs: logs, mem_peak: 0, canceled_by: None, success: true, cached_res_path: None, token: token.to_string() }).await.unwrap() + } else { + job_completed_tx.send(JobCompleted { job , result, logs: logs, mem_peak: 0, canceled_by: None, success: false, cached_res_path: None, token: token.to_string() }).await.unwrap() + } + }, Err(e) => { tracing::error!("Could not deserialize job result `{line}`: {e:?}"); job_completed_tx.send(JobCompleted { job , result: to_raw_value(&serde_json::json!({"error": format!("Could not deserialize job result `{line}`: {e:?}")})), logs: "".to_string(), mem_peak: 0, canceled_by: None, success: false, cached_res_path: None, token: token.to_string() }).await.unwrap(); diff --git a/backend/windmill-worker/src/deno_executor.rs b/backend/windmill-worker/src/deno_executor.rs index fd0ae1b3b5..746a60db3b 100644 --- a/backend/windmill-worker/src/deno_executor.rs +++ b/backend/windmill-worker/src/deno_executor.rs @@ -463,9 +463,9 @@ for await (const chunk of Deno.stdin.readable) {{ try {{ let {{ {spread} }} = JSON.parse(line) let res: any = await main(...[ {spread} ]); - console.log("wm_res:" + JSON.stringify(res ?? null, (key, value) => typeof value === 'undefined' ? null : value) + '\n'); + console.log("wm_res[success]:" + JSON.stringify(res ?? null, (key, value) => typeof value === 'undefined' ? null : value) + '\n'); }} catch (e) {{ - console.log("wm_res:" + JSON.stringify({{ error: {{ message: e.message, name: e.name, stack: e.stack, line: line }}}}) + '\n'); + console.log("wm_res[error]:" + JSON.stringify({{ message: e.message, name: e.name, stack: e.stack, line: line }}) + '\n'); }} }} if (exit) {{ diff --git a/backend/windmill-worker/src/python_executor.rs b/backend/windmill-worker/src/python_executor.rs index 19d686055b..66170aa5ab 100644 --- a/backend/windmill-worker/src/python_executor.rs +++ b/backend/windmill-worker/src/python_executor.rs @@ -925,12 +925,12 @@ for line in sys.stdin: if type(v).__name__ == 'bytes': res[k] = to_b_64(v) res_json = re.sub(replace_nan, ' null ', json.dumps(res, separators=(',', ':'), default=str).replace('\n', '')) - sys.stdout.write("wm_res:" + res_json + "\n") + sys.stdout.write("wm_res[success]:" + res_json + "\n") except BaseException as e: exc_type, exc_value, exc_traceback = sys.exc_info() tb = traceback.format_tb(exc_traceback) - err_json = json.dumps({{ "error": {{ "message": str(e), "name": e.__class__.__name__, "stack": '\n'.join(tb[1:]) }} }}, separators=(',', ':'), default=str).replace('\n', '') - sys.stdout.write("wm_res:" + err_json + "\n") + err_json = json.dumps({{ "message": str(e), "name": e.__class__.__name__, "stack": '\n'.join(tb[1:]) }}, separators=(',', ':'), default=str).replace('\n', '') + sys.stdout.write("wm_res[error]:" + err_json + "\n") sys.stdout.flush() "#, );