From 7c699598533dade9713d976d8dd90fc657ebb503 Mon Sep 17 00:00:00 2001 From: Ruben Fiszel Date: Sun, 11 May 2025 08:38:46 +0200 Subject: [PATCH] fix: improve agents workers handling of WHITELIST_ENVS --- backend/src/monitor.rs | 33 +++++--- backend/windmill-common/src/worker.rs | 106 +++++++++++++++----------- backend/windmill-worker/src/worker.rs | 28 ++----- 3 files changed, 91 insertions(+), 76 deletions(-) diff --git a/backend/src/monitor.rs b/backend/src/monitor.rs index 04359d3c4a..b7ae065f43 100644 --- a/backend/src/monitor.rs +++ b/backend/src/monitor.rs @@ -28,10 +28,10 @@ use windmill_api::{ SCIM_TOKEN, }; -#[cfg(feature = "enterprise")] -use windmill_common::ee::{jobs_waiting_alerts, worker_groups_alerts}; #[cfg(feature = "enterprise")] use windmill_common::ee::low_disk_alerts; +#[cfg(feature = "enterprise")] +use windmill_common::ee::{jobs_waiting_alerts, worker_groups_alerts}; #[cfg(feature = "oauth2")] use windmill_common::global_settings::OAUTH_SETTING; @@ -60,13 +60,9 @@ use windmill_common::{ server::load_smtp_config, tracing_init::JSON_FMT, users::truncate_token, - utils::empty_string_as_none, - utils::{now_from_db, rd_string, report_critical_error, Mode}, + utils::{empty_string_as_none, now_from_db, rd_string, report_critical_error, Mode}, worker::{ - load_worker_config, reload_custom_tags_setting, store_pull_query, - store_suspended_pull_query, update_min_version, Connection, DEFAULT_TAGS_PER_WORKSPACE, - DEFAULT_TAGS_WORKSPACES, INDEXER_CONFIG, SCRIPT_TOKEN_EXPIRY, SMTP_CONFIG, TMP_DIR, - WORKER_CONFIG, WORKER_GROUP, + load_env_vars, load_init_bash_from_env, load_whitelist_env_vars_from_env, load_worker_config, reload_custom_tags_setting, store_pull_query, store_suspended_pull_query, update_min_version, Connection, WorkerConfig, DEFAULT_TAGS_PER_WORKSPACE, DEFAULT_TAGS_WORKSPACES, INDEXER_CONFIG, SCRIPT_TOKEN_EXPIRY, SMTP_CONFIG, TMP_DIR, WORKER_CONFIG, WORKER_GROUP }, KillpillSender, BASE_URL, CRITICAL_ALERTS_ON_DB_OVERSIZE, CRITICAL_ALERT_MUTE_UI_ENABLED, CRITICAL_ERROR_CHANNELS, DB, DEFAULT_HUB_BASE_URL, HUB_BASE_URL, JOB_RETENTION_SECS, @@ -203,10 +199,23 @@ pub async fn initial_load( } Connection::Http(_) => { // TODO: reload worker config from http - WORKER_CONFIG.write().await.worker_tags = DECODED_AGENT_TOKEN - .as_ref() - .map(|x| x.tags.clone()) - .unwrap_or_default(); + let mut config = WORKER_CONFIG.write().await; + *config = WorkerConfig { + worker_tags: DECODED_AGENT_TOKEN + .as_ref() + .map(|x| x.tags.clone()) + .unwrap_or_default(), + env_vars: load_env_vars( + load_whitelist_env_vars_from_env(), + &std::collections::HashMap::new(), + ), + priority_tags_sorted: vec![], + dedicated_worker: None, + init_bash: load_init_bash_from_env(), + cache_clear: None, + additional_python_paths: None, + pip_local_dependencies: None, + }; } } } diff --git a/backend/windmill-common/src/worker.rs b/backend/windmill-common/src/worker.rs index b65b50a33a..ef4cfa0f46 100644 --- a/backend/windmill-common/src/worker.rs +++ b/backend/windmill-common/src/worker.rs @@ -382,10 +382,7 @@ fn normalize_path(path: &Path) -> PathBuf { ret } -pub fn is_allowed_file_location( - job_dir: &str, - user_defined_path: &str, -) -> error::Result { +pub fn is_allowed_file_location(job_dir: &str, user_defined_path: &str) -> error::Result { let job_dir = Path::new(job_dir); let user_path = PathBuf::from(user_defined_path); @@ -1398,18 +1395,66 @@ pub async fn load_worker_config( tracing::debug!("Custom tags priority set: {:?}", priority_tags_sorted); let env_vars_static = config.env_vars_static.unwrap_or_default().clone(); - let resolved_env_vars: HashMap = env_vars_static - .keys() - .map(|x| x.to_string()) - .chain(config.env_vars_allowlist.unwrap_or_default()) - .chain( - std::env::var("WHITELIST_ENVS") - .ok() - .map(|x| x.split(',').map(|x| x.to_string()).collect_vec()) - .unwrap_or_default() - .into_iter(), - ) - .sorted() + let resolved_env_vars: HashMap = load_env_vars( + config + .env_vars_allowlist + .unwrap_or_default() + .into_iter() + .chain(load_whitelist_env_vars_from_env()) + .chain(env_vars_static.keys().map(|x| x.to_string())), + &env_vars_static, + ); + + Ok(WorkerConfig { + worker_tags, + priority_tags_sorted, + dedicated_worker, + init_bash: config + .init_bash + .or_else(|| load_init_bash_from_env()) + .and_then(|x| if x.is_empty() { None } else { Some(x) }), + cache_clear: config.cache_clear, + pip_local_dependencies: config + .pip_local_dependencies + .or_else(|| load_pip_local_dependencies_from_env()), + additional_python_paths: config + .additional_python_paths + .or_else(|| load_additional_python_paths_from_env()), + env_vars: resolved_env_vars, + }) +} + +pub fn load_init_bash_from_env() -> Option { + std::env::var("INIT_SCRIPT") + .ok() + .and_then(|x| if x.is_empty() { None } else { Some(x) }) +} + +pub fn load_pip_local_dependencies_from_env() -> Option> { + std::env::var("PIP_LOCAL_DEPENDENCIES") + .ok() + .map(|x| x.split(',').map(|x| x.to_string()).collect_vec()) +} + +pub fn load_additional_python_paths_from_env() -> Option> { + std::env::var("ADDITIONAL_PYTHON_PATHS") + .ok() + .map(|x| x.split(':').map(|x| x.to_string()).collect_vec()) +} + +pub fn load_whitelist_env_vars_from_env() -> std::vec::IntoIter { + std::env::var("WHITELIST_ENVS") + .ok() + .map(|x| x.split(',').map(|x| x.to_string()).collect_vec()) + .unwrap_or_default() + .into_iter() +} + +pub fn load_env_vars( + iter: impl Iterator, + env_vars_static: &HashMap, +) -> HashMap { + iter.sorted() .unique() .map(|envvar_name| { ( @@ -1422,34 +1467,7 @@ pub async fn load_worker_config( }), ) }) - .collect(); - - Ok(WorkerConfig { - worker_tags, - priority_tags_sorted, - dedicated_worker, - init_bash: config - .init_bash - .or_else(|| std::env::var("INIT_SCRIPT").ok()) - .and_then(|x| if x.is_empty() { None } else { Some(x) }), - cache_clear: config.cache_clear, - pip_local_dependencies: config.pip_local_dependencies.or_else(|| { - let pip_local_dependencies = std::env::var("PIP_LOCAL_DEPENDENCIES") - .ok() - .map(|x| x.split(',').map(|x| x.to_string()).collect()); - if pip_local_dependencies == Some(vec!["".to_string()]) { - None - } else { - pip_local_dependencies - } - }), - additional_python_paths: config.additional_python_paths.or_else(|| { - std::env::var("ADDITIONAL_PYTHON_PATHS") - .ok() - .map(|x| x.split(':').map(|x| x.to_string()).collect()) - }), - env_vars: resolved_env_vars, - }) + .collect() } #[derive(Clone, PartialEq, Debug)] diff --git a/backend/windmill-worker/src/worker.rs b/backend/windmill-worker/src/worker.rs index 7199504beb..54820d52eb 100644 --- a/backend/windmill-worker/src/worker.rs +++ b/backend/windmill-worker/src/worker.rs @@ -1822,26 +1822,14 @@ async fn queue_init_bash_maybe<'c>( same_worker_tx: SameWorkerSender, worker_name: &str, ) -> anyhow::Result { - let uuid_content = match conn { - Connection::Sql(db) => { - if let Some(content) = WORKER_CONFIG.read().await.init_bash.clone() { - Some(( - push_init_job(db, content.clone(), worker_name).await?, - content, - )) - } else { - None - } - } - Connection::Http(client) => { - let init_script = std::env::var("INIT_SCRIPT"); - if init_script.is_ok() { - let content = init_script.unwrap(); - Some((queue_init_job(client, &content).await?, content)) - } else { - None - } - } + let uuid_content = if let Some(content) = WORKER_CONFIG.read().await.init_bash.clone() { + let uuid = match conn { + Connection::Sql(db) => push_init_job(db, content.clone(), worker_name).await?, + Connection::Http(client) => queue_init_job(client, &content).await?, + }; + Some((uuid, content)) + } else { + None }; if let Some((uuid, content)) = uuid_content { same_worker_tx