From 25084170d7449dfa06dd758bc3783580916cfa95 Mon Sep 17 00:00:00 2001 From: Ruben Fiszel Date: Sat, 1 Aug 2026 13:50:47 +0200 Subject: [PATCH] feat: make job subprocess oom_score_adj configurable (#10443) * feat: make job subprocess oom_score_adj configurable via JOB_OOM_SCORE_ADJ Co-Authored-By: Claude Opus 5 (1M context) * fix: warn when JOB_OOM_SCORE_ADJ leaves no gap over the worker's own score Co-Authored-By: Claude Opus 5 (1M context) * style: drop em dash from JOB_OOM_SCORE_ADJ doc comment Co-Authored-By: Claude Opus 5 (1M context) * fix: warn on any oom_score_adj gap too small to steer the OOM killer Co-Authored-By: Claude Opus 5 (1M context) --------- Co-authored-by: Claude Opus 5 (1M context) --- backend/src/main.rs | 76 +++++++++++-------- .../windmill-common/src/global_settings.rs | 1 + backend/windmill-common/src/worker.rs | 56 ++++++++++++++ backend/windmill-worker/src/handle_child.rs | 8 +- .../windmill-worker/src/python_executor.rs | 6 +- 5 files changed, 114 insertions(+), 33 deletions(-) diff --git a/backend/src/main.rs b/backend/src/main.rs index 969fd4e94d..51febad6f0 100644 --- a/backend/src/main.rs +++ b/backend/src/main.rs @@ -932,46 +932,62 @@ async fn windmill_main() -> anyhow::Result<()> { } // Lower the worker's oom_score_adj so the OOM killer strongly prefers killing - // job subprocesses (oom_score_adj=1000) over the worker itself. + // job subprocesses (oom_score_adj=JOB_OOM_SCORE_ADJ) over the worker itself. // Kubernetes sets it high for burstable QoS (e.g. 937), leaving a tiny gap vs jobs. // Requires CAP_SYS_RESOURCE to lower it; if missing, we just warn. #[cfg(any(target_os = "linux"))] - match std::fs::read_to_string("/proc/self/oom_score_adj") { - Ok(current) => { - let current = current.trim().to_string(); - let current_val = match current.parse::() { - Ok(v) => v, - Err(e) => { - tracing::warn!("Could not parse oom_score_adj '{current}': {e}"); - 0 - } - }; - if current_val > 0 { - match std::fs::write("/proc/self/oom_score_adj", "0") { - Ok(_) => { - tracing::info!( - "Lowered worker oom_score_adj from {current} to 0 \ - (jobs get 1000, gap=1000)" - ); + { + // Badness is (memory used, in permille of host RAM) + oom_score_adj, so the gap + // must exceed the worker's own footprint in permille to actually steer the kill. + // 100 covers a worker holding up to ~10% of host RAM. + const MIN_OOM_SCORE_GAP: i32 = 100; + + let job_adj = *windmill_common::worker::JOB_OOM_SCORE_ADJ; + match std::fs::read_to_string("/proc/self/oom_score_adj") { + Ok(current) => { + let current = current.trim().to_string(); + match current.parse::() { + Ok(mut worker_adj) => { + if worker_adj > 0 { + match std::fs::write("/proc/self/oom_score_adj", "0") { + Ok(_) => { + tracing::info!( + "Lowered worker oom_score_adj from {worker_adj} to 0" + ); + worker_adj = 0; + } + Err(e) => { + tracing::warn!( + "Could not lower worker oom_score_adj from {worker_adj} to 0: {e}. \ + Add CAP_SYS_RESOURCE to the container to fix this" + ); + } + } + } + let gap = job_adj - worker_adj; + if gap >= MIN_OOM_SCORE_GAP { + tracing::info!( + "Worker oom_score_adj={worker_adj}, jobs get {job_adj} (gap={gap})" + ); + } else { + tracing::warn!( + "Worker oom_score_adj={worker_adj}, jobs get {job_adj} (gap={gap}): \ + too small to reliably steer the OOM killer to the job. \ + Raise JOB_OOM_SCORE_ADJ or lower the worker's own score" + ); + } } Err(e) => { tracing::warn!( - "Could not lower worker oom_score_adj from {current} to 0: {e}. \ - Gap to jobs is only {} — OOM killer may target the worker instead. \ - Add CAP_SYS_RESOURCE to the container to fix this", - 1000 - current_val + "Could not parse worker oom_score_adj '{current}': {e}. \ + Cannot tell whether jobs (oom_score_adj={job_adj}) outrank the worker" ); } } - } else { - tracing::info!( - "Worker oom_score_adj={current} (jobs get 1000, gap={})", - 1000 - current_val - ); } - } - Err(e) => { - tracing::warn!("Could not read worker oom_score_adj: {e}"); + Err(e) => { + tracing::warn!("Could not read worker oom_score_adj: {e}"); + } } } } diff --git a/backend/windmill-common/src/global_settings.rs b/backend/windmill-common/src/global_settings.rs index 427077fdb9..6fb435cd71 100644 --- a/backend/windmill-common/src/global_settings.rs +++ b/backend/windmill-common/src/global_settings.rs @@ -298,6 +298,7 @@ pub const ENV_SETTINGS: &[&str] = &[ "GLOBAL_ERROR_HANDLER_PATH_IN_ADMINS_WORKSPACE", "MAX_WAIT_FOR_SIGINT", "MAX_WAIT_FOR_SIGTERM", + "JOB_OOM_SCORE_ADJ", "WORKER_GROUP", "SAML_METADATA", "INSTANCE_IS_DEV", diff --git a/backend/windmill-common/src/worker.rs b/backend/windmill-common/src/worker.rs index b9dd29c00a..226d74cf47 100644 --- a/backend/windmill-common/src/worker.rs +++ b/backend/windmill-common/src/worker.rs @@ -217,6 +217,33 @@ pub const CONCURRENCY_KEY_MAX_QUEUED_DEFAULT: u32 = 10_000; /// the setting is cleared or malformed. A workspace spans many keys, so this sits well above /// the per-key cap. pub const WORKSPACE_MAX_QUEUED_JOBS_DEFAULT: u32 = 20_000; +/// Default for [`JOB_OOM_SCORE_ADJ`]; also the value used when the env var is out of range or +/// unparseable. +pub const JOB_OOM_SCORE_ADJ_DEFAULT: i32 = 1000; + +/// procfs accepts -1000..=1000, but a job must never be *less* killable than the worker that +/// supervises it, so negative adjustments are rejected rather than clamped. +fn parse_job_oom_score_adj(raw: Option<&str>) -> i32 { + let Some(raw) = raw else { + return JOB_OOM_SCORE_ADJ_DEFAULT; + }; + match raw.trim().parse::() { + Ok(v) if (0..=1000).contains(&v) => v, + Ok(v) => { + tracing::warn!( + "JOB_OOM_SCORE_ADJ={v} is outside the accepted 0..=1000 range, \ + using {JOB_OOM_SCORE_ADJ_DEFAULT}" + ); + JOB_OOM_SCORE_ADJ_DEFAULT + } + Err(e) => { + tracing::warn!( + "Could not parse JOB_OOM_SCORE_ADJ='{raw}': {e}, using {JOB_OOM_SCORE_ADJ_DEFAULT}" + ); + JOB_OOM_SCORE_ADJ_DEFAULT + } + } +} lazy_static::lazy_static! { pub static ref WORKER_GROUP: String = std::env::var("WORKER_GROUP").unwrap_or_else(|_| { #[cfg(not(feature = "enterprise"))] @@ -240,6 +267,15 @@ lazy_static::lazy_static! { pub static ref LIMIT_WINDOWS_TO_1CU: bool = std::env::var("LIMIT_WINDOWS_TO_1CU").ok().is_some_and(|x| x == "1" || x == "true"); + /// `oom_score_adj` applied to job subprocesses. The kernel adds it to the process's memory + /// use expressed in permille of host RAM, so the job only reliably outranks the worker once + /// the gap between their two adjustments exceeds the worker's own footprint in permille; the + /// default maximizes that margin. Userspace OOM daemons (earlyoom, systemd-oomd, nohang) rank + /// every process on the host by the same score, so at 1000 a tiny job outranks multi-GB + /// processes and gets killed first. Lowering this trades margin over the worker for a fairer + /// ranking against everything else on the host. + pub static ref JOB_OOM_SCORE_ADJ: i32 = parse_job_oom_score_adj(std::env::var("JOB_OOM_SCORE_ADJ").ok().as_deref()); + pub static ref CGROUP_V2_PATH_RE: Regex = Regex::new(r#"(?m)^0::(/.*)$"#).unwrap(); pub static ref CGROUP_V2_CPU_RE: Regex = Regex::new(r#"(?m)^(\d+) \S+$"#).unwrap(); pub static ref CGROUP_V1_INACTIVE_FILE_RE: Regex = Regex::new(r#"(?m)^total_inactive_file (\d+)$"#).unwrap(); @@ -2441,6 +2477,26 @@ mod tests { ids.iter().map(|s| s.to_string()).collect() } + #[test] + fn test_parse_job_oom_score_adj() { + assert_eq!(parse_job_oom_score_adj(Some("300")), 300); + assert_eq!(parse_job_oom_score_adj(Some(" 0\n")), 0); + assert_eq!(parse_job_oom_score_adj(None), JOB_OOM_SCORE_ADJ_DEFAULT); + // Out of range and unparseable both fall back rather than weaken the worker's protection. + assert_eq!( + parse_job_oom_score_adj(Some("-500")), + JOB_OOM_SCORE_ADJ_DEFAULT + ); + assert_eq!( + parse_job_oom_score_adj(Some("1001")), + JOB_OOM_SCORE_ADJ_DEFAULT + ); + assert_eq!( + parse_job_oom_score_adj(Some("high")), + JOB_OOM_SCORE_ADJ_DEFAULT + ); + } + #[test] fn test_bash_sandbox_image_annotation() { // `# sandbox ` selects the container runtime and returns the image. diff --git a/backend/windmill-worker/src/handle_child.rs b/backend/windmill-worker/src/handle_child.rs index b7ac2d54e2..439a4a09a2 100644 --- a/backend/windmill-worker/src/handle_child.rs +++ b/backend/windmill-worker/src/handle_child.rs @@ -125,13 +125,17 @@ pub async fn handle_child( let pid = child.id(); #[cfg(target_os = "linux")] if let Some(pid) = pid { + let oom_score_adj = *windmill_common::worker::JOB_OOM_SCORE_ADJ; // procfs handles writes synchronously in-kernel; no fsync (it returns // EINVAL on procfs files). - match std::fs::write(format!("/proc/{pid}/oom_score_adj"), b"1000") { + match std::fs::write( + format!("/proc/{pid}/oom_score_adj"), + oom_score_adj.to_string(), + ) { Ok(()) => {} Err(e) => { tracing::error!( - "Failed to set oom_score_adj=1000 for pid {pid}: {e:#}. \ + "Failed to set oom_score_adj={oom_score_adj} for pid {pid}: {e:#}. \ OOM killer may target the worker instead of this job" ); } diff --git a/backend/windmill-worker/src/python_executor.rs b/backend/windmill-worker/src/python_executor.rs index 13e3749fbd..397367742a 100644 --- a/backend/windmill-worker/src/python_executor.rs +++ b/backend/windmill-worker/src/python_executor.rs @@ -3026,7 +3026,11 @@ pub async fn handle_python_reqs( "failed to get PID for python installation process: {}", &req ))) - .and_then(|pid| write_file(&format!("/proc/{pid}"), "oom_score_adj", "1000")) + .and_then(|pid| write_file( + &format!("/proc/{pid}"), + "oom_score_adj", + &windmill_common::worker::JOB_OOM_SCORE_ADJ.to_string(), + )) { tracing::error!( req = %req,